docs(examples): replace Langfuse wrapper with native tracing example
The examples/langfuse wrapper (everos_langfuse.py) was the interim client-side instrumentation before EverOS gained native OpenTelemetry export. Now that [observability] emits real OTLP spans, the wrapper is redundant and its faked child spans could mislead. Replace it with a minimal, dependency-light example: enable [observability] in everos.toml, run the server, and drive one add/flush/search cycle (demo.py, stdlib only) to see native traces in Langfuse. Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
parent
e85db903a2
commit
869dc67804
|
|
@ -1,3 +1,2 @@
|
|||
spans.jsonl
|
||||
__pycache__/
|
||||
.venv/
|
||||
|
|
|
|||
|
|
@ -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.<strategy>` + generation |
|
||||
| `POST /api/v1/memory/search` | retriever `everos.memory.search` → `recall` / `rank` |
|
||||
| query / recall embedding | embedding `everos.embedding` |
|
||||
| OME reflection strategies | agent `everos.ome.<strategy>` (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
|
||||
|
||||
With a server running and `[observability]` enabled:
|
||||
|
||||
```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
|
||||
|
||||
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`.
|
||||
|
|
|
|||
|
|
@ -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
|
||||
import urllib.request
|
||||
|
||||
from everos_langfuse import HTTPTransport, InstrumentedEverOS, force_flush, init_tracing
|
||||
|
||||
TS = int(time.time() * 1000)
|
||||
DAY = "20260702"
|
||||
BASE = "http://127.0.0.1:8000"
|
||||
SESSION = "langfuse_demo"
|
||||
USER = "alice"
|
||||
|
||||
|
||||
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,
|
||||
},
|
||||
)
|
||||
|
||||
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}")
|
||||
ts = int(time.time() * 1000)
|
||||
|
||||
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)")
|
||||
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"])
|
||||
|
||||
# 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"
|
||||
flush = _post("/api/v1/memory/flush", {"session_id": SESSION, "messages": []})
|
||||
print("flush ->", flush["data"])
|
||||
|
||||
# -- 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)
|
||||
print("waiting for async index sync ...")
|
||||
time.sleep(10)
|
||||
|
||||
# -- 2. boundary/flush: LLM extraction -> markdown -> index ------------
|
||||
ev.flush(session, user_id=user)
|
||||
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)")
|
||||
|
||||
# -- 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.")
|
||||
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__":
|
||||
|
|
|
|||
|
|
@ -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.<strategy>" (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.<strategy> 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()
|
||||
Loading…
Reference in New Issue