@volter/twin-langfuse
v0.1.37
Published
Local Langfuse twin — a faithful, stateful local Langfuse LLM-observability API your real `langfuse` SDK (or `@langfuse/otel`) talks to unmodified. Ingest batched/OTLP trace events into a queryable local trace store, browse the Public API + trace-explorer
Downloads
4,963
Readme
@volter/twin-langfuse
A local Langfuse twin — a faithful, stateful local Langfuse LLM-observability API your real
langfuse SDK (or @langfuse/tracing + @langfuse/otel) talks to unmodified. Point the
SDK's baseUrl at the twin and a normal langfuse.trace(...) / trace.generation(...) flushes
into a queryable local trace store; browse traces, observation trees, sessions, scores,
prompts and datasets through the real Public API (/api/public/...) or the trace-explorer
mirror. Mirror real Langfuse, simulate, and fork. Built on @volter/world-core.
It is two API families over one local event-sourced state (the @volter/world-core kernel):
- Ingestion (the SDK side) —
POST /api/public/ingestion, the batched event envelope every Langfuse SDK flushes, andPOST /api/public/otel/v1/traces, the OTLP endpoint@langfuse/otelexports to. Events fold into traces + observations + sessions + scores: atrace-createupserts by body id (that is how an SDK "updates" a trace), aspan-create/generation-createauto-creates its trace if it hasn't arrived yet,*-updateevents merge onto the existing observation, generations are priced from the project's model definitions intocostDetails/calculatedTotalCost(rolled up into the trace'stotalCost), and the batch answers 207 with per-eventsuccesses/errors— never a 4xx for a bad event, exactly as the vendor documents. Envelope ids deduplicate, so an SDK's at-least-once retry is safe. - The Public API (the app/dashboard side,
/api/public/...) — traces (list/filter/get/ delete), observations (v1 + v2), sessions, scores + score configs, prompt management with versioning and exclusive deployment labels, datasets + items + experiment runs, model definitions, comments, annotation queues, media handshakes, LLM connections, projects + API keys, a metrics query subset, and health.
Auth is HTTP Basic (publicKey:secretKey) on every route except GET /api/public/health,
which Langfuse leaves unauthenticated for load-balancer probes. No real Langfuse is ever
contacted on the serve path — it's fully local and offline. Real vendor I/O happens only in
the connector, over an injected client (a fake in tests, a real key-authed client in prod).
readOnly rejects writes (405). An unmodeled operation fails like Langfuse (4xx JSON) —
never a fabricated success.
Quick start
# serve the twin (ingestion + OTLP + Public API); prints the SDK env to copy
bunx world-langfuse serve --port 9000
# LANGFUSE_BASEURL=http://127.0.0.1:9000
# LANGFUSE_PUBLIC_KEY=pk-lf-00000000-0000-4000-8000-000000000000
# LANGFUSE_SECRET_KEY=sk-lf-00000000-0000-4000-8000-000000000000
# the trace-explorer mirror (React)
bunx world-langfuse mirror --port 9001
# offline shape conformance
bunx world-langfuse conformanceimport { Langfuse } from 'langfuse';
const langfuse = new Langfuse({
publicKey: 'pk-lf-00000000-0000-4000-8000-000000000000',
secretKey: 'sk-lf-00000000-0000-4000-8000-000000000000',
baseUrl: 'http://127.0.0.1:9000',
});
const trace = langfuse.trace({ name: 'support-agent', userId: 'ada' });
const gen = trace.generation({ name: 'answer', model: 'gpt-4o' });
gen.end({ output: 'Rayleigh scattering.', usage: { input: 100, output: 50 } });
await langfuse.flushAsync(); // → a queryable trace, priced, in the twinCoverage
The honest denominator is the capability manifest (src/langfuse-capabilities.ts),
enumerated top-down from Langfuse's own first-party OpenAPI document
(https://cloud.langfuse.com/generated/api/openapi.yml, read 2026-08-19 — 114 declared
operations across 22 tag groups), ranked core → common → niche. Every done capability is proven
by an offline verify() that drives a real ingest/create → read → assert against the twin
and checks the vendor's negative 4xx paths. The baseline test
(src/langfuse-capabilities.test.ts) asserts via assertManifestBaseline: total ≥ 50, done > 0,
0 regressions.
Coverage is partial and honest — a large majority of the read/write surface an application
integrates against is modeled; the admin, SCIM, /unstable/ dashboard-and-evaluator, and
successor-version (v2/metrics, v3/scores) surfaces are enumerated as todo at the same
per-operation granularity as the done entries, so the percentage is not inflated by a coarse
todo list.
Modeled (core → niche):
- Ingestion — the batched
/api/public/ingestionenvelope with its 207 + per-eventsuccesses/errorscontract, every declared event type (trace-create,span-create/ -update,generation-create/-update,event-create,observation-create/-update,score-create,sdk-log), envelope-id deduplication, trace upsert-by-id merge, auto-created traces for orphan observations, auto-created sessions, environment validation, usage normalization (modernusageDetails, legacyusage, OpenAI'sprompt_tokens/completion_tokens, derivedtotal), cost calculation from matching model definitions, derived latency + time-to-first-token, observation levels, Basic auth,readOnly. - OTLP —
/api/public/otel/v1/traces(OTLP/JSON only — see the caveat below): root span → trace, child spans → nested observations, attribute decoding, observation-type inference (langfuse.observation.type, else the GenAI semconv),user.id/session.idmapping,gen_ai.usage.*token mapping,STATUS_CODE_ERROR→ an ERROR-level observation, andpartialSuccess.rejectedSpansfor unusable spans. - Traces — list with the
{data, meta}envelope, filters (userId/name/sessionId/tags(ALL must match)/version/release/environment/time window),orderBy, pagination, full detail (observations + scores inline,htmlPath, derivedlatency/totalCost), delete (cascading to observations + scores), bulk delete, and a clean re-create after delete. - Observations — list + get (v1 and v2), filters by type/name/traceId/model/level/userId
(resolved through the parent trace)/start-time window, and the full
ObservationsViewshape including prices and calculated costs. - Sessions — list + get (with the session's traces), environment + time filters.
- Scores — create (NUMERIC / CATEGORICAL / BOOLEAN discrimination), list with name/source/dataType/traceId/value+operator filters, get, delete, subjects beyond a trace (session, dataset run), and score-config binding (matching name + numeric range enforced).
- Score configs — create (with categories / min-max validation), list, get, PATCH-archive.
- Prompt management — text + chat prompts, versioning (re-creating a name mints vn+1),
exclusive deployment labels (a label lives on exactly one version; creating or PATCHing
moves it),
production-then-newest default resolution, per-name meta listing with label/tag filters, validation, delete-all-versions. - Datasets + experiments — dataset upsert-by-name, get, items (create/list/get/delete,
sourceTraceIdfilter), dataset run items (the run is created on first use), run get with its items, run delete cascading to its items. - Models — custom definitions with modern
pricesmap or legacy flat scalars,matchPatternvalidation (including Postgres's(?i)inline-flag idiom the vendor documents),startDatescoping, list/get/delete, and the pricing that flows into every generation. - Comments — create against a real trace/observation/session/prompt, validation, filtered list.
- Annotation queues — create bound to real score configs, item lifecycle (create → list → get → PATCH to COMPLETED → delete), status filter, negative paths.
- Media — the presigned-upload handshake (idempotent by content hash,
uploadUrl: nullonce uploaded), get, PATCH the upload result. - LLM connections — upsert-by-provider with a masked
displaySecretKey(the secret is never echoed), list, delete. - Projects + API keys — the key's own project, key listing, minting (plaintext secret shown exactly once, and the minted pair really authenticates), revocation (the pair stops working).
- Metrics —
GET /api/public/metricswith a real aggregate subset:traces/observationsviews,count/sum/avg/min/maxover count/latency/totalCost/totalTokens/value, one dimension, filters and a time window — and a 400 for anything unsupported rather than fabricated zeros. - Mirror UI — Tracing stream, Trace detail with the observation tree (spans nested under
parents, orphans still visible) + model/token/cost badges + score chips + I/O preview,
Sessions, Scores, Prompts, Datasets, Models, Settings — all served over the same
/api/public/endpoints the API serves (one read path, so API↔UI parity cannot drift), and two committed Playwright journeys. - Connector — pull (traces/observations/sessions/scores/datasets/models/prompts) + push (scores, datasets, dataset items) over an injected executor, with confirm + idempotent re-push.
Planned (todo), beyond the per-operation gaps already enumerated in the manifest:
- Media blob bytes — the presigned-URL handshake is modeled; storing the uploaded bytes and
serving them back at that URL is
langfuse.media.blob_storage. - Data-retention expiry —
retentionDaysis stored faithfully; traces older than it do not yet drop out of reads (langfuse.projects.retention_expiry). - Rate-limit enforcement — the client half of the contract is implemented and tested
(see below); the server's 429 +
Retry-Afterislangfuse.health.rate_limit_enforcement.
The Playground's model call and LLM-as-a-judge grades come from a live model: point the app at the
openai/anthropic twins for that plane. The twin's read path folds an exact projection over
the local event log rather than reproducing ClickHouse's approximate aggregation — better for
tests, not the same engine.
Caveat: OTLP/JSON only
The OTLP endpoint accepts OTLP/JSON. @langfuse/otel's OTLP HTTP exporter can be configured
to send application/x-protobuf, which this twin does not decode — it answers a loud 400
rather than a fabricated success, and the gap is tracked as a core todo
(langfuse.otel.protobuf). Configure the exporter for JSON, or use the langfuse SDK's batched
/api/public/ingestion path, which is fully modeled.
Error envelope: modeled, not spec-derived
Langfuse's OpenAPI declares every 4xx response body as schema: {} (unspecified). The status
codes below are therefore spec-derived, but the body shape the twin serves —
{ message, error }, where error is the server error-class name
(LangfuseNotFoundError / UnauthorizedError / InvalidRequestError / MethodNotAllowedError)
— is modeled from Langfuse's server error classes and is documented here as modeled rather
than claimed as spec-conformant.
Architecture
State and execution follow the shared model and architecture. The kernel owns storage and execution mechanics; this package owns vendor semantics. Conformance tooling is a development dependency. The generated index records this package's protocol standing.
What the connector will and will not push. Traces and observations are ingested telemetry: re-emitting a locally forked trace into a real project would corrupt the customer's observability data, so they are deliberately not pushable. What is pushable is the human/eval layer an operator authors against a twin — scores, datasets and dataset items. Anything else fails loudly rather than hitting a wrong real endpoint.
Rate budget — the fail-closed backstop on live calls
liveLangfuseExecute is the one place this pack issues a live request, so every call it makes
is charged against a persistent, fail-closed spend ledger before the request goes out. Past
the ceiling, or while a Retry-After/429 cooldown is armed, it throws instead of calling. The
ledger is keyed by vendor and a hash of the credential (Langfuse's limits are per organization,
attached to the key, so it is deliberately not cwd-scoped) and persists across processes, so a
fresh process does not get a fresh allowance; a corrupt ledger counts as a full window rather
than zero spend. There is no option to disable it, and no value you can pass for budget that
yields an unguarded client — an injected budget is validated by method identity, so a subclass
or a Proxy that replaces checkBudget is refused.
Unusually for this repo, Langfuse publishes scalar limits (langfuse.com/faq/all/api-limits, read 2026-08-19), so every number here is arithmetic against that table at its lowest plan (Hobby) — a twin cannot know an operator's plan, and guessing upward is how a lockout happens:
| bucket | published (Hobby) | this declaration |
|---|---|---|
| General API | 30 req/min | 60 units / 60s at weight 2 = 30/min (also exactly the kernel fallback) |
| Deprecated reads (GET /traces, /traces/{id}, /observations, /observations/{id}, /sessions, /sessions/{id}) | 15 req/min | weight 4 = 15/min — stricter than the fallback, because the vendor documents it |
| Trace-data deletion / Metrics API | 50–100 req/day | not expressible in a 60s window → priced at the full ceiling (≤1 per minute) as a mitigation, not a guarantee |
The daily buckets are stated plainly rather than dressed up as compliance; the genuinely
protective mechanism for them is the cooldown — Langfuse's documented 429 + Retry-After
becomes a persisted client-side stop instead of a retry storm.
The mechanism is shared and vendor-agnostic — it lives in the kernel (@volter/world-core →
packages/world-core/src/rateBudget.ts); what lives here in
src/langfuse-budget.ts is this vendor's declaration (window,
ceiling, per-endpoint weights, and a reason citing the table above) plus the vendor-bound
LangfuseBudget. The rule is ratified as ../../../docs/contributing/architecture.md D8, and
the kernel module's header documents what the guard does not guarantee — read that before
trusting it.
