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.<strategy> 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) <noreply@anthropic.com>
This commit is contained in:
zhanghui 2026-07-23 20:47:39 +08:00
parent bd26f5a80f
commit 357f619c64
18 changed files with 1153 additions and 185 deletions

View File

@ -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:

View File

@ -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()

View File

@ -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,

View File

@ -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",

View File

@ -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(

View File

@ -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)

View File

@ -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):

View File

@ -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

View File

@ -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(

View File

@ -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

View File

@ -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

View File

@ -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.<name> 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.<name> 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

View File

@ -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="<prompt>")
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="<prompt>")
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"
)

View File

@ -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:

View File

@ -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"
)

View File

@ -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

View File

@ -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

View File

@ -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