diff --git a/.claude/rules/logging-observability.md b/.claude/rules/logging-observability.md index 46faf51..d7055d2 100644 --- a/.claude/rules/logging-observability.md +++ b/.claude/rules/logging-observability.md @@ -21,3 +21,11 @@ paths: - **Metrics** go through `core.observability.metrics` (Prometheus); don't invent ad-hoc counters. Histograms/counters/gauges have registry helpers. - Don't log secrets, API keys, or full memory content at `info`/above. +- **Tracing** (optional, `[otel]` extra, **off by default**): open spans with + `memory_span(...)` from `core.observability.tracing` — it stamps the Langfuse + `langfuse.*` attributes and is a no-op until `[observability] enabled`, so call + sites never branch on config. LLM / embedding token usage rides + `set_generation_usage` onto the active span (Langfuse computes cost). + Request/response content is emitted only when `capture_content` is on + (redaction hook + truncation). `request_id` is kept independent of the OTel + `trace_id`; an upstream `traceparent` header is continued when present. diff --git a/pyproject.toml b/pyproject.toml index 216b247..b9693c9 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -70,6 +70,12 @@ dependencies = [ [project.optional-dependencies] multimodal = ["everalgo-parser[svg]>=0.2.1"] # [svg] bundles cairosvg → SVG works by default +# Native OpenTelemetry tracing export. Optional — EverOS never imports these +# unless [observability] is enabled. Install with: pip install everos[otel] +otel = [ + "opentelemetry-sdk>=1.27.0", + "opentelemetry-exporter-otlp-proto-http>=1.27.0", +] [project.urls] Homepage = "https://evermind.ai" @@ -255,4 +261,8 @@ dev = [ "pre-commit>=4.0.0", "ipdb>=0.13.13", "pyinstrument>=5.0.0", + # Tracing tests must actually run (no skip-when-absent), so the optional + # [otel] stack is always present in the dev / CI environment. + "opentelemetry-sdk>=1.27.0", + "opentelemetry-exporter-otlp-proto-http>=1.27.0", ] diff --git a/src/everos/component/embedding/openai_provider.py b/src/everos/component/embedding/openai_provider.py index f756127..08d131d 100644 --- a/src/everos/component/embedding/openai_provider.py +++ b/src/everos/component/embedding/openai_provider.py @@ -23,6 +23,8 @@ from collections.abc import Sequence import openai +from everos.core.observability.tracing import set_generation_usage + from .protocol import EmbeddingServiceError @@ -94,5 +96,12 @@ class OpenAIEmbeddingProvider: ) except openai.OpenAIError as exc: raise EmbeddingServiceError(str(exc)) from exc + # Surface token usage onto the active span (e.g. everos.search.embed_query). + # No-op when tracing is off; embeddings report only input (prompt) tokens. + usage = getattr(response, "usage", None) + set_generation_usage( + model=self._model, + input_tokens=usage.prompt_tokens if usage else None, + ) # OpenAI returns ``data`` indexed by request order; truncate to ``dim``. return [list(item.embedding[: self.dim]) for item in response.data] diff --git a/src/everos/component/llm/_usage_client.py b/src/everos/component/llm/_usage_client.py new file mode 100644 index 0000000..a4bb7b7 --- /dev/null +++ b/src/everos/component/llm/_usage_client.py @@ -0,0 +1,62 @@ +"""Token-usage-recording LLM client wrapper. + +Wraps any :class:`everalgo.llm.LLMClient` and, after each ``chat`` call, +writes ``response.usage`` (+ model) onto the current OpenTelemetry span via +``set_generation_usage``. The wrapped client's behaviour is otherwise +untouched — the response is returned verbatim and all other attributes +delegate through. + +This is how token counts reach the ``everos.extract`` generation span: +everalgo's extractors own the ``chat`` call and discard the ``ChatResponse`` +(keeping only ``.content``), so usage is captured here at the client +boundary — no everalgo change required. When no span is active (e.g. the +search path, or tracing disabled), ``set_generation_usage`` no-ops. +""" + +from __future__ import annotations + +from typing import TYPE_CHECKING, Any + +from everos.core.observability.tracing import set_generation_usage + +if TYPE_CHECKING: + from everalgo.llm import ChatMessage, ChatResponse + from everalgo.llm.protocols import LLMClient + from pydantic import BaseModel + + +class UsageRecordingClient: + """LLM client proxy that records token usage onto the active span.""" + + def __init__(self, inner: LLMClient) -> None: + self._inner = inner + + async def chat( + self, + messages: list[ChatMessage], + *, + model: str | None = None, + temperature: float | None = None, + max_tokens: int | None = None, + response_format: type[BaseModel] | None = None, + **extra: Any, + ) -> ChatResponse: + response = await self._inner.chat( + messages, + model=model, + temperature=temperature, + max_tokens=max_tokens, + response_format=response_format, + **extra, + ) + usage = response.usage + set_generation_usage( + model=response.model, + input_tokens=usage.prompt_tokens if usage else None, + output_tokens=usage.completion_tokens if usage else None, + ) + return response + + def __getattr__(self, name: str) -> Any: + # Delegate everything else to the wrapped client. + return getattr(self._inner, name) diff --git a/src/everos/component/llm/client.py b/src/everos/component/llm/client.py index 846dcf1..5b0de8c 100644 --- a/src/everos/component/llm/client.py +++ b/src/everos/component/llm/client.py @@ -16,6 +16,8 @@ from everalgo.llm.protocols import LLMClient from everos.config import load_settings from everos.core.observability.logging import get_logger +from ._usage_client import UsageRecordingClient + logger = get_logger(__name__) @@ -38,7 +40,8 @@ def get_llm_client() -> LLMClient: if _llm_client is not None: return _llm_client - llm_cfg = load_settings().llm + settings = load_settings() + llm_cfg = settings.llm api_key = ( llm_cfg.api_key.get_secret_value() if llm_cfg.api_key is not None else None ) @@ -46,13 +49,18 @@ def get_llm_client() -> LLMClient: raise LLMNotConfiguredError( "LLM is required; set EVEROS_LLM__API_KEY + EVEROS_LLM__BASE_URL" ) - _llm_client = build_client( + client: LLMClient = build_client( LLMConfig( model=llm_cfg.model, api_key=api_key, base_url=llm_cfg.base_url, ) ) + # Wrap for OTel token capture only when tracing is on — keeps the + # disabled path (the default) allocation- and overhead-free. + if settings.observability.enabled: + client = UsageRecordingClient(client) + _llm_client = client logger.info("llm_client_built", model=llm_cfg.model) return _llm_client diff --git a/src/everos/config/default.toml b/src/everos/config/default.toml index 68fab81..539e68e 100644 --- a/src/everos/config/default.toml +++ b/src/everos/config/default.toml @@ -145,3 +145,21 @@ session_lock_timeout_seconds = 360.0 threshold = 0.65 time_window_days = 7.0 + +[observability] +# OpenTelemetry tracing export. Off by default; pure OTLP/HTTP, vendor-neutral +# (Langfuse, an OTel Collector, or any OTLP backend). EverOS ships no vendor SDK. +# Override via EVEROS_OBSERVABILITY__ENABLED, EVEROS_OBSERVABILITY__ENDPOINT, etc. +enabled = false +exporter = "otlp_http" # "otlp_http" | "none" +endpoint = "" # e.g. https://us.cloud.langfuse.com/api/public/otel/v1/traces +service_name = "everos" +sample_rate = 1.0 # 0.0 to 1.0 +# Privacy: false (default) = metadata only; true also emits query / extracted +# memory / .md paths as span input/output (redacted + truncated). +capture_content = false +# Recall-quality scores pushed to Langfuse (Langfuse-specific REST, off the +# OTLP stream). Only fires when langfuse_public_key/secret_key/host are set +# (via everos.toml or EVEROS_OBSERVABILITY__LANGFUSE_* — secrets, not shipped here). +emit_recall_scores = true +recall_hit_threshold = 0.6 # only meaningful for calibrated methods diff --git a/src/everos/config/settings.py b/src/everos/config/settings.py index 5ae2c8d..354d6c5 100644 --- a/src/everos/config/settings.py +++ b/src/everos/config/settings.py @@ -350,6 +350,53 @@ class KnowledgeSettings(BaseModel): search: KnowledgeSearchSettings = KnowledgeSearchSettings() +class ObservabilitySettings(BaseModel): + """``[observability]`` — OpenTelemetry tracing export. + + Off by default. When ``enabled`` is true a ``TracerProvider`` is built + once at startup and standard OTLP/HTTP spans are exported to + ``endpoint``. The signal is pure OpenTelemetry — vendor-neutral — so it + works with any OTLP backend (Langfuse, an OTel Collector, ...); EverOS + does not depend on any vendor SDK. + + ``langfuse_*`` are convenience credentials for pushing recall-quality + *scores* to Langfuse (a Langfuse-specific REST call, independent of the + OTLP span stream). Leave unset for a pure vendor-neutral OTLP export. + + Env binding: + EVEROS_OBSERVABILITY__ENABLED + EVEROS_OBSERVABILITY__EXPORTER + EVEROS_OBSERVABILITY__ENDPOINT + EVEROS_OBSERVABILITY__SERVICE_NAME + EVEROS_OBSERVABILITY__SAMPLE_RATE + EVEROS_OBSERVABILITY__LANGFUSE_PUBLIC_KEY / __LANGFUSE_SECRET_KEY + EVEROS_OBSERVABILITY__LANGFUSE_HOST + EVEROS_OBSERVABILITY__EMIT_RECALL_SCORES + EVEROS_OBSERVABILITY__RECALL_HIT_THRESHOLD + """ + + enabled: bool = False + exporter: Literal["otlp_http", "none"] = "otlp_http" + endpoint: str = "" + headers: dict[str, str] = Field(default_factory=dict) + service_name: str = "everos" + sample_rate: float = Field(default=1.0, ge=0.0, le=1.0) + # Privacy: when False (default) spans carry metadata only — no query text, + # extracted memory, or .md paths. Set True to also emit request/response + # content as span input/output (redacted + truncated). + capture_content: bool = False + + # Langfuse scores (recall-quality feedback) — optional, Langfuse-specific. + langfuse_public_key: str | None = None + langfuse_secret_key: SecretStr | None = None + langfuse_host: str | None = None + emit_recall_scores: bool = True + # ``hit`` threshold: only meaningful for calibrated-score methods + # (HYBRID LR / rerank / agentic). Not bounded to [0, 1] because raw + # BM25 scores are unbounded; tune per method on the eval side. + recall_hit_threshold: float = 0.6 + + class Settings(BaseSettings): """Top-level application settings.""" @@ -365,6 +412,7 @@ class Settings(BaseSettings): clustering: ClusteringSettings = ClusteringSettings() multimodal: MultimodalSettings = MultimodalSettings() knowledge: KnowledgeSettings = KnowledgeSettings() + observability: ObservabilitySettings = ObservabilitySettings() model_config = SettingsConfigDict( env_prefix="EVEROS_", diff --git a/src/everos/core/context/__init__.py b/src/everos/core/context/__init__.py index e69de29..3070c2d 100644 --- a/src/everos/core/context/__init__.py +++ b/src/everos/core/context/__init__.py @@ -0,0 +1,22 @@ +"""core.context — request-scoped context propagation (contextvars). + +External usage:: + + from everos.core.context import ( + get_request_id, + set_request_id, + reset_request_id, + ) +""" + +from .request import get_request_id as get_request_id +from .request import reset_request_id as reset_request_id +from .request import resolve_request_id as resolve_request_id +from .request import set_request_id as set_request_id + +__all__ = [ + "get_request_id", + "reset_request_id", + "resolve_request_id", + "set_request_id", +] diff --git a/src/everos/core/context/request.py b/src/everos/core/context/request.py new file mode 100644 index 0000000..f598536 --- /dev/null +++ b/src/everos/core/context/request.py @@ -0,0 +1,39 @@ +"""Request-scoped context propagation via ``contextvars``. + +The request id is stored in a module-level ``ContextVar`` so it survives +``await`` boundaries and is readable anywhere in the call chain (service, +infra, log processors) without being threaded through call signatures. +""" + +from __future__ import annotations + +from contextvars import ContextVar, Token + +from everos.core.observability.tracing import gen_request_id + +_request_id: ContextVar[str | None] = ContextVar("everos_request_id", default=None) + + +def get_request_id() -> str | None: + """Return the request id bound to the current context, or ``None``.""" + return _request_id.get() + + +def set_request_id(value: str | None) -> Token[str | None]: + """Bind ``value`` as the current request id; return a reset token.""" + return _request_id.set(value) + + +def reset_request_id(token: Token[str | None]) -> None: + """Restore the request id to what it was before the matching ``set``.""" + _request_id.reset(token) + + +def resolve_request_id() -> str: + """Return the propagated request id, or mint a fresh W3C-compatible one. + + Call sites that need an id (search / get managers) use this so an id + injected upstream by ``RequestIdMiddleware`` flows through to the + response, while direct / CLI callers still get a freshly minted id. + """ + return get_request_id() or gen_request_id() diff --git a/src/everos/core/lifespan/__init__.py b/src/everos/core/lifespan/__init__.py index dd75abb..b9fe386 100644 --- a/src/everos/core/lifespan/__init__.py +++ b/src/everos/core/lifespan/__init__.py @@ -19,9 +19,11 @@ External usage: from .base import LifespanProvider as LifespanProvider from .factory import build_lifespan as build_lifespan from .metrics_lifespan import MetricsLifespanProvider as MetricsLifespanProvider +from .tracing_lifespan import TracingLifespanProvider as TracingLifespanProvider __all__ = [ "LifespanProvider", "MetricsLifespanProvider", + "TracingLifespanProvider", "build_lifespan", ] diff --git a/src/everos/core/lifespan/tracing_lifespan.py b/src/everos/core/lifespan/tracing_lifespan.py new file mode 100644 index 0000000..29fc4f7 --- /dev/null +++ b/src/everos/core/lifespan/tracing_lifespan.py @@ -0,0 +1,50 @@ +"""Tracing lifespan provider. + +Builds the OpenTelemetry ``TracerProvider`` at startup (from the +``[observability]`` settings) and flushes + tears it down at shutdown. +Chassis-level (backend-agnostic), so it lives here alongside +``MetricsLifespanProvider`` rather than under the API entrypoint. + +Registered with a low ``order`` so the tracer is live before other +providers start and can themselves be traced. +""" + +from __future__ import annotations + +from fastapi import FastAPI + +from everos.config import load_settings +from everos.core.observability.logging import get_logger +from everos.core.observability.tracing import ( + init_score_sink, + init_tracing, + shutdown_score_sink, + shutdown_tracing, +) + +from .base import LifespanProvider + +logger = get_logger(__name__) + + +class TracingLifespanProvider(LifespanProvider): + """Manages the OTel tracer provider + recall-score sink over the app life.""" + + def __init__(self, order: int = 1) -> None: + super().__init__(name="tracing", order=order) + + async def startup(self, app: FastAPI) -> bool: + """Install the tracer provider + recall-score sink when configured. + + Returns True if tracing was enabled and a provider installed. + """ + settings = load_settings().observability + enabled = init_tracing(settings) + scores = init_score_sink(settings) + logger.info("tracing_lifespan_startup", enabled=enabled, scores=scores) + return enabled + + async def shutdown(self, app: FastAPI) -> None: + await shutdown_score_sink() + shutdown_tracing() + logger.info("tracing_lifespan_shutdown") diff --git a/src/everos/core/middleware/__init__.py b/src/everos/core/middleware/__init__.py index c6949c3..d298ad9 100644 --- a/src/everos/core/middleware/__init__.py +++ b/src/everos/core/middleware/__init__.py @@ -18,6 +18,7 @@ from .cors import DEFAULT_CORS_ALLOW_METHODS as DEFAULT_CORS_ALLOW_METHODS from .cors import DEFAULT_CORS_ORIGINS as DEFAULT_CORS_ORIGINS from .profile import ProfileMiddleware as ProfileMiddleware from .prometheus import PrometheusMiddleware as PrometheusMiddleware +from .request_id import RequestIdMiddleware as RequestIdMiddleware __all__ = [ "DEFAULT_CORS_ALLOW_CREDENTIALS", @@ -26,4 +27,5 @@ __all__ = [ "DEFAULT_CORS_ORIGINS", "ProfileMiddleware", "PrometheusMiddleware", + "RequestIdMiddleware", ] diff --git a/src/everos/core/middleware/request_id.py b/src/everos/core/middleware/request_id.py new file mode 100644 index 0000000..51a7c5f --- /dev/null +++ b/src/everos/core/middleware/request_id.py @@ -0,0 +1,56 @@ +"""Request-context middleware. + +Establishes per-request context at HTTP entry: + +* Mints a W3C-compatible ``request_id`` and binds it for the request's + lifetime — via ``request.state``, the ``core.context`` contextvar + (readable by service / infra), and structlog contextvars (so every log + line carries it). Echoed on the ``X-Request-Id`` response header. +* Continues an upstream **distributed trace**: if the request carries a + W3C ``traceparent`` header, our first span nests under that trace instead + of rooting a new one. Absent → we root our own trace (the common case). + +This is the single place a request id enters the system; ``extract_request_id`` +(``entrypoints/api/utils.py``) reads what this middleware sets. +""" + +from __future__ import annotations + +from collections.abc import Awaitable, Callable + +import structlog +from starlette.middleware.base import BaseHTTPMiddleware +from starlette.requests import Request +from starlette.responses import Response + +from everos.core.context import reset_request_id, set_request_id +from everos.core.observability.tracing import gen_request_id, use_traceparent + +_HEADER = "X-Request-Id" + + +class RequestIdMiddleware(BaseHTTPMiddleware): + """Assigns a request id + continues any upstream trace for the request.""" + + async def dispatch( + self, + request: Request, + call_next: Callable[[Request], Awaitable[Response]], + ) -> Response: + request_id = gen_request_id() + request.state.request_id = request_id + # Bind before ``call_next`` so the downstream task (endpoint) inherits + # the value; reset in ``finally`` so it never leaks to the next request. + token = set_request_id(request_id) + structlog.contextvars.bind_contextvars(request_id=request_id) + try: + # Continue the upstream trace when a traceparent header is present + # (no-op otherwise). Attached before call_next so downstream spans + # inherit it across the middleware task boundary. + with use_traceparent(request.headers.get("traceparent")): + response = await call_next(request) + response.headers[_HEADER] = request_id + return response + finally: + structlog.contextvars.unbind_contextvars("request_id") + reset_request_id(token) diff --git a/src/everos/core/observability/tracing/__init__.py b/src/everos/core/observability/tracing/__init__.py index 4f3f263..e437cf8 100644 --- a/src/everos/core/observability/tracing/__init__.py +++ b/src/everos/core/observability/tracing/__init__.py @@ -1,32 +1,56 @@ -"""Tracing utilities — W3C-compatible request id generation. +"""Tracing — W3C id generation + OpenTelemetry tracer lifecycle. External usage:: - from everos.core.observability.tracing import gen_request_id + from everos.core.observability.tracing import ( + gen_request_id, + get_tracer, + init_tracing, + shutdown_tracing, + force_flush, + ) + +``get_tracer`` is safe to call unconditionally: it returns a no-op tracer +until ``init_tracing`` installs a provider (and when the optional ``[otel]`` +extra is not installed), so call sites never need to branch on config. """ from __future__ import annotations -from uuid import uuid4 +from .attributes import capture_input as capture_input +from .attributes import capture_output as capture_output +from .attributes import current_trace_ids as current_trace_ids +from .attributes import current_traceparent as current_traceparent +from .attributes import memory_span as memory_span +from .attributes import set_capture_content as set_capture_content +from .attributes import set_generation_usage as set_generation_usage +from .attributes import set_redactor as set_redactor +from .attributes import use_traceparent as use_traceparent +from .ids import gen_request_id as gen_request_id +from .provider import force_flush as force_flush +from .provider import get_tracer as get_tracer +from .provider import init_tracing as init_tracing +from .provider import shutdown_tracing as shutdown_tracing +from .scores import emit_recall_scores as emit_recall_scores +from .scores import init_score_sink as init_score_sink +from .scores import shutdown_score_sink as shutdown_score_sink - -def gen_request_id() -> str: - """Generate a request id matching the W3C trace-context spec. - - Returns 32 lowercase hex characters (128-bit, no prefix) — the same - format as a W3C ``trace_id`` / OpenTelemetry trace identifier. Routes - and services that mint a fresh request id (when one wasn't injected - by upstream middleware) should call this helper rather than rolling - their own uuid / prefix format, so the id layer stays compatible - with OpenTelemetry exporters and standard APM tooling. - - Example:: - - >>> rid = gen_request_id() - >>> len(rid) - 32 - """ - return uuid4().hex - - -__all__ = ["gen_request_id"] +__all__ = [ + "capture_input", + "capture_output", + "current_trace_ids", + "current_traceparent", + "emit_recall_scores", + "force_flush", + "gen_request_id", + "get_tracer", + "init_score_sink", + "init_tracing", + "memory_span", + "set_capture_content", + "set_generation_usage", + "set_redactor", + "shutdown_score_sink", + "shutdown_tracing", + "use_traceparent", +] diff --git a/src/everos/core/observability/tracing/attributes.py b/src/everos/core/observability/tracing/attributes.py new file mode 100644 index 0000000..46d2f5a --- /dev/null +++ b/src/everos/core/observability/tracing/attributes.py @@ -0,0 +1,195 @@ +"""Span helpers implementing the Langfuse / OpenTelemetry attribute contract. + +``memory_span`` opens a span under the shared ``everos`` tracer and stamps +the ``langfuse.*`` trace/observation attributes (observation type, session / +user ids, trace metadata, tags). ``set_generation_usage`` writes the +``gen_ai.*`` model + token attributes onto the *current* span, so an LLM +client wrapper can record usage without knowing which span is active. + +Request/response content (``langfuse.observation.input/output``) is +privacy-gated: ``capture_input`` / ``capture_output`` only emit it when +``capture_content`` is on, after a redaction hook + truncation. Off by +default, so spans carry metadata only. +""" + +from __future__ import annotations + +import json +from collections.abc import Callable, Iterator, Mapping, Sequence +from contextlib import contextmanager +from typing import Any + +from .provider import get_tracer + +# ── langfuse.* trace/observation attribute keys ────────────────────────── +LF_OBSERVATION_TYPE = "langfuse.observation.type" +LF_SESSION_ID = "langfuse.session.id" +LF_USER_ID = "langfuse.user.id" +LF_TAGS = "langfuse.trace.tags" +LF_METADATA_PREFIX = "langfuse.trace.metadata." + +# ── gen_ai.* generation attribute keys (Langfuse computes cost from these) ─ +GEN_AI_REQUEST_MODEL = "gen_ai.request.model" +GEN_AI_USAGE_INPUT_TOKENS = "gen_ai.usage.input_tokens" +GEN_AI_USAGE_OUTPUT_TOKENS = "gen_ai.usage.output_tokens" + +# ── content capture (privacy-gated) ────────────────────────────────────── +LF_OBSERVATION_INPUT = "langfuse.observation.input" +LF_OBSERVATION_OUTPUT = "langfuse.observation.output" +_MAX_CONTENT_CHARS = 4096 + +DEFAULT_TAGS: tuple[str, ...] = ("everos", "memory") + +# Off by default: no request/response content leaves the process unless +# capture_content is turned on (set once at init_tracing from settings). +_capture_content = False +_redactor: Callable[[str], str] = lambda text: text # noqa: E731 - overridable hook + + +def set_capture_content(enabled: bool) -> None: + """Toggle content capture (called from init_tracing / shutdown).""" + global _capture_content + _capture_content = enabled + + +def set_redactor(redactor: Callable[[str], str] | None) -> None: + """Install a redaction hook applied to captured content; None resets it.""" + global _redactor + _redactor = redactor if redactor is not None else (lambda text: text) + + +def _prepare_content(value: Any) -> str: + """Serialize, redact, then truncate content for a span attribute.""" + text = value if isinstance(value, str) else json.dumps(value, default=str) + text = _redactor(text) + return text[:_MAX_CONTENT_CHARS] + + +def capture_input(span: Any, value: Any) -> None: + """Set ``langfuse.observation.input`` — only when capture_content is on.""" + if _capture_content and value is not None: + span.set_attribute(LF_OBSERVATION_INPUT, _prepare_content(value)) + + +def capture_output(span: Any, value: Any) -> None: + """Set ``langfuse.observation.output`` — only when capture_content is on.""" + if _capture_content and value is not None: + span.set_attribute(LF_OBSERVATION_OUTPUT, _prepare_content(value)) + + +try: + from opentelemetry import trace as _otel_trace + + _OTEL_AVAILABLE = True +except ImportError: # pragma: no cover - only without the [otel] extra + _OTEL_AVAILABLE = False + + +@contextmanager +def memory_span( + name: str, + *, + observation_type: str, + session_id: str | None = None, + user_id: str | None = None, + metadata: Mapping[str, Any] | None = None, + tags: Sequence[str] = DEFAULT_TAGS, +) -> Iterator[Any]: + """Open a span named ``name`` and stamp the langfuse.* attributes. + + Args: + name: Span name (e.g. ``everos.memory.search``). + observation_type: ``langfuse.observation.type`` — span / generation / + embedding / retriever / agent. + session_id / user_id: Grouping ids (dropped when None). + metadata: Flat mapping → ``langfuse.trace.metadata.``; None + values are dropped rather than emitted as the string "None". + tags: ``langfuse.trace.tags`` list. + """ + tracer = get_tracer("everos") + with tracer.start_as_current_span(name) as span: + span.set_attribute(LF_OBSERVATION_TYPE, observation_type) + if session_id: + span.set_attribute(LF_SESSION_ID, session_id) + if user_id: + span.set_attribute(LF_USER_ID, user_id) + if tags: + span.set_attribute(LF_TAGS, list(tags)) + for key, value in (metadata or {}).items(): + if value is not None: + span.set_attribute(f"{LF_METADATA_PREFIX}{key}", value) + yield span + + +def current_traceparent() -> str | None: + """W3C ``traceparent`` for the current span, or None when there is none. + + Captured where a request's span is active (e.g. OME enqueue) and carried + across the async / process boundary so a background strategy span can + re-attach to the originating trace. + """ + if not _OTEL_AVAILABLE: + return None + from opentelemetry.propagate import inject + + carrier: dict[str, str] = {} + inject(carrier) + return carrier.get("traceparent") + + +@contextmanager +def use_traceparent(traceparent: str | None) -> Iterator[None]: + """Attach ``traceparent`` as the current context for the block. + + Spans opened inside become children of that (remote) trace. No-op when + the traceparent is absent or OTel is not installed — the span then roots + its own trace. + """ + if not _OTEL_AVAILABLE or not traceparent: + yield + return + from opentelemetry import context as otel_context + from opentelemetry.propagate import extract + + token = otel_context.attach(extract({"traceparent": traceparent})) + try: + yield + finally: + otel_context.detach(token) + + +def current_trace_ids() -> tuple[str, str] | None: + """Return ``(trace_id_hex_032x, span_id_hex_016x)`` of the current span. + + Returns None when OTel is absent or there is no valid recording span — + the exact hex mapping Langfuse's OTLP ingestion uses for traceId / + observationId, so recall scores attach to the right observation. + """ + if not _OTEL_AVAILABLE: + return None + ctx = _otel_trace.get_current_span().get_span_context() + if not ctx.is_valid: + return None + return format(ctx.trace_id, "032x"), format(ctx.span_id, "016x") + + +def set_generation_usage( + *, + model: str | None = None, + input_tokens: int | None = None, + output_tokens: int | None = None, +) -> None: + """Record ``gen_ai.*`` model + token attributes on the current span. + + No-op when OTel is absent or there is no active recording span, so LLM + client wrappers can call it unconditionally. + """ + if not _OTEL_AVAILABLE: + return + span = _otel_trace.get_current_span() + if model is not None: + span.set_attribute(GEN_AI_REQUEST_MODEL, model) + if input_tokens is not None: + span.set_attribute(GEN_AI_USAGE_INPUT_TOKENS, input_tokens) + if output_tokens is not None: + span.set_attribute(GEN_AI_USAGE_OUTPUT_TOKENS, output_tokens) diff --git a/src/everos/core/observability/tracing/ids.py b/src/everos/core/observability/tracing/ids.py new file mode 100644 index 0000000..73f8ab8 --- /dev/null +++ b/src/everos/core/observability/tracing/ids.py @@ -0,0 +1,24 @@ +"""W3C-compatible request/trace id generation.""" + +from __future__ import annotations + +from uuid import uuid4 + + +def gen_request_id() -> str: + """Generate a request id matching the W3C trace-context spec. + + Returns 32 lowercase hex characters (128-bit, no prefix) — the same + format as a W3C ``trace_id`` / OpenTelemetry trace identifier. Routes + and services that mint a fresh request id (when one wasn't injected + by upstream middleware) should call this helper rather than rolling + their own uuid / prefix format, so the id layer stays compatible + with OpenTelemetry exporters and standard APM tooling. + + Example:: + + >>> rid = gen_request_id() + >>> len(rid) + 32 + """ + return uuid4().hex diff --git a/src/everos/core/observability/tracing/provider.py b/src/everos/core/observability/tracing/provider.py new file mode 100644 index 0000000..95630cb --- /dev/null +++ b/src/everos/core/observability/tracing/provider.py @@ -0,0 +1,178 @@ +"""OpenTelemetry tracer provider lifecycle + a no-op-safe tracer facade. + +OpenTelemetry is an optional dependency (the ``[otel]`` extra). This module +never fails to import when it is absent: the SDK imports are guarded, and +``get_tracer`` returns a no-op tracer until ``init_tracing`` installs a real +provider. + +The provider is held here (module-local) rather than on the OTel *global* +so it can be built and torn down repeatedly — in tests and across restarts — +without tripping OTel's "set global provider once" guard. Span context +propagation (parent/child nesting) still works: that rides OTel's context +vars, which are independent of which provider produced the tracer. +""" + +from __future__ import annotations + +import base64 +from collections.abc import Iterator +from contextlib import contextmanager +from typing import TYPE_CHECKING, Any + +from everos.core.observability.logging import get_logger + +logger = get_logger(__name__) + +try: + from opentelemetry.exporter.otlp.proto.http.trace_exporter import ( + OTLPSpanExporter, + ) + from opentelemetry.sdk.resources import Resource + from opentelemetry.sdk.trace import TracerProvider + from opentelemetry.sdk.trace.export import BatchSpanProcessor + from opentelemetry.sdk.trace.sampling import ParentBased, TraceIdRatioBased + + _OTEL_AVAILABLE = True +except ImportError: # pragma: no cover - only without the [otel] extra + _OTEL_AVAILABLE = False + +if TYPE_CHECKING: + from opentelemetry.sdk.trace.export import SpanProcessor + + from everos.config.settings import ObservabilitySettings + +# Our TracerProvider, deliberately kept off the OTel global (see module docstring). +_provider: Any = None + + +class _NoopSpan: + """Span stand-in used when tracing is off or OTel is not installed.""" + + def set_attribute(self, key: str, value: object) -> None: ... + + def set_attributes(self, attributes: dict[str, object]) -> None: ... + + def record_exception(self, exception: BaseException) -> None: ... + + def set_status(self, *args: object, **kwargs: object) -> None: ... + + def end(self) -> None: ... + + +class _NoopTracer: + """Tracer stand-in whose spans do nothing (zero-overhead when off).""" + + @contextmanager + def start_as_current_span(self, name: str, **kwargs: object) -> Iterator[_NoopSpan]: + yield _NoopSpan() + + +_NOOP_TRACER = _NoopTracer() + + +def _resolve_otlp_target( + settings: ObservabilitySettings, +) -> tuple[str, dict[str, str]]: + """Endpoint + headers for the OTLP exporter. + + Convenience: when ``langfuse_*`` creds are set, derive the OTLP traces + endpoint and a Basic-auth header from them — but explicit ``endpoint`` / + ``headers`` always win, so a plain vendor-neutral export is unaffected. + """ + endpoint = settings.endpoint + headers = dict(settings.headers) + pk = settings.langfuse_public_key + sk = settings.langfuse_secret_key + host = settings.langfuse_host + if host and pk and sk: + if not endpoint: + endpoint = host.rstrip("/") + "/api/public/otel/v1/traces" + if "Authorization" not in headers: + token = base64.b64encode(f"{pk}:{sk.get_secret_value()}".encode()).decode() + headers["Authorization"] = f"Basic {token}" + return endpoint, headers + + +def init_tracing( + settings: ObservabilitySettings, + *, + span_processor: SpanProcessor | None = None, +) -> bool: + """Build and install the TracerProvider. + + Args: + settings: Observability config. When ``enabled`` is false or + ``exporter`` is ``"none"``, this is a no-op. + span_processor: Injectable processor (tests pass an in-memory one). + Defaults to a ``BatchSpanProcessor`` over an OTLP/HTTP exporter. + + Returns: + True if a real provider was installed; False when disabled, the + exporter is ``none``, or the ``[otel]`` extra is not installed. + """ + global _provider + if not settings.enabled or settings.exporter == "none": + return False + if not _OTEL_AVAILABLE: + logger.warning("observability_enabled_but_otel_not_installed") + return False + + from everos import __version__ + + resource = Resource.create( + {"service.name": settings.service_name, "service.version": __version__} + ) + provider = TracerProvider( + resource=resource, + sampler=ParentBased(TraceIdRatioBased(settings.sample_rate)), + ) + endpoint, headers = _resolve_otlp_target(settings) + processor = span_processor or BatchSpanProcessor( + OTLPSpanExporter(endpoint=endpoint, headers=headers) + ) + provider.add_span_processor(processor) + _provider = provider + + from .attributes import set_capture_content + + set_capture_content(settings.capture_content) + logger.info( + "tracing_initialized", + service_name=settings.service_name, + capture_content=settings.capture_content, + ) + return True + + +def get_tracer(name: str) -> Any: + """Return a tracer for ``name`` — the real one if initialized, else no-op.""" + if _provider is None: + return _NOOP_TRACER + return _provider.get_tracer(name) + + +def force_flush(timeout_millis: int = 5000) -> None: + """Flush pending spans. No-op (and never raises) when uninitialized.""" + if _provider is None: + return + try: + _provider.force_flush(timeout_millis) + except Exception: # pragma: no cover - telemetry must never break callers + logger.warning("tracing_force_flush_failed", exc_info=True) + + +def shutdown_tracing() -> None: + """Flush + tear down the provider; safe to call when uninitialized.""" + global _provider + if _provider is None: + return + try: + _provider.force_flush() + _provider.shutdown() + except Exception: # pragma: no cover - telemetry must never break callers + logger.warning("tracing_shutdown_failed", exc_info=True) + finally: + _provider = None + from .attributes import set_capture_content + + set_capture_content(False) diff --git a/src/everos/core/observability/tracing/scores.py b/src/everos/core/observability/tracing/scores.py new file mode 100644 index 0000000..ac480eb --- /dev/null +++ b/src/everos/core/observability/tracing/scores.py @@ -0,0 +1,183 @@ +"""Non-blocking recall-score push to Langfuse. + +Recall-quality scores are a Langfuse-specific REST object (POST +``/api/public/scores``), independent of the OTLP span stream. To keep the +search request path free of any network time, scores go through a bounded +queue + a single background worker: + +- ``enqueue`` is O(1) and never blocks or raises — when the queue is full it + drops + counts (back-pressure never reaches the caller). +- the worker drains the queue and POSTs one score at a time (matching the + Langfuse scores API + the reference prototype); every send is wrapped so a + network failure only logs and the worker keeps going. + +Attaches to the originating span via ``traceId`` (OTel trace_id, 032x hex) + +``observationId`` (OTel span_id, 016x hex) — exactly the mapping Langfuse's +OTLP ingestion uses. +""" + +from __future__ import annotations + +import asyncio +import base64 +import contextlib +from collections.abc import Awaitable, Callable +from dataclasses import dataclass +from typing import TYPE_CHECKING + +from everos.core.observability.logging import get_logger + +if TYPE_CHECKING: + from everos.config.settings import ObservabilitySettings + +logger = get_logger(__name__) + +Sender = Callable[[dict], Awaitable[None]] + + +@dataclass(frozen=True) +class ScoreRecord: + """One recall score bound for the Langfuse scores API.""" + + trace_id: str + observation_id: str | None + name: str + value: float + comment: str | None + + +def _to_payload(record: ScoreRecord) -> dict: + payload: dict = { + "traceId": record.trace_id, + "name": record.name, + "value": record.value, + "dataType": "NUMERIC", + } + if record.observation_id: + payload["observationId"] = record.observation_id + if record.comment: + payload["comment"] = record.comment + return payload + + +class RecallScoreSink: + """Bounded queue + background worker that POSTs scores out-of-band.""" + + def __init__( + self, + *, + sender: Sender, + closer: Callable[[], Awaitable[None]] | None = None, + max_queue: int = 1000, + ) -> None: + self._sender = sender + self._closer = closer + self._queue: asyncio.Queue[ScoreRecord] = asyncio.Queue(maxsize=max_queue) + self._task: asyncio.Task[None] | None = None + self.dropped = 0 + + def enqueue(self, record: ScoreRecord) -> None: + """Hand a score to the worker; never blocks or raises.""" + try: + self._queue.put_nowait(record) + except asyncio.QueueFull: + self.dropped += 1 + logger.warning("recall_score_dropped_queue_full", dropped=self.dropped) + + def start(self) -> None: + self._task = asyncio.create_task(self._run()) + + async def _run(self) -> None: + while True: + record = await self._queue.get() + try: + await self._sender(_to_payload(record)) + except Exception: # telemetry must never break; log + continue + logger.warning("recall_score_send_failed", exc_info=True) + finally: + self._queue.task_done() + + async def stop(self, *, drain_timeout: float = 5.0) -> None: + """Drain pending scores (bounded by ``drain_timeout``), then tear down.""" + try: + await asyncio.wait_for(self._queue.join(), timeout=drain_timeout) + except TimeoutError: + logger.warning("recall_score_drain_timeout") + if self._task is not None: + self._task.cancel() + with contextlib.suppress(asyncio.CancelledError): + await self._task + self._task = None + if self._closer is not None: + with contextlib.suppress(Exception): + await self._closer() + + +# ── Module-level lifecycle (mirrors the tracer provider pattern) ───────── +_sink: RecallScoreSink | None = None + + +def init_score_sink(settings: ObservabilitySettings) -> bool: + """Build + start the sink when Langfuse creds + emit_recall_scores are set. + + Returns True if a sink was installed, False otherwise (disabled, scores + off, or missing creds) — in which case ``emit_recall_scores`` is a no-op. + """ + global _sink + if not settings.enabled or not settings.emit_recall_scores: + return False + pk = settings.langfuse_public_key + sk = settings.langfuse_secret_key + host = settings.langfuse_host + if not (pk and sk and host): + return False + + import httpx + + endpoint = host.rstrip("/") + "/api/public/scores" + token = base64.b64encode(f"{pk}:{sk.get_secret_value()}".encode()).decode() + auth = f"Basic {token}" + client = httpx.AsyncClient(timeout=5.0) + + async def sender(payload: dict) -> None: + resp = await client.post( + endpoint, json=payload, headers={"Authorization": auth} + ) + resp.raise_for_status() + + _sink = RecallScoreSink(sender=sender, closer=client.aclose) + _sink.start() + logger.info("recall_score_sink_started", endpoint=endpoint) + return True + + +def emit_recall_scores( + *, + trace_id: str, + observation_id: str | None, + top_score: float, + hit: bool, + method: str, +) -> None: + """Enqueue recall_top_score + recall_hit; no-op when the sink is off.""" + if _sink is None: + return + comment = f"method={method}" + _sink.enqueue( + ScoreRecord( + trace_id, observation_id, "recall_top_score", float(top_score), comment + ) + ) + _sink.enqueue( + ScoreRecord( + trace_id, observation_id, "recall_hit", 1.0 if hit else 0.0, comment + ) + ) + + +async def shutdown_score_sink() -> None: + """Drain + tear down the sink; safe when uninitialized.""" + global _sink + if _sink is not None: + await _sink.stop() + _sink = None diff --git a/src/everos/entrypoints/api/app.py b/src/everos/entrypoints/api/app.py index 1583671..16c2464 100644 --- a/src/everos/entrypoints/api/app.py +++ b/src/everos/entrypoints/api/app.py @@ -14,6 +14,7 @@ from fastapi.middleware.cors import CORSMiddleware from everos.core.lifespan import ( LifespanProvider, MetricsLifespanProvider, + TracingLifespanProvider, build_lifespan, ) from everos.core.middleware import ( @@ -23,6 +24,7 @@ from everos.core.middleware import ( DEFAULT_CORS_ORIGINS, ProfileMiddleware, PrometheusMiddleware, + RequestIdMiddleware, ) from everos.core.observability.logging import get_logger @@ -68,7 +70,8 @@ def create_app( cors_allow_methods: Allowed CORS methods (default: ``["*"]``). cors_allow_headers: Allowed CORS headers (default: ``["*"]``). lifespan_providers: Optional list of LifespanProvider; defaults to - ``[MetricsLifespanProvider(), SqliteLifespanProvider(), + ``[TracingLifespanProvider(), MetricsLifespanProvider(), + LLMLifespanProvider(), SqliteLifespanProvider(), LanceDBLifespanProvider(), CascadeLifespanProvider(), OmeLifespanProvider()]``. @@ -79,6 +82,7 @@ def create_app( if lifespan_providers is None: lifespan_providers = [ + TracingLifespanProvider(), MetricsLifespanProvider(), LLMLifespanProvider(), SqliteLifespanProvider(), @@ -113,6 +117,9 @@ def create_app( ) app.add_middleware(PrometheusMiddleware) app.add_middleware(ProfileMiddleware) + # Outermost: every request gets a request id before any other middleware + # or handler runs, so all logs + the response header carry it. + app.add_middleware(RequestIdMiddleware) # Routes. app.include_router(health.router) 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/get/manager.py b/src/everos/memory/get/manager.py index 1000cbd..472033c 100644 --- a/src/everos/memory/get/manager.py +++ b/src/everos/memory/get/manager.py @@ -21,8 +21,8 @@ import json from typing import TYPE_CHECKING, Any from everos.component.utils.datetime import to_display_tz +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 .dto import ( GetAgentCaseItem, @@ -68,7 +68,7 @@ class GetManager: # ── Public entry ───────────────────────────────────────────────── async def get(self, req: GetRequest) -> GetResponse: - request_id = gen_request_id() + request_id = resolve_request_id() descending = req.sort_order == "desc" where = compile_filters_for_get( req.filters, 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_component/test_embedding/test_usage_span.py b/tests/unit/test_component/test_embedding/test_usage_span.py new file mode 100644 index 0000000..6005607 --- /dev/null +++ b/tests/unit/test_component/test_embedding/test_usage_span.py @@ -0,0 +1,95 @@ +"""``OpenAIEmbeddingProvider`` surfaces token usage onto the active span. + +The OpenAI embeddings response carries ``usage``; the provider records it +via ``set_generation_usage`` (no contract change) so the +``everos.search.embed_query`` embedding span shows model + input tokens. +Captured via an in-memory exporter. +""" + +from __future__ import annotations + +from collections.abc import Iterator +from types import SimpleNamespace + +import pytest +from opentelemetry.sdk.trace.export import SimpleSpanProcessor +from opentelemetry.sdk.trace.export.in_memory_span_exporter import ( + InMemorySpanExporter, +) + +from everos.component.embedding.openai_provider import OpenAIEmbeddingProvider +from everos.config.settings import ObservabilitySettings +from everos.core.observability.tracing import ( + force_flush, + init_tracing, + memory_span, + shutdown_tracing, +) + + +class _FakeEmbeddings: + def __init__(self, response: object) -> None: + self._response = response + + async def create(self, *, model: str, input: list[str]) -> object: + return self._response + + +class _FakeClient: + def __init__(self, response: object) -> None: + self.embeddings = _FakeEmbeddings(response) + + +@pytest.fixture(autouse=True) +def _reset() -> Iterator[None]: + shutdown_tracing() + yield + shutdown_tracing() + + +@pytest.fixture +def captured() -> Iterator[InMemorySpanExporter]: + exporter = InMemorySpanExporter() + init_tracing( + ObservabilitySettings(enabled=True, endpoint="http://collector.invalid"), + span_processor=SimpleSpanProcessor(exporter), + ) + yield exporter + + +def _provider(response: object) -> OpenAIEmbeddingProvider: + provider = OpenAIEmbeddingProvider( + model="emb-m", api_key="k", base_url="http://x", dim=8 + ) + provider._client = _FakeClient(response) # type: ignore[assignment] + return provider + + +async def test_embed_records_usage_to_active_span( + captured: InMemorySpanExporter, +) -> None: + resp = SimpleNamespace( + data=[SimpleNamespace(embedding=[0.1] * 8)], + usage=SimpleNamespace(prompt_tokens=5, total_tokens=5), + ) + provider = _provider(resp) + with memory_span("everos.search.embed_query", observation_type="embedding"): + vec = await provider.embed("hello") + assert len(vec) == 8 + force_flush() + attrs = captured.get_finished_spans()[0].attributes + assert attrs["gen_ai.request.model"] == "emb-m" + assert attrs["gen_ai.usage.input_tokens"] == 5 + + +async def test_embed_without_usage_does_not_raise( + captured: InMemorySpanExporter, +) -> None: + resp = SimpleNamespace(data=[SimpleNamespace(embedding=[0.1] * 8)], usage=None) + provider = _provider(resp) + with memory_span("everos.search.embed_query", observation_type="embedding"): + vec = await provider.embed("hello") + assert len(vec) == 8 + force_flush() + attrs = captured.get_finished_spans()[0].attributes + assert "gen_ai.usage.input_tokens" not in attrs diff --git a/tests/unit/test_component/test_llm/test_client.py b/tests/unit/test_component/test_llm/test_client.py index eb7427e..46098e8 100644 --- a/tests/unit/test_component/test_llm/test_client.py +++ b/tests/unit/test_component/test_llm/test_client.py @@ -8,8 +8,9 @@ import pytest from pydantic import SecretStr from everos.component.llm import LLMNotConfiguredError +from everos.component.llm._usage_client import UsageRecordingClient from everos.config import Settings -from everos.config.settings import LLMSettings +from everos.config.settings import LLMSettings, ObservabilitySettings _client_mod = importlib.import_module("everos.component.llm.client") @@ -62,3 +63,42 @@ def test_returns_singleton_when_configured(monkeypatch: pytest.MonkeyPatch) -> N assert first is sentinel assert first is second + + +def _patch_settings_with_observability( + monkeypatch: pytest.MonkeyPatch, *, enabled: bool +) -> None: + cfg = Settings( + llm=LLMSettings( + model="gpt-4.1-mini", + api_key=SecretStr("sk-test"), + base_url="https://example.test", + ), + observability=ObservabilitySettings(enabled=enabled), + ) + monkeypatch.setattr(_client_mod, "load_settings", lambda: cfg) + + +def test_wraps_client_when_observability_enabled( + monkeypatch: pytest.MonkeyPatch, +) -> None: + _reset_singleton(monkeypatch) + _patch_settings_with_observability(monkeypatch, enabled=True) + sentinel = object() + monkeypatch.setattr(_client_mod, "build_client", lambda cfg: sentinel) + + client = _client_mod.get_llm_client() + + assert isinstance(client, UsageRecordingClient) + assert client._inner is sentinel + + +def test_does_not_wrap_client_when_observability_disabled( + monkeypatch: pytest.MonkeyPatch, +) -> None: + _reset_singleton(monkeypatch) + _patch_settings_with_observability(monkeypatch, enabled=False) + sentinel = object() + monkeypatch.setattr(_client_mod, "build_client", lambda cfg: sentinel) + + assert _client_mod.get_llm_client() is sentinel diff --git a/tests/unit/test_component/test_llm/test_usage_client.py b/tests/unit/test_component/test_llm/test_usage_client.py new file mode 100644 index 0000000..6583bc8 --- /dev/null +++ b/tests/unit/test_component/test_llm/test_usage_client.py @@ -0,0 +1,101 @@ +"""``UsageRecordingClient`` — wraps an LLM client and records token usage +onto the active span, without altering the response or the call. + +The extractor calls happen inside an ``everos.extract`` generation span; the +wrapper's job is to surface ``response.usage`` there (Langfuse computes cost +from model + tokens). Everything else delegates to the wrapped client. +""" + +from __future__ import annotations + +from collections.abc import Iterator +from typing import Any + +import pytest +from everalgo.llm import ChatMessage, ChatResponse, Usage +from opentelemetry.sdk.trace.export import SimpleSpanProcessor +from opentelemetry.sdk.trace.export.in_memory_span_exporter import ( + InMemorySpanExporter, +) + +from everos.component.llm._usage_client import UsageRecordingClient +from everos.config.settings import ObservabilitySettings +from everos.core.observability.tracing import ( + force_flush, + init_tracing, + memory_span, + shutdown_tracing, +) + + +class _FakeInner: + """Minimal stand-in LLM client that returns a preset response.""" + + def __init__(self, response: ChatResponse | None) -> None: + self._response = response + self.calls: list[Any] = [] + self.some_attr = 42 + + async def chat(self, messages: list[ChatMessage], **kwargs: Any) -> ChatResponse: + self.calls.append((messages, kwargs)) + assert self._response is not None + return self._response + + +@pytest.fixture(autouse=True) +def _reset() -> Iterator[None]: + shutdown_tracing() + yield + shutdown_tracing() + + +@pytest.fixture +def captured() -> Iterator[InMemorySpanExporter]: + exporter = InMemorySpanExporter() + init_tracing( + ObservabilitySettings(enabled=True, endpoint="http://collector.invalid"), + span_processor=SimpleSpanProcessor(exporter), + ) + yield exporter + + +async def test_records_usage_to_active_span(captured: InMemorySpanExporter) -> None: + resp = ChatResponse( + content="x", model="gpt-x", usage=Usage(prompt_tokens=7, completion_tokens=3) + ) + client = UsageRecordingClient(_FakeInner(resp)) + with memory_span("everos.extract", observation_type="generation"): + out = await client.chat([ChatMessage(role="user", content="hi")], model="gpt-x") + assert out is resp + force_flush() + attrs = captured.get_finished_spans()[0].attributes + assert attrs["gen_ai.request.model"] == "gpt-x" + assert attrs["gen_ai.usage.input_tokens"] == 7 + assert attrs["gen_ai.usage.output_tokens"] == 3 + + +async def test_passes_through_when_usage_absent( + captured: InMemorySpanExporter, +) -> None: + resp = ChatResponse(content="x", model="gpt-x", usage=None) + client = UsageRecordingClient(_FakeInner(resp)) + with memory_span("everos.extract", observation_type="generation"): + out = await client.chat([ChatMessage(role="user", content="hi")]) + assert out is resp + force_flush() + attrs = captured.get_finished_spans()[0].attributes + assert attrs["gen_ai.request.model"] == "gpt-x" + assert "gen_ai.usage.input_tokens" not in attrs + + +async def test_delegates_unknown_attributes() -> None: + client = UsageRecordingClient(_FakeInner(None)) + assert client.some_attr == 42 + + +async def test_chat_forwards_kwargs_to_inner() -> None: + resp = ChatResponse(content="x", model="m") + inner = _FakeInner(resp) + client = UsageRecordingClient(inner) + await client.chat([ChatMessage(role="user", content="hi")], temperature=0.5) + assert inner.calls[0][1]["temperature"] == 0.5 diff --git a/tests/unit/test_config/test_observability_settings.py b/tests/unit/test_config/test_observability_settings.py new file mode 100644 index 0000000..8939fc6 --- /dev/null +++ b/tests/unit/test_config/test_observability_settings.py @@ -0,0 +1,86 @@ +"""Unit tests for ``ObservabilitySettings`` (OpenTelemetry tracing config).""" + +from __future__ import annotations + +import os +from pathlib import Path + +import pytest +from pydantic import ValidationError + +from everos.config import Settings, load_settings +from everos.config.settings import ObservabilitySettings + + +@pytest.fixture(autouse=True) +def _isolate_env(monkeypatch: pytest.MonkeyPatch, tmp_path: Path) -> None: + for key in list(os.environ): + if key.startswith("EVEROS_"): + monkeypatch.delenv(key, raising=False) + monkeypatch.setenv("EVEROS_ROOT", str(tmp_path)) + monkeypatch.chdir(tmp_path) + load_settings.cache_clear() + + +def test_observability_defaults_are_off_and_neutral() -> None: + obs = load_settings().observability + assert obs.enabled is False + assert obs.exporter == "otlp_http" + assert obs.endpoint == "" + assert obs.headers == {} + assert obs.service_name == "everos" + assert obs.sample_rate == 1.0 + + +def test_env_overrides_observability( + monkeypatch: pytest.MonkeyPatch, tmp_path: Path +) -> None: + root = tmp_path / "r" + root.mkdir() + monkeypatch.setenv("EVEROS_OBSERVABILITY__ENABLED", "true") + monkeypatch.setenv( + "EVEROS_OBSERVABILITY__ENDPOINT", "https://otlp.example/v1/traces" + ) + monkeypatch.setenv("EVEROS_OBSERVABILITY__SERVICE_NAME", "everos-test") + s = Settings(_everos_root=root) + assert s.observability.enabled is True + assert s.observability.endpoint == "https://otlp.example/v1/traces" + assert s.observability.service_name == "everos-test" + + +def test_sample_rate_out_of_range_rejected() -> None: + with pytest.raises(ValidationError): + ObservabilitySettings(sample_rate=1.5) + with pytest.raises(ValidationError): + ObservabilitySettings(sample_rate=-0.1) + + +def test_recall_score_defaults() -> None: + obs = load_settings().observability + assert obs.langfuse_public_key is None + assert obs.langfuse_secret_key is None + assert obs.langfuse_host is None + assert obs.emit_recall_scores is True + assert obs.recall_hit_threshold == 0.6 + + +def test_secret_key_is_not_leaked_in_repr() -> None: + obs = ObservabilitySettings(langfuse_secret_key="sk-lf-supersecret") + # SecretStr masks the value in repr/str. + assert "sk-lf-supersecret" not in repr(obs) + assert obs.langfuse_secret_key is not None + assert obs.langfuse_secret_key.get_secret_value() == "sk-lf-supersecret" + + +def test_env_overrides_recall_score_fields( + monkeypatch: pytest.MonkeyPatch, tmp_path: Path +) -> None: + root = tmp_path / "r" + root.mkdir() + monkeypatch.setenv("EVEROS_OBSERVABILITY__LANGFUSE_HOST", "https://lf.example") + monkeypatch.setenv("EVEROS_OBSERVABILITY__EMIT_RECALL_SCORES", "false") + monkeypatch.setenv("EVEROS_OBSERVABILITY__RECALL_HIT_THRESHOLD", "0.8") + s = Settings(_everos_root=root) + assert s.observability.langfuse_host == "https://lf.example" + assert s.observability.emit_recall_scores is False + assert s.observability.recall_hit_threshold == 0.8 diff --git a/tests/unit/test_core/test_context.py b/tests/unit/test_core/test_context.py new file mode 100644 index 0000000..2005ab1 --- /dev/null +++ b/tests/unit/test_core/test_context.py @@ -0,0 +1,48 @@ +"""``core.context`` — request-scoped contextvar propagation. + +The request id lives in a ``ContextVar`` so it flows across ``await`` +boundaries (HTTP middleware → service → infra → logs) without being +threaded through every call signature. +""" + +from __future__ import annotations + +from everos.core.context import ( + get_request_id, + reset_request_id, + resolve_request_id, + set_request_id, +) + + +def test_get_request_id_defaults_to_none() -> None: + assert get_request_id() is None + + +def test_set_request_id_roundtrip() -> None: + token = set_request_id("abc123") + try: + assert get_request_id() == "abc123" + finally: + reset_request_id(token) + + +def test_reset_request_id_restores_previous() -> None: + token = set_request_id("first") + assert get_request_id() == "first" + reset_request_id(token) + assert get_request_id() is None + + +def test_resolve_returns_bound_id_when_present() -> None: + token = set_request_id("deadbeef" * 4) + try: + assert resolve_request_id() == "deadbeef" * 4 + finally: + reset_request_id(token) + + +def test_resolve_mints_32hex_when_absent() -> None: + rid = resolve_request_id() + assert len(rid) == 32 + assert all(c in "0123456789abcdef" for c in rid) diff --git a/tests/unit/test_core/test_lifespan/test_tracing_lifespan.py b/tests/unit/test_core/test_lifespan/test_tracing_lifespan.py new file mode 100644 index 0000000..f6bd230 --- /dev/null +++ b/tests/unit/test_core/test_lifespan/test_tracing_lifespan.py @@ -0,0 +1,92 @@ +"""``TracingLifespanProvider`` — reads [observability] and manages the +tracer provider over the app lifespan. + +Startup returns whether tracing was enabled; shutdown always flushes and +tears down without raising. Span-capture behavior itself is covered in +``test_observability/test_tracing.py``; here we assert the lifespan wiring +contract only. +""" + +from __future__ import annotations + +import os +from collections.abc import Iterator +from pathlib import Path + +import pytest +from fastapi import FastAPI + +from everos.config import load_settings +from everos.core.lifespan import TracingLifespanProvider +from everos.core.observability.tracing import shutdown_tracing + + +@pytest.fixture(autouse=True) +def _isolate(monkeypatch: pytest.MonkeyPatch, tmp_path: Path) -> Iterator[None]: + for key in list(os.environ): + if key.startswith("EVEROS_"): + monkeypatch.delenv(key, raising=False) + monkeypatch.setenv("EVEROS_ROOT", str(tmp_path)) + monkeypatch.chdir(tmp_path) + load_settings.cache_clear() + shutdown_tracing() + yield + shutdown_tracing() + + +async def test_startup_returns_true_when_enabled( + monkeypatch: pytest.MonkeyPatch, +) -> None: + monkeypatch.setenv("EVEROS_OBSERVABILITY__ENABLED", "true") + monkeypatch.setenv("EVEROS_OBSERVABILITY__ENDPOINT", "http://collector.invalid") + load_settings.cache_clear() + provider = TracingLifespanProvider() + app = FastAPI() + result = await provider.startup(app) + assert result is True + await provider.shutdown(app) # must not raise + + +async def test_startup_returns_false_when_disabled() -> None: + provider = TracingLifespanProvider() + app = FastAPI() + result = await provider.startup(app) + assert result is False + await provider.shutdown(app) # must not raise + + +def test_provider_has_low_order_to_start_first() -> None: + # Tracer must be live before other providers start. + assert TracingLifespanProvider().order < 5 # MetricsLifespanProvider is 5 + + +async def test_startup_installs_score_sink_with_creds( + monkeypatch: pytest.MonkeyPatch, +) -> None: + """With Langfuse creds + emit_recall_scores, startup installs the score + sink; shutdown drains and tears it down.""" + from pydantic import SecretStr + + import everos.core.observability.tracing.scores as scores_mod + from everos.config import Settings + from everos.config.settings import ObservabilitySettings + + obs = ObservabilitySettings( + enabled=True, + endpoint="http://collector.invalid", + langfuse_public_key="pk-lf", + langfuse_secret_key=SecretStr("sk-lf"), + langfuse_host="https://lf.example", + ) + monkeypatch.setattr( + "everos.core.lifespan.tracing_lifespan.load_settings", + lambda: Settings(observability=obs), + ) + provider = TracingLifespanProvider() + app = FastAPI() + await provider.startup(app) + try: + assert scores_mod._sink is not None + finally: + await provider.shutdown(app) + assert scores_mod._sink is None diff --git a/tests/unit/test_core/test_middleware/test_request_id.py b/tests/unit/test_core/test_middleware/test_request_id.py new file mode 100644 index 0000000..e91253c --- /dev/null +++ b/tests/unit/test_core/test_middleware/test_request_id.py @@ -0,0 +1,113 @@ +"""``RequestIdMiddleware`` — mints a W3C-compatible request id per request, +propagates it to the endpoint via both ``request.state`` and the +``core.context`` contextvar, and echoes it on the ``X-Request-Id`` +response header. + +The contextvar assertion is the load-bearing one: it proves the id set in +the middleware crosses Starlette's ``BaseHTTPMiddleware`` task boundary and +is visible to downstream handlers / loggers. +""" + +from __future__ import annotations + +from collections.abc import AsyncIterator + +import pytest +from fastapi import FastAPI, Request +from httpx import ASGITransport, AsyncClient + +from everos.core.context import get_request_id +from everos.core.middleware import RequestIdMiddleware + + +def _build_app() -> FastAPI: + app = FastAPI() + app.add_middleware(RequestIdMiddleware) + + @app.get("/echo") + async def echo(request: Request) -> dict[str, str | None]: + return { + "from_state": getattr(request.state, "request_id", None), + "from_contextvar": get_request_id(), + } + + return app + + +@pytest.fixture +async def client() -> AsyncIterator[AsyncClient]: + app = _build_app() + async with AsyncClient( + transport=ASGITransport(app=app), base_url="http://test" + ) as c: + yield c + + +async def test_sets_request_id_on_state_and_response_header( + client: AsyncClient, +) -> None: + resp = await client.get("/echo") + assert resp.status_code == 200 + rid = resp.headers["x-request-id"] + assert len(rid) == 32 # W3C trace-id shape (gen_request_id) + assert resp.json()["from_state"] == rid + + +async def test_request_id_visible_to_endpoint_via_contextvar( + client: AsyncClient, +) -> None: + resp = await client.get("/echo") + body = resp.json() + assert body["from_contextvar"] is not None + assert body["from_contextvar"] == resp.headers["x-request-id"] + + +async def test_each_request_gets_distinct_id(client: AsyncClient) -> None: + r1 = await client.get("/echo") + r2 = await client.get("/echo") + assert r1.headers["x-request-id"] != r2.headers["x-request-id"] + + +async def test_inbound_traceparent_continues_upstream_trace() -> None: + """When the request carries a W3C traceparent header, our first span + continues that upstream trace (distributed tracing). Absent → own root.""" + 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, + memory_span, + shutdown_tracing, + ) + + exporter = InMemorySpanExporter() + shutdown_tracing() + init_tracing( + ObservabilitySettings(enabled=True, endpoint="http://collector.invalid"), + span_processor=SimpleSpanProcessor(exporter), + ) + app = FastAPI() + app.add_middleware(RequestIdMiddleware) + + @app.get("/s") + async def s() -> dict[str, str]: + with memory_span("everos.memory.search", observation_type="retriever"): + pass + return {"ok": "1"} + + upstream = "1234567890abcdef1234567890abcdef" + tp = f"00-{upstream}-1111111111111111-01" + try: + async with AsyncClient( + transport=ASGITransport(app=app), base_url="http://test" + ) as c: + await c.get("/s", headers={"traceparent": tp}) + force_flush() + span = exporter.get_finished_spans()[0] + assert format(span.context.trace_id, "032x") == upstream + finally: + shutdown_tracing() diff --git a/tests/unit/test_core/test_observability/test_scores.py b/tests/unit/test_core/test_observability/test_scores.py new file mode 100644 index 0000000..6cd80f7 --- /dev/null +++ b/tests/unit/test_core/test_observability/test_scores.py @@ -0,0 +1,89 @@ +"""``RecallScoreSink`` — non-blocking recall-score push. + +The sink is the piece that guarantees the Langfuse scores REST call never +touches the search request path: ``enqueue`` is O(1) and never blocks / raises +(drops + counts when full); a background worker drains and sends. The network +``sender`` is injected so these tests assert the queue/worker contract offline. +""" + +from __future__ import annotations + +import asyncio + +from everos.core.observability.tracing.scores import RecallScoreSink, ScoreRecord + + +async def test_worker_sends_payload_in_langfuse_shape() -> None: + sent: list[dict] = [] + done = asyncio.Event() + + async def sender(payload: dict) -> None: + sent.append(payload) + done.set() + + sink = RecallScoreSink(sender=sender, max_queue=10) + sink.start() + sink.enqueue( + ScoreRecord( + trace_id="tid", + observation_id="oid", + name="recall_top_score", + value=0.8, + comment="method=hybrid", + ) + ) + await asyncio.wait_for(done.wait(), timeout=1.0) + await sink.stop() + + assert sent[0] == { + "traceId": "tid", + "observationId": "oid", + "name": "recall_top_score", + "value": 0.8, + "dataType": "NUMERIC", + "comment": "method=hybrid", + } + + +async def test_enqueue_never_blocks_or_raises_when_full() -> None: + async def slow(_: dict) -> None: + await asyncio.sleep(10) + + sink = RecallScoreSink(sender=slow, max_queue=1) + # No worker started → queue fills. enqueue must stay non-blocking. + sink.enqueue(ScoreRecord("t", "o", "n", 1.0, None)) # fills the single slot + sink.enqueue(ScoreRecord("t", "o", "n", 1.0, None)) # dropped, no raise + assert sink.dropped == 1 + + +async def test_sender_failure_does_not_crash_worker() -> None: + calls: list[dict] = [] + second = asyncio.Event() + + async def flaky(payload: dict) -> None: + calls.append(payload) + if len(calls) == 1: + raise RuntimeError("boom") + second.set() + + sink = RecallScoreSink(sender=flaky, max_queue=10) + sink.start() + sink.enqueue(ScoreRecord("t", "o", "n1", 1.0, None)) # sender raises + sink.enqueue(ScoreRecord("t", "o", "n2", 2.0, None)) # worker must survive + await asyncio.wait_for(second.wait(), timeout=1.0) + await sink.stop() + assert len(calls) == 2 # first failed but worker kept going + + +async def test_stop_drains_pending() -> None: + sent: list[dict] = [] + + async def sender(payload: dict) -> None: + sent.append(payload) + + sink = RecallScoreSink(sender=sender, max_queue=10) + sink.start() + for i in range(3): + sink.enqueue(ScoreRecord("t", "o", f"n{i}", float(i), None)) + await sink.stop() # should drain the queue before returning + assert len(sent) == 3 diff --git a/tests/unit/test_core/test_observability/test_span_attributes.py b/tests/unit/test_core/test_observability/test_span_attributes.py new file mode 100644 index 0000000..ee3a89f --- /dev/null +++ b/tests/unit/test_core/test_observability/test_span_attributes.py @@ -0,0 +1,138 @@ +"""``memory_span`` / ``set_generation_usage`` — the langfuse.* + gen_ai.* +attribute contract applied to spans. + +Captured via an in-memory exporter so the emitted attribute keys/values +are asserted directly against §4 of the implementation plan. +""" + +from __future__ import annotations + +from collections.abc import Iterator + +import pytest +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 ( + capture_input, + capture_output, + force_flush, + init_tracing, + memory_span, + set_capture_content, + set_generation_usage, + set_redactor, + shutdown_tracing, +) + + +@pytest.fixture(autouse=True) +def _reset() -> Iterator[None]: + shutdown_tracing() + yield + shutdown_tracing() + + +@pytest.fixture +def captured() -> Iterator[InMemorySpanExporter]: + exporter = InMemorySpanExporter() + init_tracing( + ObservabilitySettings(enabled=True, endpoint="http://collector.invalid"), + span_processor=SimpleSpanProcessor(exporter), + ) + yield exporter + + +def test_memory_span_sets_langfuse_attributes( + captured: InMemorySpanExporter, +) -> None: + with memory_span( + "everos.memory.search", + observation_type="retriever", + session_id="s1", + user_id="u1", + metadata={"app_id": "a", "project_id": "p", "agent_id": None}, + ): + pass + force_flush() + span = captured.get_finished_spans()[0] + attrs = span.attributes + assert span.name == "everos.memory.search" + assert attrs["langfuse.observation.type"] == "retriever" + assert attrs["langfuse.session.id"] == "s1" + assert attrs["langfuse.user.id"] == "u1" + assert attrs["langfuse.trace.metadata.app_id"] == "a" + assert attrs["langfuse.trace.metadata.project_id"] == "p" + # None-valued metadata is dropped, not emitted as "None". + assert "langfuse.trace.metadata.agent_id" not in attrs + assert list(attrs["langfuse.trace.tags"]) == ["everos", "memory"] + + +def test_set_generation_usage_annotates_current_span( + captured: InMemorySpanExporter, +) -> None: + with memory_span("everos.extract", observation_type="generation"): + set_generation_usage(model="gpt-x", input_tokens=11, output_tokens=22) + force_flush() + attrs = captured.get_finished_spans()[0].attributes + assert attrs["gen_ai.request.model"] == "gpt-x" + assert attrs["gen_ai.usage.input_tokens"] == 11 + assert attrs["gen_ai.usage.output_tokens"] == 22 + + +def test_set_generation_usage_outside_span_is_noop() -> None: + # No active span → must not raise (and nothing to record). + set_generation_usage(model="x", input_tokens=1, output_tokens=2) + + +def test_content_dropped_when_capture_off(captured: InMemorySpanExporter) -> None: + # Default: capture_content off → no observation.input/output emitted. + with memory_span("everos.extract", observation_type="generation") as span: + capture_input(span, {"query": "sensitive query"}) + capture_output(span, "secret memory text") + force_flush() + attrs = captured.get_finished_spans()[0].attributes + assert "langfuse.observation.input" not in attrs + assert "langfuse.observation.output" not in attrs + + +def test_content_emitted_when_capture_on(captured: InMemorySpanExporter) -> None: + set_capture_content(True) + try: + with memory_span("everos.memory.search", observation_type="retriever") as span: + capture_input(span, {"query": "hello"}) + capture_output(span, "world") + finally: + set_capture_content(False) + force_flush() + attrs = captured.get_finished_spans()[0].attributes + import json + + assert json.loads(attrs["langfuse.observation.input"]) == {"query": "hello"} + assert attrs["langfuse.observation.output"] == "world" + + +def test_content_redacted_and_truncated(captured: InMemorySpanExporter) -> None: + calls: list[str] = [] + + def redact(text: str) -> str: + calls.append(text) + return text.replace("SECRET", "***") + + set_redactor(redact) + set_capture_content(True) + try: + with memory_span("everos.extract", observation_type="generation") as span: + capture_output(span, "SECRET " + "a" * 6000) + finally: + set_capture_content(False) + set_redactor(None) + force_flush() + val = captured.get_finished_spans()[0].attributes["langfuse.observation.output"] + assert "SECRET" not in val + assert "***" in val + assert len(val) <= 4096 # truncated + assert calls # redaction hook was invoked diff --git a/tests/unit/test_core/test_observability/test_tracing.py b/tests/unit/test_core/test_observability/test_tracing.py new file mode 100644 index 0000000..5e3518a --- /dev/null +++ b/tests/unit/test_core/test_observability/test_tracing.py @@ -0,0 +1,170 @@ +"""``core.observability.tracing`` — provider lifecycle + tracer facade. + +Spans are captured via an in-memory exporter (a ``SimpleSpanProcessor`` +injected into ``init_tracing``) so assertions run offline, with no OTLP +endpoint. The provider is kept off the OTel *global* on purpose — the +module holds its own reference — so tests can init / shutdown repeatedly +without hitting OTel's set-global-once guard. +""" + +from __future__ import annotations + +from collections.abc import Iterator + +import pytest +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, + get_tracer, + init_tracing, + shutdown_tracing, +) + + +@pytest.fixture(autouse=True) +def _reset_tracing() -> Iterator[None]: + """Ensure each test starts and ends with no provider installed.""" + shutdown_tracing() + yield + shutdown_tracing() + + +@pytest.fixture +def captured_spans() -> Iterator[InMemorySpanExporter]: + exporter = InMemorySpanExporter() + settings = ObservabilitySettings(enabled=True, endpoint="http://collector.invalid") + init_tracing(settings, span_processor=SimpleSpanProcessor(exporter)) + yield exporter + + +def test_disabled_tracer_is_noop_and_never_raises() -> None: + # No init_tracing → no provider → get_tracer returns a no-op tracer. + tracer = get_tracer("everos.test") + with tracer.start_as_current_span("everos.noop") as span: + span.set_attribute("k", "v") # must not raise + + +def test_init_returns_false_when_disabled() -> None: + assert init_tracing(ObservabilitySettings(enabled=False)) is False + + +def test_init_returns_true_when_enabled(captured_spans: InMemorySpanExporter) -> None: + # captured_spans fixture already called init_tracing(enabled=True). + tracer = get_tracer("x") + with tracer.start_as_current_span("s"): + pass + force_flush() + assert len(captured_spans.get_finished_spans()) == 1 + + +def test_span_captured_when_enabled(captured_spans: InMemorySpanExporter) -> None: + tracer = get_tracer("everos.test") + with tracer.start_as_current_span("everos.memory.search"): + pass + force_flush() + names = [s.name for s in captured_spans.get_finished_spans()] + assert "everos.memory.search" in names + + +def test_child_span_nests_under_parent( + captured_spans: InMemorySpanExporter, +) -> None: + tracer = get_tracer("everos.test") + with ( + tracer.start_as_current_span("parent"), + tracer.start_as_current_span("child"), + ): + pass + force_flush() + spans = {s.name: s for s in captured_spans.get_finished_spans()} + assert spans["child"].parent is not None + assert spans["child"].parent.span_id == spans["parent"].context.span_id + + +def test_resource_carries_service_name( + captured_spans: InMemorySpanExporter, +) -> None: + tracer = get_tracer("x") + with tracer.start_as_current_span("s"): + pass + force_flush() + span = captured_spans.get_finished_spans()[0] + assert span.resource.attributes["service.name"] == "everos" + + +async def test_child_spans_nest_across_asyncio_gather( + captured_spans: InMemorySpanExporter, +) -> None: + """OTel context must survive the asyncio.gather task boundary: child + spans started inside gathered coroutines nest under the parent span, + not as siblings/roots. This is the mechanism SearchManager.search + relies on (plan cross-cutting note #1).""" + import asyncio + + tracer = get_tracer("everos.test") + + async def _child(name: str) -> None: + with tracer.start_as_current_span(name): + await asyncio.sleep(0) + + with tracer.start_as_current_span("parent"): + await asyncio.gather(_child("child_a"), _child("child_b")) + force_flush() + + spans = {s.name: s for s in captured_spans.get_finished_spans()} + parent_span_id = spans["parent"].context.span_id + trace_id = spans["parent"].context.trace_id + for child in ("child_a", "child_b"): + assert spans[child].parent is not None + assert spans[child].parent.span_id == parent_span_id + assert spans[child].context.trace_id == trace_id + + +def test_resolve_otlp_target_derives_from_langfuse_creds() -> None: + import base64 + + from pydantic import SecretStr + + from everos.core.observability.tracing.provider import _resolve_otlp_target + + settings = ObservabilitySettings( + enabled=True, + langfuse_public_key="pk", + langfuse_secret_key=SecretStr("sk"), + langfuse_host="https://us.cloud.langfuse.com", + ) + endpoint, headers = _resolve_otlp_target(settings) + assert endpoint == "https://us.cloud.langfuse.com/api/public/otel/v1/traces" + assert headers["Authorization"] == "Basic " + base64.b64encode(b"pk:sk").decode() + + +def test_resolve_otlp_target_explicit_values_win() -> None: + from pydantic import SecretStr + + from everos.core.observability.tracing.provider import _resolve_otlp_target + + settings = ObservabilitySettings( + enabled=True, + endpoint="http://explicit/v1/traces", + headers={"Authorization": "Basic explicit"}, + langfuse_public_key="pk", + langfuse_secret_key=SecretStr("sk"), + langfuse_host="https://us.cloud.langfuse.com", + ) + endpoint, headers = _resolve_otlp_target(settings) + assert endpoint == "http://explicit/v1/traces" + assert headers["Authorization"] == "Basic explicit" + + +def test_resolve_otlp_target_no_langfuse_returns_as_is() -> None: + from everos.core.observability.tracing.provider import _resolve_otlp_target + + settings = ObservabilitySettings(enabled=True, endpoint="http://x/v1/traces") + endpoint, headers = _resolve_otlp_target(settings) + assert endpoint == "http://x/v1/traces" + assert "Authorization" not in headers 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 diff --git a/uv.lock b/uv.lock index 5620515..94a8541 100644 --- a/uv.lock +++ b/uv.lock @@ -595,11 +595,17 @@ dependencies = [ multimodal = [ { name = "everalgo-parser", extra = ["svg"] }, ] +otel = [ + { name = "opentelemetry-exporter-otlp-proto-http" }, + { name = "opentelemetry-sdk" }, +] [package.dev-dependencies] dev = [ { name = "import-linter" }, { name = "ipdb" }, + { name = "opentelemetry-exporter-otlp-proto-http" }, + { name = "opentelemetry-sdk" }, { name = "pre-commit" }, { name = "pyinstrument" }, { name = "pytest" }, @@ -625,6 +631,8 @@ requires-dist = [ { name = "jieba", specifier = ">=0.42.1,<1.0" }, { name = "lancedb", specifier = ">=0.13.0" }, { name = "openai", specifier = ">=1.0.0" }, + { name = "opentelemetry-exporter-otlp-proto-http", marker = "extra == 'otel'", specifier = ">=1.27.0" }, + { name = "opentelemetry-sdk", marker = "extra == 'otel'", specifier = ">=1.27.0" }, { name = "portalocker", specifier = ">=2.8.2" }, { name = "prometheus-client", specifier = ">=0.20.0" }, { name = "pydantic", specifier = ">=2.7.1" }, @@ -639,12 +647,14 @@ requires-dist = [ { name = "watchdog", specifier = ">=4.0.0" }, { name = "watchfiles", specifier = ">=0.21.0" }, ] -provides-extras = ["multimodal"] +provides-extras = ["multimodal", "otel"] [package.metadata.requires-dev] dev = [ { name = "import-linter", specifier = ">=2.0" }, { name = "ipdb", specifier = ">=0.13.13" }, + { name = "opentelemetry-exporter-otlp-proto-http", specifier = ">=1.27.0" }, + { name = "opentelemetry-sdk", specifier = ">=1.27.0" }, { name = "pre-commit", specifier = ">=4.0.0" }, { name = "pyinstrument", specifier = ">=5.0.0" }, { name = "pytest", specifier = ">=8.4.0" }, @@ -688,6 +698,18 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/81/47/dd9a212ef6e343a6857485ffe25bba537304f1913bdbed446a23f7f592e1/filelock-3.29.0-py3-none-any.whl", hash = "sha256:96f5f6344709aa1572bbf631c640e4ebeeb519e08da902c39a001882f30ac258", size = 39812, upload-time = "2026-04-19T15:39:08.752Z" }, ] +[[package]] +name = "googleapis-common-protos" +version = "1.75.0" +source = { registry = "https://pypi.org/simple" } +dependencies = [ + { name = "protobuf" }, +] +sdist = { url = "https://files.pythonhosted.org/packages/b5/c8/f439cffde755cffa462bfbb156278fa6f9d09119719af9814b858fd4f81f/googleapis_common_protos-1.75.0.tar.gz", hash = "sha256:53a062ff3c32552fbd62c11fe23768b78e4ddf0494d5e5fd97d3f4689c75fbbd", size = 151035, upload-time = "2026-05-07T08:04:49.423Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/e7/c8/e2645aa8ed02fd4c7a2f59d68783b65b1f3cbdfe39a6308e156509d1fee8/googleapis_common_protos-1.75.0-py3-none-any.whl", hash = "sha256:961ed60399c457ceb0ee8f285a84c870aabc9c6a832b9d37bb281b5bebde43ed", size = 300631, upload-time = "2026-05-07T08:03:30.345Z" }, +] + [[package]] name = "greenlet" version = "3.5.3" @@ -1339,6 +1361,87 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/9d/1c/5d43735b2553baae2a5e899dcbcd0670a86930d993184d72ca909bf11c9b/openai-2.36.0-py3-none-any.whl", hash = "sha256:143f6194b548dbc2c921af1f1b03b9f14c85fed8a75b5b516f5bcc11a2a50c63", size = 1302361, upload-time = "2026-05-07T17:33:15.063Z" }, ] +[[package]] +name = "opentelemetry-api" +version = "1.44.0" +source = { registry = "https://pypi.org/simple" } +dependencies = [ + { name = "typing-extensions" }, +] +sdist = { url = "https://files.pythonhosted.org/packages/ee/8b/aa9e2d8b8dfa7c946f7dec5d1f8f6ba8eca062f43509a06bdb5ce93d26c0/opentelemetry_api-1.44.0.tar.gz", hash = "sha256:67647e5e9566edcf421166fdf022b3537f818635daa852b289e34604dc6fb33a", size = 72406, upload-time = "2026-07-16T15:25:32.678Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/ca/6f/a04e900f465ff3221ccc395522503e2d10e79fa21f2723c8e177aae1e0d1/opentelemetry_api-1.44.0-py3-none-any.whl", hash = "sha256:94b98c893a91b88657eaac1e3ba89618cdb85be6918196705354f34728b2cdef", size = 60018, upload-time = "2026-07-16T15:25:11.657Z" }, +] + +[[package]] +name = "opentelemetry-exporter-otlp-proto-common" +version = "1.44.0" +source = { registry = "https://pypi.org/simple" } +dependencies = [ + { name = "opentelemetry-proto" }, +] +sdist = { url = "https://files.pythonhosted.org/packages/61/09/4d717852c1cf3f854b76c7110a5d00883bc3c99288b9b0dbcbeb9e306eb6/opentelemetry_exporter_otlp_proto_common-1.44.0.tar.gz", hash = "sha256:dc87a5a5bc58f149a56d1547e4691588fa12994cdc3bc039a694ccb3375862ac", size = 20202, upload-time = "2026-07-16T15:25:37.658Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/5e/71/65fd9d54c10b860f87c045ccee1264cab7011268895d3528818a29c1172a/opentelemetry_exporter_otlp_proto_common-1.44.0-py3-none-any.whl", hash = "sha256:9a9fe61bba73d802904bc989f1d6b4a7b1ee40f06c40e98d6f85af65aaebb694", size = 17045, upload-time = "2026-07-16T15:25:18.201Z" }, +] + +[[package]] +name = "opentelemetry-exporter-otlp-proto-http" +version = "1.44.0" +source = { registry = "https://pypi.org/simple" } +dependencies = [ + { name = "googleapis-common-protos" }, + { name = "opentelemetry-api" }, + { name = "opentelemetry-exporter-otlp-proto-common" }, + { name = "opentelemetry-proto" }, + { name = "opentelemetry-sdk" }, + { name = "requests" }, + { name = "typing-extensions" }, +] +sdist = { url = "https://files.pythonhosted.org/packages/1a/87/95e2a5aaa795b4e2260d74e16df2d5541deb2ea9de010bcd615f4dee2654/opentelemetry_exporter_otlp_proto_http-1.44.0.tar.gz", hash = "sha256:c633d7270ad6b57cd4cfbe8b0007a9e2e7c0cb50bd6c50fe2a7b245f721a09d8", size = 25806, upload-time = "2026-07-16T15:25:39.162Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/cd/d0/fdeb1a98d8d3a6205f5f297c51b4a9bfe65126ab60339669bbe3dd54c2e2/opentelemetry_exporter_otlp_proto_http-1.44.0-py3-none-any.whl", hash = "sha256:838592fce774c1c8bb7b9a0a7facbfa82e17be5a8a4e94cef10cb84ae026bae3", size = 21850, upload-time = "2026-07-16T15:25:20.006Z" }, +] + +[[package]] +name = "opentelemetry-proto" +version = "1.44.0" +source = { registry = "https://pypi.org/simple" } +dependencies = [ + { name = "protobuf" }, +] +sdist = { url = "https://files.pythonhosted.org/packages/64/01/40ac4ae9a149263cc52c2cee200ddd80cb6d8db1a4610abf8eabce0fe771/opentelemetry_proto-1.44.0.tar.gz", hash = "sha256:c547a79c2f8c0c515d31509154682e5921c7cfd5ca67b70e1f9266e2c3e103f3", size = 46488, upload-time = "2026-07-16T15:25:45.34Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/d1/7c/8be563d68e93bbefa5c8affb82ddcff91b3ad858ce49957ba7b16fd3e0ab/opentelemetry_proto-1.44.0-py3-none-any.whl", hash = "sha256:898b155a0e1557afd867478fb6158e8122a46329ca0bb8dc53cc55e98f017f56", size = 72483, upload-time = "2026-07-16T15:25:28.429Z" }, +] + +[[package]] +name = "opentelemetry-sdk" +version = "1.44.0" +source = { registry = "https://pypi.org/simple" } +dependencies = [ + { name = "opentelemetry-api" }, + { name = "opentelemetry-semantic-conventions" }, + { name = "typing-extensions" }, +] +sdist = { url = "https://files.pythonhosted.org/packages/5d/77/a6592cbc7c8d9bcc9d6757a9df45e04a7c585e3e6e7a13456da522b21109/opentelemetry_sdk-1.44.0.tar.gz", hash = "sha256:cebe7f65dc12f26ead75c6064de12fd2a9052e5060c0272d402cfa203aae123b", size = 208624, upload-time = "2026-07-16T15:25:46.078Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/e7/23/ff077e61886ee020a17ce9c8b6fa11c601c8d8345b09ea24f605445df62a/opentelemetry_sdk-1.44.0-py3-none-any.whl", hash = "sha256:df081c4c6bcfdb1211e3e86140376792643128a25f8d72d1d27675936e7e96ad", size = 137221, upload-time = "2026-07-16T15:25:29.534Z" }, +] + +[[package]] +name = "opentelemetry-semantic-conventions" +version = "0.65b0" +source = { registry = "https://pypi.org/simple" } +dependencies = [ + { name = "opentelemetry-api" }, + { name = "typing-extensions" }, +] +sdist = { url = "https://files.pythonhosted.org/packages/8f/73/0cbdebcb4cf545fdd328da14f5137e37d0770c3f26185e478b0d15d94f50/opentelemetry_semantic_conventions-0.65b0.tar.gz", hash = "sha256:f9b2b81e9d5b64f11bc952075e7e9c7fb0aab075c7fd1c46d597f1b919852d60", size = 148774, upload-time = "2026-07-16T15:25:46.902Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/a6/0e/49df70d9b81fb5cbae4bbf2a49d865b09bcbcbc4eb53f5851b1027738d78/opentelemetry_semantic_conventions-0.65b0-py3-none-any.whl", hash = "sha256:1cacde7b0ad306f84c5ef08c3dbe1bbaf20165bba6f8bff43b670e555a086bcb", size = 204645, upload-time = "2026-07-16T15:25:30.688Z" }, +] + [[package]] name = "packaging" version = "26.2" @@ -1505,6 +1608,21 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/84/03/0d3ce49e2505ae70cf43bc5bb3033955d2fc9f932163e84dc0779cc47f48/prompt_toolkit-3.0.52-py3-none-any.whl", hash = "sha256:9aac639a3bbd33284347de5ad8d68ecc044b91a762dc39b7c21095fcd6a19955", size = 391431, upload-time = "2025-08-27T15:23:59.498Z" }, ] +[[package]] +name = "protobuf" +version = "7.35.1" +source = { registry = "https://pypi.org/simple" } +sdist = { url = "https://files.pythonhosted.org/packages/da/01/9ef0afd7999eb9badb3a768b4aedd78c86d4c65cfaf1958ab276199e76b4/protobuf-7.35.1.tar.gz", hash = "sha256:ce115a26fe0c39a2c29973d914d327e516a6455464489fe3cd1e51a1b354f81a", size = 458717, upload-time = "2026-06-11T21:55:40.257Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/10/03/8aeeb7458d22546bf64b5250ca1daeb5ff757d900e8e4a7476c6f0db843e/protobuf-7.35.1-cp310-abi3-macosx_10_9_universal2.whl", hash = "sha256:24f857477359a85c0c235261b8ba905fd51b2562f4a64ca1df5473f29850cbf6", size = 433226, upload-time = "2026-06-11T21:55:31.719Z" }, + { url = "https://files.pythonhosted.org/packages/37/4b/dfb89eb0e652a1ff073c39a59fb5e3a83cfe9b57a2c83fa6d78270101767/protobuf-7.35.1-cp310-abi3-manylinux2014_aarch64.whl", hash = "sha256:11d6b0ec246892d85215b0a13ca6e0233cf5284b68f0ac02646427f4ff88a799", size = 328847, upload-time = "2026-06-11T21:55:34.035Z" }, + { url = "https://files.pythonhosted.org/packages/0f/58/dc12f2cd484951524af6e3382c785869b9b3fb5e52ee95ae23add53ee8f9/protobuf-7.35.1-cp310-abi3-manylinux2014_s390x.whl", hash = "sha256:b73f9489a4b8b1c9cb1f8ed951c736392592edb24b9d6819f36d2e10b171d5b4", size = 344030, upload-time = "2026-06-11T21:55:34.941Z" }, + { url = "https://files.pythonhosted.org/packages/e4/be/5b3cfe508bfab6761414ff944e3366eb13be4fd71efcd69450f89ba39f43/protobuf-7.35.1-cp310-abi3-manylinux2014_x86_64.whl", hash = "sha256:74758715c53d7158fb76caf4f0cfdacc5329a4b1bb994f865d6cf302d413a1c4", size = 327130, upload-time = "2026-06-11T21:55:35.921Z" }, + { url = "https://files.pythonhosted.org/packages/d8/bc/6d6c7ba8709c85f8f2c390b2b118d6fb08a783676a572271851bf45a7d22/protobuf-7.35.1-cp310-abi3-win32.whl", hash = "sha256:353652e4efd0bca5b5fc2656abf8307ef351f0cf938c9eba09f0e09c20a25c30", size = 428945, upload-time = "2026-06-11T21:55:37.034Z" }, + { url = "https://files.pythonhosted.org/packages/0a/19/8d0cb6f20a1ef7b18f1c8986ad5783f22f84cce39c6ce9a6e645ea55192e/protobuf-7.35.1-cp310-abi3-win_amd64.whl", hash = "sha256:230a75ddfc2de4806e56696ce9640c1cdfdb6543b7cfce98d42a4c0a0e7bdb87", size = 439996, upload-time = "2026-06-11T21:55:38.123Z" }, + { url = "https://files.pythonhosted.org/packages/19/c7/5f7c636ec43e0c545e28d1f1db71990108306f7bdcb89f069ba97e428e7f/protobuf-7.35.1-py3-none-any.whl", hash = "sha256:4bc97768d8fe4ad6743c8a19403e314511ed9f6d13205b687e52421c023ac1b9", size = 171659, upload-time = "2026-06-11T21:55:39.155Z" }, +] + [[package]] name = "psutil" version = "7.2.2"