diff --git a/src/everos/component/embedding/openai_provider.py b/src/everos/component/embedding/openai_provider.py index 08d131d..1e0f436 100644 --- a/src/everos/component/embedding/openai_provider.py +++ b/src/everos/component/embedding/openai_provider.py @@ -23,7 +23,7 @@ from collections.abc import Sequence import openai -from everos.core.observability.tracing import set_generation_usage +from everos.core.observability.tracing import memory_span, set_generation_usage from .protocol import EmbeddingServiceError @@ -88,20 +88,23 @@ class OpenAIEmbeddingProvider: async def _embed_chunk(self, chunk: list[str]) -> list[list[float]]: """One ``/embeddings`` call, semaphore-guarded.""" - async with self._semaphore: - try: - response = await self._client.embeddings.create( - model=self._model, - input=chunk, - ) - except openai.OpenAIError as exc: - raise EmbeddingServiceError(str(exc)) from exc - # Surface token usage onto the active span (e.g. everos.search.embed_query). - # No-op when tracing is off; embeddings report only input (prompt) tokens. - usage = getattr(response, "usage", None) - set_generation_usage( - model=self._model, - input_tokens=usage.prompt_tokens if usage else None, - ) - # OpenAI returns ``data`` indexed by request order; truncate to ``dim``. - return [list(item.embedding[: self.dim]) for item in response.data] + # Wrap in an EMBEDDING-typed span so token usage lands on an embedding + # observation (which Langfuse can price) rather than the enclosing + # retriever span. No-op when tracing is off. + with memory_span("everos.embedding", observation_type="embedding"): + async with self._semaphore: + try: + response = await self._client.embeddings.create( + model=self._model, + input=chunk, + ) + except openai.OpenAIError as exc: + raise EmbeddingServiceError(str(exc)) from exc + # Embeddings report only input (prompt) tokens. + usage = getattr(response, "usage", None) + set_generation_usage( + model=self._model, + input_tokens=usage.prompt_tokens if usage else None, + ) + # OpenAI returns ``data`` indexed by request order; truncate to ``dim``. + return [list(item.embedding[: self.dim]) for item in response.data] diff --git a/src/everos/memory/extract/pipeline/user_memory.py b/src/everos/memory/extract/pipeline/user_memory.py index a58ee1b..e36532c 100644 --- a/src/everos/memory/extract/pipeline/user_memory.py +++ b/src/everos/memory/extract/pipeline/user_memory.py @@ -13,12 +13,14 @@ Run inside ``service.memorize`` via ``asyncio.gather`` alongside from __future__ import annotations +from pathlib import Path from typing import TYPE_CHECKING from everalgo.types import MemCell as AlgoMemCell from everalgo.user_memory import EpisodeExtractor from everos.component.utils.datetime import from_timestamp, to_iso_format +from everos.config import resolve_root from everos.core.observability.logging import get_logger from everos.core.observability.tracing import capture_output, memory_span from everos.memory import Episode, IngestResult, PipelineOutcome @@ -36,6 +38,14 @@ logger = get_logger(__name__) _TRACK = "user_memory" +def _root_relative(path: str) -> str: + """Path relative to the memory root, for telemetry (never the host abspath).""" + try: + return str(Path(path).relative_to(resolve_root())) + except ValueError: + return path + + class UserMemoryPipeline: """Per-sender Episode extraction on a list of pre-cut MemCells.""" @@ -152,8 +162,10 @@ class UserMemoryPipeline: ) ) md_paths.append(md_path) - # Written .md path (only when capture_content is on). - capture_output(persist_span, md_path) + # Written .md path, memory-root-relative (only when + # capture_content is on) — never leak the host absolute + # path to the telemetry backend. + capture_output(persist_span, _root_relative(md_path)) await self._engine.emit( EpisodeExtracted( memcell_id=memcell_id, diff --git a/src/everos/memory/search/agentic.py b/src/everos/memory/search/agentic.py index aab634a..1b79e8a 100644 --- a/src/everos/memory/search/agentic.py +++ b/src/everos/memory/search/agentic.py @@ -202,19 +202,17 @@ async def search_episodes_agentic( # 6. cluster_scoped: narrows hybrid_full to top-K cluster member expansions. async def cluster_scoped(q: str, _k: int) -> list[Candidate]: - with memory_span( - "everos.search.recall", - observation_type="retriever", - metadata={"phase": "agentic_cluster_scoped"}, - ): - return await acluster_retrieve( - q, - base_retrieve=hybrid_full, - base_candidates=_CLUSTER_BASE_CANDIDATES, - clusters=clusters, - all_docs=all_docs, - cluster_top_k=_CLUSTER_TOP_K, - ) + # No recall span here: cluster_scoped delegates to hybrid_full, which + # owns the recall span (and does the embedding). Wrapping it too + # produced a duplicate, same-name everos.search.recall nested in itself. + return await acluster_retrieve( + q, + base_retrieve=hybrid_full, + base_candidates=_CLUSTER_BASE_CANDIDATES, + clusters=clusters, + all_docs=all_docs, + cluster_top_k=_CLUSTER_TOP_K, + ) # 7. Cross-encoder rerank fn (2-arg RerankFn, no internal truncation). rerank_fn = build_rerank_fn( diff --git a/src/everos/memory/search/manager.py b/src/everos/memory/search/manager.py index 85c8e0f..14c1d27 100644 --- a/src/everos/memory/search/manager.py +++ b/src/everos/memory/search/manager.py @@ -42,6 +42,7 @@ from everos.core.context import resolve_request_id from everos.core.observability.logging import get_logger from everos.core.observability.tracing import ( capture_input, + capture_output, current_trace_ids, emit_recall_scores, memory_span, @@ -230,6 +231,17 @@ class SearchManager: unprocessed_messages=unprocessed, ) + # Returned hits (ids only) — content, so only when capture_content + # is on. This is the point of a search trace: what came back. + capture_output( + span, + { + "episodes": [e.id for e in data.episodes], + "agent_cases": [c.id for c in data.agent_cases], + "agent_skills": [s.id for s in data.agent_skills], + }, + ) + # Recall-quality signal on the span (always on; the Langfuse # scores push is separate and gated on creds — see the score sink). # `hit` is only meaningful for calibrated-score methods; leave it diff --git a/src/everos/service/memorize.py b/src/everos/service/memorize.py index d65fdfe..9fe806b 100644 --- a/src/everos/service/memorize.py +++ b/src/everos/service/memorize.py @@ -30,6 +30,7 @@ from pydantic import BaseModel from everos.component.llm import get_llm_client from everos.config import load_settings +from everos.core.context import resolve_request_id from everos.core.observability.logging import get_logger from everos.core.observability.tracing import memory_span from everos.core.persistence import MemoryRoot @@ -176,7 +177,11 @@ async def memorize( span_name, observation_type="span", session_id=session_id, - metadata={"mode": mode, "is_final": is_final}, + metadata={ + "mode": mode, + "is_final": is_final, + "request_id": resolve_request_id(), + }, ): async with asyncio.timeout(settings.memorize.session_lock_timeout_seconds): async with get_session_lock(session_id): @@ -197,15 +202,19 @@ async def _memorize_locked( ) -> MemorizeResult: """Inner critical section — runs under the per-session lock.""" ingested = await ingest_process(payload) - boundary = await prepare_cells( - ingested, - mode=mode, - is_final=is_final, - llm_client=get_llm_client(), - prompt_loader=_get_prompt_loader(), - hard_token_limit=boundary_cfg.hard_token_limit, - hard_msg_limit=boundary_cfg.hard_msg_limit, - ) + # Boundary detection runs an LLM (everalgo detect_boundaries). Wrap it in a + # generation span so its token usage lands on a GENERATION observation + # (costed by Langfuse) instead of the SPAN-typed request root (dropped). + with memory_span("everos.memcell.boundary", observation_type="generation"): + boundary = await prepare_cells( + ingested, + mode=mode, + is_final=is_final, + llm_client=get_llm_client(), + prompt_loader=_get_prompt_loader(), + hard_token_limit=boundary_cfg.hard_token_limit, + hard_msg_limit=boundary_cfg.hard_msg_limit, + ) if not boundary.cells: # Nothing went past the boundary stage — no pipelines to dispatch. diff --git a/tests/unit/test_memory/test_search/test_manager.py b/tests/unit/test_memory/test_search/test_manager.py index 88ae50c..f352575 100644 --- a/tests/unit/test_memory/test_search/test_manager.py +++ b/tests/unit/test_memory/test_search/test_manager.py @@ -1182,6 +1182,27 @@ async def test_search_captures_query_when_content_on(_search_spans: Any) -> None assert json.loads(attrs["langfuse.observation.input"])["query"] == "hi" +async def test_search_captures_returned_hits_when_content_on( + _search_spans: Any, +) -> None: + """capture_content on → the search span records the returned hit ids + (episodes/agent_cases/agent_skills), not just the query.""" + from everos.core.observability.tracing import set_capture_content + + set_capture_content(True) + try: + mgr = _build_manager(episode_sparse=[_episode_row("ep_1")]) + await mgr.search(_user_req(method=SearchMethod.KEYWORD)) + finally: + set_capture_content(False) + import json + + attrs = _span_index(_search_spans)["everos.memory.search"].attributes + out = json.loads(attrs["langfuse.observation.output"]) + assert out["episodes"] == ["ep_1"] + assert out["agent_cases"] == [] and out["agent_skills"] == [] + + async def test_search_omits_query_when_content_off(_search_spans: Any) -> None: mgr = _build_manager(episode_sparse=[_episode_row("ep_1")]) await mgr.search(_user_req(method=SearchMethod.KEYWORD))