diff --git a/README.md b/README.md index 58c1da8..0d9abfb 100644 --- a/README.md +++ b/README.md @@ -62,6 +62,33 @@ export LANGFUSE_USER_ID="your-user-id" If both `LANGFUSE_PUBLIC_KEY` and `LANGFUSE_SECRET_KEY` are set, the plugin uses environment variables instead of reading the config file. Optional values can be supplied either way. +## Redacting tool data + +Add `redaction` to `~/.config/opencode/opencode-langfuse.json` (also when credentials come from environment variables): + +```json +{ + "publicKey": "pk-lf-...", + "secretKey": "sk-lf-...", + "redaction": { + "tools": [ + { "name": "read", "path": "**/.env*", "input": "redact" }, + { + "name": "read", + "path": "**/.env.example", + "input": "as-is", + "output": "as-is" + }, + { "name": "bash", "output": "redact" } + ] + } +} +``` + +Rules match the tool name exactly (or `*` for every tool). An optional `path` glob matches the tool input's `filePath`, `path`, or `filename` (with `*` matching within a directory and `**` across directories). The last matching rule wins for each path; when multiple path arguments are present, a redaction on any of them wins. `input` and `output` accept `"as-is"` or `"redact"`; a matching rule defaults to sending the input as-is and replacing the output with `"[REDACTED]"`. Without a matching rule, data is sent as-is. Redaction covers tool observations and tool results copied into later generation inputs in both OpenCode versions. To hide the file path or other tool arguments too, set `"input": "redact"`. + +These rules do not redact text that the user or assistant independently includes in messages or reasoning. If a tool call's path is unavailable (including results without their calls), any potentially matching path rule redacts its input and/or output. If the file cannot be read or parsed, tracing stops rather than silently ignoring potentially configured redaction rules. Restart OpenCode after changing the file. + ## Contributing See the [contributing guide](./CONTRIBUTING.md). diff --git a/src/langfuse.ts b/src/langfuse.ts index b909fd2..8d8d145 100644 --- a/src/langfuse.ts +++ b/src/langfuse.ts @@ -13,23 +13,32 @@ import { NodeTracerProvider } from "@opentelemetry/sdk-trace-node"; import { Context as EffectContext, Effect } from "effect"; import { PLUGIN_VERSION } from "./version.js"; +import { TelemetryRedaction, type RedactionConfig } from "./redaction.js"; export class LangfuseClient { readonly baseUrl: string; readonly forceFlush: Effect.Effect; private readonly traceState: LangfuseTraceState; + private readonly redaction: TelemetryRedaction; constructor(input: { baseUrl: string; traceState: LangfuseTraceState; forceFlush: Effect.Effect; + redaction?: RedactionConfig; }) { this.baseUrl = input.baseUrl; this.traceState = input.traceState; this.forceFlush = input.forceFlush; + this.redaction = new TelemetryRedaction(input.redaction); + } + + private serialize(value: unknown, sessionID?: string) { + return JSON.stringify(this.redaction.sanitize(value, sessionID)); } clearTraceState() { + this.redaction.clear(); this.traceState.assistantParts.clear(); this.traceState.abortedSessions.clear(); this.traceState.tracedEventIds.clear(); @@ -50,6 +59,7 @@ export class LangfuseClient { } clearSessionTraceState(sessionID: string) { + this.redaction.clearSession(sessionID); const sessionMessageIds = new Set(); for (const [messageID, parts] of this.traceState.assistantParts) { @@ -236,11 +246,24 @@ export class LangfuseClient { "session.id": input.sessionID, ...(input.input === undefined ? {} - : { "langfuse.observation.input": JSON.stringify(input.input) }), + : { + "langfuse.observation.input": this.serialize( + input.input, + input.sessionID, + ), + }), ...(input.output === undefined ? {} - : { "langfuse.observation.output": JSON.stringify(input.output) }), - "langfuse.observation.metadata": JSON.stringify(input.metadata), + : { + "langfuse.observation.output": this.serialize( + input.output, + input.sessionID, + ), + }), + "langfuse.observation.metadata": this.serialize( + input.metadata, + input.sessionID, + ), }, startTime: new Date(input.timestamp), }); @@ -443,7 +466,10 @@ export class LangfuseClient { "langfuse.observation.model.name": input.model.id, ...(generationInput ? { - "langfuse.observation.input": JSON.stringify(generationInput), + "langfuse.observation.input": this.serialize( + generationInput, + input.sessionID, + ), } : {}), "langfuse.observation.metadata": JSON.stringify({ @@ -573,7 +599,10 @@ export class LangfuseClient { "langfuse.observation.type": "agent", "langfuse.internal.is_app_root": !parentSpan, "session.id": input.sessionID, - "langfuse.observation.input": JSON.stringify([formattedMessage]), + "langfuse.observation.input": this.serialize( + [formattedMessage], + input.sessionID, + ), "langfuse.observation.metadata": JSON.stringify({ messageID: input.messageID, agent: input.agent, @@ -613,7 +642,10 @@ export class LangfuseClient { attributes: { "langfuse.observation.type": "event", "session.id": input.sessionID, - "langfuse.observation.input": JSON.stringify([formattedMessage]), + "langfuse.observation.input": this.serialize( + [formattedMessage], + input.sessionID, + ), "langfuse.observation.metadata": JSON.stringify({ messageID: input.messageID, agent: input.agent, @@ -660,6 +692,12 @@ export class LangfuseClient { tool: string; args: Record; }) { + this.redaction.register( + input.callID, + input.tool, + input.args, + input.sessionID, + ); this.traceState.toolMessageIdsByCallId.set(input.callID, input.messageID); const parts = @@ -714,7 +752,7 @@ export class LangfuseClient { if (input.mode !== "compaction") { turn?.span.setAttribute( "langfuse.observation.output", - JSON.stringify(output), + this.serialize(output, input.sessionID), ); } const activeStep = this.traceState.activeGenerationSteps.get( @@ -731,7 +769,7 @@ export class LangfuseClient { step.span.setAttribute("langfuse.observation.model.name", input.modelID); step.span.setAttribute( "langfuse.observation.output", - JSON.stringify(output), + this.serialize(output, input.sessionID), ); step.span.setAttribute( "langfuse.observation.usage_details", @@ -795,10 +833,16 @@ export class LangfuseClient { "langfuse.observation.model.name": input.modelID, ...(generationInput ? { - "langfuse.observation.input": JSON.stringify(generationInput), + "langfuse.observation.input": this.serialize( + generationInput, + input.sessionID, + ), } : {}), - "langfuse.observation.output": JSON.stringify(output), + "langfuse.observation.output": this.serialize( + output, + input.sessionID, + ), "langfuse.observation.usage_details": JSON.stringify({ input: input.tokens.input, output: input.tokens.output, @@ -976,6 +1020,13 @@ export class LangfuseClient { return; } + this.redaction.register( + input.callID, + input.tool, + input.args, + input.sessionID, + ); + this.ensureGenerationParent(input.sessionID); this.withObservationParent( @@ -985,7 +1036,10 @@ export class LangfuseClient { attributes: { "langfuse.observation.type": "tool", "session.id": input.sessionID, - "langfuse.observation.input": JSON.stringify(input.args), + "langfuse.observation.input": this.serialize( + this.redaction.toolInput(input.tool, input.args), + input.sessionID, + ), "langfuse.observation.metadata": JSON.stringify({ callID: input.callID, tool: input.tool, @@ -1041,7 +1095,13 @@ export class LangfuseClient { span.setAttribute( "langfuse.observation.output", - JSON.stringify({ title: input.title, output: input.output }), + this.serialize( + { + title: this.redaction.toolOutput(input.callID, input.title), + output: this.redaction.toolOutput(input.callID, input.output), + }, + input.sessionID, + ), ); span.setAttribute( "langfuse.observation.metadata", @@ -1106,13 +1166,18 @@ export class LangfuseClient { span.setAttribute( "langfuse.observation.output", - JSON.stringify({ error: input.error }), + this.serialize( + { error: this.redaction.toolOutput(input.callID, input.error) }, + observation.sessionID, + ), ); span.setStatus({ code: SpanStatusCode.ERROR, - message: input.error, + message: String(this.redaction.toolOutput(input.callID, input.error)), + }); + span.recordException({ + message: String(this.redaction.toolOutput(input.callID, input.error)), }); - span.recordException({ message: input.error }); span.end(new Date(input.completed)); this.rememberToolResult({ sessionID: observation.sessionID, @@ -1146,7 +1211,10 @@ export class LangfuseClient { "session.id": sessionID, ...(generationInput ? { - "langfuse.observation.input": JSON.stringify(generationInput), + "langfuse.observation.input": this.serialize( + generationInput, + sessionID, + ), } : {}), }, @@ -1733,6 +1801,7 @@ export const createLangfuseClient = (input: { userId?: string; serviceName?: string; opencodeVersion?: string; + redaction?: RedactionConfig; }) => Effect.gen(function* () { const tracerName = "opencode-langfuse-plugin"; @@ -1802,5 +1871,6 @@ export const createLangfuseClient = (input: { baseUrl: input.baseUrl, traceState, forceFlush: Effect.tryPromise(() => processor.forceFlush()), + redaction: input.redaction, }); }); diff --git a/src/redaction.ts b/src/redaction.ts new file mode 100644 index 0000000..43e5c37 --- /dev/null +++ b/src/redaction.ts @@ -0,0 +1,235 @@ +export type RedactionConfig = { + tools?: readonly { + readonly name: string; + readonly path?: string; + readonly input?: "as-is" | "redact"; + readonly output?: "as-is" | "redact"; + }[]; +}; + +type Policy = { input: boolean; output: boolean }; + +const REDACTED = "[REDACTED]"; +const isRecord = (value: unknown): value is Record => + typeof value === "object" && value !== null && !Array.isArray(value); + +export class TelemetryRedaction { + private readonly calls = new Map(); + + constructor(private readonly config: RedactionConfig = {}) {} + + private policy(tool: string, args: unknown): Policy { + let decodedArgs: unknown = args; + if (typeof args === "string") { + try { + decodedArgs = JSON.parse(args); + } catch { + // Unknown input: apply path-scoped redaction conservatively below. + } + } + const pathKeys = ["filePath", "path", "filename"]; + const argsRecord = isRecord(decodedArgs) ? decodedArgs : undefined; + const paths = + argsRecord !== undefined + ? pathKeys + .map((key) => argsRecord[key]) + .filter((value): value is string => typeof value === "string") + : []; + const unknownPath = + paths.length === 0 || + (argsRecord !== undefined && + pathKeys.some( + (key) => key in argsRecord && typeof argsRecord[key] !== "string", + )); + const policies = (paths.length > 0 ? paths : [undefined]).map((path) => { + let policy: Policy = { input: false, output: false }; + for (const rule of this.config.tools ?? []) { + if (rule.name !== "*" && rule.name !== tool) { + continue; + } + if (rule.path !== undefined) { + if (path === undefined) { + continue; + } + const pattern = rule.path.replace( + /\*\*\/|\*\*|\*|[.+?^${}()|[\]\\]/g, + (match) => + match === "**/" + ? "(?:.*/)?" + : match === "**" + ? ".*" + : match === "*" + ? "[^/]*" + : `\\${match}`, + ); + if (!new RegExp(`^${pattern}$`).test(path)) { + continue; + } + } + policy = { + input: (rule.input ?? "as-is") === "redact", + output: (rule.output ?? "redact") === "redact", + }; + } + return policy; + }); + return { + input: + policies.some((policy) => policy.input) || + (unknownPath && this.unknownPathRedaction(tool).input), + output: + policies.some((policy) => policy.output) || + (unknownPath && this.unknownPathRedaction(tool).output), + }; + } + + private unknownPathRedaction(tool: string): Policy { + const rules = (this.config.tools ?? []).filter( + (rule) => + (rule.name === "*" || rule.name === tool) && rule.path !== undefined, + ); + return { + input: rules.some((rule) => rule.input === "redact"), + output: rules.some((rule) => (rule.output ?? "redact") === "redact"), + }; + } + + register(callID: string, tool: string, args: unknown, sessionID?: string) { + if (args === undefined && this.calls.has(callID)) { + return; + } + this.calls.set(callID, { + ...this.policy(tool, args), + sessionID: sessionID ?? this.calls.get(callID)?.sessionID, + }); + } + + clear() { + this.calls.clear(); + } + + clearSession(sessionID: string) { + for (const [callID, policy] of this.calls) { + if (policy.sessionID === sessionID) { + this.calls.delete(callID); + } + } + } + + toolInput(tool: string, args: unknown) { + return this.policy(tool, args).input ? REDACTED : args; + } + + toolOutput(callID: string, output: unknown) { + return this.calls.get(callID)?.output === true ? REDACTED : output; + } + + sanitize(value: unknown, sessionID?: string): unknown { + // Register calls before walking results: a history snapshot can contain + // tool results from earlier steps as well as their original tool calls. + const discover = (node: unknown): void => { + if (typeof node !== "object" || node === null) { + return; + } + if (Array.isArray(node)) { + node.forEach(discover); + return; + } + if (!isRecord(node)) { + return; + } + const object = node; + if (typeof object.id === "string" && typeof object.name === "string") { + if (object.type === "tool-call") { + this.register(object.id, object.name, object.input, sessionID); + } else if (typeof object.arguments === "string") { + try { + this.register( + object.id, + object.name, + JSON.parse(object.arguments), + sessionID, + ); + } catch { + this.register(object.id, object.name, object.arguments, sessionID); + } + } + } + Object.values(object).forEach(discover); + }; + discover(value); + + const visit = (node: unknown): unknown => { + if (typeof node !== "object" || node === null) { + return node; + } + if (Array.isArray(node)) { + return node.map(visit); + } + if (!isRecord(node)) { + return node; + } + const object = node; + const callID = + typeof object.tool_call_id === "string" + ? object.tool_call_id + : typeof object.toolCallId === "string" + ? object.toolCallId + : typeof object.callID === "string" + ? object.callID + : typeof object.id === "string" && + (object.type === "tool-result" || + object.type === "tool" || + object.type === "tool-call" || + (typeof object.name === "string" && "arguments" in object)) + ? object.id + : undefined; + const resultTool = + (object.type === "tool-result" || object.role === "tool") && + typeof object.name === "string" + ? object.name + : undefined; + const namedPolicy = + resultTool !== undefined + ? this.policy(resultTool, undefined) + : undefined; + const policy = + (callID !== undefined ? this.calls.get(callID) : undefined) ?? + namedPolicy; + const isResult = object.role === "tool" || object.type === "tool-result"; + const isCall = typeof object.name === "string" && "arguments" in object; + if (policy?.output === true && isResult) { + return { + ...Object.fromEntries( + Object.entries(object).filter(([key]) => + [ + "role", + "type", + "id", + "name", + "tool_call_id", + "toolCallId", + ].includes(key), + ), + ), + ...(object.role === "tool" ? { content: REDACTED } : {}), + ...(object.type === "tool-result" ? { result: REDACTED } : {}), + }; + } + if (policy?.input === true && (isCall || object.type === "tool-call")) { + return { + ...Object.fromEntries( + Object.entries(object).filter(([key]) => + ["type", "id", "name", "namespace"].includes(key), + ), + ), + ...(isCall ? { arguments: REDACTED } : { input: REDACTED }), + }; + } + return Object.fromEntries( + Object.entries(object).map(([key, item]) => [key, visit(item)]), + ); + }; + return visit(value); + } +} diff --git a/src/runtime.ts b/src/runtime.ts index dd842ad..6ec620e 100644 --- a/src/runtime.ts +++ b/src/runtime.ts @@ -11,17 +11,33 @@ import { } from "effect"; import { createLangfuseClient, type LangfuseClient } from "./langfuse.js"; +import type { RedactionConfig } from "./redaction.js"; + +const RedactionConfigSchema = Schema.Struct({ + tools: Schema.optional( + Schema.Array( + Schema.Struct({ + name: Schema.NonEmptyString, + path: Schema.optional(Schema.NonEmptyString), + input: Schema.optional(Schema.Literal("as-is", "redact")), + output: Schema.optional(Schema.Literal("as-is", "redact")), + }), + ), + ), +}); -const LangfuseCredentialsSchema = Schema.Struct({ - publicKey: Schema.NonEmptyString, - secretKey: Schema.NonEmptyString, +const LangfuseConfigSchema = Schema.Struct({ + publicKey: Schema.optional(Schema.NonEmptyString), + secretKey: Schema.optional(Schema.NonEmptyString), baseUrl: Schema.optional(Schema.NonEmptyString), environment: Schema.optional(Schema.NonEmptyString), userId: Schema.optional(Schema.NonEmptyString), serviceName: Schema.optional(Schema.NonEmptyString), + redaction: Schema.optional(RedactionConfigSchema), +}); +const RedactionOnlySchema = Schema.Struct({ + redaction: Schema.optional(RedactionConfigSchema), }); - -type LangfuseCredentials = typeof LangfuseCredentialsSchema.Type; class MissingLangfuseCredentials extends Data.TaggedError( "MissingLangfuseCredentials", @@ -31,61 +47,52 @@ class ChangedLangfuseConfiguration extends Data.TaggedError( "ChangedLangfuseConfiguration", )<{ readonly message: string }> {} -const loadLangfuseCredentials = Effect.gen(function* () { - const publicKey = process.env.LANGFUSE_PUBLIC_KEY; - const secretKey = process.env.LANGFUSE_SECRET_KEY; - - if ( - publicKey !== undefined && - publicKey !== "" && - secretKey !== undefined && - secretKey !== "" - ) { - return { - publicKey, - secretKey, - baseUrl: process.env.LANGFUSE_BASE_URL ?? process.env.LANGFUSE_BASEURL, - environment: process.env.LANGFUSE_ENVIRONMENT, - userId: process.env.LANGFUSE_USER_ID, - serviceName: process.env.LANGFUSE_SERVICE_NAME, - } satisfies LangfuseCredentials; - } - - const configPath = join( - homedir(), - ".config", - "opencode", - "opencode-langfuse.json", - ); +class InvalidLangfuseConfiguration extends Data.TaggedError( + "InvalidLangfuseConfiguration", +)<{ readonly message: string }> {} - const credentials = yield* Effect.tryPromise({ - try: async () => - Schema.decodeUnknownSync(Schema.parseJson(LangfuseCredentialsSchema))( - await readFile(configPath, "utf8"), - ), +const loadLangfuseConfig = (useEnvironmentCredentials: boolean) => + Effect.tryPromise({ + try: async (): Promise => { + const path = join( + homedir(), + ".config", + "opencode", + "opencode-langfuse.json", + ); + let contents: string; + try { + contents = await readFile(path, "utf8"); + } catch (error) { + if ( + typeof error === "object" && + error !== null && + "code" in error && + error.code === "ENOENT" + ) { + return {}; + } + throw error; + } + const document = Schema.decodeUnknownSync(Schema.parseJson())(contents); + if ( + typeof document === "object" && + document !== null && + "redaction" in document + ) { + Schema.decodeUnknownSync(RedactionConfigSchema, { + onExcessProperty: "error", + })(document.redaction); + } + return useEnvironmentCredentials + ? Schema.decodeUnknownSync(RedactionOnlySchema)(document) + : Schema.decodeUnknownSync(LangfuseConfigSchema)(document); + }, catch: () => - new MissingLangfuseCredentials({ - message: "Missing Langfuse credentials", + new InvalidLangfuseConfiguration({ + message: "Invalid Langfuse configuration", }), - }).pipe( - Effect.mapError( - () => - new MissingLangfuseCredentials({ - message: "Missing Langfuse credentials", - }), - ), - ); - - if (!credentials.publicKey || !credentials.secretKey) { - return yield* Effect.fail( - new MissingLangfuseCredentials({ - message: "Missing Langfuse credentials", - }), - ); - } - - return credentials; -}); + }); /** * opencode disposes and re-creates plugin instances inside the same process @@ -113,10 +120,34 @@ declare global { export const createLangfuseRuntime = (input: { opencodeVersion?: string }) => Effect.gen(function* () { - const credentials = yield* loadLangfuseCredentials; + const useEnvironmentCredentials = + process.env.LANGFUSE_PUBLIC_KEY !== undefined && + process.env.LANGFUSE_PUBLIC_KEY !== "" && + process.env.LANGFUSE_SECRET_KEY !== undefined && + process.env.LANGFUSE_SECRET_KEY !== ""; + const config = yield* loadLangfuseConfig(useEnvironmentCredentials); + const credentials = useEnvironmentCredentials + ? { + publicKey: process.env.LANGFUSE_PUBLIC_KEY, + secretKey: process.env.LANGFUSE_SECRET_KEY, + baseUrl: + process.env.LANGFUSE_BASE_URL ?? process.env.LANGFUSE_BASEURL, + environment: process.env.LANGFUSE_ENVIRONMENT, + userId: process.env.LANGFUSE_USER_ID, + serviceName: process.env.LANGFUSE_SERVICE_NAME, + } + : config; + const { publicKey, secretKey } = credentials; + if (publicKey === undefined || secretKey === undefined) { + return yield* Effect.fail( + new MissingLangfuseCredentials({ + message: "Missing Langfuse credentials", + }), + ); + } const clientInput = { - publicKey: credentials.publicKey, - secretKey: credentials.secretKey, + publicKey, + secretKey, baseUrl: credentials.baseUrl ?? process.env.LANGFUSE_BASE_URL ?? @@ -129,6 +160,7 @@ export const createLangfuseRuntime = (input: { opencodeVersion?: string }) => userId: credentials.userId ?? process.env.LANGFUSE_USER_ID, serviceName: credentials.serviceName ?? process.env.LANGFUSE_SERVICE_NAME, opencodeVersion: input.opencodeVersion, + redaction: (config.redaction ?? {}) satisfies RedactionConfig, } satisfies Parameters[0]; const cacheKey = JSON.stringify(clientInput); const sharedClientState = diff --git a/test/integration/redaction.test.ts b/test/integration/redaction.test.ts new file mode 100644 index 0000000..ac6d338 --- /dev/null +++ b/test/integration/redaction.test.ts @@ -0,0 +1,261 @@ +import { expect, test } from "vitest"; + +import { TelemetryRedaction } from "../../src/redaction.js"; + +test("redacts matching file reads in tool spans and subsequent generation history", () => { + const redaction = new TelemetryRedaction({ + tools: [ + { name: "read", path: "**/.env*", input: "redact" }, + { + name: "read", + path: "**/.env.example", + input: "as-is", + output: "as-is", + }, + ], + }); + redaction.register("secret", "read", { filePath: "/app/.env" }); + redaction.register("public", "read", { filePath: "/app/.env.example" }); + + expect(redaction.toolInput("read", { filePath: "/app/.env" })).toBe( + "[REDACTED]", + ); + expect(redaction.toolOutput("secret", "PASSWORD=abc")).toBe("[REDACTED]"); + expect(redaction.toolOutput("public", "example")).toBe("example"); + + const history = [ + { + role: "assistant", + tool_calls: [ + { id: "secret", name: "read", arguments: '{"filePath":"/app/.env"}' }, + ], + }, + { role: "tool", tool_call_id: "secret", content: "PASSWORD=abc" }, + { role: "tool", tool_call_id: "public", content: "example" }, + { + type: "tool-result", + toolCallId: "secret", + output: { text: "PASSWORD=abc" }, + }, + ]; + const serialized = JSON.stringify(redaction.sanitize(history)); + expect(serialized).not.toContain("PASSWORD=abc"); + expect(serialized).not.toContain('/app/.env"'); + expect(serialized).toContain("example"); + expect(history[1]).toHaveProperty("content", "PASSWORD=abc"); +}); + +test("unmatched tools stay as-is and earlier tool calls in a snapshot are discovered", () => { + const redaction = new TelemetryRedaction({ + tools: [{ name: "read", path: "**/secrets/**" }], + }); + const snapshot = [ + { + role: "assistant", + tool_calls: [ + { + id: "one", + name: "read", + arguments: '{"filePath":"/repo/secrets/token"}', + }, + ], + }, + { role: "tool", tool_call_id: "one", content: "sensitive" }, + { role: "tool", tool_call_id: "two", content: "ordinary" }, + ]; + expect(redaction.sanitize(snapshot)).toEqual([ + snapshot[0], + { ...snapshot[1], content: "[REDACTED]" }, + snapshot[2], + ]); +}); + +test("redacts real OpenCode 2 tool-call inputs and tool-result values", () => { + const redaction = new TelemetryRedaction({ + tools: [{ name: "read", path: "**/.private", input: "redact" }], + }); + const snapshot = { + messages: [ + { + role: "assistant", + content: [ + { + type: "tool-call", + id: "call-1", + name: "read", + input: { filePath: "/app/.private" }, + }, + ], + }, + { + role: "tool", + content: [ + { + type: "tool-result", + id: "call-1", + name: "read", + result: { type: "text", value: "PASSWORD=abc" }, + }, + ], + }, + ], + }; + + expect(redaction.sanitize(snapshot)).toEqual({ + messages: [ + { + role: "assistant", + content: [{ ...snapshot.messages[0].content[0], input: "[REDACTED]" }], + }, + { + role: "tool", + content: [{ ...snapshot.messages[1].content[0], result: "[REDACTED]" }], + }, + ], + }); + expect(snapshot.messages[1].content[0]).toHaveProperty( + "result.value", + "PASSWORD=abc", + ); +}); + +test("removes call policies when their session ends", () => { + const redaction = new TelemetryRedaction({ + tools: [{ name: "read", output: "redact" }], + }); + redaction.register("call-1", "read", {}, "session-1"); + redaction.register("call-2", "read", {}, "session-2"); + redaction.clearSession("session-1"); + + expect(redaction.toolOutput("call-1", "visible")).toBe("visible"); + expect(redaction.toolOutput("call-2", "secret")).toBe("[REDACTED]"); +}); + +test("redacts a named tool result even when its call was not observed", () => { + const redaction = new TelemetryRedaction({ + tools: [ + { name: "read", path: "**/.env" }, + { name: "read", path: "**/.env.example", output: "as-is" }, + ], + }); + const result = { + role: "tool", + content: [ + { + type: "tool-result", + id: "unknown-call", + name: "read", + result: { type: "text", value: "PASSWORD=secret" }, + }, + ], + }; + const expected = { + role: "tool", + content: [ + { + type: "tool-result", + id: "unknown-call", + name: "read", + result: "[REDACTED]", + }, + ], + }; + expect(redaction.sanitize(result)).toEqual(expected); + expect( + new TelemetryRedaction({ tools: [{ name: "read" }] }).sanitize(result), + ).toEqual(expected); +}); + +test("checks every path argument before sending a tool result as-is", () => { + const redaction = new TelemetryRedaction({ + tools: [ + { name: "read", path: "**/.env*", input: "redact" }, + { + name: "read", + path: "**/.env.example", + input: "as-is", + output: "as-is", + }, + ], + }); + const args = { + path: "/repo/.env.example", + filePath: "/repo/.env", + }; + redaction.register("sensitive", "read", args); + + expect(redaction.toolInput("read", args)).toBe("[REDACTED]"); + expect(redaction.toolOutput("sensitive", "PASSWORD=secret")).toBe( + "[REDACTED]", + ); +}); + +test("redacts OpenCode 2 calls whose input is a JSON string", () => { + const redaction = new TelemetryRedaction({ + tools: [{ name: "read", path: "**/.env", input: "redact" }], + }); + const snapshot = { + messages: [ + { + type: "tool-call", + id: "json-call", + name: "read", + input: '{"filePath":"/repo/.env"}', + }, + { + type: "tool-result", + id: "json-call", + name: "read", + result: { type: "text", value: "PASSWORD=secret" }, + }, + ], + }; + + expect(redaction.sanitize(snapshot)).toEqual({ + messages: [ + { ...snapshot.messages[0], input: "[REDACTED]" }, + { ...snapshot.messages[1], result: "[REDACTED]" }, + ], + }); +}); + +test("a known call without a path does not suppress path-scoped redaction", () => { + const redaction = new TelemetryRedaction({ + tools: [{ name: "read", path: "**/.env", input: "redact" }], + }); + redaction.register("unknown-input", "read", undefined); + + expect(redaction.toolOutput("unknown-input", "PASSWORD=secret")).toBe( + "[REDACTED]", + ); + expect( + redaction.sanitize({ + type: "tool-result", + id: "unknown-input", + name: "read", + result: { type: "text", value: "PASSWORD=secret" }, + }), + ).toEqual({ + type: "tool-result", + id: "unknown-input", + name: "read", + result: "[REDACTED]", + }); +}); + +test("matches literal regex characters in file paths", () => { + const redaction = new TelemetryRedaction({ + tools: [{ name: "read", path: "**/secret?.txt" }], + }); + redaction.register("literal-question", "read", { + filePath: "/repo/secret?.txt", + }); + redaction.register("different-file", "read", { + filePath: "/repo/secrett.txt", + }); + + expect(redaction.toolOutput("literal-question", "private")).toBe( + "[REDACTED]", + ); + expect(redaction.toolOutput("different-file", "public")).toBe("public"); +}); diff --git a/test/integration/runtime.test.ts b/test/integration/runtime.test.ts index f847d3a..62bca22 100644 --- a/test/integration/runtime.test.ts +++ b/test/integration/runtime.test.ts @@ -1,4 +1,4 @@ -import { mkdtemp, rm } from "node:fs/promises"; +import { mkdir, mkdtemp, rm, writeFile } from "node:fs/promises"; import { join } from "node:path"; import { Effect } from "effect"; @@ -41,6 +41,75 @@ afterEach(async () => { }); describe("Langfuse runtime", () => { + test("loads redaction rules alongside environment credentials", async () => { + const createClient = vi.fn(() => + Effect.succeed({ id: "configured-client" }), + ); + vi.doMock("../../src/langfuse.js", () => ({ + createLangfuseClient: createClient, + })); + temporaryHome = await mkdtemp(join(process.cwd(), ".test-runtime-")); + const directory = join(temporaryHome, ".config", "opencode"); + await mkdir(directory, { recursive: true }); + await writeFile( + join(directory, "opencode-langfuse.json"), + JSON.stringify({ + publicKey: "", + secretKey: false, + redaction: { tools: [{ name: "read", output: "redact" }] }, + }), + ); + process.env.HOME = temporaryHome; + process.env.LANGFUSE_PUBLIC_KEY = "pk-test"; + process.env.LANGFUSE_SECRET_KEY = "sk-test"; + + const { createLangfuseRuntime } = await import("../../src/runtime.js"); + await Effect.runPromise(createLangfuseRuntime({})); + expect(createClient).toHaveBeenCalledWith( + expect.objectContaining({ + redaction: { tools: [{ name: "read", output: "redact" }] }, + }), + ); + }); + + test("rejects invalid redaction rules even with environment credentials", async () => { + const createClient = vi.fn(() => Effect.succeed({ id: "unused" })); + vi.doMock("../../src/langfuse.js", () => ({ + createLangfuseClient: createClient, + })); + temporaryHome = await mkdtemp(join(process.cwd(), ".test-runtime-")); + const directory = join(temporaryHome, ".config", "opencode"); + await mkdir(directory, { recursive: true }); + await writeFile( + join(directory, "opencode-langfuse.json"), + JSON.stringify({ + redaction: { tools: [{ name: "read", output: "typo" }] }, + }), + ); + process.env.HOME = temporaryHome; + process.env.LANGFUSE_PUBLIC_KEY = "pk-test"; + process.env.LANGFUSE_SECRET_KEY = "sk-test"; + + const { createLangfuseRuntime } = await import("../../src/runtime.js"); + const error = await Effect.runPromise( + Effect.flip(createLangfuseRuntime({})), + ); + expect(error).toMatchObject({ _tag: "InvalidLangfuseConfiguration" }); + expect(createClient).not.toHaveBeenCalled(); + + await writeFile( + join(directory, "opencode-langfuse.json"), + JSON.stringify({ redaction: { tool: [{ name: "read" }] } }), + ); + const misspelledKey = await Effect.runPromise( + Effect.flip(createLangfuseRuntime({})), + ); + expect(misspelledKey).toMatchObject({ + _tag: "InvalidLangfuseConfiguration", + }); + expect(createClient).not.toHaveBeenCalled(); + }); + test("creates one shared client for concurrent initialization", async () => { const client = { id: "shared-client" }; let clientCreations = 0; diff --git a/test/integration/v1.test.ts b/test/integration/v1.test.ts index d5010d8..ab9c40c 100644 --- a/test/integration/v1.test.ts +++ b/test/integration/v1.test.ts @@ -1,10 +1,13 @@ import { readFileSync } from "node:fs"; +import { mkdir, mkdtemp, rm, writeFile } from "node:fs/promises"; import { createServer } from "node:http"; import type { IncomingHttpHeaders, Server } from "node:http"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; import { trace } from "@opentelemetry/api"; import LangfusePlugin from "@langfuse/opencode-observability-plugin/v1"; -import { Schema } from "effect"; +import { Effect, Schema, SynchronizedRef } from "effect"; import { afterAll, afterEach, @@ -131,6 +134,7 @@ const packageJson = Schema.decodeUnknownSync( const packageVersion = packageJson.version; const originalEnvironment = { + home: process.env.HOME, publicKey: process.env.LANGFUSE_PUBLIC_KEY, secretKey: process.env.LANGFUSE_SECRET_KEY, baseUrl: process.env.LANGFUSE_BASE_URL, @@ -147,6 +151,7 @@ let plugin: Plugin; let collectorStatus = 200; let hooksDisposed = false; let collectorBaseUrl: string; +let temporaryHome: string; let toolListCalls = 0; let toolListShouldFail = false; @@ -520,6 +525,18 @@ const disposeHooks = async () => { }; beforeAll(async () => { + temporaryHome = await mkdtemp(join(tmpdir(), "langfuse-redaction-")); + const configDirectory = join(temporaryHome, ".config", "opencode"); + await mkdir(configDirectory, { recursive: true }); + await writeFile( + join(configDirectory, "opencode-langfuse.json"), + JSON.stringify({ + redaction: { + tools: [{ name: "read", path: "**/.private", input: "redact" }], + }, + }), + ); + process.env.HOME = temporaryHome; server = createServer((request, response) => { const chunks: Buffer[] = []; @@ -604,7 +621,9 @@ afterAll(async () => { }); }); } finally { + await rm(temporaryHome, { recursive: true, force: true }); for (const [name, value] of Object.entries({ + HOME: originalEnvironment.home, LANGFUSE_PUBLIC_KEY: originalEnvironment.publicKey, LANGFUSE_SECRET_KEY: originalEnvironment.secretKey, LANGFUSE_BASE_URL: originalEnvironment.baseUrl, @@ -622,6 +641,199 @@ afterAll(async () => { }); describe("built plugin", { concurrent: false }, () => { + test("exports redacted tool spans and later generation history", async () => { + const sessionID = "redaction-session"; + const secret = "PASSWORD=not-for-export"; + const path = "/repo/.private"; + await sendUserMessage({ + sessionID, + messageID: "redaction-user-1", + text: "Read the file", + started: startedAt, + }); + await startGeneration({ + id: "redaction-step-1", + sessionID, + assistantMessageID: "redaction-assistant-1", + started: startedAt + 100, + }); + await hooks["tool.execute.before"]?.( + { sessionID, callID: "redaction-call", tool: "read" }, + { args: { filePath: path } }, + ); + await hooks["tool.execute.after"]?.( + { + sessionID, + callID: "redaction-call", + tool: "read", + args: { filePath: path }, + }, + { title: "file", output: secret, metadata: {} }, + ); + await completeGeneration({ + sessionID, + userMessageID: "redaction-user-1", + assistantMessageID: "redaction-assistant-1", + started: startedAt + 100, + completed: startedAt + 200, + }); + storeMessage(sessionID, { + info: { id: "redaction-assistant-1", role: "assistant", sessionID }, + parts: [ + { + id: "redaction-tool-part", + sessionID, + messageID: "redaction-assistant-1", + type: "tool", + callID: "redaction-call", + tool: "read", + state: { + status: "completed", + input: { filePath: path }, + output: secret, + title: "file", + metadata: {}, + time: { start: startedAt + 120, end: startedAt + 150 }, + }, + }, + ], + }); + await sendUserMessage({ + sessionID, + messageID: "redaction-user-2", + text: "What did you find?", + started: startedAt + 300, + }); + await startGeneration({ + id: "redaction-step-2", + sessionID, + assistantMessageID: "redaction-assistant-2", + started: startedAt + 400, + }); + await completeGeneration({ + sessionID, + userMessageID: "redaction-user-2", + assistantMessageID: "redaction-assistant-2", + started: startedAt + 400, + completed: startedAt + 500, + text: "Done", + }); + const { requests: exportedRequests, spans } = await flushSession(sessionID); + const tool = getSessionSpan(spans, "read", sessionID); + expect(getJsonAttribute(tool, "langfuse.observation.input")).toBe( + "[REDACTED]", + ); + expect(getJsonAttribute(tool, "langfuse.observation.output")).toEqual({ + title: "[REDACTED]", + output: "[REDACTED]", + }); + const generations = spans.filter( + (span) => span.name === "opencode.generation", + ); + expect(generations).toHaveLength(2); + expect( + JSON.stringify( + getJsonAttribute(generations[1], "langfuse.observation.input"), + ), + ).toContain('"content":"[REDACTED]"'); + expect( + JSON.stringify(exportedRequests.map((request) => request.body)), + ).not.toContain(secret); + expect( + JSON.stringify(exportedRequests.map((request) => request.body)), + ).not.toContain(path); + }); + + test("exports redacted OpenCode 2 context snapshots", async () => { + const sharedState = globalThis.langfuseOpencodeRuntimeState; + if (sharedState === undefined) { + throw new Error("Expected the shared Langfuse runtime"); + } + const state = await Effect.runPromise(SynchronizedRef.get(sharedState)); + if (state._tag !== "Ready") { + throw new Error("Expected a running Langfuse client"); + } + const langfuse = state.client; + const sessionID = "redaction-v2-session"; + const secret = "PASSWORD=v2-not-for-export"; + langfuse.traceUserPrompt({ + sessionID, + messageID: "redaction-v2-user", + content: [{ type: "text", text: "Check the file" }], + }); + langfuse.setGenerationInputSnapshot(sessionID, { + system: [], + messages: [ + { + role: "assistant", + content: [ + { + type: "tool-call", + id: "redaction-v2-call", + name: "read", + input: { filePath: "/repo/.private" }, + }, + ], + }, + { + role: "tool", + content: [ + { + type: "tool-result", + id: "redaction-v2-call", + name: "read", + result: { type: "text", value: secret }, + }, + ], + }, + ], + tools: [], + }); + langfuse.startActiveGenerationStep({ + sessionID, + assistantMessageID: "redaction-v2-assistant", + agent: "build", + model: { id: "test-model", providerID: "test-provider" }, + started: startedAt, + }); + langfuse.traceGeneration({ + sessionID, + messageID: "redaction-v2-assistant", + parentID: "redaction-v2-user", + modelID: "test-model", + providerID: "test-provider", + mode: "build", + created: startedAt, + completed: startedAt + 100, + cost: 0, + tokens: { + input: 1, + output: 1, + reasoning: 0, + cache: { read: 0, write: 0 }, + }, + output: [{ role: "assistant", content: "Done" }], + }); + + const requestCount = requests.length; + await Effect.runPromise(langfuse.forceFlush); + const exportedRequests = requests.slice(requestCount); + const generation = getSessionSpan( + exportedRequests.flatMap(getSpans), + "opencode.generation", + sessionID, + ); + const input = getJsonAttribute(generation, "langfuse.observation.input"); + expect(JSON.stringify(input)).toContain('"result":"[REDACTED]"'); + expect(JSON.stringify(input)).toContain('"input":"[REDACTED]"'); + expect( + JSON.stringify(exportedRequests.map((request) => request.body)), + ).not.toContain(secret); + expect( + JSON.stringify(exportedRequests.map((request) => request.body)), + ).not.toContain("/repo/.private"); + }); + test("resolves the OpenCode 1 package entrypoint", () => { expect(typeof LangfusePlugin).toBe("function"); });