From b417f4aa6a6fd0691e56caf2145c63bb3373cb97 Mon Sep 17 00:00:00 2001 From: Daryll Netscher Date: Mon, 14 Sep 2026 14:18:32 +0700 Subject: [PATCH] feat: improve trace completeness --- src/index.ts | 202 ++++++++++++++++++++++++++++++------- src/telemetry.ts | 114 +++++++++++++++++++++ test/completeness.test.ts | 111 ++++++++++++++++++++ test/system-prompt.test.ts | 3 +- 4 files changed, 391 insertions(+), 39 deletions(-) create mode 100644 src/telemetry.ts create mode 100644 test/completeness.test.ts diff --git a/src/index.ts b/src/index.ts index 5a8f290..778f029 100644 --- a/src/index.ts +++ b/src/index.ts @@ -14,6 +14,8 @@ import { } from "@langfuse/tracing"; import { type SpanContext, TraceFlags } from "@opentelemetry/api"; import { AlwaysOnSampler, NodeTracerProvider } from "@opentelemetry/sdk-trace-node"; +import { captureProviderPayload, contextRegistry, maskTelemetry, memoryBlockHashes, type ParentEnvelope } from "./telemetry.ts"; +export { captureProviderPayload } from "./telemetry.ts"; const EXTENSION_NAME = "@langfuse/pi-observability-plugin"; const EXTENSION_VERSION = "0.1.2"; @@ -35,6 +37,8 @@ const EMIT_IMAGE_MEDIA = (() => { return raw ? !["false", "0"].includes(raw) : true; })(); +// pi subagents are child processes that get process.env from the parent. pi +// itself has no trace propagation, so these variables connect the traces. const ENV_PARENT_TRACE_ID = "LANGFUSE_PI_PARENT_TRACE_ID"; const ENV_PARENT_SPAN_ID = "LANGFUSE_PI_PARENT_SPAN_ID"; const ENV_PARENT_SESSION_ID = "LANGFUSE_PI_PARENT_SESSION_ID"; @@ -109,6 +113,13 @@ export function loadConfig(): LangfuseConfig | undefined { }; } +/** + * Persistent config lives at `/langfuse.json` (usually + * `~/.pi/agent/langfuse.json`, keep it chmod 600) so plain `pi` in any project + * is traced without exporting env vars. Environment variables override the + * file for ad-hoc runs. Literal values only (no env interpolation, no + * command execution — a config file must not be able to run code). + */ function readConfigFile(): Partial> { try { const path = join(getAgentDir(), "langfuse.json"); @@ -152,27 +163,32 @@ export function truncateText(text: string): { text: string; meta: TruncationMeta const SECRET_REDACTION_MARK = "[redacted-langfuse-secret]"; const CYCLE_MARK = "[circular-ref]"; const LANGFUSE_KEY_TOKEN = String.raw`\b[sp]k-lf-[\w-]+\b`; +type TelemetryValue = boolean | number | string | null | undefined | TelemetryValue[] | { [key: string]: TelemetryValue }; function escapeRegExpLiteral(text: string): string { return text.replace(/[.*+?^${}()|[\]\\]/g, "\\$&"); } -export function createSecretRedactor(...extraSecrets: string[]): (value: unknown) => unknown { +/** Redact Langfuse API keys and the given literal secrets from arbitrarily shaped payloads; cycles collapse to a marker so the result survives JSON serialization. */ +export function createSecretRedactor(...extraSecrets: string[]): (value: unknown) => TelemetryValue { const alternatives = extraSecrets.filter((s) => s.length > 0).map(escapeRegExpLiteral); alternatives.push(LANGFUSE_KEY_TOKEN); const pattern = new RegExp(alternatives.join("|"), "g"); - const walk = (value: unknown, ancestors: readonly object[]): unknown => { + const walk = (value: unknown, ancestors: readonly object[]): TelemetryValue => { if (typeof value === "string") return value.replace(pattern, SECRET_REDACTION_MARK); - if (value === null || typeof value !== "object") return value; + if (value === null || typeof value === "boolean" || typeof value === "number") return value; + if (typeof value !== "object") return undefined; if (ancestors.includes(value)) return CYCLE_MARK; const chain = [...ancestors, value]; if (Array.isArray(value)) { - const items: unknown[] = []; + const items: TelemetryValue[] = []; for (const item of value) items.push(walk(item, chain)); return items; } - const fields: Record = {}; + const fields: { [key: string]: TelemetryValue } = {}; for (const [key, field] of Object.entries(value)) { + // defineProperty, not assignment: a key literally named "__proto__" + // must stay a data key instead of mutating the clone's prototype. Object.defineProperty(fields, key, { value: walk(field, chain), enumerable: true, @@ -188,16 +204,15 @@ export function createSecretRedactor(...extraSecrets: string[]): (value: unknown const redactLangfuseKeys = createSecretRedactor(); export function readSystemPrompt(ctx: { getSystemPrompt?: () => string | undefined }): string | undefined { - let raw: unknown; try { - raw = ctx.getSystemPrompt?.(); + const prompt = ctx.getSystemPrompt?.(); + return typeof prompt === "string" && prompt.trim() ? redactLangfuseKeys(prompt) as string : undefined; } catch { return undefined; } - if (typeof raw !== "string" || !raw.trim()) return undefined; - return redactLangfuseKeys(raw) as string; } +/** Extract plain text from a pi message content array. */ export function extractText(content: unknown): string { if (typeof content === "string") return content; if (!Array.isArray(content)) return ""; @@ -270,6 +285,12 @@ export function toDataUri(image: PiImagePart): string | undefined { return `data:${image.mimeType};base64,${data}`; } +/** + * Returns `text` unchanged when there are no images. With images it returns the + * OpenAI-style content parts (the text, then one `image_url` per image) that the + * Langfuse UI shows as a picture. Never truncate the result: a cut data URI is + * uploaded as a corrupt file. + */ export function toMultimodalContent( text: string, images: readonly PiImagePart[] | undefined, @@ -335,6 +356,10 @@ export function buildCostDetails(usage: PiUsage): Record | undef const details: Record = { total: cost.total }; if (cost.input > 0) details.input = cost.input; if (cost.output > 0) { + // pi prices every output token of a call at one rate (its tier selection + // reads only input-side tokens), so the reasoning share is exactly + // proportional. Deriving the non-reasoning bucket by subtraction keeps the + // two buckets summing bit-for-bit to the total pi reported. const { reasoning, canSplit } = resolveReasoningSplit(usage); if (canSplit) { const reasoningCost = cost.output * (reasoning / usage.output); @@ -378,7 +403,7 @@ function createRuntime( baseUrl: config.baseUrl, environment: config.environment, release: config.release, - mask: ({ data }) => redactSecrets(data), + mask: ({ data }) => maskTelemetry(redactSecrets(data)), shouldExportSpan: ({ otelSpan }) => isLangfuseSpan(otelSpan), }); // Trace fields are stamped here rather than on the root span: Langfuse reads @@ -420,7 +445,10 @@ interface PromptState { sawError: boolean; userText: string; turnImages: PiImagePart[]; + sessionId: string; systemPrompt?: string; + memoryInjection?: Record; + memoryRetrievals: Array<{ retrievalId: string; observationId: string }>; } const DEBUG = process.env.PI_LANGFUSE_DEBUG === "true"; @@ -446,9 +474,12 @@ export default function (pi: ExtensionAPI) { let gitBranch: string | undefined; let sessionHadImages = false; let fallbackTurnCounter = 0; - let lastPromptText = ""; let compactionStartedAt: Date | undefined; - const inheritedParent = readInheritedParent(); + const environmentParent = readInheritedParent(); + let inheritedParent = environmentParent; + let explicitParent: ParentEnvelope | undefined; + let isChildSession = Boolean(environmentParent); + let lastContext: ExtensionContext | undefined; const inheritedParentEnv: Record = { [ENV_PARENT_TRACE_ID]: process.env[ENV_PARENT_TRACE_ID], [ENV_PARENT_SPAN_ID]: process.env[ENV_PARENT_SPAN_ID], @@ -458,9 +489,17 @@ export default function (pi: ExtensionAPI) { const ensureRuntime = (): Runtime => { if (!runtime || runtime.shutdown) runtime = createRuntime(config, () => traceAttributes); + // In-process children share the Langfuse SDK's provider pointer. + setLangfuseTracerProvider(runtime.provider); return runtime; }; + const readExplicitParent = (): ParentEnvelope | undefined => { + let parent: ParentEnvelope | undefined; + pi.events.emit("pi:trace-parent-request", { reply: (value: ParentEnvelope) => { parent = value; } }); + return parent; + }; + const resolveTurnNumber = (ctx: ExtensionContext, promptText: string): number => { try { const entries = ctx.sessionManager.getEntries() as Array<{ @@ -506,6 +545,21 @@ export default function (pi: ExtensionAPI) { // Do not publish a root that the sampler dropped. Child spans would point // to a trace with no exported root. if (!(ctx.traceFlags & TraceFlags.SAMPLED)) return; + contextRegistry().set(sessionId, Object.freeze({ + traceId: ctx.traceId, spanId: ctx.spanId, + rootSessionId: inheritedParent?.sessionId ?? sessionId, + parentSessionId: explicitParent?.parentSessionId, + depth: inheritedParent?.depth ?? 0, + })); + if (isChildSession) { + const parent = explicitParent?.parentSessionId ? contextRegistry().get(explicitParent.parentSessionId) : undefined; + if (parent && explicitParent && parent.spanId === explicitParent.spanId && explicitParent.runId && explicitParent.childIndex !== undefined) { + contextRegistry().set(explicitParent.parentSessionId!, Object.freeze({...parent, + tracedChildren: [...new Set([...(parent.tracedChildren ?? []), `${explicitParent.runId}:${explicitParent.childIndex}`])], + })); + } + return; + } process.env[ENV_PARENT_TRACE_ID] = ctx.traceId; process.env[ENV_PARENT_SPAN_ID] = ctx.spanId; process.env[ENV_PARENT_SESSION_ID] = sessionId; @@ -515,6 +569,10 @@ export default function (pi: ExtensionAPI) { // A child must not attach to a turn that has ended. The inherited values // stay valid, so put them back instead of a plain delete. const withdrawParentContext = () => { + if (state && contextRegistry().get(state.sessionId)?.spanId === state.root.otelSpan.spanContext().spanId) { + contextRegistry().delete(state.sessionId); + } + if (isChildSession) return; for (const [key, value] of Object.entries(inheritedParentEnv)) { if (value === undefined) delete process.env[key]; else process.env[key] = value; @@ -573,6 +631,7 @@ export default function (pi: ExtensionAPI) { }; pi.on("session_start", async (_event, ctx) => { + lastContext = ctx; debug("session_start"); if (ctx.hasUI) ctx.ui.setStatus("langfuse", "langfuse ✓"); try { @@ -583,23 +642,35 @@ export default function (pi: ExtensionAPI) { } }); - pi.on("before_agent_start", (event, ctx) => { + const beginPrompt = (event: { prompt: string; images?: unknown }, ctx: ExtensionContext, trigger: string) => { + lastContext = ctx; debug("before_agent_start"); ensureRuntime(); - debug("runtime ready"); + debug("runtime ready"); + // A previous prompt that never settled (e.g. rapid re-prompt) is closed + // rather than leaked. if (state) finalizeRoot({ cancelled: true }); + explicitParent = readExplicitParent(); + // Missing explicit IDs mean orphan, not permission to borrow sibling env. + inheritedParent = explicitParent ? readInheritedParent({ + [ENV_PARENT_TRACE_ID]: explicitParent.traceId, + [ENV_PARENT_SPAN_ID]: explicitParent.spanId, + [ENV_PARENT_SESSION_ID]: explicitParent.rootSessionId ?? explicitParent.parentSessionId, + [ENV_PARENT_DEPTH]: String(explicitParent.depth ?? 1), + }) : environmentParent; + isChildSession = Boolean(explicitParent || inheritedParent); + const sessionId = ctx.sessionManager.getSessionId(); const turnNumber = resolveTurnNumber(ctx, event.prompt); const promptImages = extractImages(event.images); const { text: promptText, meta: userMeta } = truncateText(event.prompt); const userText = [promptText, ...promptImages.map(describeImage)].filter(Boolean).join("\n"); - lastPromptText = userText; - const isSubagent = Boolean(inheritedParent); + const isSubagent = isChildSession; traceAttributes = { ...(isSubagent ? {} : { [LangfuseOtelSpanAttributes.TRACE_NAME]: TRACE_NAME }), [LangfuseOtelSpanAttributes.TRACE_SESSION_ID]: isSubagent - ? (inheritedParent!.sessionId ?? sessionId) + ? (inheritedParent?.sessionId ?? explicitParent?.rootSessionId ?? sessionId) : sessionId, [LangfuseOtelSpanAttributes.TRACE_TAGS]: BASE_TAGS, ...(config.userId ? { [LangfuseOtelSpanAttributes.TRACE_USER_ID]: config.userId } : {}), @@ -615,12 +686,20 @@ export default function (pi: ExtensionAPI) { extension_version: EXTENSION_VERSION, session_id: sessionId, turn_number: turnNumber, + trigger, + parent_link_source: explicitParent ? (inheritedParent ? "explicit" : "unavailable") : inheritedParent ? "environment" : "none", cwd: ctx.cwd, user_text_meta: userMeta, ...(gitBranch ? { git_branch: gitBranch } : {}), ...(ctx.model ? { model: ctx.model.id, provider: ctx.model.provider } : {}), ...(isSubagent - ? { pi_subagent: true, subagent_depth: inheritedParent!.depth, parent_session_id: inheritedParent!.sessionId } + ? { pi_subagent: true, subagent_depth: explicitParent?.depth ?? inheritedParent?.depth, + parent_session_id: explicitParent?.parentSessionId ?? inheritedParent?.sessionId, + root_session_id: inheritedParent?.sessionId ?? explicitParent?.rootSessionId, + parent_trace_id: inheritedParent?.spanContext.traceId, + parent_observation_id: inheritedParent?.spanContext.spanId, + run_id: explicitParent?.runId, agent: explicitParent?.agent, child_index: explicitParent?.childIndex, + parent_tool_call_id: explicitParent?.parentToolCallId, source_run_id: explicitParent?.sourceRunId } : {}), }, }, @@ -636,47 +715,88 @@ export default function (pi: ExtensionAPI) { sawError: false, userText, turnImages: [...promptImages], + sessionId, + memoryRetrievals: [], }; publishParentContext(root, sessionId); debug("root created, turn", turnNumber); - }); + }; + pi.on("before_agent_start", (event, ctx) => beginPrompt(event, ctx, "before_agent_start")); pi.on("agent_start", (_event, ctx) => { - if (!state) return; + lastContext = ctx; + if (!state) { + const latest = [...ctx.sessionManager.getBranch()].reverse().find(entry => entry.type === "message"); + beginPrompt({ prompt: latest?.type === "message" ? extractText((latest.message as {content?: unknown}).content) : "" }, ctx, "agent_start"); + } const systemPrompt = readSystemPrompt(ctx); - if (!systemPrompt) return; + if (!state || !systemPrompt) return; state.systemPrompt = systemPrompt; + try { state.root.update({ metadata: { system_prompt: systemPrompt } }); } catch { /* tracing is best effort */ } + }); + + pi.events.on("hindsight:retrieval", (raw: unknown) => { + // Telemetry is optional: observer failures must not alter recall behavior. try { - state.root.update({ metadata: { system_prompt: systemPrompt } }); - } catch {} - debug("system prompt captured", systemPrompt.length); + if (!raw || typeof raw !== "object" || !lastContext) return; + const event = raw as Record; + if (event.version !== 1 || !["retrieval", "injection"].includes(String(event.phase))) return; + if (event.sessionId && event.sessionId !== lastContext.sessionManager.getSessionId()) return; + const standalone = !state; + if (!state) beginPrompt({prompt:""}, lastContext, "memory-event"); + if (!state) return; + ensureRuntime(); + const { results, query, ...details } = event; + const date = new Date(String(event.startedAt)); + const filters = event.filters && typeof event.filters === "object" ? event.filters as Record : {}; + const attributes = { + input: event.phase === "retrieval" ? {query, bank_id:event.bankId, tag_groups:event.tagGroups, filters:event.filters, budget:event.budget ?? filters.budget, max_tokens:event.maxTokens ?? filters.maxTokens} : undefined, + output: event.phase === "retrieval" ? {results, kept_ids:event.keptIds, injected_ids:event.injectedIds} : {injected:event.injected, rendered_hash:event.renderedHash, rendered_length:event.renderedLength}, + level: event.status === "error" || event.status === "timeout" ? "WARNING" as const : "DEFAULT" as const, + statusMessage: typeof event.error === "string" ? event.error : undefined, + metadata: {...details, bank_id:event.bankId, retrieval_id:event.retrievalId, context_id:event.contextId}, + }; + const timing = Number.isFinite(date.getTime()) ? {startTime:date} : {}; + const observation = event.phase === "retrieval" + ? state.root.startObservation("Hindsight Recall", attributes, {asType:"retriever", ...timing}) + : state.root.startObservation("Hindsight Context Injection", attributes, {asType:"event", ...timing}); + observation.end(); + if (event.phase === "retrieval" && typeof event.retrievalId === "string") state.memoryRetrievals.push({retrievalId:event.retrievalId, observationId:observation.otelSpan.spanContext().spanId}); + if (event.phase === "injection") state.memoryInjection = event; + if (standalone) { finalizeRoot({cancelled:false}); void flush(); } + } catch { /* Tracing must never break the session. */ } }); - pi.on("before_provider_request", (_event, ctx) => { + pi.on("before_provider_request", (event, ctx) => { + if (!state) beginPrompt({ prompt: "" }, ctx, "provider-fallback"); if (!state) return; + ensureRuntime(); + // A new provider request while one is open means the previous HTTP + // attempt was retried/superseded — close it instead of leaking it. if (state.openGeneration && !state.openGeneration.finished) { const gen = state.openGeneration; gen.obs.update({ level: "WARNING", statusMessage: "Superseded by provider retry", metadata: { superseded: true } }); gen.obs.end(); } const index = ++state.generationCount; - const baseInput = - index === 1 - ? { role: "user", content: lastPromptText } - : state.pendingToolResults.length - ? { role: "tool", tool_results: state.pendingToolResults } - : undefined; - const generationInput = state.systemPrompt - ? [{ role: "system", content: state.systemPrompt }, ...(baseInput ? [baseInput] : [])] - : baseInput; + const assembled = captureProviderPayload(event.payload); + const memoryHashes = memoryBlockHashes(event.payload); const obs = state.root.startObservation( GENERATION_PREFIX, { - input: generationInput, + input: assembled.input, model: ctx.model?.id, metadata: { assistant_index: index - 1, + assembled_context: assembled.meta, + memory_block_hashes: memoryHashes, + memory_context_id: state.memoryInjection?.contextId, + memory_retrieval_ids: state.memoryInjection?.retrievalIds, + memory_retrieval_observation_ids: state.memoryRetrievals.filter(item => + (state?.memoryInjection?.retrievalIds as string[] | undefined)?.includes(item.retrievalId)).map(item => item.observationId), + memory_injection_verified: state.memoryInjection?.injected === true + ? memoryHashes.includes(String(state.memoryInjection.renderedHash)) : undefined, ...(ctx.model ? { provider: ctx.model.provider } : {}), }, }, @@ -754,6 +874,7 @@ export default function (pi: ExtensionAPI) { pi.on("tool_execution_start", (event) => { if (!state) return; + ensureRuntime(); const redactedArgs = redactLangfuseKeys(event.args) as Record; const serializedArgs = safeStringify(redactedArgs); const hasDataUri = /data:[^;,]{0,100};base64,/.test(serializedArgs); @@ -789,14 +910,20 @@ export default function (pi: ExtensionAPI) { if (!open) return; state.openTools.delete(event.toolCallId); - const result = event.result as { content?: unknown; usage?: PiUsage } | undefined; + const result = event.result as { content?: unknown; usage?: PiUsage; details?: {runId?: string; results?: Array<{index?: number}>} } | undefined; const images = extractImages(result?.content); const rawOutput = renderContentWithImageMarkers(result?.content) || safeStringify(result?.content); const { text: outText, meta: outMeta } = truncateText(rawOutput); if (event.isError) state.sawError = true; state.turnImages.push(...images); - if (result?.usage) { + const childResults = result?.details?.results; + const tracedChildren = contextRegistry().get(state.sessionId)?.tracedChildren ?? []; + const childUsageAlreadyTraced = event.toolName === "subagent" && Boolean(childResults?.length) + && childResults!.every(child => child.index !== undefined && tracedChildren.includes(`${result?.details?.runId}:${child.index}`)); + if (childUsageAlreadyTraced) open.obs.update({metadata:{usage_source:"linked_child_generations", aggregate_usage_omitted:true}}); + if (result?.usage && !childUsageAlreadyTraced) { + ensureRuntime(); const usageObs = startObservation( TOOL_USAGE_OBSERVATION_NAME, { @@ -854,6 +981,7 @@ export default function (pi: ExtensionAPI) { }, startedAt?: Date, ): void => { + ensureRuntime(); if (state) { const obs = startObservation(name, attributes, { asType: "generation", diff --git a/src/telemetry.ts b/src/telemetry.ts new file mode 100644 index 0000000..6209144 --- /dev/null +++ b/src/telemetry.ts @@ -0,0 +1,114 @@ +import { createHash } from "node:crypto"; + +const REDACTED = "[redacted-secret]"; +const SECRET_FIELD = /^(?:authorization|proxy-authorization|cookie|set-cookie|api[-_]?key|secret[-_]?key|secret|password|passwd|client[-_]?secret|(?:api|auth|access|refresh|id)[-_]?token|token|credentials|connection[-_]?string|database[-_]?url|private[-_]?key|encrypted_content)$/i; +type SanitizedTelemetry = boolean | number | string | null | undefined | SanitizedTelemetry[] | { [key: string]: SanitizedTelemetry }; + +/** Copy only; never mutate the provider request or memory result. */ +export function sanitizeTelemetry(value: unknown, secrets: string[] = [], ancestors: object[] = [], omitBinary = true): SanitizedTelemetry { + if (typeof value === "string") { + let text = value; + for (const secret of secrets) if (secret.length >= 8) text = text.split(secret).join(REDACTED); + if (omitBinary) text = text.replace(/data:[^;,\s]{1,100};base64,[A-Za-z0-9+/\s]+=*/g, "[binary data omitted]"); + return text + .replace(/\b(Bearer|Basic)\s+[A-Za-z0-9._~+\/-]+=*/gi, `$1 ${REDACTED}`) + .replace(/\b[sp]k-lf-[\w-]+\b/g, "[redacted-langfuse-secret]") + .replace(/\b(?:sk-(?:proj-|ant-)?[\w-]{8,}|pk-lf-[\w-]+|gh[pousr]_[\w]{16,})\b/g, REDACTED) + .replace(/([?&](?:api[_-]?key|token|access_token|signature|x-amz-signature)=)[^&\s"'<>]+/gi, `$1${REDACTED}`) + .replace(/(\b[a-z][a-z\d+.-]*:\/\/)[^/\s:@]+:[^@\s/]+@/gi, `$1[redacted-credentials]@`) + .replace(/(\b[A-Z_]*(?:API[_-]?KEY|SECRET|PASSWORD|ACCESS_TOKEN|REFRESH_TOKEN)["']?\s*[:=]\s*["']?)[^\s"'`,;}]+/gi, `$1${REDACTED}`); + } + if (value === null || typeof value === "boolean" || typeof value === "number") return value; + if (typeof value !== "object") return undefined; + if (ancestors.includes(value)) return "[circular-ref]"; + if (ancestors.length >= 80) return "[depth limit]"; + const next = [...ancestors, value]; + if (Array.isArray(value)) return value.map(item => sanitizeTelemetry(item, secrets, next, omitBinary)); + // SAFETY: value is non-null object; only own enumerable properties are copied below. + const record = value as Record; + if (omitBinary && record.type === "image" && (record.data || record.source)) { + return { type:"image", mimeType:sanitizeTelemetry(record.mimeType, secrets, next, omitBinary), omitted:true }; + } + const fields: { [key: string]: SanitizedTelemetry } = {}; + for (const [key, item] of Object.entries(record)) { + const inlineBinary = omitBinary && (key === "file_data" || (key === "data" && + (record.mimeType || record.mime_type || record.media_type || record.format || record.type === "base64"))); + fields[key] = SECRET_FIELD.test(key) ? REDACTED : inlineBinary ? "[binary omitted]" : sanitizeTelemetry(item, secrets, next, omitBinary); + } + return fields; +} + +export function environmentSecrets(): string[] { + return Object.entries(process.env) + .filter(([key, value]) => /(?:KEY|SECRET|PASSWORD|TOKEN)$/.test(key) && (value?.length ?? 0) >= 8) + .map(([, value]) => value!); +} + +export function maskTelemetry(data: unknown): SanitizedTelemetry { + const secrets = environmentSecrets(); + if (typeof data === "string") { + try { return JSON.stringify(sanitizeTelemetry(JSON.parse(data), secrets, [], false)); } + catch { return sanitizeTelemetry(data, secrets, [], false); } + } + return sanitizeTelemetry(data, secrets, [], false); +} + +export function captureProviderPayload(payload: unknown, limit = Number(process.env.PI_LANGFUSE_CONTEXT_MAX_CHARS ?? 1000000)) { + const request = sanitizeTelemetry(payload, environmentSecrets()); + const serialized = JSON.stringify(request) ?? "null"; + const cap = Number.isFinite(limit) && limit >= 1000 ? Math.floor(limit) : 1000000; + const truncated = serialized.length > cap; + const captured = truncated + ? { truncated:true, head:serialized.slice(0, Math.floor(cap / 2)), tail:serialized.slice(-Math.floor(cap / 2)) } + : request; + const input = !truncated && captured && typeof captured === "object" && !Array.isArray(captured) + && Array.isArray((captured as { messages?: unknown }).messages) + ? (captured as { messages: unknown[] }).messages + : captured; + return { + input, + meta: { + capture_source:"before_provider_request", redacted:true, binary_omitted:true, + truncated, original_characters:serialized.length, captured_characters:Math.min(serialized.length, cap), + sha256:createHash("sha256").update(serialized).digest("hex"), request:captured, + }, + }; +} + +export interface ParentEnvelope { + traceId?: string; + spanId?: string; + rootSessionId?: string; + parentSessionId?: string; + depth?: number; + runId?: string; + agent?: string; + childIndex?: number; + parentToolCallId?: string; + sourceRunId?: string; + tracedChildren?: readonly string[]; +} +export const CONTEXT_REGISTRY = Symbol.for("pi.langfuse.contexts.v1"); +export function contextRegistry(): Map { + // SAFETY: this process-local symbol is owned exclusively by this extension integration. + const global = globalThis as unknown as Record; + if (!(global[CONTEXT_REGISTRY] instanceof Map)) global[CONTEXT_REGISTRY] = new Map(); + return global[CONTEXT_REGISTRY] as Map; +} + +/** Hash the actual memory blocks present in a request, without persisting raw text. */ +export function memoryBlockHashes(payload: unknown): string[] { + const hashes = new Set(); + const seen = new Set(); + const walk = (value: unknown) => { + if (typeof value === "string") { + if (value.includes("[\s\S]*?<\/hindsight-mental-models>\s*)?[\s\S]*?<\/hindsight-memory>|[\s\S]*?<\/hindsight-mental-models>/g)) { + hashes.add(createHash("sha256").update(match[0]).digest("hex")); + } + } else if (value && typeof value === "object" && !seen.has(value)) { + seen.add(value); for (const item of Object.values(value)) walk(item); + } + }; + walk(payload); return [...hashes]; +} diff --git a/test/completeness.test.ts b/test/completeness.test.ts new file mode 100644 index 0000000..d95d9ff --- /dev/null +++ b/test/completeness.test.ts @@ -0,0 +1,111 @@ +import assert from "node:assert/strict"; +import { EventEmitter } from "node:events"; +import { createHash } from "node:crypto"; +import { mkdtempSync } from "node:fs"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import { describe, it } from "node:test"; +import extension, * as plugin from "../src/index.ts"; +import { maskTelemetry, memoryBlockHashes } from "../src/telemetry.ts"; +import { startCaptureServer, waitForRequests } from "./helpers.ts"; + +function harness(sessionId: string, parent?: Record) { + const handlers = new Map(); const bus = new EventEmitter(); + if (parent) bus.on("pi:trace-parent-request", ({reply}) => reply(parent)); + const ctx: any = {cwd: tmpdir(), hasUI:false, model:{id:"fixture",provider:"fixture"}, sessionManager:{getSessionId:()=>sessionId,getEntries:()=>[],getBranch:()=>[]}}; + const pi: any = {on:(name:string, fn:Function)=>handlers.set(name,[...(handlers.get(name)??[]),fn]), events:{on:(n:string,f:any)=>{bus.on(n,f);return()=>bus.off(n,f);},emit:(n:string,x:any)=>bus.emit(n,x)},exec:async()=>({code:0,stdout:"test\n"})}; + extension(pi); + return {ctx,bus,emit:async(name:string,event:any={})=>{for(const fn of handlers.get(name)??[])await fn({type:name,...event},ctx);}}; +} +async function withCapture(fn: (capture:any)=>Promise) { + const capture=await startCaptureServer(); const saved={...process.env}; + process.env.PI_CODING_AGENT_DIR=mkdtempSync(join(tmpdir(),"pi-tracing-config-")); + Object.assign(process.env,{LANGFUSE_PUBLIC_KEY:"pk-lf-test",LANGFUSE_SECRET_KEY:"sk-lf-test",LANGFUSE_BASE_URL:`http://127.0.0.1:${capture.port}`,LANGFUSE_TRACING_ENABLED:"true"}); + for(const key of Object.keys(process.env))if(key.startsWith("LANGFUSE_PI_PARENT_"))delete process.env[key]; + try{await fn(capture);}finally{for(const key of Object.keys(process.env))if(!(key in saved))delete process.env[key];Object.assign(process.env,saved);capture.close();} +} +const metadata=(s:any,k:string)=>s.attrs[`langfuse.observation.metadata.${k}`]; + +describe("tracing completeness regressions",()=>{ + it("omits provider-native inline image, audio and document bytes",()=>{ + const payload={contents:[{parts:[{inlineData:{mimeType:"image/png",data:"GEMINI_BINARY_SENTINEL"}}]}],audio:{type:"input_audio",input_audio:{data:"AUDIO_BINARY_SENTINEL",format:"wav"}},doc:{type:"document",source:{type:"base64",media_type:"application/pdf",data:"DOCUMENT_BINARY_SENTINEL"}},file:{type:"input_file",file_data:"FILE_BINARY_SENTINEL"}}; + const captured=JSON.stringify((plugin as any).captureProviderPayload(payload)); + for(const value of ["GEMINI_BINARY_SENTINEL","AUDIO_BINARY_SENTINEL","DOCUMENT_BINARY_SENTINEL","FILE_BINARY_SENTINEL"])assert.ok(!captured.includes(value)); + }); + it("verifies combined mental-model and recalled-memory injections",()=>{ + const mental="model"; + const memory="fact"; + for(const text of [mental,memory,`${mental}\n\n${memory}`]){ + const digest=createHash("sha256").update(text).digest("hex"); + assert.ok(memoryBlockHashes({messages:[{content:text}]}).includes(digest)); + assert.ok(memoryBlockHashes({messages:[{content:`prefix\n${text}\nsuffix`}]}).includes(digest)); + } + }); + it("preserves upstream opt-in image uploads outside assembled-request capture",()=>{ + const root={role:"user",content:[{type:"image_url",image_url:{url:"data:image/png;base64,AAAA"}}]}; + assert.ok(JSON.stringify(maskTelemetry(root)).includes("data:image/png;base64,AAAA")); + assert.ok(!JSON.stringify((plugin as any).captureProviderPayload(root)).includes("data:image/png;base64,AAAA")); + }); + it("redacts credentials embedded in prose and URLs, not only JSON fields",()=>{ + const result=(plugin as any).captureProviderPayload({messages:[{role:"user",content:'API_KEY=fixture-plain-secret\nAuthorization: Basic Zml4dHVyZTpzZWNyZXQ=\npostgresql://user:fixture-db-password@localhost/db'}],token:"fixture-token-value"}); + for(const secret of ['fixture-plain-secret','Zml4dHVyZTpzZWNyZXQ=','fixture-db-password','fixture-token-value'])assert.ok(!JSON.stringify(result).includes(secret),`leaked ${secret}`); + }); + it("does not double-count a linked native child's usage as a second model call",async()=>withCapture(async capture=>{ + const parent=harness("usage-parent");await parent.emit("before_agent_start",{prompt:"parent"}); + await parent.emit("tool_execution_start",{toolCallId:"t1",toolName:"subagent",args:{agent:"scout"}}); + const registry=(globalThis as any)[Symbol.for("pi.langfuse.contexts.v1")]; + const child=harness("usage-child",{...registry.get("usage-parent"),parentSessionId:"usage-parent",depth:1,runId:"usage-run",agent:"scout",childIndex:0}); + await child.emit("before_agent_start",{prompt:"child"});await child.emit("before_provider_request",{payload:{messages:[]}}); + const usage={input:10,output:5,cacheRead:0,cacheWrite:0,cost:{input:0.01,output:0.01,total:0.02}}; + await child.emit("message_end",{message:{role:"assistant",content:[{type:"text",text:"done"}],usage}});await child.emit("agent_settled"); + await parent.emit("tool_execution_end",{toolCallId:"t1",toolName:"subagent",result:{content:[],usage,details:{runId:"usage-run",results:[{index:0}]}}});await parent.emit("agent_settled"); + await waitForRequests(capture,2);assert.equal(capture.spans().filter((s:any)=>s.name==="Tool LLM Usage").length,0); + for(const h of [child,parent])await h.emit("session_shutdown",{reason:"quit"}); + })); + it("captures assembled payload without mutation, preserving history/system/tools and masking secrets",()=>{ + const payload={model:"test",messages:[{role:"system",content:"system instructions"},{role:"user",content:"earlier user"},{role:"assistant",content:"earlier reply"},{role:"user",content:"new turn"}],tools:[{name:"read",parameters:{type:"object"}}],credentials:{api_key:"fixture-private-value"},header:"Bearer fixture-bearer-value",binary:{type:"image",data:"AAAA",mimeType:"image/png"}}; + const before=structuredClone(payload);const capture=(plugin as any).captureProviderPayload; + assert.equal(typeof capture,"function","provider payload capture must exist"); + const result=capture(payload);assert.deepEqual(payload,before); + assert.equal(result.input.length,4);assert.deepEqual(result.meta.request.tools,payload.tools); + assert.ok(!JSON.stringify(result).includes("fixture-private-value"));assert.ok(!JSON.stringify(result).includes("fixture-bearer-value"));assert.ok(!JSON.stringify(result).includes('"data":"AAAA"')); + assert.equal(result.meta.truncated,false); + }); + it("captures payload tail beyond old 20K cap and labels intentional limits",()=>{ + const capture=(plugin as any).captureProviderPayload;assert.equal(typeof capture,"function"); + const result=capture({messages:[{role:"system",content:"a".repeat(30000)+"TAIL_SENTINEL"}]}); + assert.ok(JSON.stringify(result.input).includes("TAIL_SENTINEL"));assert.equal(result.meta.truncated,false); + assert.equal(capture({content:"x".repeat(1500)},1000).meta.truncated,true); + }); + it("traces autonomous agent_start without before_agent_start and includes actual payload",async()=>withCapture(async capture=>{ + const h=harness("autonomous-session");await h.emit("agent_start"); + const payload={model:"fixture",messages:[{role:"system",content:"ASSEMBLED_SYSTEM"},{role:"user",content:"AUTONOMOUS_RESULT"}],tools:[]}; + await h.emit("before_provider_request",{payload});await h.emit("message_end",{message:{role:"assistant",content:[{type:"text",text:"Done"}]}});await h.emit("agent_settled"); + await waitForRequests(capture,1);const spans=capture.spans();assert.equal(spans.filter((s:any)=>s.name==="LLM Call").length,1); + const gen=spans.find((s:any)=>s.name==="LLM Call");assert.ok(String(gen.attrs["langfuse.observation.input"]).includes("ASSEMBLED_SYSTEM")); + assert.equal(metadata(spans.find((s:any)=>s.name==="Conversational Turn"),"trigger"),"agent_start"); + await h.emit("session_shutdown",{reason:"quit"}); + })); + it("keeps in-process siblings under the true parent and preserves parent attribution after children",async()=>withCapture(async capture=>{ + const parent=harness("parent-session");await parent.emit("before_agent_start",{prompt:"parent",images:[]}); + const registry=(globalThis as any)[Symbol.for("pi.langfuse.contexts.v1")];assert.ok(registry instanceof Map,"session-scoped parent registry missing"); + const envelope={...registry.get("parent-session"),parentSessionId:"parent-session",depth:1,runId:"run-1",agent:"reviewer"}; + const originalSpan=process.env.LANGFUSE_PI_PARENT_SPAN_ID; + const a=harness("child-a",{...envelope,childIndex:0});await a.emit("before_agent_start",{prompt:"child-a"}); + const b=harness("child-b",{...envelope,childIndex:1});await b.emit("before_agent_start",{prompt:"child-b"}); + assert.equal(process.env.LANGFUSE_PI_PARENT_SPAN_ID,originalSpan,"children must not overwrite parent env"); + for(const h of [b,a,parent]){await h.emit("before_provider_request",{payload:{messages:[{role:"user",content:"test"}]}});await h.emit("message_end",{message:{role:"assistant",content:[{type:"text",text:"done"}]}});await h.emit("agent_settled");} + await waitForRequests(capture,3);const spans=capture.spans();const root=spans.find((s:any)=>s.name==="Conversational Turn");const children=spans.filter((s:any)=>s.name==="Subagent Turn");assert.equal(children.length,2); + for(const c of children){assert.equal(c.parentSpanId,root.spanId);assert.equal(c.traceId,root.traceId);assert.equal(metadata(c,"run_id"),"run-1");assert.equal(metadata(c,"agent"),"reviewer");assert.equal(c.attrs["session.id"],"parent-session");} + assert.ok(spans.filter((s:any)=>s.name==="LLM Call").every((s:any)=>s.attrs["session.id"]==="parent-session")); + for(const h of [a,b,parent])await h.emit("session_shutdown",{reason:"quit"}); + })); + it("exports correlated RETRIEVER events and injection linkage on the following generation",async()=>withCapture(async capture=>{ + const h=harness("memory-session");await h.emit("before_agent_start",{prompt:"recall decision"}); + const now=new Date().toISOString();h.bus.emit("hindsight:retrieval",{version:1,phase:"retrieval",mode:"automatic",sessionId:"memory-session",retrievalId:"r1",contextId:"c1",bankId:"fixture-bank",query:"exact scoped query",tagGroups:[{all:[{tags:["project:fixture"],match:"any_strict"}]}],startedAt:now,endedAt:now,durationMs:1,status:"success",cache:"miss",keptIds:["f1"],injectedIds:["f1"],results:[{id:"f1",text:"fixture memory"}]}); + h.bus.emit("hindsight:retrieval",{version:1,phase:"injection",mode:"automatic",sessionId:"memory-session",contextId:"c1",retrievalIds:["r1"],startedAt:now,endedAt:now,status:"success",cache:"miss",injected:true,renderedHash:"fixture-hash",renderedLength:14}); + await h.emit("before_provider_request",{payload:{messages:[{role:"user",content:"question"}]}});await h.emit("message_end",{message:{role:"assistant",content:[{type:"text",text:"done"}]}});await h.emit("agent_settled");await waitForRequests(capture,1); + const spans=capture.spans();const retriever=spans.find((s:any)=>s.attrs["langfuse.observation.type"]==="retriever");assert.ok(retriever);assert.equal(metadata(retriever,"bank_id"),"fixture-bank"); + const gen=spans.find((s:any)=>s.name==="LLM Call");assert.equal(metadata(gen,"memory_context_id"),"c1");assert.ok(String(metadata(gen,"memory_retrieval_ids")).includes("r1"));await h.emit("session_shutdown",{reason:"quit"}); + })); +}); diff --git a/test/system-prompt.test.ts b/test/system-prompt.test.ts index 9a13cb9..ce89f51 100644 --- a/test/system-prompt.test.ts +++ b/test/system-prompt.test.ts @@ -66,8 +66,7 @@ describe("integration: system prompt", () => { const [head, ...rest] = input as Array<{ role?: string; content?: unknown }>; assert.equal(head?.role, "system", `generation ${i} must start with a system message`); assert.equal(head?.content, systemPrompt, `generation ${i} must carry the captured prompt`); - assert.equal(rest.length, 1, `generation ${i} keeps exactly one base message`); - assert.equal(rest[0]!.role, i === 0 ? "user" : "tool"); + assert.ok(rest.some((message) => message.role === (i === 0 ? "user" : "tool")), `generation ${i} retains its active message`); } } finally { capture.close();