diff --git a/.env.example b/.env.example index 26746e0a..42f29b55 100644 --- a/.env.example +++ b/.env.example @@ -35,6 +35,18 @@ REMEMBERSTACK_OPENROUTER_MAX_COMPLETION_TOKENS=32000 # reasoning model. # REMEMBERSTACK_OPENROUTER_REASONING_EFFORT_MAP={"z-ai/glm-4.7-flash":"none","openai/gpt-5.6-luna":"high"} +# OBSERVABILITY (optional; empty or unset keeps all exporters disabled). +# Sentry-protocol error tracking works with Sentry, GlitchTip, and Bugsink. +# The environment defaults to REMEMBERSTACK_SELFHOST_DEPLOYMENT_SLUG and the +# error-event sample rate defaults to 1.0. +# REMEMBERSTACK_SENTRY_DSN= +# REMEMBERSTACK_SENTRY_ENVIRONMENT= +# REMEMBERSTACK_SENTRY_SAMPLE_RATE=1.0 +# LoCoMo answer/judge tracing activates only when all three values are non-empty. +# LANGFUSE_PUBLIC_KEY= +# LANGFUSE_SECRET_KEY= +# LANGFUSE_HOST=https://cloud.langfuse.com + # Optional benchmark/deployment model overrides. Keep explicit model IDs for a # reproducible run; do not use a rotating router such as openrouter/free. # REMEMBERSTACK_E2_EXTRACT_MODEL=nvidia/nemotron-3-super-120b-a12b:free diff --git a/Dockerfile b/Dockerfile index bc364349..f68a5078 100644 --- a/Dockerfile +++ b/Dockerfile @@ -17,11 +17,11 @@ RUN addgroup --system app \ --home /var/lib/rememberstack --no-create-home app \ && mkdir -p /var/lib/rememberstack/forget-manifests \ && chown -R app:app /var/lib/rememberstack \ - && uv sync --locked --no-dev --extra server --no-install-project + && uv sync --locked --no-dev --extra observability --extra server --no-install-project COPY src ./src -RUN uv sync --locked --no-dev --extra server +RUN uv sync --locked --no-dev --extra observability --extra server # Provenance: the exact source revision baked into this image. The benchmark # harness compares it against the revision it prepared with, so a run can never diff --git a/benchmarks/locomo/runner.py b/benchmarks/locomo/runner.py index 560f5b87..c9c5616a 100644 --- a/benchmarks/locomo/runner.py +++ b/benchmarks/locomo/runner.py @@ -8,16 +8,21 @@ from decimal import Decimal import hashlib import json +import logging import os from pathlib import Path import subprocess import sys import time from typing import Final +from typing import TYPE_CHECKING from uuid import UUID from pydantic import BaseModel +from pydantic import SecretStr from pydantic import ValidationError +from pydantic_settings import BaseSettings +from pydantic_settings import SettingsConfigDict from benchmarks.locomo.dataset import DATASET_COMMIT from benchmarks.locomo.dataset import DATASET_SHA256 @@ -76,6 +81,12 @@ from rememberstack.surfaces.sdk import MemoryApiError from rememberstack.surfaces.sdk import MemoryClient +_logger = logging.getLogger(__name__) + +if TYPE_CHECKING: + from benchmarks.locomo.tracing import LocomoTracer + from benchmarks.locomo.tracing import QuestionTrace + _RUN_FILE: Final = "run.json" _MANIFEST_FILE: Final = "manifest.json" _DOCUMENTS_FILE: Final = "documents.json" @@ -91,6 +102,33 @@ class ExecutionGuardError(BenchmarkRunError): """A remote stage lacks an exact execution/cost/isolation acknowledgement.""" +class _LangfuseActivationSettings(BaseSettings): + """Standard Langfuse bindings used only to decide whether to load the shim.""" + + model_config = SettingsConfigDict(env_prefix="LANGFUSE_", extra="ignore") + + public_key: SecretStr | None = None + secret_key: SecretStr | None = None + host: str | None = None + + def configured_values(self) -> tuple[str, str, str] | None: + """Return credentials only when all three explicit opt-in values are set.""" + public_key = ( + "" + if self.public_key is None + else self.public_key.get_secret_value().strip() + ) + secret_key = ( + "" + if self.secret_key is None + else self.secret_key.get_secret_value().strip() + ) + host = (self.host or "").strip() + if not public_key or not secret_key or not host: + return None + return public_key, secret_key, host + + def prepare_run(*, dataset_path: Path, tier: str, output: Path) -> RunConfiguration: """Validate, fingerprint, and render a local run without remote calls.""" dataset = load_dataset(dataset_path) @@ -389,19 +427,47 @@ def answer_sample( for record in context.state.ingests.values() if record.sample_id == sample_id } - for question in remaining: - record = _answer_one( - question=question, - client=client, - provider=provider, - tools=tools, - doc_sessions=doc_sessions, - state=context.state, - max_agent_calls=max_agent_calls, - max_evaluator_cost_usd=max_evaluator_cost_usd, - ) - context.state.answers[question.item_id] = record - _save_state(run_dir=run_dir, state=context.state) + tracer = _configured_langfuse_tracer(context=context) + try: + for question in remaining: + if tracer is None: + record = _answer_one( + question=question, + client=client, + provider=provider, + tools=tools, + doc_sessions=doc_sessions, + state=context.state, + max_agent_calls=max_agent_calls, + max_evaluator_cost_usd=max_evaluator_cost_usd, + ) + else: + with tracer.question( + item_id=question.item_id, question=question.question, stage="answer" + ) as question_trace: + record = _answer_one( + question=question, + client=client, + provider=provider, + tools=tools, + doc_sessions=doc_sessions, + state=context.state, + max_agent_calls=max_agent_calls, + max_evaluator_cost_usd=max_evaluator_cost_usd, + question_trace=question_trace, + ) + if question_trace is not None: + question_trace.finish_answer( + final_answer=record.generated_answer, + failure_kind=( + None if record.failure is None else record.failure.kind + ), + ) + context.state.answers[question.item_id] = record + _save_state(run_dir=run_dir, state=context.state) + finally: + if tracer is not None: + tracer.flush() return tuple(context.state.answers[question.item_id] for question in questions) @@ -447,25 +513,47 @@ def judge_sample( _require_cost_ceiling( spent=context.state.evaluator_cost_usd, ceiling=max_evaluator_cost_usd ) - for question in questions: - if question.item_id in context.state.judges: - continue - answer = context.state.answers[question.item_id] - if answer.failure is not None: - judge = JudgeRecord( - item_id=question.item_id, label="WRONG", model_called=False - ) - else: - judge = _judge_one( - question=question, - answer=answer, - provider=provider, - state=context.state, - max_judge_calls=max_judge_calls, - max_evaluator_cost_usd=max_evaluator_cost_usd, - ) - context.state.judges[question.item_id] = judge - _save_state(run_dir=run_dir, state=context.state) + tracer = _configured_langfuse_tracer(context=context) + try: + for question in questions: + if question.item_id in context.state.judges: + continue + answer = context.state.answers[question.item_id] + if tracer is None: + judge = _judge_answer( + question=question, + answer=answer, + provider=provider, + state=context.state, + max_judge_calls=max_judge_calls, + max_evaluator_cost_usd=max_evaluator_cost_usd, + ) + else: + with tracer.question( + item_id=question.item_id, question=question.question, stage="judge" + ) as question_trace: + judge = _judge_answer( + question=question, + answer=answer, + provider=provider, + state=context.state, + max_judge_calls=max_judge_calls, + max_evaluator_cost_usd=max_evaluator_cost_usd, + question_trace=question_trace, + ) + if question_trace is not None: + question_trace.finish_judge( + final_answer=answer.generated_answer, + verdict=judge.label, + failure_kind=( + None if judge.failure is None else judge.failure.kind + ), + ) + context.state.judges[question.item_id] = judge + _save_state(run_dir=run_dir, state=context.state) + finally: + if tracer is not None: + tracer.flush() return tuple(context.state.judges[question.item_id] for question in questions) @@ -939,6 +1027,32 @@ def _require_sample_ingested(*, context: _RunContext, sample_id: str) -> None: ) +def _configured_langfuse_tracer(*, context: _RunContext) -> LocomoTracer | None: + """Load the optional observer only when all standard bindings are non-empty.""" + configured = _LangfuseActivationSettings.model_validate({}).configured_values() + if configured is None: + return None + public_key, secret_key, host = configured + from benchmarks.locomo.tracing import create_langfuse_tracer + + configuration = context.configuration + run_identity = ( + f"{configuration.protocol_fingerprint}:" + f"{configuration.repository_revision}:" + f"{configuration.prepared_at.isoformat()}" + ) + try: + return create_langfuse_tracer( + public_key=public_key, + secret_key=secret_key, + host=host, + run_identity=run_identity, + ) + except Exception: + _logger.warning("optional Langfuse tracer initialization failed", exc_info=True) + return None + + def _answer_one( *, question: LoCoMoQuestion, @@ -949,6 +1063,7 @@ def _answer_one( state: RunState, max_agent_calls: int, max_evaluator_cost_usd: Decimal, + question_trace: QuestionTrace | None = None, ) -> AnswerRecord: """Let a bounded agent choose ordinary public recipes, then answer.""" tool_names = {tool.name for tool in tools} @@ -968,6 +1083,11 @@ def _answer_one( prompt = render_answer_agent_prompt( question=question.question, tools=tools, trace=tuple(trace) ) + agent_observation = ( + None + if question_trace is None + else question_trace.start_agent_call(model=ANSWER_AGENT_MODEL) + ) started = time.monotonic_ns() try: response = provider.generate( @@ -977,6 +1097,11 @@ def _answer_one( response_type=AnswerAgentStep, ) except ProviderAccountingError as error: + call_latency_ms = _elapsed_ms(started) + if agent_observation is not None: + agent_observation.finish( + usage=None, latency_ms=call_latency_ms, outcome="accounting_error" + ) return _failed_answer( question=question, kind="accounting", @@ -984,7 +1109,7 @@ def _answer_one( retrieval_latency_ms=tool_latency_ms, retrieval_succeeded=bool(trace), agent_call_count=len(usages) + 1, - reader_latency_ms=agent_latency_ms + _elapsed_ms(started), + reader_latency_ms=agent_latency_ms + call_latency_ms, claims=_claims_from_trace( trace=tuple(trace), doc_sessions=doc_sessions ), @@ -992,6 +1117,11 @@ def _answer_one( usages=tuple(usages), ) except ValidationError as error: + call_latency_ms = _elapsed_ms(started) + if agent_observation is not None: + agent_observation.finish( + usage=None, latency_ms=call_latency_ms, outcome="invalid_response" + ) return _failed_answer( question=question, kind="invalid_response", @@ -999,7 +1129,7 @@ def _answer_one( retrieval_latency_ms=tool_latency_ms, retrieval_succeeded=bool(trace), agent_call_count=len(usages) + 1, - reader_latency_ms=agent_latency_ms + _elapsed_ms(started), + reader_latency_ms=agent_latency_ms + call_latency_ms, claims=_claims_from_trace( trace=tuple(trace), doc_sessions=doc_sessions ), @@ -1007,9 +1137,16 @@ def _answer_one( usages=tuple(usages), ) except OpenRouterProviderError as error: + call_latency_ms = _elapsed_ms(started) if error.usage is not None: usages.append(error.usage) state.evaluator_cost_usd += error.usage.cost_usd + if agent_observation is not None: + agent_observation.finish( + usage=error.usage, + latency_ms=call_latency_ms, + outcome="provider_error", + ) return _failed_answer( question=question, kind="reader", @@ -1019,17 +1156,29 @@ def _answer_one( agent_call_count=( len(usages) if error.usage is not None else len(usages) + 1 ), - reader_latency_ms=agent_latency_ms + _elapsed_ms(started), + reader_latency_ms=agent_latency_ms + call_latency_ms, claims=_claims_from_trace( trace=tuple(trace), doc_sessions=doc_sessions ), tool_calls=tuple(trace), usages=tuple(usages), ) - agent_latency_ms += _elapsed_ms(started) + call_latency_ms = _elapsed_ms(started) + agent_latency_ms += call_latency_ms usages.append(response.usage) state.evaluator_cost_usd += response.usage.cost_usd + step = response.output + if agent_observation is not None and step.action == "tool": + agent_observation.finish( + usage=response.usage, latency_ms=call_latency_ms, outcome="tool" + ) if state.evaluator_cost_usd > max_evaluator_cost_usd: + if agent_observation is not None and step.action == "answer": + agent_observation.finish( + usage=response.usage, + latency_ms=call_latency_ms, + outcome="accounting_error", + ) return _failed_answer( question=question, kind="accounting", @@ -1047,9 +1196,14 @@ def _answer_one( tool_calls=tuple(trace), usages=tuple(usages), ) - step = response.output if step.action == "answer": if not trace: + if agent_observation is not None: + agent_observation.finish( + usage=response.usage, + latency_ms=call_latency_ms, + outcome="invalid_response", + ) return _failed_answer( question=question, kind="invalid_response", @@ -1063,6 +1217,12 @@ def _answer_one( ) answer = step.answer or "" if len(answer.split()) > 6: + if agent_observation is not None: + agent_observation.finish( + usage=response.usage, + latency_ms=call_latency_ms, + outcome="invalid_response", + ) return _failed_answer( question=question, kind="invalid_response", @@ -1078,6 +1238,13 @@ def _answer_one( usages=tuple(usages), ) claims = _claims_from_trace(trace=tuple(trace), doc_sessions=doc_sessions) + if agent_observation is not None: + agent_observation.finish( + usage=response.usage, + latency_ms=call_latency_ms, + outcome="answer", + final_answer=answer, + ) return AnswerRecord( item_id=question.item_id, sample_id=question.sample_id, @@ -1128,21 +1295,31 @@ def _answer_one( tool_calls=tuple(trace), usages=tuple(usages), ) - tool_started = time.monotonic_ns() # Validated by the model's own validator for tool steps, so this cannot # raise here; trailing junk (observed: a sentence period after the # closing brace at temperature 0) is recorded on the trace row. step_arguments, step_trailing = step.parsed_arguments() + tool_observation = ( + None + if question_trace is None + else question_trace.start_tool_call( + name=step.tool_name or "", arguments=step_arguments + ) + ) + tool_started = time.monotonic_ns() try: envelope = client.run_recipe( name=step.tool_name or "", arguments=step_arguments ) except MemoryApiError as error: + failed_latency = _elapsed_ms(tool_started) + if tool_observation is not None: + tool_observation.finish(latency_ms=failed_latency, outcome="api_error") return _failed_answer( question=question, kind="tool", message=str(error), - retrieval_latency_ms=tool_latency_ms + _elapsed_ms(tool_started), + retrieval_latency_ms=tool_latency_ms + failed_latency, retrieval_succeeded=False, agent_call_count=len(usages), reader_latency_ms=agent_latency_ms, @@ -1153,6 +1330,8 @@ def _answer_one( usages=tuple(usages), ) latency = _elapsed_ms(tool_started) + if tool_observation is not None: + tool_observation.finish(latency_ms=latency, outcome="succeeded") tool_latency_ms += latency trace.append( ToolCallRecord( @@ -1177,6 +1356,30 @@ def _answer_one( ) +def _judge_answer( + *, + question: LoCoMoQuestion, + answer: AnswerRecord, + provider: ModelProviderPort, + state: RunState, + max_judge_calls: int, + max_evaluator_cost_usd: Decimal, + question_trace: QuestionTrace | None = None, +) -> JudgeRecord: + """Return a local wrong for answer failures or invoke the configured judge.""" + if answer.failure is not None: + return JudgeRecord(item_id=question.item_id, label="WRONG", model_called=False) + return _judge_one( + question=question, + answer=answer, + provider=provider, + state=state, + max_judge_calls=max_judge_calls, + max_evaluator_cost_usd=max_evaluator_cost_usd, + question_trace=question_trace, + ) + + def _judge_one( *, question: LoCoMoQuestion, @@ -1185,6 +1388,7 @@ def _judge_one( state: RunState, max_judge_calls: int, max_evaluator_cost_usd: Decimal, + question_trace: QuestionTrace | None = None, ) -> JudgeRecord: """Invoke the judge once; every call failure becomes a visible wrong.""" called = sum(record.model_called for record in state.judges.values()) @@ -1193,6 +1397,11 @@ def _judge_one( _require_cost_before_call( spent=state.evaluator_cost_usd, ceiling=max_evaluator_cost_usd ) + judge_observation = ( + None + if question_trace is None + else question_trace.start_judge_call(model=JUDGE_MODEL) + ) started = time.monotonic_ns() try: response = provider.generate( @@ -1208,40 +1417,69 @@ def _judge_one( response_type=JudgeOutput, ) except ProviderAccountingError as error: + latency_ms = _elapsed_ms(started) + if judge_observation is not None: + judge_observation.finish( + usage=None, latency_ms=latency_ms, outcome="accounting_error" + ) return JudgeRecord( item_id=question.item_id, label="WRONG", model_called=True, - latency_ms=_elapsed_ms(started), + latency_ms=latency_ms, failure=_failure(kind="accounting", message=str(error)), ) except OpenRouterProviderError as error: + latency_ms = _elapsed_ms(started) if error.usage is not None: state.evaluator_cost_usd += error.usage.cost_usd + if judge_observation is not None: + judge_observation.finish( + usage=error.usage, + latency_ms=latency_ms, + outcome="provider_error", + verdict="WRONG", + ) return JudgeRecord( item_id=question.item_id, label="WRONG", model_called=True, usage=error.usage, - latency_ms=_elapsed_ms(started), + latency_ms=latency_ms, failure=_failure(kind="judge", message=str(error)), ) except ValidationError as error: + latency_ms = _elapsed_ms(started) + if judge_observation is not None: + judge_observation.finish( + usage=None, + latency_ms=latency_ms, + outcome="invalid_response", + verdict="WRONG", + ) return JudgeRecord( item_id=question.item_id, label="WRONG", model_called=True, - latency_ms=_elapsed_ms(started), + latency_ms=latency_ms, failure=_failure(kind="judge", message=str(error)), ) state.evaluator_cost_usd += response.usage.cost_usd + latency_ms = _elapsed_ms(started) if state.evaluator_cost_usd > max_evaluator_cost_usd: + if judge_observation is not None: + judge_observation.finish( + usage=response.usage, + latency_ms=latency_ms, + outcome="accounting_error", + verdict="WRONG", + ) return JudgeRecord( item_id=question.item_id, label="WRONG", model_called=True, usage=response.usage, - latency_ms=_elapsed_ms(started), + latency_ms=latency_ms, failure=_failure( kind="accounting", message=( @@ -1250,12 +1488,19 @@ def _judge_one( ), ), ) + if judge_observation is not None: + judge_observation.finish( + usage=response.usage, + latency_ms=latency_ms, + outcome="judged", + verdict=response.output.label, + ) return JudgeRecord( item_id=question.item_id, label=response.output.label, model_called=True, usage=response.usage, - latency_ms=_elapsed_ms(started), + latency_ms=latency_ms, ) diff --git a/benchmarks/locomo/tracing.py b/benchmarks/locomo/tracing.py new file mode 100644 index 00000000..0976cbf8 --- /dev/null +++ b/benchmarks/locomo/tracing.py @@ -0,0 +1,381 @@ +"""Optional Langfuse observer for LoCoMo answer and judge stages. + +This module is imported only after all three Langfuse environment bindings are +non-empty. It never receives rendered prompts, source chunks, tool results, gold +answers, or failure messages. +""" + +from __future__ import annotations + +from collections.abc import Iterator +from contextlib import AbstractContextManager +from contextlib import contextmanager +from importlib import import_module +import logging +import sys +from typing import Protocol +from typing import Self + +from rememberstack.model import ProviderCallUsage + +_logger = logging.getLogger(__name__) + + +class _Observation(Protocol): + """The observation operations used by the benchmark shim.""" + + def update(self, **values: object) -> Self: + """Attach bounded metadata to the observation.""" + ... + + def end(self) -> Self: + """End a manually created observation.""" + ... + + +class _LangfuseClient(Protocol): + """The dynamically loaded Langfuse client surface used here.""" + + def create_trace_id(self, *, seed: str | None = None) -> str: + """Create a W3C trace identifier, deterministically when seeded.""" + ... + + def start_as_current_observation( + self, **values: object + ) -> AbstractContextManager[_Observation]: + """Create the active root observation for one stage.""" + ... + + def start_observation(self, **values: object) -> _Observation: + """Create one child observation under the active question root.""" + ... + + def flush(self) -> None: + """Flush the stage's pending observations.""" + ... + + +class _LangfuseModule(Protocol): + """The optional module constructor used by the shim.""" + + def Langfuse(self, **values: object) -> _LangfuseClient: # noqa: N802 + """Build one explicitly configured client.""" + ... + + +def create_langfuse_tracer( + *, public_key: str, secret_key: str, host: str, run_identity: str +) -> LocomoTracer: + """Load Langfuse lazily and create one stage-shared tracing client.""" + module = _load_langfuse() + client = module.Langfuse( + public_key=public_key, secret_key=secret_key, base_url=host + ) + return LocomoTracer(client=client, run_identity=run_identity) + + +class LocomoTracer: + """Create deterministic per-question traces and flush them as a batch.""" + + def __init__(self, *, client: _LangfuseClient, run_identity: str) -> None: + """Bind a client and stable prepared-run identity.""" + self._client = client + self._run_identity = run_identity + + @contextmanager + def question( + self, *, item_id: str, question: str, stage: str + ) -> Iterator[QuestionTrace | None]: + """Open one best-effort answer/judge root on a deterministic trace.""" + try: + trace_id = self._client.create_trace_id( + seed=f"{self._run_identity}:{item_id}" + ) + manager = self._client.start_as_current_observation( + name=f"locomo.{stage}", + as_type="agent" if stage == "answer" else "span", + trace_context={"trace_id": trace_id}, + input={"question": question}, + metadata={"item_id": item_id, "stage": stage}, + ) + root = manager.__enter__() + except Exception: + _logger.warning( + "optional Langfuse question observation start failed", exc_info=True + ) + yield None + return + + trace = QuestionTrace(client=self._client, root=root) + try: + yield trace + finally: + exception_details = sys.exc_info() + try: + trace.end_started() + finally: + try: + # Ignore a truthy result: tracing must not suppress an exception + # raised by the answer or judge protocol. + manager.__exit__(*exception_details) + except Exception: + _logger.warning( + "optional Langfuse question observation end failed", + exc_info=True, + ) + + def flush(self) -> None: + """Best-effort flush observations at the end of a protocol stage.""" + try: + self._client.flush() + except Exception: + _logger.warning("optional Langfuse flush failed", exc_info=True) + + +class QuestionTrace: + """Bounded child-span writer for one question trace.""" + + def __init__(self, *, client: _LangfuseClient, root: _Observation) -> None: + """Bind the active question root.""" + self._client = client + self._root = root + self._agent_calls = 0 + self._started: list[ModelCall | ToolCall] = [] + + def start_agent_call(self, *, model: str) -> ModelCall | None: + """Best-effort start one generation without recording its prompt.""" + self._agent_calls += 1 + return self._start_model_call( + name="locomo.answer-agent", + model=model, + metadata={"call_index": self._agent_calls}, + ) + + def start_tool_call( + self, *, name: str, arguments: dict[str, object] + ) -> ToolCall | None: + """Best-effort start one content-free tool observation.""" + try: + observation = self._client.start_observation( + name="locomo.tool", + as_type="tool", + input=_arguments_summary(arguments=arguments), + metadata={"tool_name": name}, + ) + except Exception: + _logger.warning( + "optional Langfuse tool observation start failed", exc_info=True + ) + return None + call = ToolCall(observation=observation) + self._started.append(call) + return call + + def start_judge_call(self, *, model: str) -> ModelCall | None: + """Best-effort start one judge generation without its prompt.""" + return self._start_model_call(name="locomo.judge", model=model) + + def _start_model_call( + self, *, name: str, model: str, metadata: dict[str, object] | None = None + ) -> ModelCall | None: + """Start and retain one generation so the root can always close it.""" + try: + values: dict[str, object] = { + "name": name, + "as_type": "generation", + "model": model, + } + if metadata is not None: + values["metadata"] = metadata + observation = self._client.start_observation(**values) + except Exception: + _logger.warning( + "optional Langfuse model observation start failed", exc_info=True + ) + return None + call = ModelCall(observation=observation) + self._started.append(call) + return call + + def end_started(self) -> None: + """End every child observation, including unfinished error paths.""" + for observation in reversed(self._started): + try: + observation.end() + except Exception: + _logger.warning( + "optional Langfuse child observation cleanup failed", exc_info=True + ) + + def finish_answer( + self, *, final_answer: str | None, failure_kind: str | None + ) -> None: + """Best-effort finish the answer root with the allowed answer body.""" + values: dict[str, object] = { + "metadata": { + "outcome": "failed" if failure_kind is not None else "answered", + "failure_kind": failure_kind, + } + } + if final_answer is not None: + values["output"] = {"final_answer": final_answer} + self._update_root(values=values) + + def finish_judge( + self, *, final_answer: str | None, verdict: str, failure_kind: str | None + ) -> None: + """Best-effort finish the judge root with answer and bounded verdict.""" + values: dict[str, object] = { + "metadata": {"verdict": verdict, "failure_kind": failure_kind} + } + if final_answer is not None: + values["output"] = {"final_answer": final_answer} + self._update_root(values=values) + + def _update_root(self, *, values: dict[str, object]) -> None: + """Update the root without allowing the observer to affect protocol flow.""" + try: + self._root.update(**values) + except Exception: + _logger.warning( + "optional Langfuse question observation update failed", exc_info=True + ) + + +class ModelCall: + """Finish one generation with usage, cost, latency, and bounded outcome.""" + + def __init__(self, *, observation: _Observation) -> None: + """Retain the manual child observation.""" + self._observation = observation + self._ended = False + + def finish( + self, + *, + usage: ProviderCallUsage | None, + latency_ms: int, + outcome: str, + final_answer: str | None = None, + verdict: str | None = None, + ) -> None: + """Best-effort update and end the generation exactly once.""" + if self._ended: + return + try: + metadata: dict[str, object] = {"latency_ms": latency_ms, "outcome": outcome} + if usage is not None: + metadata["provider_latency_ms"] = usage.latency_ms + if verdict is not None: + metadata["verdict"] = verdict + values: dict[str, object] = {"metadata": metadata} + if usage is not None: + values.update(_usage_values(usage=usage)) + if final_answer is not None: + values["output"] = {"final_answer": final_answer} + self._observation.update(**values) + except Exception: + _logger.warning( + "optional Langfuse model observation update failed", exc_info=True + ) + finally: + self.end() + + def end(self) -> None: + """Best-effort end an unfinished generation exactly once.""" + if self._ended: + return + self._ended = True + try: + self._observation.end() + except Exception: + _logger.warning( + "optional Langfuse model observation end failed", exc_info=True + ) + + +class ToolCall: + """Finish one tool span without sending its response body.""" + + def __init__(self, *, observation: _Observation) -> None: + """Retain the manual child observation.""" + self._observation = observation + self._ended = False + + def finish(self, *, latency_ms: int, outcome: str) -> None: + """Best-effort finish one tool span with bounded metadata.""" + if self._ended: + return + try: + self._observation.update( + output={"outcome": outcome}, metadata={"latency_ms": latency_ms} + ) + except Exception: + _logger.warning( + "optional Langfuse tool observation update failed", exc_info=True + ) + finally: + self.end() + + def end(self) -> None: + """Best-effort end an unfinished tool span exactly once.""" + if self._ended: + return + self._ended = True + try: + self._observation.end() + except Exception: + _logger.warning( + "optional Langfuse tool observation end failed", exc_info=True + ) + + +def _usage_values(*, usage: ProviderCallUsage) -> dict[str, object]: + """Map provider accounting to Langfuse's generation fields.""" + values: dict[str, object] = { + "model": usage.model_name, + "usage_details": { + "input": usage.tokens_in, + "output": usage.tokens_out, + "total": usage.tokens_in + usage.tokens_out, + }, + } + values["cost_details"] = {"total": float(usage.cost_usd)} + return values + + +def _arguments_summary(*, arguments: dict[str, object]) -> dict[str, object]: + """Describe argument shapes without copying keys or values.""" + shapes = [ + _value_shape(value=value) + for _, value in sorted(arguments.items(), key=lambda item: item[0]) + ] + return {"argument_count": len(arguments), "value_shapes": shapes} + + +def _value_shape(*, value: object) -> dict[str, object]: + """Return a content-free type/size summary for one argument.""" + if isinstance(value, str): + return {"type": "string", "length": len(value)} + if isinstance(value, dict): + return {"type": "object", "field_count": len(value)} + if isinstance(value, (list, tuple)): + return {"type": "array", "length": len(value)} + if value is None: + return {"type": "null"} + if isinstance(value, bool): + return {"type": "boolean"} + if isinstance(value, (int, float)): + return {"type": "number"} + return {"type": "other"} + + +def _load_langfuse() -> _LangfuseModule: + """Resolve the optional dependency only after environment opt-in.""" + try: + module = import_module("langfuse") + except ModuleNotFoundError as error: + raise RuntimeError( + "Langfuse tracing requires rememberstack[observability]" + ) from error + return module # type: ignore[return-value] diff --git a/compose.yaml b/compose.yaml index 4ebba275..aa20fa57 100644 --- a/compose.yaml +++ b/compose.yaml @@ -24,6 +24,12 @@ x-app: &app REMEMBERSTACK_SELFHOST_DEPLOYMENT_SLUG: ${REMEMBERSTACK_SELFHOST_DEPLOYMENT_SLUG} REMEMBERSTACK_SELFHOST_DEPLOYMENT_NAME: ${REMEMBERSTACK_SELFHOST_DEPLOYMENT_NAME} REMEMBERSTACK_SELFHOST_API_PORT: "8000" + REMEMBERSTACK_SENTRY_DSN: ${REMEMBERSTACK_SENTRY_DSN:-} + REMEMBERSTACK_SENTRY_ENVIRONMENT: ${REMEMBERSTACK_SENTRY_ENVIRONMENT:-} + REMEMBERSTACK_SENTRY_SAMPLE_RATE: ${REMEMBERSTACK_SENTRY_SAMPLE_RATE:-} + LANGFUSE_PUBLIC_KEY: ${LANGFUSE_PUBLIC_KEY:-} + LANGFUSE_SECRET_KEY: ${LANGFUSE_SECRET_KEY:-} + LANGFUSE_HOST: ${LANGFUSE_HOST:-} # D79: STRUCTURER now names only the anchor-proposal fallback seat. REMEMBERSTACK_STRUCTURER_MODEL: ${REMEMBERSTACK_STRUCTURER_MODEL:-openai/gpt-5.6-luna} REMEMBERSTACK_SKELETON_CHECK_MODEL: ${REMEMBERSTACK_SKELETON_CHECK_MODEL:-z-ai/glm-4.7-flash} diff --git a/plan/designs/orchestration_design.md b/plan/designs/orchestration_design.md index 4edb4cfd..b2c6d95c 100644 --- a/plan/designs/orchestration_design.md +++ b/plan/designs/orchestration_design.md @@ -213,6 +213,8 @@ system: metric). The deployment operator or cloud product chooses collection, retention, dashboards, and alerts. +Instance provisioning — Langfuse self-host versus cloud, and GlitchTip versus Better Stack — is +an infrastructure decision tracked separately from the engine. Those consumers must derive their view from this state/telemetry rather than becoming another authority for pipeline truth (D60/D61). diff --git a/pyproject.toml b/pyproject.toml index e9b3ec7f..807e381c 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -62,6 +62,10 @@ benchmark = [ "nltk>=3.9", "regex>=2024.11.6", ] +observability = [ + "langfuse>=4.14.1,<5", + "sentry-sdk==2.66.1", +] [project.scripts] remember = "rememberstack.surfaces.cli:main" diff --git a/src/rememberstack/adapters/selfhost/__init__.py b/src/rememberstack/adapters/selfhost/__init__.py index 05608c5e..a3089c15 100644 --- a/src/rememberstack/adapters/selfhost/__init__.py +++ b/src/rememberstack/adapters/selfhost/__init__.py @@ -17,6 +17,7 @@ from rememberstack.adapters.selfhost.queue import SelfHostTaskQueue from rememberstack.adapters.selfhost.queue import SelfHostWorkerLoop from rememberstack.adapters.selfhost.queue import TokenBucket +from rememberstack.adapters.selfhost.telemetry import FanoutTelemetry from rememberstack.adapters.selfhost.telemetry import JsonLineTelemetry from rememberstack.adapters.selfhost.watcher import LocalDirectoryWatcher @@ -25,6 +26,7 @@ __all__ = ( "LanceChunkIndex", + "FanoutTelemetry", "LocalFSForgetManifestStore", "LocalGitRepository", "LocalDirectoryWatcher", diff --git a/src/rememberstack/adapters/selfhost/telemetry.py b/src/rememberstack/adapters/selfhost/telemetry.py index 8ca2b0f7..14adc677 100644 --- a/src/rememberstack/adapters/selfhost/telemetry.py +++ b/src/rememberstack/adapters/selfhost/telemetry.py @@ -1,12 +1,16 @@ """Simple JSON-lines telemetry for self-hosted process logs.""" import json +import logging import sys from threading import Lock import traceback from typing import TextIO from rememberstack.model import TelemetryEvent +from rememberstack.ports.telemetry import TelemetryPort + +_logger = logging.getLogger(__name__) class JsonLineTelemetry: @@ -43,3 +47,35 @@ def _write(self, *, payload: dict[str, object]) -> None: with self._lock: self._stream.write(line + "\n") self._stream.flush() + + +class FanoutTelemetry: + """Send events to one authoritative local sink and optional remote sinks.""" + + def __init__(self, *, sinks: tuple[TelemetryPort, ...]) -> None: + """Retain a non-empty set whose first sink is authoritative.""" + if not sinks: + raise ValueError("fanout telemetry requires at least one sink") + self._sinks = sinks + + def export_event(self, *, event: TelemetryEvent) -> None: + """Export locally, then isolate every optional sink failure.""" + self._sinks[0].export_event(event=event) + for sink in self._sinks[1:]: + try: + sink.export_event(event=event) + except Exception: + _logger.warning("optional telemetry event export failed", exc_info=True) + + def export_exception( + self, *, event: TelemetryEvent, exception: BaseException + ) -> None: + """Export locally, then isolate every optional sink failure.""" + self._sinks[0].export_exception(event=event, exception=exception) + for sink in self._sinks[1:]: + try: + sink.export_exception(event=event, exception=exception) + except Exception: + _logger.warning( + "optional telemetry exception export failed", exc_info=True + ) diff --git a/src/rememberstack/adapters/sentry.py b/src/rememberstack/adapters/sentry.py new file mode 100644 index 00000000..75a01e2f --- /dev/null +++ b/src/rememberstack/adapters/sentry.py @@ -0,0 +1,174 @@ +"""Opt-in, metadata-only Sentry-protocol telemetry. + +The vendor SDK is resolved dynamically only after a self-host entrypoint has +validated a non-empty DSN. Importing RememberStack never imports ``sentry_sdk``. +""" + +from __future__ import annotations + +from importlib import import_module +import logging +from threading import Lock +from typing import Any +from typing import Protocol +from typing import Self + +from rememberstack.model import TelemetryEvent + +_logger = logging.getLogger(__name__) + + +class _Scope(Protocol): + """The small Sentry scope surface used by the adapter.""" + + def __enter__(self) -> Self: + """Enter an event-local scope.""" + ... + + def __exit__( + self, + exception_type: type[BaseException] | None, + exception: BaseException | None, + traceback: object | None, + ) -> bool | None: + """Leave an event-local scope.""" + ... + + def set_tag(self, key: str, value: str) -> None: + """Attach one searchable metadata tag.""" + ... + + +class _SentrySdk(Protocol): + """The dynamically loaded Sentry SDK surface used here.""" + + def init(self, **options: object) -> object: + """Initialize the process-global SDK client.""" + ... + + def new_scope(self) -> _Scope: + """Create an isolated scope for one captured exception.""" + ... + + def capture_exception(self, error: BaseException) -> object: + """Capture one real exception object.""" + ... + + +_INITIALIZE_LOCK = Lock() +_INITIALIZED_SDK: _SentrySdk | None = None + + +def initialize_sentry( + *, dsn: str, environment: str, sample_rate: float, sdk: _SentrySdk | None = None +) -> SentryTelemetry: + """Initialize Sentry once and return its worker telemetry sink. + + Request bodies, breadcrumbs, local variables, user data, exception messages, + and ad-hoc extras are all removed. Events retain exception type/stack + metadata plus the explicit worker routing tags. + """ + global _INITIALIZED_SDK + + with _INITIALIZE_LOCK: + if _INITIALIZED_SDK is None: + resolved = sdk or _load_sdk() + resolved.init( + dsn=dsn, + environment=environment, + sample_rate=sample_rate, + traces_sample_rate=0.0, + send_default_pii=False, + include_local_variables=False, + max_request_body_size="never", + max_breadcrumbs=0, + before_send=_metadata_only_event, + ) + _INITIALIZED_SDK = resolved + return SentryTelemetry(sdk=_INITIALIZED_SDK) + + +class SentryTelemetry: + """Capture only worker exceptions; ordinary state telemetry stays local.""" + + def __init__(self, *, sdk: _SentrySdk) -> None: + """Bind the initialized SDK without exposing it to worker code.""" + self._sdk = sdk + + def export_event(self, *, event: TelemetryEvent) -> None: + """Ignore ordinary events; PostgreSQL and JSON telemetry remain authoritative.""" + + def export_exception( + self, *, event: TelemetryEvent, exception: BaseException + ) -> None: + """Capture the exception with only stage, lane, and processing tags.""" + attributes = {attribute.name: attribute.value for attribute in event.attributes} + try: + with self._sdk.new_scope() as scope: + for name in ("stage", "lane", "processing_id"): + value = attributes.get(name) + if value is not None: + scope.set_tag(name, str(value)) + self._sdk.capture_exception(exception) + except Exception: + _logger.warning("optional Sentry exception capture failed", exc_info=True) + + +def _load_sdk() -> _SentrySdk: + """Load the optional SDK only after a DSN opted the process in.""" + try: + module = import_module("sentry_sdk") + except ModuleNotFoundError as error: + raise RuntimeError( + "REMEMBERSTACK_SENTRY_DSN requires rememberstack[observability]" + ) from error + return module # type: ignore[return-value] + + +def _metadata_only_event( + event: dict[str, Any], hint: dict[str, Any] +) -> dict[str, Any] | None: + """Remove runtime text and user/request state just before transport.""" + log_record = hint.get("log_record") + if getattr(log_record, "name", None) == "rememberstack.workers.base": + # Worker failures are explicitly captured with route tags after their + # ledger transition. Drop LoggingIntegration's earlier untagged copy. + return None + # The current capture surface intentionally admits only SDK-owned envelope + # fields plus the exception type/stack scrubbed below. Every caller-owned + # top-level text field known to that surface is removed here; extending the + # integrations requires extending this list before the new data is enabled. + for key in ( + "breadcrumbs", + "extra", + "logentry", + "message", + "request", + "threads", + "user", + ): + event.pop(key, None) + exceptions = event.get("exception") + if isinstance(exceptions, dict): + values = exceptions.get("values") + if isinstance(values, list): + for value in values: + if not isinstance(value, dict): + continue + value["value"] = "[redacted]" + _strip_stack_runtime_text(value.get("stacktrace")) + return event + + +def _strip_stack_runtime_text(stacktrace: object) -> None: + """Keep stack coordinates while removing locals and source-line text.""" + if not isinstance(stacktrace, dict): + return + frames = stacktrace.get("frames") + if not isinstance(frames, list): + return + for frame in frames: + if not isinstance(frame, dict): + continue + for key in ("context_line", "post_context", "pre_context", "vars"): + frame.pop(key, None) diff --git a/src/rememberstack/profiles/selfhost.py b/src/rememberstack/profiles/selfhost.py index 4d7e28c5..b833f272 100644 --- a/src/rememberstack/profiles/selfhost.py +++ b/src/rememberstack/profiles/selfhost.py @@ -13,6 +13,7 @@ from alembic import command from alembic.config import Config from pydantic import Field +from pydantic import SecretStr from pydantic_settings import BaseSettings from pydantic_settings import SettingsConfigDict import sqlalchemy @@ -38,6 +39,7 @@ from fastapi import FastAPI from rememberstack.adapters.selfhost import SelfHostWorkerLoop + from rememberstack.ports.telemetry import TelemetryPort from rememberstack.workers import StageHandler _SUPPORTED_WORKER_STAGES = ( @@ -96,6 +98,25 @@ class SelfHostSettings(BaseSettings): worker_session_s: float = Field(default=3_600.0, gt=0) +class SentrySettings(BaseSettings): + """Strictly opt-in self-host error-tracking settings.""" + + model_config = SettingsConfigDict( + env_prefix="REMEMBERSTACK_SENTRY_", env_ignore_empty=True, extra="ignore" + ) + + dsn: SecretStr | None = None + environment: str | None = None + sample_rate: float = Field(default=1.0, ge=0.0, le=1.0) + + def configured_dsn(self) -> str | None: + """Return a non-empty DSN only when error tracking is explicitly enabled.""" + if self.dsn is None: + return None + value = self.dsn.get_secret_value().strip() + return value or None + + class _FreshDeploymentReadiness: """Fail closed if a fresh quickstart sees portable forget history. @@ -132,6 +153,7 @@ def __init__( corpusfs_store: MinIOObjectStore, snapshot_store: MinIOObjectStore, model_provider: OpenRouterModelProvider, + error_telemetry: TelemetryPort | None = None, ) -> None: """Retain one dependency graph for an API, setup, or worker process.""" self._settings = settings @@ -141,9 +163,10 @@ def __init__( self._corpusfs_store = corpusfs_store self._snapshot_store = snapshot_store self._model_provider = model_provider + self._error_telemetry = error_telemetry @classmethod - def from_settings(cls) -> Self: + def from_settings(cls, *, error_telemetry: TelemetryPort | None = None) -> Self: """Load every external value through its typed settings boundary.""" profile_settings = SelfHostSettings.model_validate({}) minio_settings = MinIOSettings.model_validate({}) @@ -167,6 +190,7 @@ def from_settings(cls) -> Self: model_provider=OpenRouterModelProvider( settings=OpenRouterSettings.model_validate({}) ), + error_telemetry=error_telemetry, ) def close(self) -> None: @@ -287,6 +311,7 @@ def healthz() -> dict[str, str]: def worker_loop(self, *, stage: PipelineStage) -> SelfHostWorkerLoop: """Build one continuous route's ordinary LISTEN/NOTIFY worker loop.""" + from rememberstack.adapters.selfhost import FanoutTelemetry from rememberstack.adapters.selfhost import JsonLineTelemetry from rememberstack.adapters.selfhost import SelfHostTaskQueue from rememberstack.adapters.selfhost import SelfHostWorkerLoop @@ -302,12 +327,18 @@ def worker_loop(self, *, stage: PipelineStage) -> SelfHostWorkerLoop: registry = HandlerRegistry() registry.register(stage=stage, handler=self._handler(stage=stage)) ledger = WorkLedger(engine=self._engine, settings=WorkLedgerSettings()) + local_telemetry = JsonLineTelemetry() + telemetry = ( + local_telemetry + if self._error_telemetry is None + else FanoutTelemetry(sinks=(local_telemetry, self._error_telemetry)) + ) return SelfHostWorkerLoop( worker=Worker( ledger=ledger, registry=registry, queue=SelfHostTaskQueue(ledger=ledger), - telemetry=JsonLineTelemetry(), + telemetry=telemetry, ), deployment_id=self._settings.deployment_id, stage=stage, @@ -501,7 +532,13 @@ def _handler(self, *, stage: PipelineStage) -> StageHandler: def create_api() -> FastAPI: - """Uvicorn factory for the self-host API process.""" + """Uvicorn factory that initializes process-global API error tracking. + + The API process has no worker telemetry fanout, so the returned Sentry sink + is intentionally unused after its process-global SDK initialization. + """ + settings = SelfHostSettings.model_validate({}) + _initialize_error_tracking(command="api", deployment_slug=settings.deployment_slug) return SelfHostProfile.from_settings().api() @@ -533,7 +570,10 @@ def main(argv: list[str] | None = None) -> int: access_log=True, ) return 0 - profile = SelfHostProfile.from_settings() + error_telemetry = _initialize_error_tracking( + command=args.command, deployment_slug=settings.deployment_slug + ) + profile = SelfHostProfile.from_settings(error_telemetry=error_telemetry) try: if args.command == "setup": profile.setup() @@ -547,6 +587,24 @@ def main(argv: list[str] | None = None) -> int: profile.close() +def _initialize_error_tracking( + *, command: str, deployment_slug: str +) -> TelemetryPort | None: + """Initialize the optional Sentry sink only for long-lived profile entrypoints.""" + if command not in {"api", "setup", "worker"}: + return None + settings = SentrySettings.model_validate({}) + dsn = settings.configured_dsn() + if dsn is None: + return None + from rememberstack.adapters.sentry import initialize_sentry + + environment = (settings.environment or "").strip() or deployment_slug + return initialize_sentry( + dsn=dsn, environment=environment, sample_rate=settings.sample_rate + ) + + def _psycopg_url() -> str: """Remove SQLAlchemy's driver suffix for psycopg's native connection parser.""" url = make_url(load_database_settings().sqlalchemy_url()) diff --git a/src/tests/adapters/test_sentry.py b/src/tests/adapters/test_sentry.py new file mode 100644 index 00000000..fe40e6a9 --- /dev/null +++ b/src/tests/adapters/test_sentry.py @@ -0,0 +1,179 @@ +"""Opt-in Sentry initialization, redaction, and worker-tag proofs.""" + +from __future__ import annotations + +from datetime import datetime +from datetime import UTC +import json +import logging +from typing import Self +from uuid import UUID + +import pytest + +from rememberstack.model import TelemetryAttribute +from rememberstack.model import TelemetryEvent + + +class _FakeScope: + """Record the event-local tags attached before capture.""" + + def __init__(self, *, sdk: _FakeSentrySdk) -> None: + self._sdk = sdk + self.tags: dict[str, str] = {} + + def __enter__(self) -> Self: + self._sdk.scopes.append(self) + return self + + def __exit__( + self, + exception_type: type[BaseException] | None, + exception: BaseException | None, + traceback: object | None, + ) -> None: + del exception_type, exception, traceback + + def set_tag(self, key: str, value: str) -> None: + self.tags[key] = value + + +class _FakeSentrySdk: + """In-memory stand-in for the optional SDK transport.""" + + def __init__(self) -> None: + self.init_calls: list[dict[str, object]] = [] + self.scopes: list[_FakeScope] = [] + self.exceptions: list[BaseException] = [] + + def init(self, **options: object) -> object: + self.init_calls.append(options) + return object() + + def new_scope(self) -> _FakeScope: + return _FakeScope(sdk=self) + + def capture_exception(self, error: BaseException) -> object: + self.exceptions.append(error) + return object() + + +def test_sentry_initializes_once_and_captures_only_worker_metadata() -> None: + """A configured sink keeps tags and the real exception while stripping text.""" + from rememberstack.adapters import sentry as sentry_adapter + + sentry_adapter._INITIALIZED_SDK = None + sdk = _FakeSentrySdk() + telemetry = sentry_adapter.initialize_sentry( + dsn="https://public@example.test/1", + environment="deployment-slug", + sample_rate=0.25, + sdk=sdk, + ) + sentry_adapter.initialize_sentry( + dsn="https://other.example.test/2", + environment="ignored", + sample_rate=1.0, + sdk=sdk, + ) + + processing_id = UUID("62000000-0000-0000-0000-000000000001") + event = TelemetryEvent( + name="worker.run", + occurred_at=datetime(2026, 7, 29, tzinfo=UTC), + attributes=( + TelemetryAttribute(name="stage", value="extract_claims"), + TelemetryAttribute(name="lane", value="steady"), + TelemetryAttribute(name="processing_id", value=str(processing_id)), + TelemetryAttribute(name="outcome", value="retry_scheduled"), + ), + ) + exception = RuntimeError("private prompt and completion body") + telemetry.export_exception(event=event, exception=exception) + + assert len(sdk.init_calls) == 1 + options = sdk.init_calls[0] + assert options["environment"] == "deployment-slug" + assert options["sample_rate"] == 0.25 + assert options["traces_sample_rate"] == 0.0 + assert options["send_default_pii"] is False + assert options["include_local_variables"] is False + assert options["max_request_body_size"] == "never" + assert options["max_breadcrumbs"] == 0 + assert sdk.exceptions == [exception] + assert sdk.scopes[-1].tags == { + "stage": "extract_claims", + "lane": "steady", + "processing_id": str(processing_id), + } + + before_send = options["before_send"] + assert callable(before_send) + scrubbed = before_send( + { + "message": "private prompt", + "request": {"data": "private completion"}, + "breadcrumbs": {"values": [{"message": "private chunk"}]}, + "extra": {"prompt": "private prompt"}, + "user": {"email": "person@example.test"}, + "exception": { + "values": [ + { + "type": "RuntimeError", + "value": "private completion", + "stacktrace": { + "frames": [ + { + "filename": "worker.py", + "function": "handle", + "lineno": 10, + "context_line": "prompt = private", + "pre_context": ["private"], + "post_context": ["private"], + "vars": {"prompt": "private"}, + } + ] + }, + } + ] + }, + }, + {}, + ) + assert scrubbed is not None + encoded = json.dumps(scrubbed) + assert "private" not in encoded + assert "person@example.test" not in encoded + assert '"type": "RuntimeError"' in encoded + assert '"filename": "worker.py"' in encoded + + class _WorkerLog: + name = "rememberstack.workers.base" + + assert before_send({"message": "duplicate"}, {"log_record": _WorkerLog()}) is None + + +def test_sentry_capture_failure_is_logged_and_suppressed( + caplog: pytest.LogCaptureFixture, +) -> None: + """A remote capture failure cannot replace the worker's recorded outcome.""" + from rememberstack.adapters.sentry import SentryTelemetry + + class _RaisingSentrySdk(_FakeSentrySdk): + def capture_exception(self, error: BaseException) -> object: + del error + raise RuntimeError("remote capture unavailable") + + caplog.set_level(logging.WARNING, logger="rememberstack.adapters.sentry") + telemetry = SentryTelemetry(sdk=_RaisingSentrySdk()) + + telemetry.export_exception( + event=TelemetryEvent( + name="worker.run", + occurred_at=datetime(2026, 7, 29, tzinfo=UTC), + attributes=(), + ), + exception=RuntimeError("authoritative worker failure"), + ) + + assert "optional Sentry exception capture failed" in caplog.text diff --git a/src/tests/adapters/test_telemetry.py b/src/tests/adapters/test_telemetry.py index b3c43e43..8cbd10f2 100644 --- a/src/tests/adapters/test_telemetry.py +++ b/src/tests/adapters/test_telemetry.py @@ -4,10 +4,27 @@ from datetime import UTC from io import StringIO import json +from typing import cast +from uuid import UUID +from rememberstack.adapters.selfhost import FanoutTelemetry from rememberstack.adapters.selfhost import JsonLineTelemetry +from rememberstack.model import ClaimedWork +from rememberstack.model import NonRetryableHandlerError +from rememberstack.model import PipelineStage +from rememberstack.model import ProcessingLane +from rememberstack.model import ProcessingTarget +from rememberstack.model import RunResultOutcome from rememberstack.model import TelemetryAttribute from rememberstack.model import TelemetryEvent +from rememberstack.spine import WorkLedger +from rememberstack.workers import HandlerOutcome +from rememberstack.workers import HandlerRegistry +from rememberstack.workers import RunResult +from rememberstack.workers import Worker + +_DEPLOYMENT_ID = UUID("74000000-0000-0000-0000-000000000001") +_PROCESSING_ID = UUID("74000000-0000-0000-0000-000000000002") def _event() -> TelemetryEvent: @@ -46,3 +63,93 @@ def test_json_lines_flushes_one_event_per_line() -> None: "name": "outcome", "value": "dead_lettered", } + + +class _MemoryLedger: + """Small row-backed ledger double for the worker's committed failure path.""" + + def __init__(self) -> None: + self.rows: list[dict[str, object]] = [ + { + "processing_id": _PROCESSING_ID, + "status": "pending", + "attempts": 0, + "last_error": None, + } + ] + self._claimed = False + + def claim_one(self, **_: object) -> ClaimedWork | None: + if self._claimed: + return None + self._claimed = True + self.rows[0]["status"] = "running" + self.rows[0]["attempts"] = 1 + return ClaimedWork( + processing_id=_PROCESSING_ID, + deployment_id=_DEPLOYMENT_ID, + target_kind=ProcessingTarget.DOCUMENT, + target_id=UUID("74000000-0000-0000-0000-000000000003"), + stage=PipelineStage.CONVERT, + component_version="convert-test", + content_hash="content-test", + lane=ProcessingLane.STEADY, + attempt=1, + payload=None, + ) + + def fail( + self, *, processing_id: UUID, error: str, retryable: bool + ) -> datetime | None: + assert processing_id == _PROCESSING_ID + assert retryable is False + self.rows[0]["status"] = "dead_letter" + self.rows[0]["last_error"] = error + return None + + +class _PermanentFailure: + def handle(self, *, work: ClaimedWork, meter: object) -> HandlerOutcome: + del work, meter + raise NonRetryableHandlerError("authoritative handler failure") + + +class _RaisingTelemetry: + def export_event(self, **_: object) -> None: + raise RuntimeError("optional sink unavailable") + + def export_exception(self, **_: object) -> None: + raise RuntimeError("optional sink unavailable") + + +def test_worker_result_and_ledger_survive_raising_optional_fanout_sink() -> None: + """A second sink cannot alter the result after the authoritative row commit.""" + ledger = _MemoryLedger() + registry = HandlerRegistry() + registry.register(stage=PipelineStage.CONVERT, handler=_PermanentFailure()) + stream = StringIO() + telemetry = FanoutTelemetry( + sinks=(JsonLineTelemetry(stream=stream), _RaisingTelemetry()) + ) + + result = Worker( + ledger=cast(WorkLedger, ledger), registry=registry, telemetry=telemetry + ).run_one( + deployment_id=_DEPLOYMENT_ID, + stage=PipelineStage.CONVERT, + lane=ProcessingLane.STEADY, + ) + + assert result == RunResult( + processing_id=_PROCESSING_ID, outcome=RunResultOutcome.DEAD_LETTERED + ) + ledger_row = ledger.rows[0] + assert ledger_row["processing_id"] == _PROCESSING_ID + assert ledger_row["status"] == "dead_letter" + assert ledger_row["attempts"] == 1 + assert "NonRetryableHandlerError: authoritative handler failure" in str( + ledger_row["last_error"] + ) + local_row = json.loads(stream.getvalue()) + assert local_row["attributes"][5] == {"name": "outcome", "value": "dead_lettered"} + assert local_row["exception"]["message"] == "authoritative handler failure" diff --git a/src/tests/benchmarks/test_locomo_runner.py b/src/tests/benchmarks/test_locomo_runner.py index 16324f2b..826ab843 100644 --- a/src/tests/benchmarks/test_locomo_runner.py +++ b/src/tests/benchmarks/test_locomo_runner.py @@ -8,6 +8,9 @@ import hashlib import json from pathlib import Path +import sys +from types import ModuleType +from typing import Self from typing import TypeVar from uuid import UUID @@ -239,6 +242,311 @@ def test_staged_mock_run_checks_readiness_and_resumes( assert summary.answer_agent_calls == 2 +class _FakeLangfuseObservation: + """Capture one fake observation's updates and end state.""" + + def __init__(self, *, started: dict[str, object]) -> None: + self.started = started + self.updates: list[dict[str, object]] = [] + self.ended = False + + def update(self, **values: object) -> Self: + self.updates.append(values) + return self + + def end(self) -> Self: + self.ended = True + return self + + +class _FakeLangfuseContext: + """Context manager for a fake root observation.""" + + def __init__(self, *, observation: _FakeLangfuseObservation) -> None: + self._observation = observation + + def __enter__(self) -> _FakeLangfuseObservation: + return self._observation + + def __exit__( + self, + exception_type: type[BaseException] | None, + exception: BaseException | None, + traceback: object | None, + ) -> None: + del exception_type, exception, traceback + self._observation.end() + + +class _FakeLangfuseClient: + """Fake transport exposing the Langfuse tracing surface used by the shim.""" + + def __init__(self) -> None: + self.observations: list[_FakeLangfuseObservation] = [] + self.flushes = 0 + + def create_trace_id(self, *, seed: str | None = None) -> str: + return hashlib.sha256((seed or "").encode()).hexdigest()[:32] + + def start_as_current_observation(self, **values: object) -> _FakeLangfuseContext: + observation = self._start(values=values) + return _FakeLangfuseContext(observation=observation) + + def start_observation(self, **values: object) -> _FakeLangfuseObservation: + return self._start(values=values) + + def flush(self) -> None: + self.flushes += 1 + + def _start(self, *, values: dict[str, object]) -> _FakeLangfuseObservation: + observation = _FakeLangfuseObservation(started=values) + self.observations.append(observation) + return observation + + +def test_langfuse_fake_transport_is_observer_only_and_content_bounded( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch +) -> None: + """Configured traces emit every call while persisted protocol output is identical.""" + _patch_prepared_inputs(monkeypatch=monkeypatch) + monkeypatch.setattr(runner, "_elapsed_ms", lambda _started: 1) + for name in ("LANGFUSE_PUBLIC_KEY", "LANGFUSE_SECRET_KEY", "LANGFUSE_HOST"): + monkeypatch.delenv(name, raising=False) + monkeypatch.delitem(sys.modules, "benchmarks.locomo.tracing", raising=False) + monkeypatch.delitem(sys.modules, "langfuse", raising=False) + + plain_dir = tmp_path / "plain" + traced_dir = tmp_path / "traced" + for run_dir in (plain_dir, traced_dir): + prepare_run( + dataset_path=tmp_path / "synthetic.json", tier="smoke", output=run_dir + ) + + raw_clients: list[httpx.Client] = [] + + def execute_run(*, run_dir: Path, provider: FakeModelProvider) -> None: + raw_client = httpx.Client( + base_url="http://memory.test", transport=httpx.MockTransport(_run_transport) + ) + raw_clients.append(raw_client) + client = MemoryClient(client=raw_client) + ingest_sample( + run_dir=run_dir, + sample_id="conv-test", + max_documents=1, + execute=True, + isolated_deployment_confirmation="conv-test", + client=client, + provider=_PreflightProvider(), + ) + answer_sample( + run_dir=run_dir, + sample_id="conv-test", + max_questions=1, + max_agent_calls=9, + max_evaluator_cost_usd=Decimal("1"), + execute=True, + client=client, + provider=provider, + ) + judge_sample( + run_dir=run_dir, + sample_id="conv-test", + max_judge_calls=1, + max_evaluator_cost_usd=Decimal("1"), + execute=True, + provider=provider, + ) + + try: + execute_run( + run_dir=plain_dir, + provider=FakeModelProvider(generate_router=_private_tool_answer_and_judge), + ) + immutable_before = { + name: (traced_dir / name).read_bytes() + for name in ("run.json", "manifest.json", "documents.json") + } + fake_client = _FakeLangfuseClient() + constructor_calls: list[dict[str, object]] = [] + module = ModuleType("langfuse") + + def construct_langfuse(**values: object) -> _FakeLangfuseClient: + constructor_calls.append(values) + return fake_client + + module.Langfuse = construct_langfuse # type: ignore[attr-defined] + monkeypatch.setitem(sys.modules, "langfuse", module) + monkeypatch.setenv("LANGFUSE_PUBLIC_KEY", "public-test-key") + monkeypatch.setenv("LANGFUSE_SECRET_KEY", "secret-test-key") + monkeypatch.setenv("LANGFUSE_HOST", "https://langfuse.test") + execute_run( + run_dir=traced_dir, + provider=FakeModelProvider(generate_router=_private_tool_answer_and_judge), + ) + finally: + for raw_client in raw_clients: + raw_client.close() + + assert json.loads((plain_dir / "state.json").read_text()) == json.loads( + (traced_dir / "state.json").read_text() + ) + assert immutable_before == { + name: (traced_dir / name).read_bytes() for name in immutable_before + } + assert len(constructor_calls) == 2 + assert fake_client.flushes == 2 + names = [observation.started["name"] for observation in fake_client.observations] + assert names == [ + "locomo.answer", + "locomo.answer-agent", + "locomo.tool", + "locomo.answer-agent", + "locomo.judge", + "locomo.judge", + ] + roots = [ + observation + for observation in fake_client.observations + if observation.started["name"] in {"locomo.answer", "locomo.judge"} + and "trace_context" in observation.started + ] + assert len(roots) == 2 + assert roots[0].started["trace_context"] == roots[1].started["trace_context"] + assert all(observation.ended for observation in fake_client.observations) + wire = json.dumps( + [ + {"started": observation.started, "updates": observation.updates} + for observation in fake_client.observations + ], + sort_keys=True, + ) + assert "Alpha lives in Prague." not in wire + assert "PRIVATE_TOOL_ARGUMENT_BODY" not in wire + assert "TOOL TRACE SO FAR" not in wire + assert "Gold answer" not in wire + assert '"question": "Where?"' in wire + assert '"final_answer": "Prague"' in wire + assert '"verdict": "CORRECT"' in wire + assert '"usage_details"' in wire + assert '"cost_details"' in wire + + +class _RaisingLangfuseObservation(_FakeLangfuseObservation): + """Record cleanup attempts while raising from update and end.""" + + def update(self, **values: object) -> Self: + super().update(**values) + raise RuntimeError("observation update unavailable") + + def end(self) -> Self: + self.ended = True + raise RuntimeError("observation end unavailable") + + +class _RaisingLangfuseClient(_FakeLangfuseClient): + """Fail across start, finish, context cleanup, and flush lifecycle points.""" + + def __init__(self) -> None: + super().__init__() + self.trace_ids = 0 + + def create_trace_id(self, *, seed: str | None = None) -> str: + self.trace_ids += 1 + if self.trace_ids > 1: + raise RuntimeError("observation start unavailable") + return super().create_trace_id(seed=seed) + + def flush(self) -> None: + self.flushes += 1 + raise RuntimeError("observation flush unavailable") + + def _start(self, *, values: dict[str, object]) -> _FakeLangfuseObservation: + observation = _RaisingLangfuseObservation(started=values) + self.observations.append(observation) + return observation + + +def test_raising_langfuse_lifecycle_cannot_change_outputs_or_state( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch +) -> None: + """Answer/judge outputs and checkpoints ignore every observer failure.""" + _patch_prepared_inputs(monkeypatch=monkeypatch) + monkeypatch.setattr(runner, "_elapsed_ms", lambda _started: 1) + for name in ("LANGFUSE_PUBLIC_KEY", "LANGFUSE_SECRET_KEY", "LANGFUSE_HOST"): + monkeypatch.delenv(name, raising=False) + monkeypatch.delitem(sys.modules, "benchmarks.locomo.tracing", raising=False) + monkeypatch.delitem(sys.modules, "langfuse", raising=False) + + plain_dir = tmp_path / "plain-errors" + traced_dir = tmp_path / "traced-errors" + for run_dir in (plain_dir, traced_dir): + prepare_run( + dataset_path=tmp_path / "synthetic.json", tier="smoke", output=run_dir + ) + + raw_clients: list[httpx.Client] = [] + + def execute_run(*, run_dir: Path) -> tuple[tuple[object, ...], tuple[object, ...]]: + raw_client = httpx.Client( + base_url="http://memory.test", transport=httpx.MockTransport(_run_transport) + ) + raw_clients.append(raw_client) + client = MemoryClient(client=raw_client) + provider = FakeModelProvider(generate_router=_private_tool_answer_and_judge) + ingest_sample( + run_dir=run_dir, + sample_id="conv-test", + max_documents=1, + execute=True, + isolated_deployment_confirmation="conv-test", + client=client, + provider=_PreflightProvider(), + ) + answers = answer_sample( + run_dir=run_dir, + sample_id="conv-test", + max_questions=1, + max_agent_calls=9, + max_evaluator_cost_usd=Decimal("1"), + execute=True, + client=client, + provider=provider, + ) + judges = judge_sample( + run_dir=run_dir, + sample_id="conv-test", + max_judge_calls=1, + max_evaluator_cost_usd=Decimal("1"), + execute=True, + provider=provider, + ) + return answers, judges + + try: + plain_outputs = execute_run(run_dir=plain_dir) + raising_client = _RaisingLangfuseClient() + module = ModuleType("langfuse") + module.Langfuse = lambda **_: raising_client # type: ignore[attr-defined] + monkeypatch.setitem(sys.modules, "langfuse", module) + monkeypatch.setenv("LANGFUSE_PUBLIC_KEY", "public-test-key") + monkeypatch.setenv("LANGFUSE_SECRET_KEY", "secret-test-key") + monkeypatch.setenv("LANGFUSE_HOST", "https://langfuse.test") + traced_outputs = execute_run(run_dir=traced_dir) + finally: + for raw_client in raw_clients: + raw_client.close() + + assert traced_outputs == plain_outputs + assert json.loads((traced_dir / "state.json").read_text()) == json.loads( + (plain_dir / "state.json").read_text() + ) + assert raising_client.flushes == 2 + assert raising_client.trace_ids == 2 + assert raising_client.observations + assert all(observation.ended for observation in raising_client.observations) + + def test_readiness_flag_cannot_hide_an_incomplete_pipeline_report( tmp_path: Path, monkeypatch: pytest.MonkeyPatch ) -> None: @@ -393,6 +701,25 @@ def _tool_answer_and_judge(prompt: str, type_name: str) -> dict[str, object]: return _tool_then_answer(prompt, type_name) +def _private_tool_answer_and_judge(prompt: str, type_name: str) -> dict[str, object]: + """Use a sentinel argument body that tracing must summarize, never copy.""" + if type_name == "JudgeOutput": + return {"label": "CORRECT"} + if "TOOL TRACE SO FAR:\n[]" in prompt: + return { + "action": "tool", + "tool_name": "claims_verbatim", + "arguments_json": '{"query": "PRIVATE_TOOL_ARGUMENT_BODY"}', + "answer": None, + } + return { + "action": "answer", + "tool_name": None, + "arguments_json": "{}", + "answer": "Prague", + } + + def _question() -> LoCoMoQuestion: return LoCoMoQuestion( item_id="conv-test/qa/0000", diff --git a/src/tests/packaging/test_client_wheel.py b/src/tests/packaging/test_client_wheel.py index 698e1a0c..939fb190 100644 --- a/src/tests/packaging/test_client_wheel.py +++ b/src/tests/packaging/test_client_wheel.py @@ -195,6 +195,7 @@ def _assert_dependency_split(*, wheel: Path) -> None: "benchmark", "connectors-watched-directory", "k", + "observability", "server", } assert any(requirement.startswith("sqlalchemy") for requirement in requirements) diff --git a/src/tests/profiles/test_selfhost_profile.py b/src/tests/profiles/test_selfhost_profile.py index 37a2add7..f27ea7f9 100644 --- a/src/tests/profiles/test_selfhost_profile.py +++ b/src/tests/profiles/test_selfhost_profile.py @@ -2,11 +2,14 @@ from pathlib import Path import re +import subprocess +import sys import pytest from rememberstack.model import PipelineStage from rememberstack.profiles.selfhost import _expected_components +from rememberstack.profiles.selfhost import _initialize_error_tracking from rememberstack.profiles.selfhost import _model_bindings from rememberstack.profiles.selfhost import _SUPPORTED_WORKER_STAGES @@ -55,6 +58,84 @@ def test_compose_wires_the_exact_supported_worker_set_and_projection_job() -> No assert composed_stages == _SUPPORTED_WORKER_STAGES assert 'profiles: ["operations"]' in compose assert 'command: ["project", "--plane", "all"]' in compose + for name in ( + "REMEMBERSTACK_SENTRY_DSN", + "REMEMBERSTACK_SENTRY_ENVIRONMENT", + "REMEMBERSTACK_SENTRY_SAMPLE_RATE", + "LANGFUSE_PUBLIC_KEY", + "LANGFUSE_SECRET_KEY", + "LANGFUSE_HOST", + ): + assert f"{name}: ${{{name}:-}}" in compose + + +def test_observability_imports_are_absent_without_environment_opt_in( + monkeypatch: pytest.MonkeyPatch, +) -> None: + """Default profile/benchmark imports do not load either optional SDK or shim.""" + for name in ( + "REMEMBERSTACK_SENTRY_DSN", + "REMEMBERSTACK_SENTRY_ENVIRONMENT", + "REMEMBERSTACK_SENTRY_SAMPLE_RATE", + "LANGFUSE_PUBLIC_KEY", + "LANGFUSE_SECRET_KEY", + "LANGFUSE_HOST", + ): + monkeypatch.delenv(name, raising=False) + result = subprocess.run( + [ + sys.executable, + "-c", + ( + "import sys;" + " import rememberstack.profiles.selfhost;" + " import benchmarks.locomo.runner;" + " forbidden = {" + "'sentry_sdk', 'langfuse'," + " 'rememberstack.adapters.sentry'," + " 'benchmarks.locomo.tracing'};" + " loaded = forbidden.intersection(sys.modules);" + " assert not loaded, sorted(loaded)" + ), + ], + cwd=_ROOT, + check=False, + capture_output=True, + text=True, + ) + assert result.returncode == 0, result.stderr + + +def test_sentry_environment_defaults_to_deployment_slug( + monkeypatch: pytest.MonkeyPatch, +) -> None: + """Self-host startup passes the deployment slug and default sample rate.""" + from rememberstack.adapters import sentry as sentry_adapter + + marker = object() + calls: list[dict[str, object]] = [] + + def initialize(**values: object) -> object: + calls.append(values) + return marker + + monkeypatch.setattr(sentry_adapter, "initialize_sentry", initialize) + monkeypatch.setenv("REMEMBERSTACK_SENTRY_DSN", "https://public@example.test/1") + monkeypatch.delenv("REMEMBERSTACK_SENTRY_ENVIRONMENT", raising=False) + monkeypatch.setenv("REMEMBERSTACK_SENTRY_SAMPLE_RATE", "") + + telemetry = _initialize_error_tracking( + command="worker", deployment_slug="customer-memory" + ) + + assert telemetry is marker + assert calls == [ + { + "dsn": "https://public@example.test/1", + "environment": "customer-memory", + "sample_rate": 1.0, + } + ] def test_compose_forwards_the_dedicated_summary_seat( diff --git a/uv.lock b/uv.lock index 352b3d78..f752e368 100644 --- a/uv.lock +++ b/uv.lock @@ -54,6 +54,15 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/da/35/f2287558c17e29fafc8ef3daf819bb9834061cfa43bff8014f7df7f63bdc/anyio-4.14.2-py3-none-any.whl", hash = "sha256:9f505dda5ac9f0c8309b5e8bd445a8c2bf7246f3ce950121e45ea15bc41d1494", size = 125813 }, ] +[[package]] +name = "backoff" +version = "2.2.1" +source = { registry = "https://pypi.org/simple" } +sdist = { url = "https://files.pythonhosted.org/packages/47/d7/5bbeb12c44d7c4f2fb5b56abce497eb5ed9f34d85701de869acedd602619/backoff-2.2.1.tar.gz", hash = "sha256:03f829f5bb1923180821643f8753b0502c3b682293992485b0eef2807afa5cba", size = 17001 } +wheels = [ + { url = "https://files.pythonhosted.org/packages/df/73/b6e24bd22e6720ca8ee9a85a0c4a2971af8497d8f3193fa05390cbd46e09/backoff-2.2.1-py3-none-any.whl", hash = "sha256:63579f9a0628e06278f7e47b7d7d5b6ce20dc65c5e96a6f3ca99a6adca0396e8", size = 15148 }, +] + [[package]] name = "beautifulsoup4" version = "4.15.0" @@ -300,6 +309,18 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/e8/2d/d2a548598be01649e2d46231d151a6c56d10b964d94043a335ae56ea2d92/flatbuffers-25.12.19-py2.py3-none-any.whl", hash = "sha256:7634f50c427838bb021c2d66a3d1168e9d199b0607e6329399f04846d42e20b4", size = 26661 }, ] +[[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 } +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 }, +] + [[package]] name = "greenlet" version = "3.5.3" @@ -598,6 +619,25 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/d9/5d/8ca165f1386caf6c4d1c515afd52f345b66432264eecfdfb7fd33eefd9af/lancedb-0.34.0-cp39-abi3-win_amd64.whl", hash = "sha256:51cbc11808f9e3332819b9367c975b3a888541447a8e7bea09c57c852a279153", size = 63530726 }, ] +[[package]] +name = "langfuse" +version = "4.14.1" +source = { registry = "https://pypi.org/simple" } +dependencies = [ + { name = "backoff" }, + { name = "httpx" }, + { name = "opentelemetry-api" }, + { name = "opentelemetry-exporter-otlp-proto-http" }, + { name = "opentelemetry-sdk" }, + { name = "packaging" }, + { name = "pydantic" }, + { name = "wrapt" }, +] +sdist = { url = "https://files.pythonhosted.org/packages/a1/61/bb7fa00f5334a670c9aae1f053bc725654a0ca5435fa17d0cb103df8c978/langfuse-4.14.1.tar.gz", hash = "sha256:576641820ae79aeca71453c8eff3c8c35d53f7e41729c918f59eb81f7499faf9", size = 381977 } +wheels = [ + { url = "https://files.pythonhosted.org/packages/e6/e4/6c305af16f1ce344271dd829383b7af586530cf60b55af1aaf274ca37ba9/langfuse-4.14.1-py3-none-any.whl", hash = "sha256:07d19f16338b8e21f8e5996b7e6c3ed150ee582fbaa6275ac9eeea297093f4be", size = 669448 }, +] + [[package]] name = "magika" version = "0.6.2" @@ -876,6 +916,87 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/b7/f6/2bac21f722aa45d876d4a51f26bd0ef30e704068a3cd5021a5a7cd784271/onnxruntime-1.27.0-cp314-cp314t-manylinux_2_27_x86_64.manylinux_2_28_x86_64.whl", hash = "sha256:370d211e1ceeac4cd5f45301655463ac59e27cdc74d9f7aeb2d19ff4b7a76715", size = 18670781 }, ] +[[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 } +wheels = [ + { url = "https://files.pythonhosted.org/packages/ca/6f/a04e900f465ff3221ccc395522503e2d10e79fa21f2723c8e177aae1e0d1/opentelemetry_api-1.44.0-py3-none-any.whl", hash = "sha256:94b98c893a91b88657eaac1e3ba89618cdb85be6918196705354f34728b2cdef", size = 60018 }, +] + +[[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 } +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 }, +] + +[[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 } +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 }, +] + +[[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 } +wheels = [ + { url = "https://files.pythonhosted.org/packages/d1/7c/8be563d68e93bbefa5c8affb82ddcff91b3ad858ce49957ba7b16fd3e0ab/opentelemetry_proto-1.44.0-py3-none-any.whl", hash = "sha256:898b155a0e1557afd867478fb6158e8122a46329ca0bb8dc53cc55e98f017f56", size = 72483 }, +] + +[[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 } +wheels = [ + { url = "https://files.pythonhosted.org/packages/e7/23/ff077e61886ee020a17ce9c8b6fa11c601c8d8345b09ea24f605445df62a/opentelemetry_sdk-1.44.0-py3-none-any.whl", hash = "sha256:df081c4c6bcfdb1211e3e86140376792643128a25f8d72d1d27675936e7e96ad", size = 137221 }, +] + +[[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 } +wheels = [ + { url = "https://files.pythonhosted.org/packages/a6/0e/49df70d9b81fb5cbae4bbf2a49d865b09bcbcbc4eb53f5851b1027738d78/opentelemetry_semantic_conventions-0.65b0-py3-none-any.whl", hash = "sha256:1cacde7b0ad306f84c5ef08c3dbe1bbaf20165bba6f8bff43b670e555a086bcb", size = 204645 }, +] + [[package]] name = "packaging" version = "26.2" @@ -1300,6 +1421,10 @@ benchmark = [ { name = "nltk" }, { name = "regex" }, ] +observability = [ + { name = "langfuse" }, + { name = "sentry-sdk" }, +] server = [ { name = "alembic" }, { name = "boto3" }, @@ -1345,6 +1470,7 @@ requires-dist = [ { name = "httpx", specifier = ">=0.28.1" }, { name = "ladybug", marker = "extra == 'server'", specifier = ">=0.18.2" }, { name = "lancedb", marker = "extra == 'server'", specifier = ">=0.34.0" }, + { name = "langfuse", marker = "extra == 'observability'", specifier = ">=4.14.1,<5" }, { name = "markdown-it-py", marker = "extra == 'server'", specifier = ">=4.2.0" }, { name = "markitdown", marker = "extra == 'server'", specifier = ">=0.1.6" }, { name = "nltk", marker = "extra == 'benchmark'", specifier = ">=3.9" }, @@ -1353,6 +1479,7 @@ requires-dist = [ { name = "pydantic", specifier = ">=2.11" }, { name = "pydantic-settings", specifier = ">=2.10" }, { name = "regex", marker = "extra == 'benchmark'", specifier = ">=2024.11.6" }, + { name = "sentry-sdk", marker = "extra == 'observability'", specifier = "==2.66.1" }, { name = "sqlalchemy", marker = "extra == 'server'", specifier = ">=2.0" }, { name = "uvicorn", marker = "extra == 'server'", specifier = ">=0.37" }, ] @@ -1445,6 +1572,19 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/24/23/e84c64ad0e8bc59cd1b2ef98def848deff0ef3456c542afe74d51e9e8c85/s3transfer-0.19.1-py3-none-any.whl", hash = "sha256:d5fd7005ee39307455ad5f310b5ea67f4b1960d7fed5b3671ee50c249de675de", size = 90072 }, ] +[[package]] +name = "sentry-sdk" +version = "2.66.1" +source = { registry = "https://pypi.org/simple" } +dependencies = [ + { name = "certifi" }, + { name = "urllib3" }, +] +sdist = { url = "https://files.pythonhosted.org/packages/7f/6f/d59cad0889d15fde85254cf58e701484de3f3f0406003b3197746910b19b/sentry_sdk-2.66.1.tar.gz", hash = "sha256:f882fb08710c5f8bfc603aafa3e901b384009a19cc3f76a572b863392ee81cdc", size = 940543 } +wheels = [ + { url = "https://files.pythonhosted.org/packages/89/d3/726bd88f0eece09ddf431bea4c9191c18e7a8d070b854eb0014d447712ee/sentry_sdk-2.66.1-py3-none-any.whl", hash = "sha256:86002793161d9a95ef04bdd8d442e9bfece5d989b755f05d6360215094a7aff6", size = 505555 }, +] + [[package]] name = "six" version = "1.17.0" @@ -1580,3 +1720,67 @@ sdist = { url = "https://files.pythonhosted.org/packages/a2/65/b7c6c443ccc58678c wheels = [ { url = "https://files.pythonhosted.org/packages/45/ec/dbb7e5a6b91f86bfb9eb7d2988a2730907b6a729875b949c7f022e8b88fa/uvicorn-0.51.0-py3-none-any.whl", hash = "sha256:5d38af6cd620f2ae3849fb44fd4879e0890aa1febe8d47eb355fb45d93fe6a5b", size = 73219 }, ] + +[[package]] +name = "wrapt" +version = "2.3.0" +source = { registry = "https://pypi.org/simple" } +sdist = { url = "https://files.pythonhosted.org/packages/2b/b0/c1f5a970721f06b85c0cd5142e0ff8fe067708abd779b0c4f4be7d61d09f/wrapt-2.3.0.tar.gz", hash = "sha256:681a2d0eefd721998f90642762b8e75c2159ec531b20ad5e437245ea7b06a107", size = 131509 } +wheels = [ + { url = "https://files.pythonhosted.org/packages/5b/4a/d17a0fad1bf1c5f2c887ff71fef75654141b0880bff71d157d955b5bec3a/wrapt-2.3.0-cp312-cp312-macosx_10_13_x86_64.whl", hash = "sha256:0a45ffae742ce91a16e11cb6c7cd71e7f9994f3cbd283b962ab093f5c6dcf525", size = 82139 }, + { url = "https://files.pythonhosted.org/packages/6e/55/51b92daaf6defb57f4dc56bdcce985400f75c6984a03ca5e78ccac717028/wrapt-2.3.0-cp312-cp312-macosx_11_0_arm64.whl", hash = "sha256:69e477046f2237ef0bc6547544ee73008dc764ca26eff44f09e976d221b34d5d", size = 82723 }, + { url = "https://files.pythonhosted.org/packages/28/7f/cfd9bc4b1f5e424eeea83d0493e43f3b1b02707ce8e50c47945873982bd5/wrapt-2.3.0-cp312-cp312-manylinux1_x86_64.manylinux_2_28_x86_64.manylinux_2_5_x86_64.whl", hash = "sha256:5d221a6e6ddd302b8397433184e96b59f259f50024b854db1c411a881586b6b8", size = 172381 }, + { url = "https://files.pythonhosted.org/packages/cb/89/ff7814f6eb6856b479946117d1138a2fbb46cdb6b1f379db359056c69743/wrapt-2.3.0-cp312-cp312-manylinux2014_aarch64.manylinux_2_17_aarch64.manylinux_2_28_aarch64.whl", hash = "sha256:392158c9a7f2ab1b8699418bfc0fe6f83548788c418b27d7bf2019ad3405cebb", size = 174120 }, + { url = "https://files.pythonhosted.org/packages/12/1e/8eded8615d39e3ce81f626937a3a87b280a2a86239a2bf14a4b4bb345034/wrapt-2.3.0-cp312-cp312-manylinux_2_31_riscv64.manylinux_2_39_riscv64.whl", hash = "sha256:e5301c35cf75655eb33498f2bd6ae8703ca19940e3167dc9cdf740c712a39c60", size = 163035 }, + { url = "https://files.pythonhosted.org/packages/35/ea/a0af2d9da62897af2a055484920de05dade30d2ba2c0d65cbdea875d3d8b/wrapt-2.3.0-cp312-cp312-musllinux_1_2_aarch64.whl", hash = "sha256:418f54bb09d1762db02c7009b4051149893af3153a87f92d70356703c11eea02", size = 171887 }, + { url = "https://files.pythonhosted.org/packages/7e/dd/63cd4c864c65ef4906df64bd2d378f4a62b54f28063f282dfb3bf93caead/wrapt-2.3.0-cp312-cp312-musllinux_1_2_riscv64.whl", hash = "sha256:1598becd30f8f2777d18564064eb4f4dbe1ab0e05a8f09786d0ef505ac782bf3", size = 161113 }, + { url = "https://files.pythonhosted.org/packages/ca/ee/82f1fc9e431b5c2c5a6d201aa865dbeae3984c311c6d11a185f0c8367cf6/wrapt-2.3.0-cp312-cp312-musllinux_1_2_x86_64.whl", hash = "sha256:3da470536bf9645143323dd41b32db55c6f4304ad382094c1a1da8a92061e10d", size = 170530 }, + { url = "https://files.pythonhosted.org/packages/37/a5/5dc590e863a419930d988f8b7ca3e75a6befcfb10b6003b3a152f3d5f732/wrapt-2.3.0-cp312-cp312-win32.whl", hash = "sha256:fb8e2e6704a1e0b1b989546c69e2688371ef4a07fa5f61bde3eb6211186f5ac1", size = 78323 }, + { url = "https://files.pythonhosted.org/packages/51/f9/4a6925a07951df56394f7e6ebe14f69f1c5ef9d87aa63e0839acf15aa63a/wrapt-2.3.0-cp312-cp312-win_amd64.whl", hash = "sha256:cdc021cb0b62471d6aac7f2bd92f3b4658073775f9ee7fcd325c511129e7bcc8", size = 81180 }, + { url = "https://files.pythonhosted.org/packages/a8/4f/8b5de0395b2a72216751d41c9861df6facaeb611b619d8810ed2b3b23eb2/wrapt-2.3.0-cp312-cp312-win_arm64.whl", hash = "sha256:67bfe2485f50368c3fcd2275fc1fd100e350d601e0058921a7c82678a465aeab", size = 80155 }, + { url = "https://files.pythonhosted.org/packages/8e/6e/0f88a072483e76b881e3fdcd6b6ffb4a5791002514fe541e72b1b73c859a/wrapt-2.3.0-cp313-cp313-macosx_10_13_x86_64.whl", hash = "sha256:0d3fb71e65b001adfc42684522eeccd9c21d8ba679945abc993439567b66e59f", size = 81960 }, + { url = "https://files.pythonhosted.org/packages/d7/ff/b7e2776e7c294075eb712cc9ef573d1b818f393006d09787262b8fc871c4/wrapt-2.3.0-cp313-cp313-macosx_11_0_arm64.whl", hash = "sha256:51a7a4181c1295774812271fbcd7c909df372bc25579d4ed9eb875caaf0ae86f", size = 82435 }, + { url = "https://files.pythonhosted.org/packages/d8/90/343bb5d0f1f9669bc252a6073f085b4abf862511bd5c9c9eaec754341f1d/wrapt-2.3.0-cp313-cp313-manylinux1_x86_64.manylinux_2_28_x86_64.manylinux_2_5_x86_64.whl", hash = "sha256:9045917809c63fdf7abe3a2ceaed3d670b8ee4500ddd9291192d30aeb34467c5", size = 170350 }, + { url = "https://files.pythonhosted.org/packages/59/f8/13b79a392930bd0dd6b86cbfbfe1c40944110456e1dc6d809e5c46ece904/wrapt-2.3.0-cp313-cp313-manylinux2014_aarch64.manylinux_2_17_aarch64.manylinux_2_28_aarch64.whl", hash = "sha256:54ca1d5573f69b5fe1d74f1f65799c68015e82f685efec9fd8cfa40a094c44d0", size = 170022 }, + { url = "https://files.pythonhosted.org/packages/b2/fc/4f1b6918f5290db959d6e0c07f77385d87cede29c39c9cf8f145e9c82954/wrapt-2.3.0-cp313-cp313-manylinux_2_31_riscv64.manylinux_2_39_riscv64.whl", hash = "sha256:242b60c21e30866e6a2fa606c612b47c553fa60c0eaeeeb7797fb842ac0ce609", size = 161043 }, + { url = "https://files.pythonhosted.org/packages/01/e1/45d3cf74414780bdff6d0380467e003f6eb0f028b6c9403db868dbc7209c/wrapt-2.3.0-cp313-cp313-musllinux_1_2_aarch64.whl", hash = "sha256:e3f3d7ec0a51fbfe00d3aef047641ff2c58b25565b4717fc1f90e050be01cba8", size = 168576 }, + { url = "https://files.pythonhosted.org/packages/f3/73/2fa58dd97f191c997755e2c6d569a68f0c433db4e4b36099bdd7227b6cac/wrapt-2.3.0-cp313-cp313-musllinux_1_2_riscv64.whl", hash = "sha256:261f53870cd4fb2bf38f9f972c56c728fd224cb7c65721307de59d9e7e6741ae", size = 159140 }, + { url = "https://files.pythonhosted.org/packages/29/a8/08a56e2000a8816d449dcbad8c8b081697acbbd490821ceca0f9d8e8d20c/wrapt-2.3.0-cp313-cp313-musllinux_1_2_x86_64.whl", hash = "sha256:8159ec0b0cb7608175eb150de94c19e34f4d47ac655f5ca9baf45df6b688ffd3", size = 169263 }, + { url = "https://files.pythonhosted.org/packages/9e/d4/354e1725e35a73b2af4fa70a3e024c7a5d1bf1802dfb862dcb668aae0253/wrapt-2.3.0-cp313-cp313-win32.whl", hash = "sha256:10461884b3014fbfc8eb7d09a93c5f246363e6711d9d881f95eb8c27fdef049f", size = 78241 }, + { url = "https://files.pythonhosted.org/packages/6c/7e/34c87fa2174848dfee820322aaa318bab08913998ccecc8d2f57b4ad4639/wrapt-2.3.0-cp313-cp313-win_amd64.whl", hash = "sha256:ac870cc97b73bb00ac353329e9559a4bebc47c4c86792ed9b23b58c15b6ad838", size = 81113 }, + { url = "https://files.pythonhosted.org/packages/11/86/fcc9a530579e008c9478bb565a6cdfbfd33536660f069c8b91a6607c5050/wrapt-2.3.0-cp313-cp313-win_arm64.whl", hash = "sha256:a65e8db2b4e90c2e7ade931086351c98ef420bf7a94ee08c95ac8a3cbbc43579", size = 80182 }, + { url = "https://files.pythonhosted.org/packages/96/50/3864848b95b28ef73e17551fc8dccbff2628a834f52cf26a57f9c419fb83/wrapt-2.3.0-cp313-cp313t-macosx_10_13_x86_64.whl", hash = "sha256:fd1f2f557dd3491fe75905e578f4db967393d40d1a8f468edc4d40ac7f2d5944", size = 83921 }, + { url = "https://files.pythonhosted.org/packages/3b/4c/3d1921a60c3e8c71c540ff136e6a47a1fbccf7f671e818394889f7871d9c/wrapt-2.3.0-cp313-cp313t-macosx_11_0_arm64.whl", hash = "sha256:9f5d2aec29dfc76c37e23897dee92766a3fd4f3bff3ae7fc9c6b4bf37d8c1360", size = 84412 }, + { url = "https://files.pythonhosted.org/packages/fa/1a/4a796ff7adb26ada6d4b758c94d47a38320b085e7099afc088efbbcdb006/wrapt-2.3.0-cp313-cp313t-manylinux1_x86_64.manylinux_2_28_x86_64.manylinux_2_5_x86_64.whl", hash = "sha256:646d20d413ffcd1b0a2f700076e2d0252d872dcb7754860a73e45a59ea883614", size = 207168 }, + { url = "https://files.pythonhosted.org/packages/1d/3e/d7777776806c579b761bac2f91721dda9f04c7a1b380213c5935cc750ae6/wrapt-2.3.0-cp313-cp313t-manylinux2014_aarch64.manylinux_2_17_aarch64.manylinux_2_28_aarch64.whl", hash = "sha256:379f670f45b7bb8993edd9f6fc36c6cc65edb81cffa0b504be34acb0303fff0a", size = 214351 }, + { url = "https://files.pythonhosted.org/packages/63/27/2d64d394df7bf181955b3bb562bf33c4492fb4be113f53071106d43ad8b5/wrapt-2.3.0-cp313-cp313t-manylinux_2_31_riscv64.manylinux_2_39_riscv64.whl", hash = "sha256:6208f302f110295d64b22a7ac96500c791bf492dce4366e622e4912b077c9687", size = 199020 }, + { url = "https://files.pythonhosted.org/packages/3e/3d/fb31d3db7d9834d265fb1a27a2adf0ddf51557c67458c97b22439ad6ae3d/wrapt-2.3.0-cp313-cp313t-musllinux_1_2_aarch64.whl", hash = "sha256:ed635a9ca4f3a5a2b900c10c69e823373bc00ebc114b459383596d3487da3570", size = 209969 }, + { url = "https://files.pythonhosted.org/packages/1f/d1/8724b5da582e62070dc9bf4d8bf1972f317297eefd7ba1f2b5c6393ccf6c/wrapt-2.3.0-cp313-cp313t-musllinux_1_2_riscv64.whl", hash = "sha256:e3b9eaa742ae7a0aaaaad4ca4b69469d757af2d6e6663ef1dadc47adec0aeb41", size = 196324 }, + { url = "https://files.pythonhosted.org/packages/0d/5c/3d9ef411149543016ee6bcf3af707f787cebd946527452b94bf122e9b7b4/wrapt-2.3.0-cp313-cp313t-musllinux_1_2_x86_64.whl", hash = "sha256:d0f7284f88f4833705132d06d3b425a43095c2cbd07c58166aac3ab646ba12a4", size = 202610 }, + { url = "https://files.pythonhosted.org/packages/13/9b/4fc042ceb757866dd4a5fc057b3b736f2b360d3703ce9f830d83dc9226e0/wrapt-2.3.0-cp313-cp313t-win32.whl", hash = "sha256:7ebb274aba688b043429eb1500ff8a76ce0cb8ac0812ca3e301f06247b8722b3", size = 79178 }, + { url = "https://files.pythonhosted.org/packages/6b/ff/b94878f8eed809ca042685276bcea9f24e8c2ca7c9653bb80bbb920a68a5/wrapt-2.3.0-cp313-cp313t-win_amd64.whl", hash = "sha256:c4bded758ad6f03b965830944a2f0bc5b2eb3767fe5a7310134315d1a6610e98", size = 82634 }, + { url = "https://files.pythonhosted.org/packages/80/fb/663e1de5332a71685a729754312d327d4cada767c36e1c5a2db4c8de49e6/wrapt-2.3.0-cp313-cp313t-win_arm64.whl", hash = "sha256:d2cc64539da63e39ffb9c7ede849b6e8ddaaf7b3876b5cfb04efd85a5f3f4eb6", size = 81387 }, + { url = "https://files.pythonhosted.org/packages/58/10/b073beaea89bc0d3670a75ff51139430a54b6af7ba7796507730634536dd/wrapt-2.3.0-cp314-cp314-macosx_10_15_x86_64.whl", hash = "sha256:ea52a0d0f08c584943d5764be0e84efa912c8da23c23e1e285ff2f5641c18fcc", size = 81978 }, + { url = "https://files.pythonhosted.org/packages/b3/31/0916d9cebf848ed3f1a0c1888faee421747df77331e4db2bc527a9a85988/wrapt-2.3.0-cp314-cp314-macosx_11_0_arm64.whl", hash = "sha256:fd85b0aa88efdb189d6ae2f35f4526943a8f091c38599c9c31478241c819e6a1", size = 82518 }, + { url = "https://files.pythonhosted.org/packages/f5/73/31c1bf0f3384062751c2094dadb314916d70aa9b6bfd26d994b4a7b393fa/wrapt-2.3.0-cp314-cp314-manylinux1_x86_64.manylinux_2_28_x86_64.manylinux_2_5_x86_64.whl", hash = "sha256:141ed6211286a9660d8d6702de598b43f0934b4f0eda16393f100a80f501d945", size = 170187 }, + { url = "https://files.pythonhosted.org/packages/ed/25/fce087d54b79b8905f3c3c9dd5f454bbd8d8acb80b960c4a6aee5b4659b3/wrapt-2.3.0-cp314-cp314-manylinux2014_aarch64.manylinux_2_17_aarch64.manylinux_2_28_aarch64.whl", hash = "sha256:2e49885a62ec4ee854d1b9e6371fda6afd219917225752abf729a3f36d4df9a5", size = 169288 }, + { url = "https://files.pythonhosted.org/packages/c7/30/0d09e6dddc6b7a7230ac77f50254b5980ab4fcd22976f72f8cc8a0404458/wrapt-2.3.0-cp314-cp314-manylinux_2_31_riscv64.manylinux_2_39_riscv64.whl", hash = "sha256:1d6159c9b2fefec02314e1332dbbbfaf960e369dfd26bcf7f8b258b5732065b3", size = 160932 }, + { url = "https://files.pythonhosted.org/packages/2c/ca/0913af0d2ec0c43865d32d615f518fea66c13c5c930e489e9b0de248e9a8/wrapt-2.3.0-cp314-cp314-musllinux_1_2_aarch64.whl", hash = "sha256:24da48596326ef8e448cfa837b454f638713d3531262375f00e5a9681682fc07", size = 169017 }, + { url = "https://files.pythonhosted.org/packages/c3/f2/3d1e47ea81b822210f5df1bf942fd90780a75c055243d569b664529dea88/wrapt-2.3.0-cp314-cp314-musllinux_1_2_riscv64.whl", hash = "sha256:cd3a2edf0427013736b8127955cec62608c56e53ea47e82812ea32059cda407f", size = 159065 }, + { url = "https://files.pythonhosted.org/packages/43/a5/ef2066ced8e5fca204e2b361e9708e36555b40949c583d997ea3b590817d/wrapt-2.3.0-cp314-cp314-musllinux_1_2_x86_64.whl", hash = "sha256:4fa0df3bff4e7ce45759f33fd39335fe2f60477bb9ecf7b8aa41e7d07ee36a23", size = 168821 }, + { url = "https://files.pythonhosted.org/packages/d5/e1/016104650d4e572fa91506eb396b3dd8efbccc9284fdc1c9479c3d21db28/wrapt-2.3.0-cp314-cp314-win32.whl", hash = "sha256:2935d5454b3f179a29b12cf390ee47246740ba2c3a7545b1b46ba31a5f2a4a0b", size = 78700 }, + { url = "https://files.pythonhosted.org/packages/3d/97/6fdc20a9f2ca304748b3f0819cbf377d55260562777bf0b615431bc3c181/wrapt-2.3.0-cp314-cp314-win_amd64.whl", hash = "sha256:cc2cea812e5cb179a796b766747e7d3b21088760d8deb95676d482b8c8e6fa7d", size = 81422 }, + { url = "https://files.pythonhosted.org/packages/5e/a4/9cbd53bf05746bea2c392af39cb052427a8ec95cbd494d930733d8f44681/wrapt-2.3.0-cp314-cp314-win_arm64.whl", hash = "sha256:22cc5c0a717bd4da87018ae0bffd4c19c6fb679d3ff357216ba566ab26c76cab", size = 80639 }, + { url = "https://files.pythonhosted.org/packages/43/bb/6c5e4a0f66ea0d2b2dd267e8dd05a0014eea56840b3c8595d40b0a5d1f91/wrapt-2.3.0-cp314-cp314t-macosx_10_15_x86_64.whl", hash = "sha256:a6b5984cd65dd639546f0eb4b8eacf1c31cb2fe9fb5c27bffe240987cdb2cf84", size = 84030 }, + { url = "https://files.pythonhosted.org/packages/6a/eb/a1aedf03283bc9cbf8a1783995ddc54e3c5a86878f19002d2c428494f4c5/wrapt-2.3.0-cp314-cp314t-macosx_11_0_arm64.whl", hash = "sha256:c88abcf53daef80e01a75c7530e727fa6e2c1888fe83e3dcdba4c96216a1f5c7", size = 84419 }, + { url = "https://files.pythonhosted.org/packages/63/61/50d511c0dc5105563849e86daa3e16ac7feef699f79fb05af45ea70107d5/wrapt-2.3.0-cp314-cp314t-manylinux1_x86_64.manylinux_2_28_x86_64.manylinux_2_5_x86_64.whl", hash = "sha256:85de890ff968196e92dd1ae73a9fb8970495e7650a457b1c9ef0ac3dd550bce2", size = 207171 }, + { url = "https://files.pythonhosted.org/packages/3f/59/9b538cf7795217e810699d16bc88b96a830d9b5c403eb2ec2db6b5f2ae81/wrapt-2.3.0-cp314-cp314t-manylinux2014_aarch64.manylinux_2_17_aarch64.manylinux_2_28_aarch64.whl", hash = "sha256:50f416b74d092bb9f41b424e90dd457f365f7ba4b11de62a23679769a21bd85c", size = 214329 }, + { url = "https://files.pythonhosted.org/packages/b3/28/9935d62b1499e5c8b3d191e99ba4eb31ca237a0b699142011a837e9dc7ea/wrapt-2.3.0-cp314-cp314t-manylinux_2_31_riscv64.manylinux_2_39_riscv64.whl", hash = "sha256:39febbee6d77301d31da6996b152ce52452da7c7ef72aba10c2fa976dff9c295", size = 199079 }, + { url = "https://files.pythonhosted.org/packages/2b/01/4446b80fa2ffa47a3449b250d004ba1c1937f07f64a179608fec735df866/wrapt-2.3.0-cp314-cp314t-musllinux_1_2_aarch64.whl", hash = "sha256:93513bec052c6cd987f9f580c3df068c8bc4ebae6543736be3ca7ec5959cafcd", size = 209992 }, + { url = "https://files.pythonhosted.org/packages/d4/07/56f26c9f9979586a021e8148747004aba4498f49458c90b0502969b904e1/wrapt-2.3.0-cp314-cp314t-musllinux_1_2_riscv64.whl", hash = "sha256:729126e667da34d251b8ebf8a45ef0c5ddadc21542b3d6e1abf4259ece6508df", size = 196334 }, + { url = "https://files.pythonhosted.org/packages/8b/41/6d7bcc895b0f28b2250e10908f060687b9165429dcd7f22ddb3d4c031b74/wrapt-2.3.0-cp314-cp314t-musllinux_1_2_x86_64.whl", hash = "sha256:626b69db2021aa01671ec7bbc9740e558522bd44c18cf2ce69bf3d666a014109", size = 202644 }, + { url = "https://files.pythonhosted.org/packages/cd/25/7860927edba06b758b8852a6f02e832be715563c67a6795d94350bc81099/wrapt-2.3.0-cp314-cp314t-win32.whl", hash = "sha256:629d73378082c00a8173031f9fb30a3ac6abbc894a5bfdfae71fabc60642d501", size = 79685 }, + { url = "https://files.pythonhosted.org/packages/c4/0f/270bafe92fde3b069a39bc01e39ee79340895b335640df861d43d2a51885/wrapt-2.3.0-cp314-cp314t-win_amd64.whl", hash = "sha256:42869085687f0aefd57c0f636c3f9354f8ffb321a8ba9cb52d19beb796e561c5", size = 83104 }, + { url = "https://files.pythonhosted.org/packages/55/b3/af176d79a8515a8a720eccdad9a96f6e31a30abf2865430c8c42adf2fd13/wrapt-2.3.0-cp314-cp314t-win_arm64.whl", hash = "sha256:b1e5aa486e269b00ed35e64771c7d0ab8096cfd2643405ca8cd60ebedc099a51", size = 81774 }, + { url = "https://files.pythonhosted.org/packages/00/39/3daf9f47be208606586de4568ba6713db53ebc8fd7a575aea1fe57983b69/wrapt-2.3.0-py3-none-any.whl", hash = "sha256:d8c7ed08477429752b8c44991f40ad7838b18332a160698740a6bfbc10d998a2", size = 61866 }, +] diff --git a/website/src/app/docs/deployment/page.mdx b/website/src/app/docs/deployment/page.mdx index 0ce30575..594b9f92 100644 --- a/website/src/app/docs/deployment/page.mdx +++ b/website/src/app/docs/deployment/page.mdx @@ -49,6 +49,30 @@ It also contains an OpenRouter placeholder. Replace it before processing a corpus: the complete pipeline makes extraction and embedding calls. Pin explicit model IDs for reproducible work; do not use a rotating free-model router. +## Optional observability + +Observability is strictly opt-in. The self-host image includes the +`rememberstack[observability]` extra, but no exporter initializes unless its +environment configuration is non-empty. + +Set `REMEMBERSTACK_SENTRY_DSN` to send metadata-only error events to a +Sentry-protocol service such as Sentry, GlitchTip, or Bugsink. +`REMEMBERSTACK_SENTRY_ENVIRONMENT` defaults to the deployment slug, and +`REMEMBERSTACK_SENTRY_SAMPLE_RATE` defaults to `1.0`. Request bodies, local +variables, breadcrumbs, PII, prompt/completion text, and exception messages are +not sent. Caught worker failures carry only their stage, lane, and processing ID +as routing tags; PostgreSQL and the existing JSON telemetry remain authoritative +for retry and ledger state. + +LoCoMo answer and judge stages create Langfuse traces only when +`LANGFUSE_PUBLIC_KEY`, `LANGFUSE_SECRET_KEY`, and `LANGFUSE_HOST` are all set. +Run the benchmark with both optional groups, for example +`uv run --extra benchmark --extra observability python -m benchmarks.locomo`. +The observer records per-call model/accounting metadata, content-free tool +argument shapes, final answers, and verdicts, then flushes at stage end. It never +sends source chunks, documents, rendered prompts, tool response bodies, or gold +answers. + Check the API and the deployment-seeded recipe registry: ```bash