import type { LangfuseRuntime, LangfuseScoreClient, PendingScore } from "./types.js"; import { state } from "./state.js"; import { randomUUID } from "node:crypto"; let runtime: LangfuseRuntime | null = null; let registeredContextManager: OtelContextManager | null = null; const activeSessions = new Set(); let lastRuntimeError: { scope: string; message: string; timestamp: Date } | null = null; type FallbackObservationType = "SPAN" | "GENERATION"; interface OtelContextManager { enable(): OtelContextManager; disable(): void; } interface OtelContextApi { setGlobalContextManager(contextManager: OtelContextManager): boolean; } type AsyncHooksContextManagerCtor = new () => OtelContextManager; interface RestFallbackTrace { id: string; timestamp: string; name: string; input?: unknown; output?: unknown; sessionId?: string; metadata?: Record; } interface RestFallbackObservation { id: string; traceId: string; type: FallbackObservationType; name: string; startTime: string; endTime?: string; parentObservationId?: string; input?: unknown; output?: unknown; metadata?: Record; model?: string; modelParameters?: Record; usageDetails?: Record; costDetails?: Record; level?: "DEBUG" | "DEFAULT" | "WARNING" | "ERROR"; statusMessage?: string; completionStartTime?: string; } interface RestFallbackStore { trace?: RestFallbackTrace; observations: RestFallbackObservation[]; observationById: Map; attempted: boolean; } const OTEL_VISIBILITY_TIMEOUT_MS = 1_500; const OTEL_VISIBILITY_POLL_INTERVAL_MS = 200; const DEFAULT_SHUTDOWN_STEP_TIMEOUT_MS = 2_000; const DEFAULT_LANGFUSE_REQUEST_TIMEOUT_SECONDS = 5; const DEFAULT_SCORE_FLUSH_AT = 10; const DEFAULT_SCORE_FLUSH_INTERVAL_MS = 1_000; const MAX_SCORE_QUEUE_SIZE = 100_000; const MAX_SCORE_BATCH_SIZE = 100; let shutdownStepTimeoutMs = DEFAULT_SHUTDOWN_STEP_TIMEOUT_MS; function nowIso() { return new Date().toISOString(); } function resolvePositiveEnvNumber(name: string, fallback: number, integer = false): number { const parsed = Number(process.env[name]); if (!Number.isFinite(parsed) || parsed <= 0) { return fallback; } return integer ? Math.floor(parsed) : parsed; } function delay(ms: number, signal?: AbortSignal) { return new Promise((resolve, reject) => { if (signal?.aborted) { reject(signal.reason); return; } const timeout = setTimeout(resolve, ms); signal?.addEventListener("abort", () => { clearTimeout(timeout); reject(signal.reason); }, { once: true }); }); } function debugLog(message: string) { if (process.env.PI_LANGFUSE_DEBUG === "1" || process.env.PI_LANGFUSE_DEBUG === "true") { console.log(message); } } export function ensureOtelContextManager( contextApi: OtelContextApi, AsyncHooksContextManager: AsyncHooksContextManagerCtor, ): boolean { if (registeredContextManager) { return true; } const contextManager = new AsyncHooksContextManager().enable(); if (contextApi.setGlobalContextManager(contextManager)) { registeredContextManager = contextManager; return true; } contextManager.disable(); return false; } function rememberRuntimeError(scope: string, error: unknown) { lastRuntimeError = { scope, message: error instanceof Error ? error.message : String(error), timestamp: new Date(), }; } export function getLastRuntimeError(): { scope: string; message: string; timestamp: Date } | null { return lastRuntimeError; } async function withShutdownDeadline(label: string, startOperation: () => Promise | undefined, deadline: number): Promise { const remainingMs = deadline - Date.now(); if (remainingMs <= 0) { debugLog(`📊 Langfuse: Skipped ${label}; shutdown deadline elapsed`); return undefined; } const operation = startOperation(); if (!operation) { return undefined; } let timeout: NodeJS.Timeout | undefined; try { return await Promise.race([ operation, new Promise((resolve) => { timeout = setTimeout(() => { debugLog(`📊 Langfuse: ${label} timed out; shutdown deadline elapsed`); resolve(undefined); }, remainingMs); }), ]); } finally { if (timeout) { clearTimeout(timeout); } } } function getRuntimeConfig(rt: LangfuseRuntime) { return rt.runtimeConfig ?? state.config; } function ingestionHeaders(rt: LangfuseRuntime): Record { const config = getRuntimeConfig(rt); if (!config) { throw new Error("Langfuse runtime config is unavailable"); } const auth = Buffer.from(`${config.publicKey}:${config.secretKey}`).toString("base64"); return { Authorization: `Basic ${auth}`, "Content-Type": "application/json", }; } async function ingestBatch(rt: LangfuseRuntime, batch: unknown[], signal: AbortSignal): Promise { const config = getRuntimeConfig(rt); if (!config) { throw new Error("Langfuse runtime config is unavailable"); } const response = await fetch(`${config.host.replace(/\/$/, "")}/api/public/ingestion`, { method: "POST", headers: ingestionHeaders(rt), body: JSON.stringify({ batch }), signal, }); if (!response.ok) { throw new Error(`Langfuse ingestion failed with HTTP ${response.status}`); } const text = await response.text(); if (!text) { return []; } const responseBody = JSON.parse(text) as { errors?: unknown[] }; return Array.isArray(responseBody.errors) ? responseBody.errors : []; } async function flushPendingScores(rt: LangfuseRuntime, signal: AbortSignal): Promise { const pendingScores = rt.pendingScores; if (!pendingScores || pendingScores.length === 0) { return; } while (pendingScores.length > 0) { const scores = pendingScores.slice(0, MAX_SCORE_BATCH_SIZE); try { const errors = await ingestBatch( rt, scores.map((score) => ({ type: "score-create", id: randomUUID(), timestamp: nowIso(), body: score, })), signal, ); pendingScores.splice(0, scores.length); if (errors.length > 0) { rememberRuntimeError("score ingestion", new Error(JSON.stringify(errors))); console.warn("📊 Langfuse: Score ingestion reported errors", errors); } } catch (error) { if ((error as { name?: string }).name !== "AbortError") { rememberRuntimeError("score ingestion", error); console.warn("📊 Langfuse: Failed to flush scores", error); } return; } } } function clearScoreFlushTimer(rt: LangfuseRuntime) { if (rt.scoreFlushTimer) { clearTimeout(rt.scoreFlushTimer); rt.scoreFlushTimer = undefined; } } function scheduleScoreFlush(rt: LangfuseRuntime) { if ( rt.scoreFlushStopped || rt.scoreFlushTimer || rt.scoreFlushPromise || !rt.pendingScores?.length ) { return; } rt.scoreFlushTimer = setTimeout(() => { rt.scoreFlushTimer = undefined; void startScoreFlush(rt); }, rt.scoreFlushIntervalMs ?? DEFAULT_SCORE_FLUSH_INTERVAL_MS); rt.scoreFlushTimer.unref?.(); } function startScoreFlush(rt: LangfuseRuntime): Promise { if (rt.scoreFlushPromise) { return rt.scoreFlushPromise; } clearScoreFlushTimer(rt); const controller = new AbortController(); const timeout = setTimeout( () => controller.abort(new DOMException("Langfuse score request timed out", "AbortError")), rt.scoreRequestTimeoutMs ?? DEFAULT_LANGFUSE_REQUEST_TIMEOUT_SECONDS * 1_000, ); timeout.unref?.(); rt.scoreFlushController = controller; const promise = flushPendingScores(rt, controller.signal).finally(() => { clearTimeout(timeout); if (rt.scoreFlushPromise === promise) { rt.scoreFlushPromise = undefined; } if (rt.scoreFlushController === controller) { rt.scoreFlushController = undefined; } scheduleScoreFlush(rt); }); rt.scoreFlushPromise = promise; return promise; } function stopScoreFlush(rt: LangfuseRuntime) { rt.scoreFlushStopped = true; clearScoreFlushTimer(rt); rt.scoreFlushController?.abort( new DOMException("Langfuse score flushing stopped", "AbortError"), ); } function toIso(value: unknown): string | undefined { if (!value) { return undefined; } if (value instanceof Date) { return value.toISOString(); } if (typeof value === "string") { return value; } return undefined; } function mergeMetadata(current: Record | undefined, next: Record | undefined) { return next ? { ...(current ?? {}), ...next } : current; } function applyObservationUpdate(record: RestFallbackObservation, body: Record | undefined) { if (!body) { return; } if ("input" in body) record.input = body.input; if ("output" in body) record.output = body.output; if ("metadata" in body && body.metadata && typeof body.metadata === "object") { record.metadata = mergeMetadata(record.metadata, body.metadata as Record); } if (typeof body.model === "string") record.model = body.model; if (body.modelParameters && typeof body.modelParameters === "object") { record.modelParameters = body.modelParameters as Record; } if (body.usageDetails && typeof body.usageDetails === "object") { record.usageDetails = body.usageDetails as Record; } if (body.costDetails && typeof body.costDetails === "object") { record.costDetails = body.costDetails as Record; } if (typeof body.level === "string") record.level = body.level as RestFallbackObservation["level"]; if (typeof body.statusMessage === "string") record.statusMessage = body.statusMessage; const completionStartTime = toIso(body.completionStartTime); if (completionStartTime) record.completionStartTime = completionStartTime; } function applyTraceUpdate(store: RestFallbackStore, body: Record | undefined) { if (!store.trace || !body) { return; } if ("input" in body) store.trace.input = body.input; if ("output" in body) store.trace.output = body.output; if ("metadata" in body && body.metadata && typeof body.metadata === "object") { store.trace.metadata = mergeMetadata(store.trace.metadata, body.metadata as Record); } } function observationType(asType?: string): FallbackObservationType { return asType === "generation" ? "GENERATION" : "SPAN"; } function wrapObservation( observation: any, store: RestFallbackStore, name: string, body: Record | undefined, asType?: string, parentObservationId?: string, ): any { const id = observation.id || randomUUID(); const traceId = observation.traceId || store.trace?.id || randomUUID(); const metadata = body?.metadata && typeof body.metadata === "object" ? body.metadata as Record : undefined; const record: RestFallbackObservation = { id, traceId, name, type: observationType(asType), startTime: nowIso(), parentObservationId, metadata: mergeMetadata(metadata, asType && asType !== "generation" && asType !== "span" ? { langfuseObservationType: asType } : undefined), }; applyObservationUpdate(record, body); store.observations.push(record); store.observationById.set(id, record); if (!parentObservationId && !store.trace) { store.trace = { id: traceId, timestamp: record.startTime, name, input: body?.input, sessionId: typeof metadata?.sessionId === "string" ? metadata.sessionId : state.currentSessionId || undefined, metadata, }; } return { ...observation, id, traceId, update(updateBody?: Record) { applyObservationUpdate(record, updateBody); if (!parentObservationId) { applyTraceUpdate(store, updateBody); } const updated = observation.update(updateBody); return updated === observation ? this : updated; }, end(endBody?: Record) { if (endBody && typeof endBody === "object") { applyObservationUpdate(record, endBody); if (!parentObservationId) { applyTraceUpdate(store, endBody); } } record.endTime = nowIso(); return observation.end(); }, startObservation(childName: string, childBody?: Record, options?: { asType?: string }) { const child = observation.startObservation(childName, childBody, options); return wrapObservation(child, store, childName, childBody, options?.asType, id); }, setTraceIO(traceBody?: { input?: unknown; output?: unknown }) { applyTraceUpdate(store, traceBody); return observation.setTraceIO?.(traceBody); }, }; } async function traceExists(rt: LangfuseRuntime, traceId: string, signal: AbortSignal): Promise { const config = getRuntimeConfig(rt); if (!config) { return false; } try { const response = await fetch( `${config.host.replace(/\/$/, "")}/api/public/traces/${encodeURIComponent(traceId)}`, { headers: ingestionHeaders(rt), signal, }, ); if (response.status === 404) { return false; } if (!response.ok) { throw new Error(`Langfuse trace visibility check failed with HTTP ${response.status}`); } return true; } catch (error) { if (signal.aborted) { throw error; } return false; } } async function waitForTraceVisibility(rt: LangfuseRuntime, traceId: string, signal: AbortSignal): Promise { const deadline = Date.now() + OTEL_VISIBILITY_TIMEOUT_MS; while (true) { if (await traceExists(rt, traceId, signal)) { return true; } const remainingMs = deadline - Date.now(); if (remainingMs <= 0) { return false; } await delay(Math.min(OTEL_VISIBILITY_POLL_INTERVAL_MS, remainingMs), signal); } } function eventTimestamp(record: { endTime?: string; startTime?: string; timestamp?: string }) { return record.endTime ?? record.startTime ?? record.timestamp ?? nowIso(); } async function fallbackToRestIngestion(rt: LangfuseRuntime, signal: AbortSignal) { const store = rt.restFallback as RestFallbackStore | undefined; if (!store?.trace || store.attempted) { return; } store.attempted = true; if (await waitForTraceVisibility(rt, store.trace.id, signal)) { return; } const trace = store.trace; const batch: any[] = [ { type: "trace-create", id: randomUUID(), timestamp: eventTimestamp(trace), body: { id: trace.id, timestamp: trace.timestamp, name: trace.name, input: trace.input, output: trace.output, sessionId: trace.sessionId, metadata: trace.metadata, }, }, ]; for (const observation of store.observations) { const body = { id: observation.id, traceId: observation.traceId, name: observation.name, startTime: observation.startTime, endTime: observation.endTime, input: observation.input, output: observation.output, metadata: observation.metadata, level: observation.level, statusMessage: observation.statusMessage, parentObservationId: observation.parentObservationId, ...(observation.type === "GENERATION" ? { completionStartTime: observation.completionStartTime, model: observation.model, modelParameters: observation.modelParameters, usageDetails: observation.usageDetails, costDetails: observation.costDetails, } : {}), }; batch.push({ type: observation.type === "GENERATION" ? "generation-create" : "span-create", id: randomUUID(), timestamp: eventTimestamp(observation), body, }); } const errors = await ingestBatch(rt, batch, signal); if (errors.length > 0) { rememberRuntimeError("REST fallback ingestion", new Error(JSON.stringify(errors))); console.warn("📊 Langfuse: REST fallback ingestion reported errors", errors); } else { debugLog(`📊 Langfuse: OTel trace ${trace.id} was not visible; wrote fallback trace via REST ingestion`); } } export async function getRuntime(): Promise { if (!state.config) { throw new Error("Langfuse config is not set"); } // Track the current session as a runtime consumer. // Multiple sessions can share the same runtime; shutdown is deferred // until the last session releases it. const sessionId = state.currentSessionId; if (sessionId) { activeSessions.add(sessionId); } if (!runtime) { const [ { BasicTracerProvider }, { context }, { AsyncHooksContextManager }, { LangfuseSpanProcessor }, tracing, { LangfuseClient }, ] = await Promise.all([ import("@opentelemetry/sdk-trace-base"), import("@opentelemetry/api"), import("@opentelemetry/context-async-hooks"), import("@langfuse/otel"), import("@langfuse/tracing"), import("@langfuse/client"), ]); const restFallback: RestFallbackStore = { observations: [], observationById: new Map(), attempted: false, }; try { ensureOtelContextManager(context, AsyncHooksContextManager); const scoreFlushAt = resolvePositiveEnvNumber("LANGFUSE_FLUSH_AT", DEFAULT_SCORE_FLUSH_AT, true); const scoreFlushIntervalMs = resolvePositiveEnvNumber("LANGFUSE_FLUSH_INTERVAL", DEFAULT_SCORE_FLUSH_INTERVAL_MS / 1_000) * 1_000; const scoreRequestTimeoutMs = resolvePositiveEnvNumber("LANGFUSE_TIMEOUT", DEFAULT_LANGFUSE_REQUEST_TIMEOUT_SECONDS) * 1_000; const spanProcessor = new LangfuseSpanProcessor({ publicKey: state.config.publicKey, secretKey: state.config.secretKey, baseUrl: state.config.host, }); const tracerProvider = new BasicTracerProvider({ spanProcessors: [spanProcessor] }); tracing.setLangfuseTracerProvider(tracerProvider); runtime = { startObservation: ((name: string, body?: Record, options?: { asType?: string }) => { const observation = (tracing as any).startObservation(name, body, options); return wrapObservation(observation, restFallback, name, body, options?.asType); }) as unknown as LangfuseRuntime["startObservation"], propagateAttributes: tracing.propagateAttributes as unknown as LangfuseRuntime["propagateAttributes"], scoreClient: new LangfuseClient({ publicKey: state.config.publicKey, secretKey: state.config.secretKey, baseUrl: state.config.host, }) as LangfuseScoreClient, spanProcessor, tracerProvider, clearTracerProvider: () => tracing.setLangfuseTracerProvider(null), restFallback, pendingScores: [], scoreFlushAt, scoreFlushIntervalMs, scoreRequestTimeoutMs, scoreFlushStopped: false, runtimeConfig: { publicKey: state.config.publicKey, secretKey: state.config.secretKey, host: state.config.host, }, }; lastRuntimeError = null; } catch (e) { rememberRuntimeError("runtime init", e); throw e; } } return runtime as LangfuseRuntime; } function doShutdownRuntime(): Promise { return (async () => { if (!runtime) { return; } const rt = runtime; runtime = null; const deadline = Date.now() + shutdownStepTimeoutMs; const controller = new AbortController(); const abortTimeout = setTimeout(() => controller.abort(), shutdownStepTimeoutMs); stopScoreFlush(rt); try { await withShutdownDeadline( "Active score flush", () => rt.scoreFlushPromise, deadline, ); await withShutdownDeadline("OTel force flush", () => rt.tracerProvider?.forceFlush?.(), deadline); await withShutdownDeadline( "REST fallback ingestion", () => fallbackToRestIngestion(rt, controller.signal), deadline, ); await flushPendingScores(rt, controller.signal); await withShutdownDeadline("Langfuse score flush", () => rt.scoreClient.flush?.(), deadline); await withShutdownDeadline("Langfuse client shutdown", () => rt.scoreClient.shutdown?.(), deadline); await withShutdownDeadline("OTel tracer shutdown", () => rt.tracerProvider?.shutdown?.(), deadline); } catch (e) { rememberRuntimeError("runtime shutdown", e); console.warn("📊 Langfuse: Failed to flush/shutdown cleanly", e); } finally { clearTimeout(abortTimeout); clearScoreFlushTimer(rt); rt.scoreFlushController?.abort(); rt.scoreFlushController = undefined; rt.scoreFlushPromise = undefined; if (!runtime) { rt.clearTracerProvider?.(); } } })(); } /** * Release the current session's reference to the Langfuse runtime. * Only actually shuts down the runtime when the last session releases it. * Accepts an optional sessionId for use outside of withSession (e.g. deferred callbacks). */ export async function shutdownRuntime(sessionId?: string): Promise { const sid = sessionId ?? state.currentSessionId; if (sid) { activeSessions.delete(sid); } // Still have active sessions — keep the runtime alive. if (activeSessions.size > 0) { return; } await doShutdownRuntime(); } /** * Force-shutdown the Langfuse runtime regardless of active session references. * Used when the user manually reconfigures (e.g. /langfuse-setup) and needs * a fresh runtime with new credentials. */ export async function forceShutdownRuntime(): Promise { activeSessions.clear(); await doShutdownRuntime(); } export function __setRuntimeForTest(rt: LangfuseRuntime | null, timeoutMs = DEFAULT_SHUTDOWN_STEP_TIMEOUT_MS): void { if (runtime && runtime !== rt) { stopScoreFlush(runtime); } runtime = rt; if (rt) { rt.scoreFlushStopped = false; } shutdownStepTimeoutMs = timeoutMs; activeSessions.clear(); } export async function sendScore(name: string, value: number, options: { traceId?: string; observationId?: string } = {}) { try { const rt = await getRuntime(); const score: PendingScore = { name, value, dataType: name === "session_had_errors" || name === "tool_is_error" ? "BOOLEAN" : "NUMERIC", traceId: options.traceId, observationId: options.observationId, sessionId: options.traceId ? undefined : state.currentSessionId || undefined, ...(process.env.LANGFUSE_TRACING_ENVIRONMENT ? { environment: process.env.LANGFUSE_TRACING_ENVIRONMENT } : {}), }; if (!rt.pendingScores) { return; } if (rt.pendingScores.length >= MAX_SCORE_QUEUE_SIZE) { const error = new Error( `Langfuse score queue is full (${MAX_SCORE_QUEUE_SIZE}); dropping score`, ); rememberRuntimeError("score queue", error); console.warn(`📊 Langfuse: ${error.message}`); return; } rt.pendingScores.push(score); if (rt.pendingScores.length >= (rt.scoreFlushAt ?? DEFAULT_SCORE_FLUSH_AT)) { void startScoreFlush(rt); } else { scheduleScoreFlush(rt); } } catch (e) { rememberRuntimeError(`score ${name}`, e); console.warn(`📊 Langfuse: Failed to send score ${name}`, e); } }