diff --git a/.github/workflows/docs.yml b/.github/workflows/docs.yml index 0d232459..f9cf3a39 100644 --- a/.github/workflows/docs.yml +++ b/.github/workflows/docs.yml @@ -1,29 +1,14 @@ name: Docs +# Runs on EVERY pull request / push to main — no `paths:` filter. When `links` +# is a required status check, path-filtering it silently blocks code-only PRs: +# the workflow is skipped, the required check never reports, and the merge box +# waits forever. `make docs-check` validates the whole doc tree regardless of +# what a PR touched and takes only seconds, so it is cheap to always report. on: pull_request: - paths: - - "**/*.md" - - ".github/ISSUE_TEMPLATE/**" - - ".github/PULL_REQUEST_TEMPLATE.md" - - ".github/workflows/docs.yml" - - ".claude/skills/**/*.md" - - "CLAUDE.md" - - "CONTRIBUTING.md" - - "scripts/check_docs.py" - - "scripts/check_github_contributor_docs.py" push: branches: [main] - paths: - - "**/*.md" - - ".github/ISSUE_TEMPLATE/**" - - ".github/PULL_REQUEST_TEMPLATE.md" - - ".github/workflows/docs.yml" - - ".claude/skills/**/*.md" - - "CLAUDE.md" - - "CONTRIBUTING.md" - - "scripts/check_docs.py" - - "scripts/check_github_contributor_docs.py" permissions: contents: read diff --git a/examples/langfuse/.gitignore b/examples/langfuse/.gitignore index ab5bda63..670a9362 100644 --- a/examples/langfuse/.gitignore +++ b/examples/langfuse/.gitignore @@ -1,3 +1,2 @@ -spans.jsonl __pycache__/ .venv/ diff --git a/examples/langfuse/README.md b/examples/langfuse/README.md index fc87ee0b..fdf8b341 100644 --- a/examples/langfuse/README.md +++ b/examples/langfuse/README.md @@ -1,64 +1,73 @@ -# EverOS × Langfuse (OpenTelemetry) +# EverOS × Langfuse (native OpenTelemetry) -Trace EverOS memory operations — writes, LLM extraction, recall with quality -scores, and reflection — into [Langfuse](https://langfuse.com) as OpenTelemetry -spans, so an agent's memory layer becomes visible and evaluable next to the rest -of its traces. +EverOS emits OpenTelemetry spans for its own memory operations — write, memcell +boundary + episode extraction (LLM), search with recall-quality scores, and OME +reflection — and exports them over OTLP to any backend, including +[Langfuse](https://langfuse.com). There is **no wrapper and no extra +instrumentation code**: enable it in config and the traces appear. -This is a thin, dependency-light wrapper (pure OpenTelemetry SDK, no Langfuse -package dependency). The same spans work with Langfuse Cloud, self-hosted -Langfuse, or any other OTLP backend. +## Enable -## Files +1. Install the optional OpenTelemetry extra: -- `everos_langfuse.py` — the instrumentation wrapper (`init_tracing`, - `InstrumentedEverOS`, `HTTPTransport`, recall-score push). -- `demo.py` — a runnable end-to-end example. Ships a mock transport, so it runs - with **no EverOS server required**; set `EVEROS_BASE_URL` to trace a real one. + ```bash + pip install "everos[otel]" + ``` -## Span model +2. Add `[observability]` to your `everos.toml`. The Langfuse keys derive the + OTLP endpoint and auth automatically: + + ```toml + [observability] + enabled = true + langfuse_public_key = "pk-lf-..." + langfuse_secret_key = "sk-lf-..." + langfuse_host = "https://us.cloud.langfuse.com" # EU: https://cloud.langfuse.com + # capture_content = true # opt-in: also record query / extracted memory text + ``` + + Container/CI equivalent via env vars: `EVEROS_OBSERVABILITY__ENABLED=true`, + `EVEROS_OBSERVABILITY__LANGFUSE_PUBLIC_KEY=...`, and so on. + +3. Run EverOS normally: + + ```bash + everos server start + ``` + +Off by default — with `enabled = false` (or the `otel` extra absent) there is +zero tracing overhead. + +## What you get | EverOS operation | Langfuse observation | | --- | --- | -| `POST /api/v1/memory/add` | span `everos.memory.add` | -| `POST /api/v1/memory/flush` → extraction | span + generation `everos.extract` (model + tokens) | +| `POST /api/v1/memory/add` · `flush` | span `everos.memory.add` / `everos.memory.flush` | +| memcell boundary detection (LLM) | generation `everos.memcell.boundary` (model + tokens) | +| episode extraction (LLM) | generation `everos.extract` | | markdown persistence | span `everos.persist.markdown` | -| async index sync | span `everos.cascade.index` (separate correlated trace) | -| `POST /api/v1/memory/search` | retriever `everos.memory.search` | -| ↳ embedding / hybrid recall / rerank | embedding / retriever / span | -| `POST /api/v1/ome/trigger` | agent `everos.ome.` + generation | +| `POST /api/v1/memory/search` | retriever `everos.memory.search` → `recall` / `rank` | +| query / recall embedding | embedding `everos.embedding` | +| OME reflection strategies | agent `everos.ome.` (linked to the triggering request's trace) | -`langfuse.session.id` / `langfuse.user.id` are set on every span; recall quality -is pushed as Langfuse scores (`recall_top_score`, `recall_hit`). +`langfuse.session.id` / `langfuse.user.id` group the traces; recall quality is +pushed as Langfuse scores (`recall_top_score` always, and `recall_hit` for +calibrated methods — HYBRID / AGENTIC). Query and memory text are captured only +when `capture_content = true`. -## Run +## Try it -```bash -pip install opentelemetry-sdk opentelemetry-exporter-otlp requests - -export LANGFUSE_PUBLIC_KEY="pk-lf-..." -export LANGFUSE_SECRET_KEY="sk-lf-..." -export LANGFUSE_HOST="https://us.cloud.langfuse.com" # EU: https://cloud.langfuse.com +With a server running and `[observability]` enabled: +```bash python demo.py ``` -- With no keys set, `demo.py` still runs against the built-in mock and writes a - local `spans.jsonl` (offline inspection) — nothing is sent anywhere. -- With Langfuse keys set, the same spans and recall scores flow into your - Langfuse project. Open **Tracing → Traces** (filter by tag `everos` / `memory`). -- To trace a real deployment, set `EVEROS_BASE_URL` to a running EverOS server - (see the [EverOS quickstart](../../README.md)); the instrumentation is identical. - -## Privacy - -Spans carry non-sensitive metadata (latency, token counts, model names, scores) -by default. Capturing raw query or memory content as span input/output is opt-in. -The demo uses synthetic data, and its `public_traces` flag (safe only for -synthetic data) marks the resulting traces as publicly shareable. +It drives one add → flush → search (keyword / hybrid / agentic) cycle against +`http://127.0.0.1:8000` using only the standard library, then tells you to open +Langfuse → **Tracing** filtered to `session.id = langfuse_demo`. ## Learn more -- Langfuse OpenTelemetry docs: https://langfuse.com/integrations/native/opentelemetry -- Native, opt-in instrumentation inside EverOS core is planned; this wrapper is - the interim path and mirrors the same span model. +- Langfuse OpenTelemetry: https://langfuse.com/integrations/native/opentelemetry +- Config reference: the `[observability]` block in `src/everos/config/default.toml`. diff --git a/examples/langfuse/demo.py b/examples/langfuse/demo.py index 49809907..ecf85420 100644 --- a/examples/langfuse/demo.py +++ b/examples/langfuse/demo.py @@ -1,218 +1,92 @@ -"""End-to-end demo: EverOS memory operations traced into Langfuse. +"""Minimal EverOS x Langfuse demo — native OpenTelemetry tracing. -Replays one realistic memory lifecycle — ingest -> extraction -> recall (with -an updated fact winning over a stale one) -> agent-skill recall -> reflection — -through the instrumentation in everos_langfuse.py. +EverOS emits OTel spans for its own memory operations when ``[observability]`` +is enabled; this script contains **no instrumentation code**. It just drives a +running server (add -> flush -> search) so the traces the server produces show +up in your Langfuse project. -Two modes, same code path: - * offline (default) — spans land in ./spans.jsonl for offline inspection - * live — set LANGFUSE_PUBLIC_KEY / LANGFUSE_SECRET_KEY / - LANGFUSE_HOST and the exact same spans + recall - scores also flow into your Langfuse project. +Prereqs (see README.md): + 1. pip install "everos[otel]" + 2. configure [observability] in everos.toml with your Langfuse keys + 3. everos server start # defaults to http://127.0.0.1:8000 -The MockEverOSTransport returns responses in the exact envelope/shape of the -EverOS HTTP API v1 (see EverOS docs/api.md); swap in HTTPTransport to run -against a real `pip install everos` server — the instrumentation is identical. +Then: python demo.py """ from __future__ import annotations +import json import time -import uuid - -from everos_langfuse import HTTPTransport, InstrumentedEverOS, force_flush, init_tracing - -TS = int(time.time() * 1000) -DAY = "20260702" - - -def _envelope(data: dict, detail: dict | None = None) -> dict: - resp = {"request_id": uuid.uuid4().hex, "data": data} - if detail: - resp["_detail"] = detail # server-side facts the spans describe - return resp - - -class MockEverOSTransport: - """Faithful mock of the EverOS HTTP API v1 (response envelope + field - shapes from docs/api.md), so the demo runs without provider keys.""" - - def __init__(self): - self.buffer: list[dict] = [] - - def __call__(self, path: str, payload: dict) -> dict: - if path == "/api/v1/memory/add": - self.buffer.extend(payload["messages"]) - time.sleep(0.012) - return _envelope({"message_count": len(payload["messages"]), - "status": "accumulated"}) - - if path == "/api/v1/memory/flush": - buffered, self.buffer = self.buffer, [] - time.sleep(0.01) - return _envelope( - {"status": "extracted"}, - detail={ - "model": "gpt-4.1-mini", - "buffered_messages": [m["content"] for m in buffered], - "memory_cell": { - "episode_id": "alice_ep_%s_001" % DAY, - "subject": "Alice's routines and recent move", - "summary": ("Alice climbs in Yosemite every spring, bikes to " - "work, and recently moved from SOMA to Oakland; " - "her go-to coffee used to be Blue Bottle in SOMA."), - "atomic_facts": [ - "Alice climbs in Yosemite every spring.", - "Alice bikes to work most days.", - "Alice moved from SOMA to Oakland in June 2026.", - "Alice's favorite coffee shop was Blue Bottle in SOMA.", - ], - }, - "usage": {"input": 642, "output": 187}, - "md_files": ["memory/alice/episodic/2026-07-02-alice-routines.md"], - "rows_indexed": 5, - "index_lag_ms": 512, - "extract_s": 0.42, - }, - ) - - if path == "/api/v1/memory/search": - q = payload["query"].lower() - if "live" in q: # conflict-resolution showcase: fresh fact outranks stale - ranked = [ - {"id": "alice_af_%s_003" % DAY, - "content": "Alice moved from SOMA to Oakland in June 2026.", - "score": 0.81}, - {"id": "alice_af_%s_004" % DAY, - "content": "Alice's favorite coffee shop was Blue Bottle in SOMA.", - "score": 0.34}, - ] - elif "sport" in q or "outdoor" in q: - ranked = [ - {"id": "alice_af_%s_001" % DAY, - "content": "Alice climbs in Yosemite every spring.", "score": 0.86}, - {"id": "alice_af_%s_002" % DAY, - "content": "Alice bikes to work most days.", "score": 0.72}, - ] - elif payload.get("agent_id"): # agent track: cases + skills - ranked = [ - {"id": "raven_case_%s_007" % DAY, - "content": "Case: flaky LanceDB test fixed by pinning fsync " - "before rename and retrying open with backoff.", - "score": 0.74}, - {"id": "raven_skill_retry_backoff", - "content": "Skill: wrap flaky IO in retry-with-backoff; verify " - "with 3 consecutive green runs.", - "score": 0.69}, - ] - else: # deliberate miss: query about something never stored - ranked = [ - {"id": "alice_af_%s_002" % DAY, - "content": "Alice bikes to work most days.", "score": 0.31}, - ] - - time.sleep(0.01) - if payload.get("agent_id"): - data = {"episodes": [], "profiles": [], - "agent_cases": [r for r in ranked if "case" in r["id"]], - "agent_skills": [r for r in ranked if "skill" in r["id"]], - "unprocessed_messages": []} - else: - data = {"episodes": [{ - "id": "alice_ep_%s_001" % DAY, - "user_id": payload.get("user_id"), - "session_id": "sess-cafe-chat-001", - "summary": "Alice's routines and recent move", - "score": ranked[0]["score"], - "atomic_facts": ranked, - }], - "profiles": [], "agent_cases": [], "agent_skills": [], - "unprocessed_messages": []} - return _envelope(data, detail={ - "embed_model": "Qwen/Qwen3-Embedding-4B", "embed_tokens": 11, - "rerank_model": "Qwen/Qwen3-Reranker-4B", - "candidates": 24, "ranked": ranked, - "embed_s": 0.028, "recall_s": 0.019, "rerank_s": 0.047, - }) - - if path == "/api/v1/ome/trigger": - time.sleep(0.01) - return _envelope( - {"status": "ok", "name": payload["name"]}, - detail={ - "model": "gpt-4.1-mini", - "episodes_in": ["alice_ep_%s_001" % DAY], - "consolidated": { - "profile_update": "home_location: SOMA -> Oakland (2026-06)", - "episodes_merged": 1, - }, - "usage": {"input": 918, "output": 141}, - "reflect_s": 0.31, - }, - ) +import urllib.request + +BASE = "http://127.0.0.1:8000" +SESSION = "langfuse_demo" +USER = "alice" - raise ValueError(f"unknown path {path}") + +def _post(path: str, body: dict) -> dict: + req = urllib.request.Request( + BASE + path, + data=json.dumps(body).encode(), + headers={"content-type": "application/json"}, + method="POST", + ) + with urllib.request.urlopen(req) as resp: + return json.load(resp) def main() -> None: - live = init_tracing(service_name="everos", spans_jsonl="spans.jsonl") - print(f"[demo] tracing initialised — live Langfuse export: {live}") - - import os - if os.getenv("EVEROS_BASE_URL"): - transport = HTTPTransport(os.environ["EVEROS_BASE_URL"]) - print(f"[demo] using real EverOS server at {os.environ['EVEROS_BASE_URL']}") - else: - transport = MockEverOSTransport() - print("[demo] using MockEverOSTransport (EverOS HTTP API v1 shapes)") - - # public_traces=True: demo data is synthetic (fictional "Alice"), so the - # resulting traces are safe to share as public Langfuse trace URLs. - ev = InstrumentedEverOS(transport, public_traces=True) - session, user = "sess-cafe-chat-001", "alice" - - # -- 1. write path: ingest a conversation ------------------------------ - ev.add(session, [ - {"sender_id": user, "role": "user", "timestamp": TS, - "content": "I love climbing in Yosemite every spring."}, - {"sender_id": user, "role": "user", "timestamp": TS + 10, - "content": "My favorite coffee shop is Blue Bottle in SOMA."}, - {"sender_id": user, "role": "user", "timestamp": TS + 20, - "content": "I bike to work most days."}, - ], user_id=user) - ev.add(session, [ - {"sender_id": user, "role": "user", "timestamp": TS + 30, - "content": "Oh — actually I moved from SOMA to Oakland last month."}, - ], user_id=user) - - # -- 2. boundary/flush: LLM extraction -> markdown -> index ------------ - ev.flush(session, user_id=user) - - # -- 3. read path: recall with quality scores --------------------------- - r1 = ev.search("What outdoor sports does Alice do?", user_id=user, - session_id=session) - r2 = ev.search("Where does Alice live now?", user_id=user, session_id=session) - r3 = ev.search("What are Alice's favorite books?", user_id=user, - session_id=session) # deliberate low-quality recall - # agent-memory track (cases / skills) — the Raven angle - r4 = ev.search("How did we fix the flaky LanceDB test last time?", - agent_id="raven-dev-agent", session_id="raven-run-042") - - # -- 4. self-evolution: offline reflection ------------------------------ - ev.trigger_ome("reflect_episodes", user_id=user, session_id=session) - - force_flush() - time.sleep(0.5) - - print("\n[demo] traces emitted:") - for label, r in [("recall: sports", r1), ("recall: moved city", r2), - ("recall: miss (books)", r3), ("recall: agent skill", r4)]: - print(f" - {label:24s} trace_id={r['_trace_id']} " - f"scores_pushed={r['_scores_pushed']}") - print("\n[demo] spans also written to spans.jsonl (offline copy)") - if live: - print("[demo] open your Langfuse project -> Traces; " - "scores 'recall_top_score' / 'recall_hit' attached to searches.") + ts = int(time.time() * 1000) + + add = _post( + "/api/v1/memory/add", + { + "session_id": SESSION, + "messages": [ + { + "message_id": "m1", + "role": "user", + "content": "Moved our vector store to LanceDB to fix index bloat.", + "timestamp": ts, + "sender_id": USER, + }, + { + "message_id": "m2", + "role": "assistant", + "content": "Noted — LanceDB with compaction keeps it compact.", + "timestamp": ts + 1000, + "sender_id": "assistant", + }, + ], + }, + ) + print("add ->", add["data"]) + + flush = _post("/api/v1/memory/flush", {"session_id": SESSION, "messages": []}) + print("flush ->", flush["data"]) + + print("waiting for async index sync ...") + time.sleep(10) + + for method in ("keyword", "hybrid", "agentic"): + resp = _post( + "/api/v1/memory/search", + { + "user_id": USER, + "query": "which vector database did we move to and why", + "method": method, + "top_k": 5, + "filters": {"session_id": SESSION}, + }, + ) + hits = len(resp["data"].get("episodes", [])) + print(f"search[{method}] -> {hits} hit(s)") + + print( + f"\nOpen Langfuse -> Tracing and filter session.id = {SESSION} " + "to see the traces (add / flush / search, with token usage and " + "recall-quality scores)." + ) if __name__ == "__main__": diff --git a/examples/langfuse/everos_langfuse.py b/examples/langfuse/everos_langfuse.py deleted file mode 100644 index 104f0610..00000000 --- a/examples/langfuse/everos_langfuse.py +++ /dev/null @@ -1,473 +0,0 @@ -"""EverOS -> Langfuse OpenTelemetry instrumentation (prototype). - -Emits EverOS memory operations as OpenTelemetry spans following Langfuse's -attribute conventions (https://langfuse.com/integrations/native/opentelemetry), -so that an agent's memory layer becomes visible — and evaluable — inside -Langfuse, next to the rest of the trace. - -Span model (mirrors EverOS's documented write/read paths): - - POST /api/v1/memory/add span "everos.memory.add" - POST /api/v1/memory/flush span "everos.memory.flush" - |- extraction (LLM) generation "everos.extract" model/tokens/cost - |- markdown persistence span "everos.persist.markdown" - |- index sync span "everos.index.sqlite+lancedb" - POST /api/v1/memory/search retriever "everos.memory.search" query/top_k -> episodes+scores - |- query embedding embedding "everos.search.embed_query" - |- hybrid recall retriever "everos.search.hybrid_recall" (BM25 + vector ANN + fusion) - |- rerank span "everos.search.rerank" scores - POST /api/v1/ome/trigger agent "everos.ome." (reflection / self-evolution) - |- consolidation (LLM) generation "everos.reflect.consolidate" - -Design notes: - * Pure OpenTelemetry SDK — no Langfuse package dependency. The same spans - can go to any OTLP backend (incl. an OpenTelemetry Collector); - Langfuse ingests them natively on /api/public/otel (HTTP/protobuf). - * `langfuse.session.id` / `langfuse.user.id` are set on EVERY span, per - Langfuse's attribute-propagation guidance. - * Recall-quality signals (fused retrieval score of the top hit, hit/miss) - are pushed as Langfuse *scores* via POST /api/public/scores, attached to - the search trace + retriever observation, so they can be plotted and - filtered in Langfuse evals. (Scores are not part of the OTel span model.) - * EverOS request-ids are already W3C trace-context format (32-hex), see - everos.core.observability.tracing — so server-side adoption is a thin, - additive layer. - -This file is written to be read: it doubles as the integration sketch for -the EverOS <> Langfuse proposal. -""" - -from __future__ import annotations - -import base64 -import json -import os -import time -from typing import Any, Callable, Optional - -import requests -from opentelemetry import trace -from opentelemetry.sdk.resources import Resource -from opentelemetry.sdk.trace import TracerProvider, ReadableSpan -from opentelemetry.sdk.trace.export import ( - BatchSpanProcessor, - SimpleSpanProcessor, - SpanExporter, - SpanExportResult, -) - -try: # OTLP/HTTP exporter (protobuf) — what Langfuse's endpoint expects - from opentelemetry.exporter.otlp.proto.http.trace_exporter import OTLPSpanExporter -except ImportError: # pragma: no cover - OTLPSpanExporter = None - -DEFAULT_LANGFUSE_HOST = "https://us.cloud.langfuse.com" - -# Attribute keys we flatten into the local JSONL dump (offline inspection) -_FLAT_KEYS = { - "langfuse.observation.type": "obs_type", - "langfuse.session.id": "session_id", - "langfuse.user.id": "user_id", - "gen_ai.request.model": "model", - "gen_ai.usage.input_tokens": "input_tokens", - "gen_ai.usage.output_tokens": "output_tokens", - "everos.search.top_score": "top_score", - "everos.search.hit": "recall_hit", - "everos.op": "op", -} - - -class JsonLinesSpanExporter(SpanExporter): - """Dump every finished span as one JSON line — a transparent, local record - of exactly what would be sent to Langfuse (handy for offline inspection).""" - - def __init__(self, path: str): - # one file per run — a deterministic offline record - self._fh = open(path, "w", encoding="utf-8") - - def export(self, spans: list[ReadableSpan]) -> SpanExportResult: - for s in spans: - ctx = s.get_span_context() - attrs = dict(s.attributes or {}) - row: dict[str, Any] = { - "trace_id": format(ctx.trace_id, "032x"), - "span_id": format(ctx.span_id, "016x"), - "parent_span_id": format(s.parent.span_id, "016x") if s.parent else "", - "name": s.name, - "start_ts": s.start_time // 1_000_000, # ms epoch - "duration_ms": round((s.end_time - s.start_time) / 1_000_000, 3), - "status": s.status.status_code.name, - } - for k, col in _FLAT_KEYS.items(): - if k in attrs: - row[col] = attrs[k] - row["attributes"] = {k: v for k, v in attrs.items()} - self._fh.write(json.dumps(row, ensure_ascii=False, default=str) + "\n") - self._fh.flush() - return SpanExportResult.SUCCESS - - def shutdown(self) -> None: - self._fh.close() - - -def init_tracing( - service_name: str = "everos", - spans_jsonl: str = "spans.jsonl", -) -> bool: - """Configure OTel. Returns True if a live Langfuse exporter is attached. - - Reads LANGFUSE_PUBLIC_KEY / LANGFUSE_SECRET_KEY / LANGFUSE_HOST from env. - Offline (no keys): spans still go to the local JSONL file, so you can - inspect exactly what would be sent to Langfuse without an account. - """ - resource = Resource.create( - { - "service.name": service_name, - "service.version": "1.1.0", # everos PyPI version this models - } - ) - provider = TracerProvider(resource=resource) - provider.add_span_processor(SimpleSpanProcessor(JsonLinesSpanExporter(spans_jsonl))) - - pk = os.getenv("LANGFUSE_PUBLIC_KEY") - sk = os.getenv("LANGFUSE_SECRET_KEY") - host = os.getenv("LANGFUSE_HOST", DEFAULT_LANGFUSE_HOST).rstrip("/") - live = bool(pk and sk and OTLPSpanExporter) - if live: - auth = base64.b64encode(f"{pk}:{sk}".encode()).decode() - exporter = OTLPSpanExporter( - endpoint=f"{host}/api/public/otel/v1/traces", - headers={ - "Authorization": f"Basic {auth}", - "x-langfuse-ingestion-version": "4", - }, - ) - provider.add_span_processor(BatchSpanProcessor(exporter)) - trace.set_tracer_provider(provider) - return live - - -def force_flush() -> None: - provider = trace.get_tracer_provider() - if hasattr(provider, "force_flush"): - provider.force_flush() - - -# -------------------------------------------------------------------------- -# Langfuse scores (recall quality) — pushed via the public API, since scores -# are first-class objects in Langfuse rather than span attributes. -# -------------------------------------------------------------------------- - - -def push_score( - trace_id: str, - name: str, - value: float, - observation_id: Optional[str] = None, - comment: Optional[str] = None, -) -> bool: - pk = os.getenv("LANGFUSE_PUBLIC_KEY") - sk = os.getenv("LANGFUSE_SECRET_KEY") - host = os.getenv("LANGFUSE_HOST", DEFAULT_LANGFUSE_HOST).rstrip("/") - if not (pk and sk): - return False - payload: dict[str, Any] = { - "traceId": trace_id, - "name": name, - "value": value, - "dataType": "NUMERIC", - } - if observation_id: - payload["observationId"] = observation_id - if comment: - payload["comment"] = comment - try: - r = requests.post( - f"{host}/api/public/scores", auth=(pk, sk), json=payload, timeout=15 - ) - return r.status_code in (200, 201, 207) - except requests.RequestException as exc: # never break the caller's flow - print(f"[everos-langfuse] score push failed ({type(exc).__name__}); " - "spans are still recorded locally") - return False - - -# -------------------------------------------------------------------------- -# Instrumented EverOS client -# -------------------------------------------------------------------------- - -Transport = Callable[[str, dict], dict] -_TRUNC = 4000 # keep span payloads bounded - - -def _j(obj: Any) -> str: - s = json.dumps(obj, ensure_ascii=False, default=str) - return s if len(s) <= _TRUNC else s[:_TRUNC] + "…" - - -def _top_score_from_data(data: dict) -> float | None: - """Best hit score across all scored result arrays in a real search - response. Each array is already sorted desc by the server, so the top - hit is derivable from the public API output alone — no server-internal - detail needed. Returns None when nothing scored came back (a miss).""" - scores = [ - float(item["score"]) - for key in ("episodes", "profiles", "agent_cases", "agent_skills") - for item in (data.get(key) or []) - if item.get("score") is not None - ] - return max(scores) if scores else None - - -class InstrumentedEverOS: - """Wraps an EverOS transport (real HTTP server or mock) and emits the - spans that the proposed server-side instrumentation would emit. - - Every public method == one EverOS API call == one Langfuse trace. - """ - - def __init__(self, transport: Transport, tracer_name: str = "everos", - public_traces: bool = False): - """public_traces: mark every trace as publicly shareable via URL - (langfuse.trace.public). Only enable for synthetic/demo data — - never for real memory content.""" - self._t = transport - self._tracer = trace.get_tracer(tracer_name) - self._public = public_traces - - # -- helpers ------------------------------------------------------------ - - def _common(self, span, *, session_id=None, user_id=None, agent_id=None, - app_id="default", project_id="default", obs_type="span", op=""): - span.set_attribute("langfuse.observation.type", obs_type) - span.set_attribute("everos.op", op) - if self._public: - span.set_attribute("langfuse.trace.public", True) - if session_id: - span.set_attribute("langfuse.session.id", session_id) - if user_id: - span.set_attribute("langfuse.user.id", user_id) - if agent_id: - span.set_attribute("langfuse.trace.metadata.agent_id", agent_id) - span.set_attribute("langfuse.trace.metadata.app_id", app_id) - span.set_attribute("langfuse.trace.metadata.project_id", project_id) - span.set_attribute("langfuse.trace.tags", ["everos", "memory"]) - - # -- write path ---------------------------------------------------------- - - def add(self, session_id: str, messages: list[dict], user_id: str | None = None, - app_id: str = "default", project_id: str = "default") -> dict: - with self._tracer.start_as_current_span("everos.memory.add") as span: - self._common(span, session_id=session_id, user_id=user_id, - app_id=app_id, project_id=project_id, op="add") - span.set_attribute("langfuse.observation.input", _j(messages)) - resp = self._t("/api/v1/memory/add", { - "session_id": session_id, "app_id": app_id, - "project_id": project_id, "messages": messages, - }) - span.set_attribute("langfuse.observation.output", _j(resp["data"])) - span.set_attribute("everos.buffer.status", resp["data"]["status"]) - return resp - - def flush(self, session_id: str, user_id: str | None = None, - app_id: str = "default", project_id: str = "default") -> dict: - """Boundary -> LLM extraction -> markdown persist -> index sync.""" - with self._tracer.start_as_current_span("everos.memory.flush") as span: - self._common(span, session_id=session_id, user_id=user_id, - app_id=app_id, project_id=project_id, op="flush") - resp = self._t("/api/v1/memory/flush", { - "session_id": session_id, "app_id": app_id, "project_id": project_id, - }) - detail = resp.get("_detail", {}) - - # Extraction (generation w/ model+tokens), markdown persist, and the - # async index trace all describe server-internal facts the HTTP API - # doesn't expose yet. Emit them with the mock / once native - # instrumentation ships; skip on a real server rather than fabricate. - # The top-level everos.memory.flush span (real latency + output) is - # always emitted. - if detail: - # 1. LLM extraction, a *generation*: model + token usage. - # EverOS does not compute cost; Langfuse derives it from - # model + usage in its model-usage views. - with self._tracer.start_as_current_span("everos.extract") as g: - self._common(g, session_id=session_id, user_id=user_id, - app_id=app_id, project_id=project_id, - obs_type="generation", op="extract") - g.set_attribute("gen_ai.request.model", detail.get("model", "gpt-4.1-mini")) - g.set_attribute("langfuse.observation.input", _j(detail.get("buffered_messages", []))) - g.set_attribute("langfuse.observation.output", _j(detail.get("memory_cell", {}))) - usage = detail.get("usage", {}) - g.set_attribute("gen_ai.usage.input_tokens", usage.get("input", 0)) - g.set_attribute("gen_ai.usage.output_tokens", usage.get("output", 0)) - time.sleep(detail.get("extract_s", 0.05)) - - # 2. Markdown persistence (atomic tmp+fsync+rename), strong consistency - with self._tracer.start_as_current_span("everos.persist.markdown") as p: - self._common(p, session_id=session_id, user_id=user_id, - app_id=app_id, project_id=project_id, op="persist") - p.set_attribute("langfuse.observation.output", - _j({"md_files": detail.get("md_files", [])})) - time.sleep(0.008) - - span.set_attribute("langfuse.observation.output", _j(resp["data"])) - - # 3. Index sync runs AFTER the API call returns, in EverOS's async - # "cascade" daemon (file watcher + debounce + entry diff -> LanceDB). - # It is therefore emitted as its OWN short-lived trace, correlated - # to the originating write by session_id, not as a child span. - # Server-internal, so mock / native only. - if detail: - with self._tracer.start_as_current_span("everos.cascade.index") as ix: - self._common(ix, session_id=session_id, user_id=user_id, - app_id=app_id, project_id=project_id, op="index") - ix.set_attribute("langfuse.observation.input", - _j({"triggered_by": "markdown change", - "correlates_to_session": session_id})) - ix.set_attribute("langfuse.observation.output", - _j({"rows_indexed": detail.get("rows_indexed", 0), - "index_lag_ms": detail.get("index_lag_ms", 500)})) - time.sleep(0.02) - - return resp - - # -- read path ----------------------------------------------------------- - - def search(self, query: str, user_id: str | None = None, agent_id: str | None = None, - top_k: int = 5, app_id: str = "default", project_id: str = "default", - session_id: str | None = None, hit_threshold: float = 0.6) -> dict: - with self._tracer.start_as_current_span("everos.memory.search") as span: - self._common(span, session_id=session_id, user_id=user_id, agent_id=agent_id, - app_id=app_id, project_id=project_id, - obs_type="retriever", op="search") - span.set_attribute("langfuse.observation.input", - _j({"query": query, "top_k": top_k, "method": "hybrid"})) - ctx = span.get_span_context() - trace_id_hex = format(ctx.trace_id, "032x") - retriever_obs_id = format(ctx.span_id, "016x") - - payload = {"query": query, "method": "hybrid", "top_k": top_k, - "app_id": app_id, "project_id": project_id} - if user_id: - payload["user_id"] = user_id - if agent_id: - payload["agent_id"] = agent_id - resp = self._t("/api/v1/memory/search", payload) - detail = resp.get("_detail", {}) - - # embed / hybrid_recall / rerank describe INTERNAL pipeline stages the - # HTTP API doesn't expose yet. Emit them with the mock / once native - # instrumentation ships; skip on a real server rather than fabricate. - if detail: - # 1. Query embedding - with self._tracer.start_as_current_span("everos.search.embed_query") as e: - self._common(e, session_id=session_id, user_id=user_id, agent_id=agent_id, - app_id=app_id, project_id=project_id, - obs_type="embedding", op="embed") - e.set_attribute("gen_ai.request.model", - detail.get("embed_model", "Qwen/Qwen3-Embedding-4B")) - e.set_attribute("langfuse.observation.input", _j(query)) - # compact output — never dump the raw vector into telemetry - e.set_attribute("langfuse.observation.output", - _j({"embedding_dims": detail.get("embed_dims", 2560)})) - e.set_attribute("gen_ai.usage.input_tokens", detail.get("embed_tokens", 0)) - time.sleep(detail.get("embed_s", 0.03)) - - # 2. Hybrid recall: single LanceDB query = BM25 + vector ANN + filter - with self._tracer.start_as_current_span("everos.search.hybrid_recall") as h: - self._common(h, session_id=session_id, user_id=user_id, agent_id=agent_id, - app_id=app_id, project_id=project_id, - obs_type="retriever", op="recall") - h.set_attribute("langfuse.observation.input", - _j({"bm25": True, "vector_ann": True, "filters": None})) - h.set_attribute("langfuse.observation.output", - _j({"candidates": detail.get("candidates", 0)})) - time.sleep(detail.get("recall_s", 0.03)) - - # 3. Rerank (cross-encoder) - with self._tracer.start_as_current_span("everos.search.rerank") as r: - self._common(r, session_id=session_id, user_id=user_id, agent_id=agent_id, - app_id=app_id, project_id=project_id, op="rerank") - r.set_attribute("langfuse.observation.metadata.rerank_model", - detail.get("rerank_model", "Qwen/Qwen3-Reranker-4B")) - r.set_attribute("langfuse.observation.output", _j(detail.get("ranked", []))) - time.sleep(detail.get("rerank_s", 0.05)) - - # Recall quality is derivable from the REAL response — every hit - # carries a fused/reranked score — so it works against a live server - # today, not just the mock. None means a miss (nothing scored). - top_score = _top_score_from_data(resp["data"]) - span.set_attribute("langfuse.observation.output", _j(resp["data"])) - if top_score is not None: - span.set_attribute("everos.search.top_score", top_score) - span.set_attribute("everos.search.hit", top_score >= hit_threshold) - else: - # Nothing scored came back: a genuine miss. Record hit so it - # still counts in recall hit-rate; no top_score (no hit to score). - span.set_attribute("everos.search.hit", False) - - # Recall-quality -> Langfuse scores (visible in evals/dashboards). - # Pushed AFTER the span closes so exporter/network time never - # inflates the measured search latency. - if top_score is not None: - pushed = push_score(trace_id_hex, "recall_top_score", top_score, - observation_id=retriever_obs_id, - comment="fused+reranked score of top memory hit") - push_score(trace_id_hex, "recall_hit", - 1.0 if top_score >= hit_threshold else 0.0, - observation_id=retriever_obs_id, - comment=f"top_score >= {hit_threshold}") - resp["_scores_pushed"] = pushed - else: - # Miss: record hit=0 so empty recalls still count in hit-rate; - # no top_score is pushed (there is no hit to score). - resp["_scores_pushed"] = push_score( - trace_id_hex, "recall_hit", 0.0, - observation_id=retriever_obs_id, - comment=f"no hit >= {hit_threshold} (empty recall)") - resp["_trace_id"] = trace_id_hex - return resp - - # -- self-evolution (OME / reflection) ------------------------------------ - - def trigger_ome(self, strategy: str = "reflect_episodes", - user_id: str | None = None, session_id: str | None = None) -> dict: - with self._tracer.start_as_current_span(f"everos.ome.{strategy}") as span: - self._common(span, session_id=session_id, user_id=user_id, - obs_type="agent", op="reflect") - span.set_attribute("langfuse.observation.input", _j({"strategy": strategy})) - resp = self._t("/api/v1/ome/trigger", {"name": strategy, "force": True}) - detail = resp.get("_detail", {}) - - # The consolidation generation (model + tokens) is server-internal; - # emit it with the mock / once native instrumentation ships, skip on - # a real server. The top-level everos.ome. agent span (real - # latency + output) is always emitted. - if detail: - with self._tracer.start_as_current_span("everos.reflect.consolidate") as g: - self._common(g, session_id=session_id, user_id=user_id, - obs_type="generation", op="consolidate") - g.set_attribute("gen_ai.request.model", detail.get("model", "gpt-4.1-mini")) - g.set_attribute("langfuse.observation.input", - _j(detail.get("episodes_in", []))) - g.set_attribute("langfuse.observation.output", - _j(detail.get("consolidated", {}))) - usage = detail.get("usage", {}) - g.set_attribute("gen_ai.usage.input_tokens", usage.get("input", 0)) - g.set_attribute("gen_ai.usage.output_tokens", usage.get("output", 0)) - time.sleep(detail.get("reflect_s", 0.08)) - - span.set_attribute("langfuse.observation.output", _j(resp["data"])) - return resp - - -class HTTPTransport: - """Real transport for a running EverOS server (pip install everos).""" - - def __init__(self, base_url: str = "http://127.0.0.1:8000"): - self.base_url = base_url.rstrip("/") - - def __call__(self, path: str, payload: dict) -> dict: - r = requests.post(f"{self.base_url}{path}", json=payload, timeout=180) - r.raise_for_status() - return r.json()