From 357f619c64c5ce03e4c25b76f19e08f0528a0588 Mon Sep 17 00:00:00 2001 From: zhanghui Date: Thu, 23 Jul 2026 20:47:39 +0800 Subject: [PATCH] feat(observability): instrument memory ops + OME trace linking Open spans at the memory hot paths (all no-op when tracing is off): - add / flush (service.memorize), extract + persist.markdown (user pipeline). - search: everos.memory.search retriever + a uniform recall / rank decomposition across keyword / vector / hybrid / agentic (manager, agentic modules, cross-encoder callbacks); query-embedding tokens land on recall. - recall quality: top_score / hit on the search span, plus recall_top_score / recall_hit pushed to Langfuse scores via the bounded-queue sink (method tagged; off the request path). - OME: everos.ome. agent span + everos.reflect.consolidate generation; a W3C traceparent captured at enqueue is threaded through the APScheduler job and re-attached in the Runner, so strategies fanned out from a request nest under that request's trace. Co-Authored-By: Claude Opus 4.8 (1M context) --- src/everos/infra/ome/_dispatch/runner.py | 22 +- src/everos/infra/ome/engine.py | 14 +- .../memory/extract/pipeline/user_memory.py | 52 ++- src/everos/memory/reflection/orchestrator.py | 21 +- src/everos/memory/search/agentic.py | 45 ++- src/everos/memory/search/agentic_agent.py | 24 +- src/everos/memory/search/callbacks.py | 19 +- src/everos/memory/search/manager.py | 341 +++++++++++------- src/everos/service/memorize.py | 24 +- .../integration/test_memorize_integration.py | 65 ++++ tests/unit/test_infra/test_ome/test_engine.py | 49 +++ tests/unit/test_infra/test_ome/test_runner.py | 104 ++++++ .../test_pipeline/test_user_memory_emits.py | 152 ++++++++ .../unit/test_memory/test_get/test_manager.py | 19 + .../test_reflection/test_orchestrator.py | 49 +++ .../test_memory/test_search/test_agentic.py | 70 ++++ .../test_memory/test_search/test_manager.py | 201 ++++++++++- tests/unit/test_service/test_memorize_span.py | 67 ++++ 18 files changed, 1153 insertions(+), 185 deletions(-) create mode 100644 tests/unit/test_service/test_memorize_span.py diff --git a/src/everos/infra/ome/_dispatch/runner.py b/src/everos/infra/ome/_dispatch/runner.py index 58b4d8c..c7a1b17 100644 --- a/src/everos/infra/ome/_dispatch/runner.py +++ b/src/everos/infra/ome/_dispatch/runner.py @@ -29,6 +29,7 @@ from structlog.contextvars import bound_contextvars from everos.component.utils.datetime import get_utc_now from everos.core.observability.logging import get_logger +from everos.core.observability.tracing import memory_span, use_traceparent from everos.infra.ome._dispatch._state import _CURRENT_STRATEGY from everos.infra.ome._stores.run_record import RunRecordStore from everos.infra.ome.decorator import StrategyMeta @@ -113,6 +114,7 @@ class Runner: *, run_id: str, max_retries_snapshot: int, + traceparent: str | None = None, ) -> None: """Execute ``meta.func(event, ctx)`` with the attempt retry loop. @@ -141,6 +143,7 @@ class Runner: event_topic=event_topic, event_payload=event_payload, max_retries_snapshot=max_retries_snapshot, + traceparent=traceparent, ) if terminated: return @@ -155,6 +158,7 @@ class Runner: event_topic: str, event_payload: str, max_retries_snapshot: int, + traceparent: str | None = None, ) -> bool: """Run one attempt; return ``True`` if a terminal state was written (success / dead-letter or persistence failure), ``False`` @@ -185,7 +189,23 @@ class Runner: try: token = _CURRENT_STRATEGY.set(meta) try: - await meta.func(event, ctx) + # Continue the triggering request's trace when a traceparent + # was carried across the APScheduler boundary; otherwise the + # agent span roots its own trace (e.g. cron / recovery). + with ( + use_traceparent(traceparent), + memory_span( + f"everos.ome.{meta.name}", + observation_type="agent", + metadata={ + "strategy": meta.name, + "run_id": current_run_id, + "attempt": attempt, + "event_topic": event_topic, + }, + ), + ): + await meta.func(event, ctx) finally: _CURRENT_STRATEGY.reset(token) except StrategyContractError as e: diff --git a/src/everos/infra/ome/engine.py b/src/everos/infra/ome/engine.py index ac9bf62..6c8c840 100644 --- a/src/everos/infra/ome/engine.py +++ b/src/everos/infra/ome/engine.py @@ -24,6 +24,7 @@ from apscheduler.triggers.interval import IntervalTrigger from everos.component.utils.datetime import get_utc_now from everos.core.observability.logging import get_logger +from everos.core.observability.tracing import current_traceparent from everos.infra.ome._background.config_reloader import ConfigReloader from everos.infra.ome._background.crash_recovery import scan_and_resume from everos.infra.ome._background.idle_scanner import IdleScanner @@ -91,12 +92,15 @@ async def _runner_entry( event_topic: str, event_payload: str, max_retries_snapshot: int, + traceparent: str = "", ) -> None: """Module-level APS jobstore callback for a single run. Looks the engine up by id and hands off to :meth:`OfflineEngine.dispatch_run`. Pickle-safe (no closures, no - bound methods captured into APS jobstore args). + bound methods captured into APS jobstore args). ``traceparent`` defaults + to "" so crash-recovered jobs enqueued before this field existed still + unpack (they simply root their own trace). """ engine = _ENGINES.get(engine_id) if engine is None: @@ -112,6 +116,7 @@ async def _runner_entry( event_topic=event_topic, event_payload=event_payload, max_retries_snapshot=max_retries_snapshot, + traceparent=traceparent, ) @@ -644,6 +649,10 @@ class OfflineEngine: else self._config.max_retries ) event_topic = type(event).topic() + # Capture the triggering request's trace context (if any) here — this + # runs synchronously in the caller's task, so its span is still active. + # Carried as a pickle-safe string across the APScheduler boundary. + traceparent = current_traceparent() or "" self._on_run_enqueued() try: self._scheduler.add_job( @@ -657,6 +666,7 @@ class OfflineEngine: event_topic, event.model_dump_json(), max_retries_snapshot, + traceparent, ], id=run_id, replace_existing=False, @@ -700,6 +710,7 @@ class OfflineEngine: event_topic: str, event_payload: str, max_retries_snapshot: int, + traceparent: str | None = None, ) -> None: """APS jobstore callback target for one strategy run. @@ -724,6 +735,7 @@ class OfflineEngine: event, run_id=run_id, max_retries_snapshot=max_retries_snapshot, + traceparent=traceparent, ) finally: self._on_run_completed() diff --git a/src/everos/memory/extract/pipeline/user_memory.py b/src/everos/memory/extract/pipeline/user_memory.py index 394fa0c..a58ee1b 100644 --- a/src/everos/memory/extract/pipeline/user_memory.py +++ b/src/everos/memory/extract/pipeline/user_memory.py @@ -20,6 +20,7 @@ from everalgo.user_memory import EpisodeExtractor from everos.component.utils.datetime import from_timestamp, to_iso_format 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 from everos.memory.events import EpisodeExtracted, UserPipelineStarted from everos.memory.prompt_slots import PromptLoader @@ -99,9 +100,23 @@ class UserMemoryPipeline: # than the per-user fan-out per the algo's docstring). Fan-out # is then md-only: every user sender owns a copy of the same # narrative under its own owner_id path. - algo_ep = await self._ep_ext.aextract( - cell, sender_id=None, prompt=episode_prompt - ) + with memory_span( + "everos.extract", + observation_type="generation", + session_id=ingested.session_id, + metadata={ + "app_id": ingested.app_id, + "project_id": ingested.project_id, + "memcell_id": memcell_id, + }, + ) as extract_span: + # Token usage is recorded onto this span by the LLM client + # wrapper when the extractor issues its chat() call. + algo_ep = await self._ep_ext.aextract( + cell, sender_id=None, prompt=episode_prompt + ) + # Extracted memory text (only when capture_content is on). + capture_output(extract_span, algo_ep.episode) for sender_id in user_senders: ep = Episode.from_algo( algo_ep, @@ -111,15 +126,24 @@ class UserMemoryPipeline: parent_id=memcell_id, ) inline, sections = _episode_to_entry_body(ep) - eid = await self._episode_writer.append_entry( - ep.owner_id, - inline=inline, - sections=sections, - app_id=ingested.app_id, - project_id=ingested.project_id, - ) - md_paths.append( - str( + with memory_span( + "everos.persist.markdown", + observation_type="span", + session_id=ingested.session_id, + metadata={ + "owner_id": ep.owner_id, + "app_id": ingested.app_id, + "project_id": ingested.project_id, + }, + ) as persist_span: + eid = await self._episode_writer.append_entry( + ep.owner_id, + inline=inline, + sections=sections, + app_id=ingested.app_id, + project_id=ingested.project_id, + ) + md_path = str( self._episode_writer.path_for( ep.owner_id, eid.date, @@ -127,7 +151,9 @@ class UserMemoryPipeline: project_id=ingested.project_id, ) ) - ) + md_paths.append(md_path) + # Written .md path (only when capture_content is on). + capture_output(persist_span, md_path) await self._engine.emit( EpisodeExtracted( memcell_id=memcell_id, diff --git a/src/everos/memory/reflection/orchestrator.py b/src/everos/memory/reflection/orchestrator.py index 1c36d1a..bb8fe12 100644 --- a/src/everos/memory/reflection/orchestrator.py +++ b/src/everos/memory/reflection/orchestrator.py @@ -29,6 +29,7 @@ import numpy as np from everos.component.utils.datetime import from_timestamp, to_iso_format from everos.core.errors import AppError from everos.core.observability.logging import get_logger +from everos.core.observability.tracing import memory_span from everos.core.persistence import MemoryRoot from everos.infra.ome.context import StrategyContext from everos.memory._partition_locks import get_partition_lock @@ -580,13 +581,19 @@ class ReflectionOrchestrator: """ algo_episodes = _to_algo_episodes(episodes) try: - if is_update: - return await self._reflect_update( - algo_episodes=algo_episodes, - episodes=episodes, - merged_entry_ids=merged_entry_ids, - ) - return await self._reflector.areflect(algo_episodes) + with memory_span( + "everos.reflect.consolidate", + observation_type="generation", + metadata={"owner_id": owner_id, "is_update": is_update}, + ): + # Token usage lands on this span via the LLM client wrapper. + if is_update: + return await self._reflect_update( + algo_episodes=algo_episodes, + episodes=episodes, + merged_entry_ids=merged_entry_ids, + ) + return await self._reflector.areflect(algo_episodes) except AppError: logger.warning( "reflection_reflector_failed", diff --git a/src/everos/memory/search/agentic.py b/src/everos/memory/search/agentic.py index 6abcf73..aab634a 100644 --- a/src/everos/memory/search/agentic.py +++ b/src/everos/memory/search/agentic.py @@ -33,6 +33,7 @@ from everalgo.types import Candidate from everos.component.utils.datetime import from_timestamp, to_timestamp_ms from everos.core.observability.logging import get_logger +from everos.core.observability.tracing import memory_span from everos.infra.persistence.sqlite import cluster_repo from everos.memory.search.callbacks import build_rerank_fn from everos.memory.search.shaper import shape_episode_from_candidate @@ -168,15 +169,20 @@ async def search_episodes_agentic( # 4. hybrid_full: RRF fusion of dense + sparse MaxSim. async def hybrid_full(q: str, k: int) -> list[Candidate]: - return await ahybrid_retrieve( - q, - dense_retrieve=_dense, - sparse_retrieve=_sparse, - top_n=k, - dense_candidates=_DENSE_CANDIDATES, - sparse_candidates=_SPARSE_CANDIDATES, - rrf_k=_HYBRID_RRF_K, - ) + with memory_span( + "everos.search.recall", + observation_type="retriever", + metadata={"phase": "agentic_hybrid"}, + ): + return await ahybrid_retrieve( + q, + dense_retrieve=_dense, + sparse_retrieve=_sparse, + top_n=k, + dense_candidates=_DENSE_CANDIDATES, + sparse_candidates=_SPARSE_CANDIDATES, + rrf_k=_HYBRID_RRF_K, + ) # 5. Load cluster snapshot + full-corpus all_docs (memcell-keyed). # Reshape metadata to the everalgo doc contract so the sufficiency / @@ -196,14 +202,19 @@ 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]: - 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, - ) + 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, + ) # 7. Cross-encoder rerank fn (2-arg RerankFn, no internal truncation). rerank_fn = build_rerank_fn( diff --git a/src/everos/memory/search/agentic_agent.py b/src/everos/memory/search/agentic_agent.py index 70ede20..012eb5e 100644 --- a/src/everos/memory/search/agentic_agent.py +++ b/src/everos/memory/search/agentic_agent.py @@ -27,6 +27,7 @@ from everalgo.rank.agentic import aagentic_retrieve from everalgo.rank.hybrid import ahybrid_retrieve from everalgo.types import Candidate +from everos.core.observability.tracing import memory_span from everos.memory.search.callbacks import build_rerank_fn from everos.memory.search.shaper import ( shape_agent_case_from_candidate, @@ -165,15 +166,20 @@ async def _run_agentic_retrieve( return await recaller.sparse_recall(q, where, limit=k) async def hybrid_full(q: str, k: int) -> list[Candidate]: - return await ahybrid_retrieve( - q, - dense_retrieve=_dense, - sparse_retrieve=_sparse, - top_n=k, - dense_candidates=_DENSE_CANDIDATES, - sparse_candidates=_SPARSE_CANDIDATES, - rrf_k=_HYBRID_RRF_K, - ) + with memory_span( + "everos.search.recall", + observation_type="retriever", + metadata={"phase": "agentic_hybrid"}, + ): + return await ahybrid_retrieve( + q, + dense_retrieve=_dense, + sparse_retrieve=_sparse, + top_n=k, + dense_candidates=_DENSE_CANDIDATES, + sparse_candidates=_SPARSE_CANDIDATES, + rrf_k=_HYBRID_RRF_K, + ) rerank_fn = build_rerank_fn(reranker, text_field=recaller.text_field) diff --git a/src/everos/memory/search/callbacks.py b/src/everos/memory/search/callbacks.py index 8d7ea88..4490698 100644 --- a/src/everos/memory/search/callbacks.py +++ b/src/everos/memory/search/callbacks.py @@ -28,6 +28,7 @@ from everalgo.rank.protocols import RerankFn, RetrieveFn from everalgo.types import Candidate from everos.component.rerank import RerankProvider +from everos.core.observability.tracing import memory_span if TYPE_CHECKING: from .recall import KindRecaller @@ -64,7 +65,12 @@ def build_rerank_fn( if not items: return [] passages = [str(c.metadata.get(text_field, "")) for c in items] - results = await provider.rerank(query, passages, instruction=instruction) + with memory_span( + "everos.search.rank", + observation_type="span", + metadata={"phase": "cross_encoder"}, + ): + results = await provider.rerank(query, passages, instruction=instruction) out: list[Candidate] = [] for r in results: if not 0 <= r.index < len(items): @@ -112,9 +118,14 @@ def build_skill_rerank_fn(provider: RerankProvider) -> RerankFn: if not items: return [] passages = [_format_skill_passage(c) for c in items] - results = await provider.rerank( - query, passages, instruction=_SKILL_RERANK_INSTRUCTION - ) + with memory_span( + "everos.search.rank", + observation_type="span", + metadata={"phase": "cross_encoder_skill"}, + ): + results = await provider.rerank( + query, passages, instruction=_SKILL_RERANK_INSTRUCTION + ) out: list[Candidate] = [] for r in results: if not 0 <= r.index < len(items): diff --git a/src/everos/memory/search/manager.py b/src/everos/memory/search/manager.py index 7633df4..3a5fc70 100644 --- a/src/everos/memory/search/manager.py +++ b/src/everos/memory/search/manager.py @@ -37,8 +37,15 @@ from everalgo.rank.fusion import rrf from everalgo.types import Candidate, RankInput from everos.component.utils.datetime import to_display_tz +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 gen_request_id +from everos.core.observability.tracing import ( + capture_input, + current_trace_ids, + emit_recall_scores, + memory_span, +) from everos.infra.persistence.sqlite import ( UnprocessedBuffer, unprocessed_buffer_repo, @@ -123,6 +130,15 @@ _MAXSIM_FACT_POOL_CAP = 2000 _UNPROCESSED_TRACK = "memorize" +def _top_score(data: SearchData) -> float: + """Max relevance score across scored result items (0.0 when empty). + + Profiles are excluded — they are a KV fetch with no query-relevance score. + """ + items = [*data.episodes, *data.agent_cases, *data.agent_skills] + return max((item.score for item in items), default=0.0) + + class SearchManager: """Orchestrates per-kind recall, fusion, and shape into the public DTO.""" @@ -152,42 +168,81 @@ class SearchManager: # ── Public entry ──────────────────────────────────────────────── async def search(self, req: SearchRequest) -> SearchResponse: - request_id = gen_request_id() - # Compile filters first: a malformed `filters` payload is a user - # input error (422) and should surface before the server-side - # component guard (500). The two steps are independent. - where = compile_filters( - req.filters, - owner_id=req.owner_id, - owner_type=req.owner_type, - app_id=req.app_id, - project_id=req.project_id, - ) - self._validate_components(req) + request_id = resolve_request_id() + with memory_span( + "everos.memory.search", + observation_type="retriever", + session_id=_extract_top_level_session_id(req.filters), + user_id=req.user_id, + metadata={ + "request_id": request_id, + "app_id": req.app_id, + "project_id": req.project_id, + "agent_id": req.agent_id, + "owner_type": req.owner_type, + "method": req.method.value, + }, + ) as span: + # Content (query text) only when capture_content is on. + capture_input( + span, + {"query": req.query, "top_k": req.top_k, "method": req.method.value}, + ) + # Compile filters first: a malformed `filters` payload is a user + # input error (422) and should surface before the server-side + # component guard (500). The two steps are independent. + where = compile_filters( + req.filters, + owner_id=req.owner_id, + owner_type=req.owner_type, + app_id=req.app_id, + project_id=req.project_id, + ) + self._validate_components(req) - if req.owner_type == "user": - episodes, profiles, unprocessed = await asyncio.gather( - self._search_episodes(req, where), - self._fetch_profile(req), - self._load_unprocessed(req), - ) - data = SearchData( - episodes=episodes, - profiles=profiles, - unprocessed_messages=unprocessed, - ) - else: # "agent" - (cases, skills), unprocessed = await asyncio.gather( - self._search_cases_and_skills(req, where), - self._load_unprocessed(req), - ) - data = SearchData( - agent_cases=cases, - agent_skills=skills, - unprocessed_messages=unprocessed, - ) + if req.owner_type == "user": + episodes, profiles, unprocessed = await asyncio.gather( + self._search_episodes(req, where), + self._fetch_profile(req), + self._load_unprocessed(req), + ) + data = SearchData( + episodes=episodes, + profiles=profiles, + unprocessed_messages=unprocessed, + ) + else: # "agent" + (cases, skills), unprocessed = await asyncio.gather( + self._search_cases_and_skills(req, where), + self._load_unprocessed(req), + ) + data = SearchData( + agent_cases=cases, + agent_skills=skills, + unprocessed_messages=unprocessed, + ) - return SearchResponse(request_id=request_id, data=data) + # Recall-quality signal on the span (always on; the Langfuse + # scores push is separate and gated on creds — see the score sink). + top_score = _top_score(data) + threshold = load_settings().observability.recall_hit_threshold + hit = top_score >= threshold + span.set_attribute("everos.search.top_score", top_score) + span.set_attribute("everos.search.hit", hit) + + # Push recall-quality scores to Langfuse out-of-band (no-op unless + # a score sink is configured); attach to this retriever span. + ids = current_trace_ids() + if ids is not None: + emit_recall_scores( + trace_id=ids[0], + observation_id=ids[1], + top_score=top_score, + hit=hit, + method=req.method.value, + ) + + return SearchResponse(request_id=request_id, data=data) # ── Unprocessed buffer ────────────────────────────────────────── @@ -267,12 +322,17 @@ class SearchManager: # ── KEYWORD / VECTOR: single-route recall ── if fusion_mode is None: - if req.method == SearchMethod.KEYWORD: - cands = await self._ep.sparse_recall( - req.query, where, limit=self._recall_limit(req.top_k) - ) - else: - cands = await self._maxsim_atomic_recall(req, where, top_k) + with memory_span( + "everos.search.recall", + observation_type="retriever", + metadata={"phase": "single_route", "method": req.method.value}, + ): + if req.method == SearchMethod.KEYWORD: + cands = await self._ep.sparse_recall( + req.query, where, limit=self._recall_limit(req.top_k) + ) + else: + cands = await self._maxsim_atomic_recall(req, where, top_k) # ``atomic_facts`` stays empty: facts come back only when the HYBRID # pipeline surfaces them with a score (see ``reshape_hybrid_output``). # Single-route recall has no per-fact score against the query, so @@ -290,43 +350,53 @@ class SearchManager: ) if fusion_mode == "hierarchy": - rrf_candidates = rrf(sparse, dense) - ep_to_parents = build_ep_to_fact_parents(rrf_candidates) - episode_to_facts = await self._fact.facts_for_episodes( - ep_to_parents, - where, - per_episode=max(top_k * 2, 20), - query_vector=query_vector, - ) - scored = heap_expand( - sparse=sparse, - dense=dense, - episode_to_facts=episode_to_facts, - top_k=top_k, - ) - episode_pool = {c.id: c for c in (*sparse, *dense)} - shaped = reshape_hybrid_output(scored, episode_pool=episode_pool) - if req.min_score is not None: - shaped = [s for s in shaped if s.score >= req.min_score] - return shaped + with memory_span( + "everos.search.rank", + observation_type="span", + metadata={"phase": "hierarchy"}, + ): + rrf_candidates = rrf(sparse, dense) + ep_to_parents = build_ep_to_fact_parents(rrf_candidates) + episode_to_facts = await self._fact.facts_for_episodes( + ep_to_parents, + where, + per_episode=max(top_k * 2, 20), + query_vector=query_vector, + ) + scored = heap_expand( + sparse=sparse, + dense=dense, + episode_to_facts=episode_to_facts, + top_k=top_k, + ) + episode_pool = {c.id: c for c in (*sparse, *dense)} + shaped = reshape_hybrid_output(scored, episode_pool=episode_pool) + if req.min_score is not None: + shaped = [s for s in shaped if s.score >= req.min_score] + return shaped # rrf / lr: standard everalgo fusion path (fallback). - output = await arank( - RankInput( - query=req.query, - memory_type=self._ep.everalgo_memory_type, # type: ignore[arg-type] - sparse_candidates=sparse, - dense_candidates=dense, - top_k=top_k, - radius=_effective_radius(req), - ), - config=RankConfig(fusion_mode=fusion_mode) - if fusion_mode != "rrf" - else DEFAULT_RANK_CONFIG, - llm=self._llm, - enable_rerank=enable_rerank, - rerank_top_k=top_k, - ) + with memory_span( + "everos.search.rank", + observation_type="span", + metadata={"phase": fusion_mode}, + ): + output = await arank( + RankInput( + query=req.query, + memory_type=self._ep.everalgo_memory_type, # type: ignore[arg-type] + sparse_candidates=sparse, + dense_candidates=dense, + top_k=top_k, + radius=_effective_radius(req), + ), + config=RankConfig(fusion_mode=fusion_mode) + if fusion_mode != "rrf" + else DEFAULT_RANK_CONFIG, + llm=self._llm, + enable_rerank=enable_rerank, + rerank_top_k=top_k, + ) ep_candidates = (_scored_as_candidate(s) for s in output.items) return [ ep @@ -363,22 +433,27 @@ class SearchManager: sparse, dense, _ = await self._recall_sparse_dense( self._case, req, where, top_k, cap=_AGENT_TOP_K_CAP ) - output = await arank( - RankInput( - query=req.query, - memory_type=self._case.everalgo_memory_type, # type: ignore[arg-type] - sparse_candidates=sparse, - dense_candidates=dense, - top_k=top_k, - radius=_effective_radius(req), - ), - config=RankConfig(fusion_mode=fusion_mode) - if fusion_mode != "rrf" - else DEFAULT_RANK_CONFIG, - llm=self._llm, - enable_rerank=enable_rerank, - rerank_top_k=top_k, - ) + with memory_span( + "everos.search.rank", + observation_type="span", + metadata={"phase": fusion_mode, "kind": "agent_case"}, + ): + output = await arank( + RankInput( + query=req.query, + memory_type=self._case.everalgo_memory_type, # type: ignore[arg-type] + sparse_candidates=sparse, + dense_candidates=dense, + top_k=top_k, + radius=_effective_radius(req), + ), + config=RankConfig(fusion_mode=fusion_mode) + if fusion_mode != "rrf" + else DEFAULT_RANK_CONFIG, + llm=self._llm, + enable_rerank=enable_rerank, + rerank_top_k=top_k, + ) case_candidates = (_scored_as_candidate(s) for s in output.items) shaped = (shape_agent_case_from_candidate(c) for c in case_candidates) return [item for item in shaped if item is not None] @@ -432,25 +507,33 @@ class SearchManager: # to the skill facade (adds the skill-only 0.4 relevance gate). # Config is ``rrf`` — ``skill_hybrid`` is an everos routing # label, not an everalgo fusion mode. - output = await arank( - RankInput( - query=req.query, - memory_type=self._skill.everalgo_memory_type, # type: ignore[arg-type] - sparse_candidates=sparse, - dense_candidates=dense, - top_k=top_k, - radius=_effective_radius(req), - ), - config=DEFAULT_RANK_CONFIG, - llm=self._llm, - enable_rerank=True, - rerank_top_k=top_k, - ) + with memory_span( + "everos.search.rank", + observation_type="span", + metadata={"phase": "skill_llm", "kind": "agent_skill"}, + ): + output = await arank( + RankInput( + query=req.query, + memory_type=self._skill.everalgo_memory_type, # type: ignore[arg-type] + sparse_candidates=sparse, + dense_candidates=dense, + top_k=top_k, + radius=_effective_radius(req), + ), + config=DEFAULT_RANK_CONFIG, + llm=self._llm, + enable_rerank=True, + rerank_top_k=top_k, + ) skill_candidates = (_scored_as_candidate(s) for s in output.items) shaped = (shape_agent_skill_from_candidate(c) for c in skill_candidates) return [item for item in shaped if item is not None] # Cross-encoder lane (default): rrf + skill-shaped cross-encoder rerank. + # The rank span is emitted inside build_skill_rerank_fn (callbacks), + # so the cross-encoder rerank is covered uniformly with the agentic + # path rather than double-wrapped here. return await search_agent_skills_hybrid( req.query, sparse=sparse, @@ -477,15 +560,20 @@ class SearchManager: *, cap: int = _DEFAULT_TOP_K_CAP, ) -> list[Candidate]: - if req.method == SearchMethod.KEYWORD: - return await recaller.sparse_recall( - req.query, where, limit=self._recall_limit(req.top_k, cap=cap) + with memory_span( + "everos.search.recall", + observation_type="retriever", + metadata={"phase": "single_route", "method": req.method.value}, + ): + if req.method == SearchMethod.KEYWORD: + return await recaller.sparse_recall( + req.query, where, limit=self._recall_limit(req.top_k, cap=cap) + ) + vector = await self._embed_query(req.query) + cands = await recaller.dense_recall( + vector, where, limit=self._recall_limit(req.top_k, cap=cap) ) - vector = await self._embed_query(req.query) - cands = await recaller.dense_recall( - vector, where, limit=self._recall_limit(req.top_k, cap=cap) - ) - return self._apply_radius(cands, _effective_radius(req)) + return self._apply_radius(cands, _effective_radius(req)) async def _recall_sparse_dense( self, @@ -504,16 +592,21 @@ class SearchManager: the query. Returns ``[]`` for ``vector`` when no embedding provider is configured. """ - vector = await self._embed_query(req.query) - limit = self._recall_limit(req.top_k, cap=cap) - sparse, dense = await asyncio.gather( - recaller.sparse_recall(req.query, where, limit=limit), - recaller.dense_recall(vector, where, limit=limit) - if vector - else _empty_candidates(), - ) - dense = self._apply_radius(dense, _effective_radius(req)) - return sparse, dense, vector + with memory_span( + "everos.search.recall", + observation_type="retriever", + metadata={"phase": "sparse_dense", "method": req.method.value}, + ): + vector = await self._embed_query(req.query) + limit = self._recall_limit(req.top_k, cap=cap) + sparse, dense = await asyncio.gather( + recaller.sparse_recall(req.query, where, limit=limit), + recaller.dense_recall(vector, where, limit=limit) + if vector + else _empty_candidates(), + ) + dense = self._apply_radius(dense, _effective_radius(req)) + return sparse, dense, vector async def _maxsim_atomic_recall( self, req: SearchRequest, where: str, top_k: int diff --git a/src/everos/service/memorize.py b/src/everos/service/memorize.py index a92a04b..03d455d 100644 --- a/src/everos/service/memorize.py +++ b/src/everos/service/memorize.py @@ -31,6 +31,7 @@ from pydantic import BaseModel from everos.component.llm import get_llm_client from everos.config import load_settings from everos.core.observability.logging import get_logger +from everos.core.observability.tracing import memory_span from everos.core.persistence import MemoryRoot from everos.infra.ome.config import OMEConfig from everos.infra.ome.engine import OfflineEngine @@ -170,14 +171,21 @@ async def memorize( boundary_cfg = settings.boundary_detection session_id = payload["session_id"] - async with asyncio.timeout(settings.memorize.session_lock_timeout_seconds): - async with get_session_lock(session_id): - return await _memorize_locked( - payload, - mode=mode, - boundary_cfg=boundary_cfg, - is_final=is_final, - ) + span_name = "everos.memory.flush" if is_final else "everos.memory.add" + with memory_span( + span_name, + observation_type="span", + session_id=session_id, + metadata={"mode": mode, "is_final": is_final}, + ): + async with asyncio.timeout(settings.memorize.session_lock_timeout_seconds): + async with get_session_lock(session_id): + return await _memorize_locked( + payload, + mode=mode, + boundary_cfg=boundary_cfg, + is_final=is_final, + ) async def _memorize_locked( diff --git a/tests/integration/test_memorize_integration.py b/tests/integration/test_memorize_integration.py index 54277ba..c591958 100644 --- a/tests/integration/test_memorize_integration.py +++ b/tests/integration/test_memorize_integration.py @@ -689,3 +689,68 @@ async def test_same_session_multi_add_concatenates( assert len(rows) == 1 # one cell from the flush ids = json.loads(rows[0]["message_ids_json"]) assert len(ids) == 6 # all 6 messages folded in + + +# --------------------------------------------------------------------------- +# Tracing: the add/flush span nests extract + persist in one trace +# --------------------------------------------------------------------------- + + +async def test_flush_produces_nested_trace( + tmp_path: Path, + memorize_env: Callable[..., Any], +) -> None: + """A real flush that extracts one Episode emits everos.memory.flush with + everos.extract + everos.persist.markdown as children of one trace.""" + from opentelemetry.sdk.trace.export import SimpleSpanProcessor + from opentelemetry.sdk.trace.export.in_memory_span_exporter import ( + InMemorySpanExporter, + ) + + from everos.config.settings import ObservabilitySettings + from everos.core.observability.tracing import ( + force_flush, + init_tracing, + shutdown_tracing, + ) + + fake = _make_fake_llm(boundary_responses=[[]]) + await memorize_env(mode="chat", fake_llm=fake) + + exporter = InMemorySpanExporter() + shutdown_tracing() + init_tracing( + ObservabilitySettings(enabled=True, endpoint="http://collector.invalid"), + span_processor=SimpleSpanProcessor(exporter), + ) + try: + payload = { + "session_id": "trace_sess", + "messages": [ + _user("hello", 1_700_000_000_000), + _assistant("hi there", 1_700_000_001_000), + ], + } + result = await memorize(payload, is_final=True) + assert result.status == "extracted" + force_flush() + finally: + shutdown_tracing() + + spans = {s.name: s for s in exporter.get_finished_spans()} + assert "everos.memory.flush" in spans + assert "everos.extract" in spans + assert "everos.persist.markdown" in spans + + root = spans["everos.memory.flush"] + extract = spans["everos.extract"] + persist = spans["everos.persist.markdown"] + + # One trace: all three share the flush span's trace id. + trace_id = root.context.trace_id + assert extract.context.trace_id == trace_id + assert persist.context.trace_id == trace_id + # flush is the root; extract / persist hang beneath it (not siblings). + assert root.parent is None + assert extract.parent is not None + assert persist.parent is not None diff --git a/tests/unit/test_infra/test_ome/test_engine.py b/tests/unit/test_infra/test_ome/test_engine.py index b46c408..682ce3f 100644 --- a/tests/unit/test_infra/test_ome/test_engine.py +++ b/tests/unit/test_infra/test_ome/test_engine.py @@ -621,3 +621,52 @@ async def test_enqueue_run_rolls_back_counter_on_add_job_failure( finally: monkeypatch.undo() await engine.stop() + + +@pytest.mark.asyncio +async def test_engine_emit_links_strategy_to_request_trace(cfg: OMEConfig) -> None: + """End-to-end: emit inside a request span → the engine captures the + traceparent at enqueue, threads it across APScheduler, and the strategy + body runs under the SAME trace as the triggering request.""" + from opentelemetry.sdk.trace.export import SimpleSpanProcessor + from opentelemetry.sdk.trace.export.in_memory_span_exporter import ( + InMemorySpanExporter, + ) + + from everos.config.settings import ObservabilitySettings + from everos.core.observability.tracing import ( + current_trace_ids, + init_tracing, + memory_span, + shutdown_tracing, + ) + + seen_tid: list[str] = [] + + @offline_strategy(name="tid_collector", trigger=Immediate(on=[_E]), emits=[]) + async def s(event: _E, ctx: StrategyContext) -> None: + ids = current_trace_ids() # inside the everos.ome.* span + seen_tid.append(ids[0] if ids else "") + + engine = OfflineEngine(config=cfg) + engine.register(s) + shutdown_tracing() + init_tracing( + ObservabilitySettings(enabled=True, endpoint="http://collector.invalid"), + span_processor=SimpleSpanProcessor(InMemorySpanExporter()), + ) + await engine.start() + try: + with memory_span("everos.memory.flush", observation_type="span") as req: + req_tid = format(req.get_span_context().trace_id, "032x") + await engine.emit(_E()) # traceparent captured at _enqueue_run here + for _ in range(50): + if seen_tid: + break + await asyncio.sleep(0.05) + finally: + await engine.stop() + shutdown_tracing() + + assert seen_tid, "strategy did not run" + assert seen_tid[0] == req_tid # same trace as the triggering request diff --git a/tests/unit/test_infra/test_ome/test_runner.py b/tests/unit/test_infra/test_ome/test_runner.py index 696de53..c7cb26a 100644 --- a/tests/unit/test_infra/test_ome/test_runner.py +++ b/tests/unit/test_infra/test_ome/test_runner.py @@ -228,3 +228,107 @@ async def test_runner_aborts_silently_when_mark_running_fails( async def _no_emit(event: BaseEvent) -> None: return None + + +@pytest.mark.asyncio +async def test_runner_emits_ome_agent_span(setup) -> None: + """Runner wraps the strategy body in an everos.ome. agent span + (its own trace — runs in an APScheduler task, no request context).""" + rec_store, sem = setup + + @offline_strategy(name="traced_strat", trigger=Immediate(on=[_E]), emits=[]) + async def s(event: _E, ctx: StrategyContext) -> None: + return None + + from opentelemetry.sdk.trace.export import SimpleSpanProcessor + from opentelemetry.sdk.trace.export.in_memory_span_exporter import ( + InMemorySpanExporter, + ) + + from everos.config.settings import ObservabilitySettings + from everos.core.observability.tracing import ( + force_flush, + init_tracing, + shutdown_tracing, + ) + + exporter = InMemorySpanExporter() + shutdown_tracing() + init_tracing( + ObservabilitySettings(enabled=True, endpoint="http://collector.invalid"), + span_processor=SimpleSpanProcessor(exporter), + ) + try: + runner = Runner( + run_record_store=rec_store, + engine_sem=sem, + emit_hook=_no_emit, + engine=MagicMock(), + ) + await runner.run(s.meta, _E(), run_id="r_trace", max_retries_snapshot=1) + force_flush() + finally: + shutdown_tracing() + + spans = {sp.name: sp for sp in exporter.get_finished_spans()} + assert "everos.ome.traced_strat" in spans + assert ( + spans["everos.ome.traced_strat"].attributes["langfuse.observation.type"] + == "agent" + ) + + +@pytest.mark.asyncio +async def test_runner_ome_span_links_to_upstream_traceparent(setup) -> None: + """Given a traceparent (captured where a request span was active), the + everos.ome. span nests under that upstream trace, not a new root.""" + rec_store, sem = setup + + @offline_strategy(name="linked_strat", trigger=Immediate(on=[_E]), emits=[]) + async def s(event: _E, ctx: StrategyContext) -> None: + return None + + from opentelemetry.sdk.trace.export import SimpleSpanProcessor + from opentelemetry.sdk.trace.export.in_memory_span_exporter import ( + InMemorySpanExporter, + ) + + from everos.config.settings import ObservabilitySettings + from everos.core.observability.tracing import ( + current_traceparent, + force_flush, + init_tracing, + memory_span, + shutdown_tracing, + ) + + exporter = InMemorySpanExporter() + shutdown_tracing() + init_tracing( + ObservabilitySettings(enabled=True, endpoint="http://collector.invalid"), + span_processor=SimpleSpanProcessor(exporter), + ) + try: + # Simulate the triggering request: capture its traceparent, then close. + with memory_span("everos.memory.flush", observation_type="span") as parent: + tp = current_traceparent() + parent_tid = parent.get_span_context().trace_id + + runner = Runner( + run_record_store=rec_store, + engine_sem=sem, + emit_hook=_no_emit, + engine=MagicMock(), + ) + await runner.run( + s.meta, _E(), run_id="r_link", max_retries_snapshot=1, traceparent=tp + ) + force_flush() + finally: + shutdown_tracing() + + ome = {sp.name: sp for sp in exporter.get_finished_spans()}[ + "everos.ome.linked_strat" + ] + assert ome.context.trace_id == parent_tid # same trace as the request + assert ome.parent is not None # child, not a fresh root diff --git a/tests/unit/test_memory/test_extract/test_pipeline/test_user_memory_emits.py b/tests/unit/test_memory/test_extract/test_pipeline/test_user_memory_emits.py index ab70c7a..f696d72 100644 --- a/tests/unit/test_memory/test_extract/test_pipeline/test_user_memory_emits.py +++ b/tests/unit/test_memory/test_extract/test_pipeline/test_user_memory_emits.py @@ -121,3 +121,155 @@ async def test_emit_episode_extracted_after_md_write() -> None: assert extracted[0].owner_id == "u1" assert extracted[0].session_id == "s1" assert extracted[0].source == "pipeline" + + +async def test_run_emits_extract_and_persist_spans() -> None: + """pipeline.run opens an everos.extract generation span (where LLM token + usage lands) and an everos.persist.markdown span around the md write.""" + from opentelemetry.sdk.trace.export import SimpleSpanProcessor + from opentelemetry.sdk.trace.export.in_memory_span_exporter import ( + InMemorySpanExporter, + ) + + from everos.config.settings import ObservabilitySettings + from everos.core.observability.tracing import ( + force_flush, + init_tracing, + shutdown_tracing, + ) + + engine = _CapturingEngine() + episode_writer = MagicMock() + episode_writer.append_entry = AsyncMock( + return_value=EntryId(prefix="ep", date=_dt.date(2026, 5, 17), seq=1) + ) + episode_writer.path_for = MagicMock(return_value="users/u1/episodes/x.md") + prompt_loader = MagicMock() + prompt_loader.load = MagicMock(return_value="") + pipeline = UserMemoryPipeline( + episode_writer=episode_writer, + prompt_loader=prompt_loader, + llm_client=MagicMock(), + engine=engine, + ) + cell = _sample_memcell() + ingested = IngestResult( + session_id="s1", + messages=[ + CanonicalMessage( + message_id="m1", + session_id="s1", + sender_id="u1", + role="user", + timestamp=_dt.datetime.fromtimestamp(1_700_000_000, tz=_dt.UTC), + text="hello", + ) + ], + ) + algo_ep = AlgoEpisode( + owner_id="u1", episode="they said hello", timestamp=1_700_000_000_000 + ) + + exporter = InMemorySpanExporter() + init_tracing( + ObservabilitySettings(enabled=True, endpoint="http://collector.invalid"), + span_processor=SimpleSpanProcessor(exporter), + ) + try: + with patch.object( + pipeline._ep_ext, "aextract", new=AsyncMock(return_value=algo_ep) + ): + await pipeline.run( + ingested=ingested, + cells=[cell], + memcell_ids=["mc_a"], + per_cell_all_senders=[["u1"]], + ) + force_flush() + spans = {s.name: s for s in exporter.get_finished_spans()} + assert "everos.extract" in spans + assert "everos.persist.markdown" in spans + assert spans["everos.extract"].attributes["langfuse.observation.type"] == ( + "generation" + ) + finally: + shutdown_tracing() + + +async def test_extract_persist_capture_content_when_on() -> None: + """With capture_content on, everos.extract carries the episode text and + everos.persist.markdown carries the written .md path; off → neither.""" + from opentelemetry.sdk.trace.export import SimpleSpanProcessor + from opentelemetry.sdk.trace.export.in_memory_span_exporter import ( + InMemorySpanExporter, + ) + + from everos.config.settings import ObservabilitySettings + from everos.core.observability.tracing import ( + force_flush, + init_tracing, + set_capture_content, + shutdown_tracing, + ) + + engine = _CapturingEngine() + episode_writer = MagicMock() + episode_writer.append_entry = AsyncMock( + return_value=EntryId(prefix="ep", date=_dt.date(2026, 5, 17), seq=1) + ) + episode_writer.path_for = MagicMock(return_value="users/u1/episodes/x.md") + prompt_loader = MagicMock() + prompt_loader.load = MagicMock(return_value="") + pipeline = UserMemoryPipeline( + episode_writer=episode_writer, + prompt_loader=prompt_loader, + llm_client=MagicMock(), + engine=engine, + ) + cell = _sample_memcell() + ingested = IngestResult( + session_id="s1", + messages=[ + CanonicalMessage( + message_id="m1", + session_id="s1", + sender_id="u1", + role="user", + timestamp=_dt.datetime.fromtimestamp(1_700_000_000, tz=_dt.UTC), + text="hello", + ) + ], + ) + algo_ep = AlgoEpisode( + owner_id="u1", episode="they said hello", timestamp=1_700_000_000_000 + ) + + exporter = InMemorySpanExporter() + init_tracing( + ObservabilitySettings(enabled=True, endpoint="http://collector.invalid"), + span_processor=SimpleSpanProcessor(exporter), + ) + set_capture_content(True) + try: + with patch.object( + pipeline._ep_ext, "aextract", new=AsyncMock(return_value=algo_ep) + ): + await pipeline.run( + ingested=ingested, + cells=[cell], + memcell_ids=["mc_a"], + per_cell_all_senders=[["u1"]], + ) + force_flush() + finally: + set_capture_content(False) + shutdown_tracing() + + spans = {s.name: s for s in exporter.get_finished_spans()} + assert spans["everos.extract"].attributes["langfuse.observation.output"] == ( + "they said hello" + ) + assert ( + spans["everos.persist.markdown"].attributes["langfuse.observation.output"] + == "users/u1/episodes/x.md" + ) diff --git a/tests/unit/test_memory/test_get/test_manager.py b/tests/unit/test_memory/test_get/test_manager.py index cb4e167..35cec6b 100644 --- a/tests/unit/test_memory/test_get/test_manager.py +++ b/tests/unit/test_memory/test_get/test_manager.py @@ -220,6 +220,25 @@ async def test_episodic_memory_populates_episodes_and_counts( assert all(item.user_id == "u1" for item in resp.data.episodes) +async def test_get_uses_propagated_request_id_when_bound( + manager: tuple[GetManager, _StubRepo, _StubRepo, _StubRepo], +) -> None: + """When a request id is bound upstream (middleware), ``get`` reuses it + instead of minting a fresh one, so the response id matches the trace.""" + from everos.core.context import reset_request_id, set_request_id + + mgr, ep, _, _ = manager + ep.rows = [_episode_row("ep_1")] + token = set_request_id("deadbeef" * 4) + try: + resp = await mgr.get( + GetRequest(user_id="u1", memory_type=GetMemoryType.EPISODE) + ) + assert resp.request_id == "deadbeef" * 4 + finally: + reset_request_id(token) + + async def test_episodic_memory_passes_where_and_sort_to_repo( manager: tuple[GetManager, _StubRepo, _StubRepo, _StubRepo], ) -> None: diff --git a/tests/unit/test_memory/test_reflection/test_orchestrator.py b/tests/unit/test_memory/test_reflection/test_orchestrator.py index c168204..a413e4a 100644 --- a/tests/unit/test_memory/test_reflection/test_orchestrator.py +++ b/tests/unit/test_memory/test_reflection/test_orchestrator.py @@ -434,3 +434,52 @@ def test_ts_to_ms_datetime() -> None: def test_ts_to_ms_int_passthrough() -> None: """int -> int passthrough.""" assert _ts_to_ms(1717200000000) == 1717200000000 + + +async def test_call_reflector_emits_consolidate_generation_span() -> None: + """_call_reflector wraps the reflector call in an everos.reflect.consolidate + generation span (token usage lands here via the LLM client wrapper).""" + from opentelemetry.sdk.trace.export import SimpleSpanProcessor + from opentelemetry.sdk.trace.export.in_memory_span_exporter import ( + InMemorySpanExporter, + ) + + from everos.config.settings import ObservabilitySettings + from everos.core.observability.tracing import ( + force_flush, + init_tracing, + shutdown_tracing, + ) + + reflector = MagicMock() + reflector.areflect = AsyncMock( + return_value=_FakeAlgoResult( + owner_id=None, episode="merged", subject="s", timestamp=1717200000000 + ) + ) + orch = _build_orchestrator(reflector=reflector) + + exporter = InMemorySpanExporter() + shutdown_tracing() + init_tracing( + ObservabilitySettings(enabled=True, endpoint="http://collector.invalid"), + span_processor=SimpleSpanProcessor(exporter), + ) + try: + result = await orch._call_reflector( + episodes=[_make_episode_row()], + merged_entry_ids=[], + is_update=False, + owner_id="u_alice", + ) + force_flush() + finally: + shutdown_tracing() + + assert result is not None + spans = {s.name: s for s in exporter.get_finished_spans()} + assert "everos.reflect.consolidate" in spans + assert ( + spans["everos.reflect.consolidate"].attributes["langfuse.observation.type"] + == "generation" + ) diff --git a/tests/unit/test_memory/test_search/test_agentic.py b/tests/unit/test_memory/test_search/test_agentic.py index 9b48ee3..9ab4896 100644 --- a/tests/unit/test_memory/test_search/test_agentic.py +++ b/tests/unit/test_memory/test_search/test_agentic.py @@ -366,3 +366,73 @@ def test_restore_shaper_metadata_reverts_bridged_fields() -> None: assert isinstance(restored["timestamp"], _dt.datetime) assert restored["timestamp"] == original assert restored["episode"] == "x" + + +async def test_agentic_emits_recall_and_rank_spans( + ep_recaller: _StubEpisodeRecaller, + fact_recaller: _StubFactRecaller, + clusters: list[Cluster], +) -> None: + """The agentic recall closures (base_retrieve) and the cross-encoder + rerank_fn emit everos.search.recall / everos.search.rank spans. Driven + by a fake aagentic_retrieve that actually invokes both callbacks — the + same way the real everalgo driver does — so the assertion is + deterministic and independent of everalgo's loop internals.""" + from opentelemetry.sdk.trace.export import SimpleSpanProcessor + from opentelemetry.sdk.trace.export.in_memory_span_exporter import ( + InMemorySpanExporter, + ) + + from everos.config.settings import ObservabilitySettings + from everos.core.observability.tracing import ( + force_flush, + init_tracing, + shutdown_tracing, + ) + + async def exercising_driver( + query: str, *, base_retrieve: Any, rerank_fn: Any, **_: Any + ) -> tuple[list[Candidate], AgenticDecision]: + cands = await base_retrieve(query, 10) # -> cluster_scoped -> hybrid_full + # Real driver reranks the recalled hits; feed a non-empty list so the + # cross-encoder rerank actually runs (empty input short-circuits). + await rerank_fn(query, cands or [_mc_candidate("mc_r", "ep_r")]) + return [], AgenticDecision(is_multi_round=False) + + async def fake_embed(q: str) -> list[float]: + return [0.1, 0.2, 0.3, 0.4] + + exporter = InMemorySpanExporter() + shutdown_tracing() + init_tracing( + ObservabilitySettings(enabled=True, endpoint="http://collector.invalid"), + span_processor=SimpleSpanProcessor(exporter), + ) + try: + with ( + patch("everos.memory.search.agentic.aagentic_retrieve", exercising_driver), + patch( + "everos.memory.search.agentic.cluster_repo.list_for_owner", + AsyncMock(return_value=clusters), + ), + ): + await search_episodes_agentic( + "What did Alice eat?", + owner_id="alice", + where="owner_id = 'alice'", + app_id="test_app", + project_id="test_proj", + episode_recaller=ep_recaller, + atomic_fact_recaller=fact_recaller, + embed_query_fn=fake_embed, + reranker=_StubReranker(), + llm=FakeLLMClient(responses=[]), + top_k=10, + ) + force_flush() + finally: + shutdown_tracing() + + names = {s.name for s in exporter.get_finished_spans()} + assert "everos.search.recall" in names + assert "everos.search.rank" in names diff --git a/tests/unit/test_memory/test_search/test_manager.py b/tests/unit/test_memory/test_search/test_manager.py index ff8817b..b396d46 100644 --- a/tests/unit/test_memory/test_search/test_manager.py +++ b/tests/unit/test_memory/test_search/test_manager.py @@ -283,6 +283,20 @@ async def test_user_keyword_returns_episodes_only() -> None: assert resp.data.profiles == [] +async def test_search_uses_propagated_request_id_when_bound() -> None: + """When a request id is bound upstream (middleware), ``search`` reuses it + instead of minting a fresh one, so the response id matches the trace.""" + from everos.core.context import reset_request_id, set_request_id + + mgr = _build_manager(episode_sparse=[_episode_row("ep_1")]) + token = set_request_id("deadbeef" * 4) + try: + resp = await mgr.search(_user_req()) + assert resp.request_id == "deadbeef" * 4 + finally: + reset_request_id(token) + + async def test_user_keyword_leaves_atomic_facts_empty() -> None: """KEYWORD never back-fills facts — only HYBRID produces relevance-scored facts. @@ -569,7 +583,9 @@ async def test_agent_hybrid_with_llm_rerank_does_not_need_reranker() -> None: class _StubReranker: """Minimal reranker stub — returns trivial scores.""" - async def rerank(self, query: str, documents: Sequence[str]) -> list[Any]: + async def rerank( + self, query: str, documents: Sequence[str], **kwargs: Any + ) -> list[Any]: from everos.component.rerank.protocol import RerankResult return [RerankResult(index=i, score=1.0) for i in range(len(documents))] @@ -936,3 +952,186 @@ async def test_agent_hybrid_llm_rerank_merges_bridged_skills_into_dense_pool( # The bridged skill inherits the matched case's score (0.85 from c1). by_id = {c.id: c for c in seen_skill_dense["dense"]} assert by_id["s_bridged"].score == pytest.approx(0.85) + + +async def test_search_emits_memory_search_span() -> None: + """search() opens an everos.memory.search retriever span carrying the + langfuse.* attribute contract (observation type / user id / metadata).""" + from opentelemetry.sdk.trace.export import SimpleSpanProcessor + from opentelemetry.sdk.trace.export.in_memory_span_exporter import ( + InMemorySpanExporter, + ) + + from everos.config.settings import ObservabilitySettings + from everos.core.observability.tracing import ( + force_flush, + init_tracing, + shutdown_tracing, + ) + + exporter = InMemorySpanExporter() + init_tracing( + ObservabilitySettings(enabled=True, endpoint="http://collector.invalid"), + span_processor=SimpleSpanProcessor(exporter), + ) + try: + mgr = _build_manager(episode_sparse=[_episode_row("ep_1")]) + await mgr.search(_user_req()) + force_flush() + spans = {s.name: s for s in exporter.get_finished_spans()} + assert "everos.memory.search" in spans + attrs = spans["everos.memory.search"].attributes + assert attrs["langfuse.observation.type"] == "retriever" + assert attrs["langfuse.user.id"] == "alice" + assert attrs["langfuse.trace.metadata.owner_type"] == "user" + assert list(attrs["langfuse.trace.tags"]) == ["everos", "memory"] + finally: + shutdown_tracing() + + +# ── Search sub-span decomposition (recall + rank phases) ──────────────── + + +@pytest.fixture +def _search_spans(): # type: ignore[no-untyped-def] + from opentelemetry.sdk.trace.export import SimpleSpanProcessor + from opentelemetry.sdk.trace.export.in_memory_span_exporter import ( + InMemorySpanExporter, + ) + + from everos.config.settings import ObservabilitySettings + from everos.core.observability.tracing import init_tracing, shutdown_tracing + + exporter = InMemorySpanExporter() + shutdown_tracing() + init_tracing( + ObservabilitySettings(enabled=True, endpoint="http://collector.invalid"), + span_processor=SimpleSpanProcessor(exporter), + ) + yield exporter + shutdown_tracing() + + +def _span_index(exporter: Any) -> dict[str, Any]: + from everos.core.observability.tracing import force_flush + + force_flush() + return {s.name: s for s in exporter.get_finished_spans()} + + +async def test_keyword_user_emits_recall_no_rank(_search_spans: Any) -> None: + mgr = _build_manager(episode_sparse=[_episode_row("ep_1")]) + await mgr.search(_user_req(method=SearchMethod.KEYWORD)) + spans = _span_index(_search_spans) + assert "everos.search.recall" in spans + assert "everos.search.rank" not in spans + # recall nests under the search span (one trace). + tid = spans["everos.memory.search"].context.trace_id + assert spans["everos.search.recall"].context.trace_id == tid + assert spans["everos.search.recall"].parent is not None + + +async def test_hybrid_user_emits_recall_and_rank(_search_spans: Any) -> None: + mgr = _build_manager( + episode_sparse=[_episode_row("ep_1")], embedding=_StubEmbedding() + ) + await mgr.search(_user_req(method=SearchMethod.HYBRID)) + spans = _span_index(_search_spans) + assert "everos.search.recall" in spans + assert "everos.search.rank" in spans + + +async def test_keyword_agent_emits_recall_no_rank(_search_spans: Any) -> None: + mgr = _build_manager(case_sparse=[_case_row("c1")], skill_sparse=[_skill_row("s1")]) + await mgr.search(_agent_req(method=SearchMethod.KEYWORD)) + spans = _span_index(_search_spans) + assert "everos.search.recall" in spans + assert "everos.search.rank" not in spans + + +async def test_hybrid_agent_emits_recall_and_rank(_search_spans: Any) -> None: + mgr = _build_manager( + case_sparse=[_case_row("c1")], + skill_sparse=[_skill_row("s1")], + embedding=_StubEmbedding(), + reranker=_StubReranker(), + ) + await mgr.search(_agent_req(method=SearchMethod.HYBRID)) + spans = _span_index(_search_spans) + assert "everos.search.recall" in spans + assert "everos.search.rank" in spans + + +async def test_search_emits_top_score_and_hit_on_span(_search_spans: Any) -> None: + """search sets everos.search.top_score (max item score) + everos.search.hit + (>= recall_hit_threshold, default 0.6) on the retriever span — always, no + Langfuse needed.""" + mgr = _build_manager(episode_sparse=[_episode_row("ep_1", score=0.75)]) + await mgr.search(_user_req(method=SearchMethod.KEYWORD)) + spans = _span_index(_search_spans) + attrs = spans["everos.memory.search"].attributes + assert attrs["everos.search.top_score"] == pytest.approx(0.75) + assert attrs["everos.search.hit"] is True + + +async def test_search_hit_false_when_below_threshold(_search_spans: Any) -> None: + mgr = _build_manager(episode_sparse=[_episode_row("ep_1", score=0.3)]) + await mgr.search(_user_req(method=SearchMethod.KEYWORD)) + spans = _span_index(_search_spans) + attrs = spans["everos.memory.search"].attributes + assert attrs["everos.search.top_score"] == pytest.approx(0.3) + assert attrs["everos.search.hit"] is False + + +async def test_search_top_score_zero_when_no_results(_search_spans: Any) -> None: + mgr = _build_manager() # no candidates + await mgr.search(_user_req(method=SearchMethod.KEYWORD)) + spans = _span_index(_search_spans) + attrs = spans["everos.memory.search"].attributes + assert attrs["everos.search.top_score"] == pytest.approx(0.0) + assert attrs["everos.search.hit"] is False + + +async def test_search_enqueues_recall_scores( + _search_spans: Any, monkeypatch: pytest.MonkeyPatch +) -> None: + """When tracing is active, search hands recall_top_score/hit to the score + sink with the retriever span's trace_id (032x) + observation_id (016x).""" + import everos.memory.search.manager as mgr_mod + + captured: dict[str, Any] = {} + + def fake_emit(**kwargs: Any) -> None: + captured.update(kwargs) + + monkeypatch.setattr(mgr_mod, "emit_recall_scores", fake_emit) + mgr = _build_manager(episode_sparse=[_episode_row("ep_1", score=0.75)]) + await mgr.search(_user_req(method=SearchMethod.KEYWORD)) + + assert captured["top_score"] == pytest.approx(0.75) + assert captured["hit"] is True + assert captured["method"] == "keyword" + assert len(captured["trace_id"]) == 32 + assert len(captured["observation_id"]) == 16 + + +async def test_search_captures_query_when_content_on(_search_spans: Any) -> None: + 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)) # query="hi" + finally: + set_capture_content(False) + import json + + attrs = _span_index(_search_spans)["everos.memory.search"].attributes + assert json.loads(attrs["langfuse.observation.input"])["query"] == "hi" + + +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)) + attrs = _span_index(_search_spans)["everos.memory.search"].attributes + assert "langfuse.observation.input" not in attrs diff --git a/tests/unit/test_service/test_memorize_span.py b/tests/unit/test_service/test_memorize_span.py new file mode 100644 index 0000000..1e7fdcc --- /dev/null +++ b/tests/unit/test_service/test_memorize_span.py @@ -0,0 +1,67 @@ +"""``memorize`` opens an everos.memory.add / everos.memory.flush span. + +The inner critical section is mocked out — this asserts only the span +wrapping + name selection (add vs flush) driven by ``is_final``. +""" + +from __future__ import annotations + +import importlib +from collections.abc import AsyncIterator, Iterator +from contextlib import asynccontextmanager +from unittest.mock import AsyncMock + +import pytest +from opentelemetry.sdk.trace.export import SimpleSpanProcessor +from opentelemetry.sdk.trace.export.in_memory_span_exporter import ( + InMemorySpanExporter, +) + +from everos.config import Settings +from everos.core.observability.tracing import ( + force_flush, + init_tracing, + shutdown_tracing, +) + +mm = importlib.import_module("everos.service.memorize") + + +@pytest.fixture(autouse=True) +def _patch(monkeypatch: pytest.MonkeyPatch) -> Iterator[InMemorySpanExporter]: + monkeypatch.setattr(mm, "load_settings", lambda: Settings()) + monkeypatch.setattr( + mm, + "_memorize_locked", + AsyncMock(return_value=mm.MemorizeResult(message_count=0, status="extracted")), + ) + + @asynccontextmanager + async def _fake_lock(session_id: str) -> AsyncIterator[None]: + yield + + monkeypatch.setattr(mm, "get_session_lock", _fake_lock) + + exporter = InMemorySpanExporter() + shutdown_tracing() + init_tracing( + Settings().observability.model_copy(update={"enabled": True}), + span_processor=SimpleSpanProcessor(exporter), + ) + yield exporter + shutdown_tracing() + + +async def test_add_emits_memory_add_span(_patch: InMemorySpanExporter) -> None: + await mm.memorize({"session_id": "s1", "messages": []}, is_final=False) + force_flush() + spans = {s.name: s for s in _patch.get_finished_spans()} + assert "everos.memory.add" in spans + assert spans["everos.memory.add"].attributes["langfuse.session.id"] == "s1" + + +async def test_flush_emits_memory_flush_span(_patch: InMemorySpanExporter) -> None: + await mm.memorize({"session_id": "s2", "messages": []}, is_final=True) + force_flush() + names = {s.name for s in _patch.get_finished_spans()} + assert "everos.memory.flush" in names