diff --git a/packages/evals/framework/harnesses/piAdapter.ts b/packages/evals/framework/harnesses/piAdapter.ts index 3d828dc25..dc10bdabe 100644 --- a/packages/evals/framework/harnesses/piAdapter.ts +++ b/packages/evals/framework/harnesses/piAdapter.ts @@ -133,12 +133,14 @@ function normalizeResult(value: unknown): { for (const block of value.content) { if (!isRecord(block)) continue; if (block.type === "text" && typeof block.text === "string") text.push(block.text); - if ( - block.type === "image" && - typeof block.data === "string" && - typeof block.mimeType === "string" - ) { - images.push({ bytes: Buffer.from(block.data, "base64"), mediaType: block.mimeType }); + if (block.type === "image" && typeof block.mimeType === "string") { + // The pi session decodes screenshots to a single Buffer when it retains + // the event; raw base64 only arrives from callers that bypass it. + if (Buffer.isBuffer(block.bytes)) { + images.push({ bytes: block.bytes, mediaType: block.mimeType }); + } else if (typeof block.data === "string") { + images.push({ bytes: Buffer.from(block.data, "base64"), mediaType: block.mimeType }); + } } } } diff --git a/packages/evals/tests/framework/piScreenshotPipeline.test.ts b/packages/evals/tests/framework/piScreenshotPipeline.test.ts new file mode 100644 index 000000000..b61065fac --- /dev/null +++ b/packages/evals/tests/framework/piScreenshotPipeline.test.ts @@ -0,0 +1,60 @@ +import { describe, expect, it } from "vitest"; +import { compactPiEvent } from "@browserbasehq/stagehand-integrations-pi-sdk"; +import { piAdapter } from "../../framework/harnesses/piAdapter.js"; + +describe("Pi screenshot evidence pipeline", () => { + it("keeps explicit omission evidence without re-decoding an over-budget screenshot", () => { + const event = compactPiEvent( + { + type: "tool_execution_end", + toolCallId: "oversize", + toolName: "screenshot", + result: { content: [{ type: "image", data: "AAAA", mimeType: "image/png" }] }, + }, + { remainingBytes: 0 }, + ); + const trajectory = piAdapter.fromHarnessResult( + { events: [event], finalAnswer: "done" }, + { id: "omitted-image", instruction: "Capture a screenshot." }, + ); + const modalities = trajectory.steps[0].agentEvidence.modalities; + expect(modalities.filter((modality) => modality.type === "image")).toEqual([]); + expect(JSON.stringify(modalities)).toContain("Screenshot omitted"); + expect(JSON.stringify(event)).not.toContain("AAAA"); + }); + it.each([false, true])("retains the screenshot after SDK compaction=%s", (compact) => { + const png = Buffer.from( + "iVBORw0KGgoAAAANSUhEUgAAAAEAAAABCAQAAAC1HAwCAAAAC0lEQVR42mP8/x8AAwMCAO+aZu0AAAAASUVORK5CYII=", + "base64", + ); + const event = { + type: "tool_execution_end", + toolCallId: "screenshot-1", + toolName: "mcp__stagehand__screenshot", + result: { + content: [ + { type: "text", text: "Screenshot captured." }, + { type: "image", data: png.toString("base64"), mimeType: "image/png" }, + ], + }, + isError: false, + }; + + const trajectory = piAdapter.fromHarnessResult( + { events: [compact ? compactPiEvent(event) : event], finalAnswer: "done" }, + { + id: "pi-screenshot-pipeline", + instruction: "Capture a screenshot.", + initUrl: "https://example.invalid", + }, + ); + + expect(trajectory.steps).toHaveLength(1); + const images = trajectory.steps[0].agentEvidence.modalities.filter( + (modality) => modality.type === "image", + ); + expect(images).toHaveLength(1); + expect(images[0]).toMatchObject({ bytes: png, mediaType: "image/png" }); + expect(event.result.content[1].data).toBe(png.toString("base64")); + }); +}); diff --git a/packages/integrations/claude-agent-sdk/src/session.ts b/packages/integrations/claude-agent-sdk/src/session.ts index f10b6a72c..6281bde9f 100644 --- a/packages/integrations/claude-agent-sdk/src/session.ts +++ b/packages/integrations/claude-agent-sdk/src/session.ts @@ -1,5 +1,6 @@ import { HarnessAdapterError, + harnessEventLogLevel, sanitizeErrorMessage, type HarnessLogger, } from "@browserbasehq/stagehand-integrations/harness"; @@ -115,7 +116,11 @@ export async function runClaudeAgentSession(input: { permissionMode: input.session.permissionMode ?? "default", settingSources: input.session.settingSources ?? [], stderr: (data: string) => { - input.logger.log({ category: "claude_code", message: data, level: 1 }); + input.logger.log({ + category: "claude_code", + message: data, + level: /\b(?:error|fatal|failed)\b/iu.test(data) ? 1 : 2, + }); }, ...(systemPrompt !== undefined && { systemPrompt }), }, @@ -277,11 +282,25 @@ export function buildClaudeCodeTranscript(messages: ClaudeSdkMessage[]): string } export function logClaudeCodeMessage(logger: HarnessLogger, message: ClaudeSdkMessage): void { + const type = String(message.type ?? "unknown"); + const level = harnessEventLogLevel(type, { + isError: + (type === "result" && message.subtype !== undefined && message.subtype !== "success") || + message.is_error === true || + (type === "user" && + isRecord(message.message) && + Array.isArray(message.message.content) && + message.message.content.some( + (block) => isRecord(block) && block.type === "tool_result" && block.is_error === true, + )), + hasContent: type === "assistant" || type === "user" || type === "result", + }); + if (level === undefined) return; const summary = summarizeClaudeCodeMessage(message); logger.log({ category: "claude_code", message: summary.message, - level: 1, + level, auxiliary: { type: { value: String(message.type ?? "unknown"), type: "string" }, ...(summary.detail && { detail: { value: summary.detail, type: "string" } }), diff --git a/packages/integrations/claude-agent-sdk/tests/eventLog.test.ts b/packages/integrations/claude-agent-sdk/tests/eventLog.test.ts new file mode 100644 index 000000000..25564ba1a --- /dev/null +++ b/packages/integrations/claude-agent-sdk/tests/eventLog.test.ts @@ -0,0 +1,48 @@ +import { describe, expect, it } from "vitest"; +import { logClaudeCodeMessage } from "../src/index.js"; + +function recordingLogger() { + const lines: Array<{ level?: number; message: string }> = []; + const push = (line: { level?: number; message: string }) => void lines.push(line); + return { lines, logger: { log: push, warn: push, error: push } }; +} + +describe("claude_code event log levels", () => { + it("drops partial stream events and demotes routine messages to debug", () => { + const { lines, logger } = recordingLogger(); + logClaudeCodeMessage(logger, { type: "stream_event", event: { type: "content_block_delta" } }); + expect(lines).toEqual([]); + + logClaudeCodeMessage(logger, { type: "system", subtype: "init" }); + logClaudeCodeMessage(logger, { + type: "assistant", + message: { content: [{ type: "text", text: "hello" }] }, + }); + logClaudeCodeMessage(logger, { type: "result", subtype: "success", result: "done" }); + expect(lines.map((line) => [line.level, line.message])).toEqual([ + [2, "system message"], + [2, "assistant: hello"], + [2, "result: success"], + ]); + }); + + it("keeps failures visible", () => { + const { lines, logger } = recordingLogger(); + logClaudeCodeMessage(logger, { type: "result", subtype: "error_max_turns", is_error: true }); + expect(lines).toEqual([ + expect.objectContaining({ level: 1, message: "result: error_max_turns" }), + ]); + }); + it.each([false, true])("classifies nested tool_result is_error=%s", (isError) => { + const { lines, logger } = recordingLogger(); + logClaudeCodeMessage(logger, { + type: "user", + message: { + content: [ + { type: "tool_result", tool_use_id: "fixture", is_error: isError, content: "result" }, + ], + }, + }); + expect(lines[0].level).toBe(isError ? 1 : 2); + }); +}); diff --git a/packages/integrations/codex-sdk/src/session.ts b/packages/integrations/codex-sdk/src/session.ts index 85f97d7dc..8e229f96e 100644 --- a/packages/integrations/codex-sdk/src/session.ts +++ b/packages/integrations/codex-sdk/src/session.ts @@ -1,5 +1,6 @@ import { HarnessAdapterError, + harnessEventLogLevel, sanitizeErrorMessage, type HarnessLogger, } from "@browserbasehq/stagehand-integrations/harness"; @@ -256,11 +257,26 @@ export function buildCodexTranscript(events: CodexEvent[]): string { } export function logCodexEvent(logger: HarnessLogger, event: CodexEvent): void { + const type = String(event.type ?? "unknown"); + const item = isRecord(event.item) ? event.item : undefined; + const level = harnessEventLogLevel(type, { + isError: + type === "turn.failed" || + type === "error" || + item?.type === "error" || + (type === "item.completed" && + (item?.status === "failed" || + (item?.type === "command_execution" && + typeof item.exit_code === "number" && + item.exit_code !== 0))), + hasContent: type === "item.completed" || type === "turn.completed", + }); + if (level === undefined) return; const summary = summarizeCodexEvent(event); logger.log({ category: "codex", message: summary.message, - level: 1, + level, auxiliary: { type: { value: String(event.type ?? "unknown"), type: "string" }, ...(summary.detail && { detail: { value: summary.detail, type: "string" } }), diff --git a/packages/integrations/codex-sdk/tests/eventLog.test.ts b/packages/integrations/codex-sdk/tests/eventLog.test.ts new file mode 100644 index 000000000..7a62a9aa8 --- /dev/null +++ b/packages/integrations/codex-sdk/tests/eventLog.test.ts @@ -0,0 +1,56 @@ +import { describe, expect, it } from "vitest"; +import { logCodexEvent } from "../src/index.js"; + +function recordingLogger() { + const lines: Array<{ level?: number; message: string }> = []; + const push = (line: { level?: number; message: string }) => void lines.push(line); + return { lines, logger: { log: push, warn: push, error: push } }; +} + +describe("codex event log levels", () => { + it("drops item updates and bare lifecycle starts, demotes completed items to debug", () => { + const { lines, logger } = recordingLogger(); + logCodexEvent(logger, { type: "item.updated", item: { type: "agent_message", text: "par" } }); + logCodexEvent(logger, { type: "thread.started", thread_id: "t1" }); + logCodexEvent(logger, { type: "turn.started" }); + logCodexEvent(logger, { + type: "item.started", + item: { type: "mcp_tool_call", server: "stagehand", tool: "run" }, + }); + expect(lines).toEqual([]); + + logCodexEvent(logger, { + type: "item.completed", + item: { type: "mcp_tool_call", server: "stagehand", tool: "run", status: "completed" }, + }); + logCodexEvent(logger, { type: "turn.completed", usage: { input_tokens: 1 } }); + expect(lines.map((line) => [line.level, line.message])).toEqual([ + [2, "mcp: stagehand.run completed"], + [2, "turn completed"], + ]); + }); + + it("keeps failures visible", () => { + const { lines, logger } = recordingLogger(); + logCodexEvent(logger, { type: "turn.failed", error: { message: "boom" } }); + logCodexEvent(logger, { + type: "item.completed", + item: { type: "mcp_tool_call", server: "stagehand", tool: "run", status: "failed" }, + }); + logCodexEvent(logger, { type: "item.completed", item: { type: "error", message: "bad" } }); + expect(lines.map((line) => line.level)).toEqual([1, 1, 1]); + }); + it.each([0, 1, -1])("classifies completed command exit_code=%s", (exitCode) => { + const { lines, logger } = recordingLogger(); + logCodexEvent(logger, { + type: "item.completed", + item: { + type: "command_execution", + status: "completed", + command: "fixture", + exit_code: exitCode, + }, + }); + expect(lines[0].level).toBe(exitCode === 0 ? 2 : 1); + }); +}); diff --git a/packages/integrations/core/src/harness/eventLog.ts b/packages/integrations/core/src/harness/eventLog.ts new file mode 100644 index 000000000..62452dda0 --- /dev/null +++ b/packages/integrations/core/src/harness/eventLog.ts @@ -0,0 +1,43 @@ +/** Log level a raw SDK event should be recorded at. */ +export type HarnessEventLogLevel = 0 | 1 | 2; + +export interface HarnessEventClassification { + /** The event reports a failure the operator should see without debug output. */ + isError?: boolean; + /** + * The event carries complete content of its own (a finished message, a tool + * result, usage). Lifecycle boundary events without content are dropped. + */ + hasContent?: boolean; +} + +/** + * Stream fragments — deltas, partial updates, chunks — that only accumulate + * into a later completed event and never carry standalone content. + */ +export function isStreamDeltaEventType(type: string): boolean { + return ( + type === "stream_event" || /(?:^|[-_.:])(?:delta|update|updated|partial|chunk)$/iu.test(type) + ); +} + +/** Lifecycle markers (`*-start`, `*.started`, `*_end`, ...) that bracket other events. */ +export function isLifecycleBoundaryEventType(type: string): boolean { + return /(?:^|[-_.:])(?:start|started|begin|end|ended|stop|stopped)$/iu.test(type); +} + +/** + * Decide how a raw SDK event lands in the harness log. The readable per-step + * trace is emitted from the normalized trajectory after the run, so raw events + * are debug material: errors stay visible at level 1, pure stream noise is + * dropped (`undefined`), and everything else is kept at level 2. + */ +export function harnessEventLogLevel( + type: string, + classification: HarnessEventClassification = {}, +): HarnessEventLogLevel | undefined { + if (classification.isError) return 1; + if (isStreamDeltaEventType(type)) return undefined; + if (isLifecycleBoundaryEventType(type) && !classification.hasContent) return undefined; + return 2; +} diff --git a/packages/integrations/core/src/harness/index.ts b/packages/integrations/core/src/harness/index.ts index 75962aa09..2f4524337 100644 --- a/packages/integrations/core/src/harness/index.ts +++ b/packages/integrations/core/src/harness/index.ts @@ -1,3 +1,4 @@ export * from "./contract.js"; export * from "./env.js"; +export * from "./eventLog.js"; export * from "./redact.js"; diff --git a/packages/integrations/core/tests/event-log.test.ts b/packages/integrations/core/tests/event-log.test.ts new file mode 100644 index 000000000..1ad03b9dd --- /dev/null +++ b/packages/integrations/core/tests/event-log.test.ts @@ -0,0 +1,43 @@ +import { describe, expect, it } from "vitest"; +import { + harnessEventLogLevel, + isLifecycleBoundaryEventType, + isStreamDeltaEventType, +} from "../src/harness/eventLog.js"; + +describe("harness event log levels", () => { + it("recognizes stream fragments across SDK vocabularies", () => { + for (const type of [ + "text-delta", + "reasoning-delta", + "tool-call-delta", + "message_update", + "tool_execution_update", + "item.updated", + "message.delta", + "stream_event", + "content_block_delta", + ]) { + expect(isStreamDeltaEventType(type), type).toBe(true); + } + expect(isStreamDeltaEventType("tool-call")).toBe(false); + expect(isStreamDeltaEventType("message_end")).toBe(false); + }); + + it("recognizes lifecycle boundaries", () => { + for (const type of ["reasoning-start", "turn.started", "message_end", "step-start"]) { + expect(isLifecycleBoundaryEventType(type), type).toBe(true); + } + expect(isLifecycleBoundaryEventType("step-finish")).toBe(false); + expect(isLifecycleBoundaryEventType("turn.completed")).toBe(false); + }); + + it("drops noise, demotes the rest, and keeps errors at level 1", () => { + expect(harnessEventLogLevel("text-delta")).toBeUndefined(); + expect(harnessEventLogLevel("reasoning-start")).toBeUndefined(); + expect(harnessEventLogLevel("message_end", { hasContent: true })).toBe(2); + expect(harnessEventLogLevel("tool-call")).toBe(2); + expect(harnessEventLogLevel("text-delta", { isError: true })).toBe(1); + expect(harnessEventLogLevel("turn.failed", { isError: true })).toBe(1); + }); +}); diff --git a/packages/integrations/deepagents-sdk/src/session.ts b/packages/integrations/deepagents-sdk/src/session.ts index 13a8c8c7e..98e39cdd9 100644 --- a/packages/integrations/deepagents-sdk/src/session.ts +++ b/packages/integrations/deepagents-sdk/src/session.ts @@ -3,6 +3,7 @@ import { fileURLToPath } from "node:url"; import path from "node:path"; import { createInterface } from "node:readline"; import { + harnessEventLogLevel, sanitizeErrorMessage, type HarnessLogger, } from "@browserbasehq/stagehand-integrations/harness"; @@ -366,11 +367,17 @@ export function buildDeepagentsTranscript(events: DeepagentsEvent[]): string { } export function logDeepagentsEvent(logger: HarnessLogger, event: DeepagentsEvent): void { + const type = String(event.type ?? "unknown"); + const level = harnessEventLogLevel(type, { + isError: type === "error" || (type === "tool_result" && event.ok === false), + hasContent: type === "assistant" || type === "tool_result" || type === "final", + }); + if (level === undefined) return; const summary = summarizeDeepagentsEvent(event); logger.log({ category: "deepagents", message: sanitizeErrorMessage(summary.message), - level: 1, + level, auxiliary: { type: { value: String(event.type ?? "unknown"), type: "string" }, ...(summary.detail && { diff --git a/packages/integrations/deepagents-sdk/tests/eventLog.test.ts b/packages/integrations/deepagents-sdk/tests/eventLog.test.ts new file mode 100644 index 000000000..1c21cb675 --- /dev/null +++ b/packages/integrations/deepagents-sdk/tests/eventLog.test.ts @@ -0,0 +1,32 @@ +import { describe, expect, it } from "vitest"; +import { logDeepagentsEvent } from "../src/index.js"; + +function recordingLogger() { + const lines: Array<{ level?: number; message: string }> = []; + const push = (line: { level?: number; message: string }) => void lines.push(line); + return { lines, logger: { log: push, warn: push, error: push } }; +} + +describe("deepagents event log levels", () => { + it("demotes routine events to debug and keeps errors visible", () => { + const { lines, logger } = recordingLogger(); + logDeepagentsEvent(logger, { type: "message_delta", text: "pa" }); + logDeepagentsEvent(logger, { type: "assistant", text: "hi" }); + logDeepagentsEvent(logger, { type: "tool_result", server: "stagehand", name: "run", ok: true }); + logDeepagentsEvent(logger, { type: "usage", input_tokens: 1 }); + logDeepagentsEvent(logger, { + type: "tool_result", + server: "stagehand", + name: "run", + ok: false, + }); + logDeepagentsEvent(logger, { type: "error", message: "boom" }); + expect(lines.map((line) => [line.level, line.message])).toEqual([ + [2, "assistant: hi"], + [2, "tool: stagehand.run ok"], + [2, "usage"], + [1, "tool: stagehand.run error"], + [1, "error: boom"], + ]); + }); +}); diff --git a/packages/integrations/eve-sdk/src/session.ts b/packages/integrations/eve-sdk/src/session.ts index 314eee0a4..8d0a3bfc2 100644 --- a/packages/integrations/eve-sdk/src/session.ts +++ b/packages/integrations/eve-sdk/src/session.ts @@ -5,6 +5,7 @@ import path from "node:path"; import { fileURLToPath } from "node:url"; import { HarnessAdapterError, + harnessEventLogLevel, sanitizeErrorMessage, type HarnessLogger, } from "@browserbasehq/stagehand-integrations/harness"; @@ -373,13 +374,23 @@ export function buildEveTranscript(events: EveEvent[]): string { } export function logEveEvent(logger: HarnessLogger, event: EveEvent): void { + const level = harnessEventLogLevel(event.type, { + isError: + event.type.endsWith(".failed") || + (event.type === "action.result" && + isRecord(event.data) && + (event.data.status !== "completed" || + (isRecord(event.data.result) && event.data.result.isError === true))), + hasContent: event.type.endsWith(".completed") || event.type === "action.result", + }); + if (level === undefined) return; const summary = summarizeEveEvent(event); const message = sanitizeErrorMessage(summary.message); const detail = summary.detail ? sanitizeErrorMessage(summary.detail) : undefined; logger.log({ category: "eve", message, - level: 1, + level, auxiliary: { type: { value: event.type, type: "string" }, ...(detail && { detail: { value: detail, type: "string" } }), diff --git a/packages/integrations/eve-sdk/tests/eventLog.test.ts b/packages/integrations/eve-sdk/tests/eventLog.test.ts new file mode 100644 index 000000000..1a1133bfd --- /dev/null +++ b/packages/integrations/eve-sdk/tests/eventLog.test.ts @@ -0,0 +1,47 @@ +import { describe, expect, it } from "vitest"; +import { logEveEvent } from "../src/index.js"; + +function recordingLogger() { + const lines: Array<{ level?: number; message: string }> = []; + const push = (line: { level?: number; message: string }) => void lines.push(line); + return { lines, logger: { log: push, warn: push, error: push } }; +} + +describe("eve event log levels", () => { + it("drops deltas and bare lifecycle markers", () => { + const { lines, logger } = recordingLogger(); + for (const type of ["message.delta", "step.started", "turn.started", "session.started"]) { + logEveEvent(logger, { type, data: {} }); + } + expect(lines).toEqual([]); + }); + + it("demotes completed steps and tool results to debug and keeps failures visible", () => { + const { lines, logger } = recordingLogger(); + logEveEvent(logger, { type: "message.completed", data: { message: "done" } }); + logEveEvent(logger, { + type: "action.result", + data: { status: "completed", result: { toolName: "run" } }, + }); + logEveEvent(logger, { type: "step.completed", data: { usage: {} } }); + logEveEvent(logger, { type: "turn.failed", data: { message: "boom" } }); + expect(lines.map((line) => [line.level, line.message])).toEqual([ + [2, "agent: done"], + [2, "tool: run completed"], + [2, "step completed"], + [1, "turn failed: boom"], + ]); + }); + it.each([ + { status: "completed", isError: false, level: 2 }, + { status: "failed", isError: false, level: 1 }, + { status: "completed", isError: true, level: 1 }, + ])("classifies action.result $status isError=$isError", ({ status, isError, level }) => { + const { lines, logger } = recordingLogger(); + logEveEvent(logger, { + type: "action.result", + data: { status, result: { toolName: "run", isError } }, + }); + expect(lines[0].level).toBe(level); + }); +}); diff --git a/packages/integrations/fx-sdk/src/session.ts b/packages/integrations/fx-sdk/src/session.ts index ee523e0b6..0c41db20a 100644 --- a/packages/integrations/fx-sdk/src/session.ts +++ b/packages/integrations/fx-sdk/src/session.ts @@ -4,6 +4,7 @@ import fsp from "node:fs/promises"; import path from "node:path"; import { HarnessAdapterError, + harnessEventLogLevel, sanitizeErrorMessage, type HarnessLogger, } from "@browserbasehq/stagehand-integrations/harness"; @@ -703,11 +704,24 @@ export function buildFxTranscript(events: FxEvent[]): string { } export function logFxEvent(logger: HarnessLogger, event: FxEvent): void { + const level = harnessEventLogLevel(event.type, { + isError: + (event.type === "stderr" && /\b(?:error|fatal|failed|panic)\b/iu.test(event.line)) || + (event.type === "tool_step" && + event.tool_results.some((result) => + /^(?:error|failed|failure)$/iu.test(result.status ?? ""), + )) || + (event.type === "ask_result" && + typeof event.ask.error === "string" && + event.ask.error.length > 0), + hasContent: true, + }); + if (level === undefined) return; const summary = summarizeFxEvent(event); logger.log({ category: "fx", message: summary.message, - level: 1, + level, auxiliary: { type: { value: event.type, type: "string" }, ...(summary.detail && { detail: { value: summary.detail, type: "string" } }), diff --git a/packages/integrations/fx-sdk/tests/eventLog.test.ts b/packages/integrations/fx-sdk/tests/eventLog.test.ts new file mode 100644 index 000000000..cfb056532 --- /dev/null +++ b/packages/integrations/fx-sdk/tests/eventLog.test.ts @@ -0,0 +1,50 @@ +import { describe, expect, it } from "vitest"; +import { logFxEvent } from "../src/index.js"; + +function recordingLogger() { + const lines: Array<{ level?: number; message: string }> = []; + const push = (line: { level?: number; message: string }) => void lines.push(line); + return { lines, logger: { log: push, warn: push, error: push } }; +} + +describe("fx event log levels", () => { + it("demotes routine events to debug and keeps failures visible", () => { + const { lines, logger } = recordingLogger(); + logFxEvent(logger, { type: "assistant", text: "hi" }); + logFxEvent(logger, { + type: "tool_step", + assistant: "", + tool_calls: [{ name: "run" }], + tool_results: [{ tool_name: "run", status: "completed" }], + }); + logFxEvent(logger, { type: "stderr", line: "loading plugins" }); + logFxEvent(logger, { type: "turn_committed", terminal_reason: "done" }); + logFxEvent(logger, { type: "stderr", line: "Error: could not reach MCP server" }); + logFxEvent(logger, { + type: "tool_step", + assistant: "", + tool_calls: [{ name: "run" }], + tool_results: [{ tool_name: "run", status: "error" }], + }); + expect(lines.map((line) => [line.level, line.message])).toEqual([ + [2, "assistant: hi"], + [2, "tools: run"], + [2, "stderr: loading plugins"], + [2, "turn committed: done"], + [1, "stderr: Error: could not reach MCP server"], + [1, "tools: run"], + ]); + }); + it("keeps failure-status tool results and structured ask errors visible", () => { + const { lines, logger } = recordingLogger(); + logFxEvent(logger, { + type: "tool_step", + assistant: "", + tool_calls: [], + tool_results: [{ status: "failure" }], + }); + logFxEvent(logger, { type: "ask_result", ask: { error: "provider unavailable" } }); + logFxEvent(logger, { type: "ask_result", ask: { output: "done" } }); + expect(lines.map((line) => line.level)).toEqual([1, 1, 2]); + }); +}); diff --git a/packages/integrations/mastra-sdk/src/session.ts b/packages/integrations/mastra-sdk/src/session.ts index 6e185bb5f..9986a31e6 100644 --- a/packages/integrations/mastra-sdk/src/session.ts +++ b/packages/integrations/mastra-sdk/src/session.ts @@ -5,6 +5,7 @@ import { createOpenAI } from "@ai-sdk/openai"; import { createOpenAICompatible } from "@ai-sdk/openai-compatible"; import { HarnessAdapterError, + harnessEventLogLevel, sanitizeErrorMessage, type HarnessLogger, } from "@browserbasehq/stagehand-integrations/harness"; @@ -265,7 +266,7 @@ export async function runMastraSession(input: { }); for await (const event of stream.fullStream) { - events.push(event); + events.push(compactMastraEvent(event)); logMastraEvent(input.logger, event); const payload = isRecord(event.payload) ? event.payload : {}; if (event.type === "tool-call") { @@ -407,12 +408,18 @@ export function buildMastraTranscript(events: MastraEvent[]): string { } export function logMastraEvent(logger: HarnessLogger, event: MastraEvent): void { - const summary = summarizeMastraEvent(event); const type = String(event.type ?? "unknown"); + const level = harnessEventLogLevel(type, { + isError: type === "error" || type === "tool-error", + hasContent: type === "tool-call" || type === "tool-result", + }); + if (level === undefined) return; + const summary = summarizeMastraEvent(event); + if (summary.detail) summary.detail = clip(summary.detail, MAX_EVENT_DETAIL_CHARS); logger.log({ category: "mastra", message: summary.message, - level: type === "text-delta" || type === "reasoning-delta" ? 2 : 1, + level, auxiliary: { type: { value: type, type: "string" }, ...(summary.detail && { detail: { value: summary.detail, type: "string" } }), @@ -445,17 +452,57 @@ export function summarizeMastraEvent(event: MastraEvent): { if (type === "reasoning-delta" && typeof payload.text === "string") { return sanitizeMastraSummary(`reasoning: ${clip(payload.text, 500)}`, payload.text); } - if (type === "finish") { + if (type === "finish" || type === "step-finish") { const stepResult = isRecord(payload.stepResult) ? payload.stepResult : undefined; const reason = String(stepResult?.reason ?? "unknown"); const output = isRecord(payload.output) ? payload.output : undefined; - return sanitizeMastraSummary(`finish: ${reason}`, safeJson(output?.usage)); + return sanitizeMastraSummary(`${type}: ${reason}`, safeJson(output?.usage)); } if (type === "error") { const message = stringifyError(payload.error) || "error"; return sanitizeMastraSummary(`error: ${clip(message, 500)}`, message); } - return sanitizeMastraSummary(`${type} event`, safeJson(event)); + const detail = safeJson(compactMastraEvent(event)); + return sanitizeMastraSummary(`${type} event`, detail); +} + +const MAX_EVENT_DETAIL_CHARS = 2_000; + +/** Chunk types the trajectory adapter and transcript read verbatim. */ +const RETAINED_MASTRA_EVENT_TYPES = new Set([ + "tool-call", + "tool-result", + "tool-error", + "text-delta", + "reasoning-delta", + "error", + "abort", +]); + +/** + * Reduce a fullStream chunk to what the session consumers need. Mastra's + * `step-finish`/`finish` payloads carry the whole request body (every prior + * message and tool result) under `metadata.request`/`response`, so retaining or + * stringifying them per step grows quadratically with the conversation. + */ +export function compactMastraEvent(event: MastraEvent): MastraEvent { + const type = String(event.type ?? ""); + if (RETAINED_MASTRA_EVENT_TYPES.has(type)) return event; + const payload = isRecord(event.payload) ? event.payload : undefined; + const output = isRecord(payload?.output) ? payload.output : undefined; + const compactPayload: Record = { + ...(payload?.stepResult !== undefined && { stepResult: payload.stepResult }), + ...(payload?.reason !== undefined && { reason: payload.reason }), + ...(payload?.toolName !== undefined && { toolName: payload.toolName }), + ...(payload?.toolCallId !== undefined && { toolCallId: payload.toolCallId }), + ...(output?.usage !== undefined && { output: { usage: output.usage } }), + }; + return { + type: event.type, + ...(event.runId !== undefined && { runId: event.runId }), + ...(event.from !== undefined && { from: event.from }), + ...(payload !== undefined && { payload: compactPayload }), + }; } export function isRecord(value: unknown): value is Record { diff --git a/packages/integrations/mastra-sdk/tests/eventLog.test.ts b/packages/integrations/mastra-sdk/tests/eventLog.test.ts new file mode 100644 index 000000000..4b5e74187 --- /dev/null +++ b/packages/integrations/mastra-sdk/tests/eventLog.test.ts @@ -0,0 +1,72 @@ +import { describe, expect, it, vi } from "vitest"; +import { logMastraEvent, buildMastraTranscript } from "../src/index.js"; + +function recordingLogger() { + const lines: Array<{ level?: number; message: string }> = []; + const push = (line: { level?: number; message: string }) => void lines.push(line); + return { lines, logger: { log: push, warn: push, error: push } }; +} + +describe("mastra event log levels", () => { + it.each(["tool-call", "tool-result", "tool-error", "error"])( + "bounds %s log details without truncating trajectory data", + (type) => { + const text = "x".repeat(10_000); + const event = { + type, + payload: { toolName: "fixture", args: { text }, result: text, error: text }, + }; + const log = vi.fn(); + logMastraEvent({ log, warn: log, error: log }, event); + expect(log.mock.calls[0][0].auxiliary.detail.value.length).toBeLessThanOrEqual(2_000); + expect(buildMastraTranscript([event])).toContain(text); + expect(event.payload.result).toBe(text); + }, + ); + it("drops stream deltas and bare start/end markers", () => { + const { lines, logger } = recordingLogger(); + for (const type of [ + "tool-call-delta", + "reasoning-delta", + "text-delta", + "reasoning-start", + "reasoning-end", + "text-start", + "text-end", + "tool-call-input-streaming-start", + "tool-call-input-streaming-end", + "step-start", + ]) { + logMastraEvent(logger, { type, payload: { text: "x" } }); + } + expect(lines).toEqual([]); + }); + + it("demotes tool calls, results and finishes to debug and keeps errors visible", () => { + const { lines, logger } = recordingLogger(); + logMastraEvent(logger, { + type: "tool-call", + payload: { toolName: "stagehand_run", args: { code: "1" } }, + }); + logMastraEvent(logger, { + type: "tool-result", + payload: { toolName: "stagehand_run", result: "ok" }, + }); + logMastraEvent(logger, { + type: "step-finish", + payload: { stepResult: { reason: "tool-calls" } }, + }); + logMastraEvent(logger, { + type: "tool-error", + payload: { toolName: "stagehand_run", error: "nope" }, + }); + logMastraEvent(logger, { type: "error", payload: { error: "fatal" } }); + expect(lines.map((line) => [line.level, line.message])).toEqual([ + [2, 'tool: stagehand_run {"code":"1"}'], + [2, "tool result: stagehand_run ok"], + [2, "step-finish: tool-calls"], + [1, "tool error: stagehand_run nope"], + [1, "error: fatal"], + ]); + }); +}); diff --git a/packages/integrations/mastra-sdk/tests/session.test.ts b/packages/integrations/mastra-sdk/tests/session.test.ts index 0b7c95faf..095b9d0ed 100644 --- a/packages/integrations/mastra-sdk/tests/session.test.ts +++ b/packages/integrations/mastra-sdk/tests/session.test.ts @@ -2,6 +2,7 @@ import { describe, expect, it, vi } from "vitest"; import { buildMastraTranscript, + compactMastraEvent, normalizeMastraModel, runMastraSession, type MastraEvent, @@ -288,23 +289,65 @@ describe("Mastra SDK session", () => { } }); - it("sanitizes the returned final text", async () => { + it("does not retain or log step-finish request bodies", async () => { + const history = "x".repeat(200_000); + const log = vi.fn(); const result = await runMastraSession({ prompt: "task", model: "gpt-5.4-mini", - logger, + logger: { ...logger, log }, sdk: fakeSdk({ events: [ + { type: "step-start", payload: { messageId: "m1", request: { body: history } } }, { - type: "text-delta", - payload: { text: "done https://x.test?apiKey=secret123" }, + type: "tool-call", + payload: { toolCallId: "1", toolName: "stagehand_run", args: { code: "1" } }, + }, + { + type: "step-finish", + payload: { + stepResult: { reason: "tool-calls", isContinued: true }, + output: { usage: { inputTokens: 7, outputTokens: 3 } }, + metadata: { + request: { body: history }, + response: { messages: [{ role: "assistant", content: history }] }, + }, + messages: { all: [{ role: "user", content: history }] }, + }, + }, + { + type: "finish", + payload: { + stepResult: { reason: "stop" }, + output: { usage: { inputTokens: 7, outputTokens: 3 } }, + metadata: { request: { body: history } }, + }, }, ], }), session: {}, }); - expect(result.finalText).toBe("done https://x.test?apiKey=[redacted]"); + expect(result.events).toHaveLength(4); + expect(result.events[1]).toMatchObject({ type: "tool-call" }); + expect(result.events[2]).toEqual({ + type: "step-finish", + payload: { + stepResult: { reason: "tool-calls", isContinued: true }, + output: { usage: { inputTokens: 7, outputTokens: 3 } }, + }, + }); + expect(JSON.stringify(result.events)).not.toContain(history.slice(0, 10_000)); + expect(result.finishReason).toBe("stop"); + expect(result.tokenUsage.inputTokens).toBe(7); + + const logged = JSON.stringify(log.mock.calls); + expect(logged).not.toContain(history.slice(0, 10_000)); + expect(logged.length).toBeLessThan(20_000); + expect(compactMastraEvent({ type: "text-delta", payload: { text: "hi" } })).toEqual({ + type: "text-delta", + payload: { text: "hi" }, + }); }); it("aborts MCP discovery and disconnects without creating an agent", async () => { @@ -365,4 +408,23 @@ describe("Mastra SDK session", () => { expect(result.status).toBe("sdk_error"); expect(result.stopReason).toBe("stop"); }); + + it("sanitizes the returned final text", async () => { + const result = await runMastraSession({ + prompt: "task", + model: "gpt-5.4-mini", + logger, + sdk: fakeSdk({ + events: [ + { + type: "text-delta", + payload: { text: "done https://x.test?apiKey=secret123" }, + }, + ], + }), + session: {}, + }); + + expect(result.finalText).toBe("done https://x.test?apiKey=[redacted]"); + }); }); diff --git a/packages/integrations/pi-sdk/README.md b/packages/integrations/pi-sdk/README.md new file mode 100644 index 000000000..5e255b183 --- /dev/null +++ b/packages/integrations/pi-sdk/README.md @@ -0,0 +1,14 @@ +# Pi SDK event retention + +The session retains screenshot evidence up to 8 MiB per image and 64 MiB across +one run. It checks the encoded size before allocating a decoded buffer. Images +within those limits remain usable screenshot bytes; images exceeding either +limit become an explicit text omission marker in the retained trajectory. The +marker contains no image payload for downstream adapters to decode again. + +These limits apply to the harness's retained evidence. They do not change the +tool result the model receives or bound Pi's own upstream message storage. +Whitespace-heavy or otherwise noncanonical base64 may be rejected conservatively. + +Normal completed-event logs are debug detail. Tool/provider failures stay visible +at level 1; log summaries redact credentials before applying their text limit. diff --git a/packages/integrations/pi-sdk/src/session.ts b/packages/integrations/pi-sdk/src/session.ts index 9a8542642..ef4e9ebb3 100644 --- a/packages/integrations/pi-sdk/src/session.ts +++ b/packages/integrations/pi-sdk/src/session.ts @@ -1,5 +1,6 @@ import { HarnessAdapterError, + harnessEventLogLevel, sanitizeErrorMessage, type HarnessLogger, } from "@browserbasehq/stagehand-integrations/harness"; @@ -161,6 +162,7 @@ export async function runPiSession(input: { }): Promise { const sdk = input.sdk ?? (await loadPiSdk({ logger: input.logger })); const events: PiEvent[] = []; + const imageBudget = { remainingBytes: MAX_PI_SESSION_IMAGE_BYTES }; let iterationError: unknown; let stopReason: string | undefined; let piSession: PiAgentSessionLike | undefined; @@ -192,7 +194,7 @@ export async function runPiSession(input: { piSession.agent.shouldStopAfterTurn = () => turns >= maxTurns; unsubscribe = piSession.subscribe((event) => { if (event.type === "message_update") return; - events.push(event); + events.push(compactPiEvent(event, imageBudget)); logPiEvent(input.logger, event); if (event.type === "turn_end") turns += 1; if (event.type === "tool_execution_end" && typeof event.toolName === "string") { @@ -302,11 +304,20 @@ export function buildPiTranscript(events: PiEvent[]): string { } export function logPiEvent(logger: HarnessLogger, event: PiEvent): void { + const type = String(event.type ?? "unknown"); + const level = harnessEventLogLevel(type, { + isError: + type === "error" || + (type === "message_end" && isRecord(event.message) && event.message.stopReason === "error") || + (type === "tool_execution_end" && event.isError === true), + hasContent: type === "message_end" || type === "tool_execution_end", + }); + if (level === undefined) return; const summary = summarizePiEvent(event); logger.log({ category: "pi", message: summary.message, - level: 1, + level, auxiliary: { type: { value: String(event.type ?? "unknown"), type: "string" }, ...(summary.detail && { detail: { value: summary.detail, type: "string" } }), @@ -318,26 +329,114 @@ export function summarizePiEvent(event: PiEvent): { message: string; detail?: st const type = String(event.type ?? "unknown"); if (type === "message_end" && isRecord(event.message)) { const text = assistantText(event.message); - const detail = text || safeJson(event.message); + const detail = text || safeJson(withoutImageData(event.message)); return { - message: sanitizeErrorMessage(`assistant: ${clip(text, 500)}`), - ...(detail && { detail: sanitizeErrorMessage(detail) }), + message: `assistant: ${clip(sanitizeErrorMessage(text), 500)}`, + ...(detail && { detail: clip(sanitizeErrorMessage(detail), MAX_EVENT_DETAIL_CHARS) }), }; } if (type.startsWith("tool_execution_")) { - const detail = safeJson(event); + const detail = safeJson(withoutImageData(event)); return { message: sanitizeErrorMessage(`${type}: ${String(event.toolName ?? "tool")}`), - ...(detail && { detail: sanitizeErrorMessage(detail) }), + ...(detail && { detail: clip(sanitizeErrorMessage(detail), MAX_EVENT_DETAIL_CHARS) }), }; } - const detail = safeJson(event); + const detail = safeJson(withoutImageData(event)); return { message: sanitizeErrorMessage(`${type} event`), - ...(detail && { detail: sanitizeErrorMessage(detail) }), + ...(detail && { detail: clip(sanitizeErrorMessage(detail), MAX_EVENT_DETAIL_CHARS) }), + }; +} + +const MAX_EVENT_DETAIL_CHARS = 20_000; +/** Retained trajectory evidence only; model-facing tool results are unchanged. */ +export const MAX_PI_IMAGE_BYTES = 8 * 1024 * 1024; +export const MAX_PI_SESSION_IMAGE_BYTES = 64 * 1024 * 1024; + +/** + * Reduce a retained pi event to what the trajectory adapter and usage + * accounting read. pi emits every message (including tool results carrying + * screenshots) both as its own `message_end` and inside `tool_execution_end`, + * so screenshots are decoded to a single Buffer on the tool event and dropped + * from non-assistant messages, which nothing downstream reads. + */ +export function compactPiEvent( + event: PiEvent, + imageBudget: { remainingBytes: number } = { remainingBytes: MAX_PI_SESSION_IMAGE_BYTES }, +): PiEvent { + if (event.type === "tool_execution_end") { + return { ...event, result: decodeImageBlocks(event.result, imageBudget) }; + } + if ( + event.type === "message_end" && + isRecord(event.message) && + event.message.role !== "assistant" + ) { + return { ...event, message: withoutImageData(event.message) }; + } + return event; +} + +function decodeImageBlocks(value: unknown, budget: { remainingBytes: number }): unknown { + if (!isRecord(value) || !Array.isArray(value.content)) return value; + return { + ...value, + content: value.content.map((block) => { + if (!isRecord(block) || block.type !== "image") { + return block; + } + const data = typeof block.data === "string" ? block.data : undefined; + const bytes = Buffer.isBuffer(block.bytes) ? block.bytes : undefined; + if (data === undefined && !bytes) return block; + // An upper bound from encoded length prevents allocating an oversized + // Buffer. Noncanonical/whitespace-heavy base64 may be rejected conservatively. + const size = + bytes?.byteLength ?? + Math.max( + 0, + Math.ceil(data!.length / 4) * 3 - + (data!.endsWith("==") ? 2 : data!.endsWith("=") ? 1 : 0), + ); + const limit = + size > MAX_PI_IMAGE_BYTES + ? `${MAX_PI_IMAGE_BYTES}-byte per-image` + : size > budget.remainingBytes + ? `${budget.remainingBytes}-byte remaining session image` + : undefined; + if (limit) { + // A text block leaves neither base64 nor Buffer for an adapter to decode. + return { + type: "text", + text: `[Screenshot omitted from retained evidence: ${size} bytes exceeds the ${limit} budget.]`, + }; + } + const retained = bytes ?? Buffer.from(data!, "base64"); + budget.remainingBytes -= retained.byteLength; + const { data: _data, bytes: _bytes, ...rest } = block; + return { ...rest, bytes: retained }; + }), }; } +/** Deep copy with image payloads replaced by a size placeholder (for logs). */ +export function withoutImageData(value: T): T { + if (Array.isArray(value)) return value.map((item) => withoutImageData(item)) as T; + if (Buffer.isBuffer(value)) return `[${value.byteLength} bytes]` as T; + if (!isRecord(value)) return value; + if (value.type === "image" && (typeof value.data === "string" || Buffer.isBuffer(value.bytes))) { + const size = + typeof value.data === "string" + ? Math.floor((value.data.length * 3) / 4) + : (value.bytes as Buffer).byteLength; + const { data: _data, bytes: _bytes, ...rest } = value; + return { ...rest, data: `[image ${size} bytes]` } as T; + } + return Object.fromEntries( + Object.entries(value).map(([key, entry]) => [key, withoutImageData(entry)]), + ) as T; +} + export function resolvePiStatus(input: { iterationError?: unknown; stopReason?: string; diff --git a/packages/integrations/pi-sdk/tests/eventLog.test.ts b/packages/integrations/pi-sdk/tests/eventLog.test.ts new file mode 100644 index 000000000..91a854a23 --- /dev/null +++ b/packages/integrations/pi-sdk/tests/eventLog.test.ts @@ -0,0 +1,50 @@ +import { describe, expect, it } from "vitest"; +import { logPiEvent } from "../src/index.js"; + +function recordingLogger() { + const lines: Array<{ level?: number; message: string }> = []; + const push = (line: { level?: number; message: string }) => void lines.push(line); + return { lines, logger: { log: push, warn: push, error: push } }; +} + +describe("pi event log levels", () => { + it("drops message updates and bare lifecycle markers", () => { + const { lines, logger } = recordingLogger(); + for (const type of [ + "message_update", + "tool_execution_update", + "message_start", + "tool_execution_start", + "agent_start", + "turn_start", + "agent_end", + "turn_end", + ]) { + logPiEvent(logger, { type, toolName: "run" }); + } + expect(lines).toEqual([]); + }); + + it("demotes completed messages and tool results to debug and keeps tool errors visible", () => { + const { lines, logger } = recordingLogger(); + logPiEvent(logger, { + type: "message_end", + message: { role: "assistant", content: [{ type: "text", text: "hi" }] }, + }); + logPiEvent(logger, { type: "tool_execution_end", toolName: "run", result: "ok" }); + logPiEvent(logger, { type: "tool_execution_end", toolName: "run", isError: true, result: "x" }); + expect(lines.map((line) => [line.level, line.message])).toEqual([ + [2, "assistant: hi"], + [2, "tool_execution_end: run"], + [1, "tool_execution_end: run"], + ]); + }); + it.each(["error", "stop"])("classifies message_end stopReason=%s", (stopReason) => { + const { lines, logger } = recordingLogger(); + logPiEvent(logger, { + type: "message_end", + message: { role: "assistant", stopReason, errorMessage: "provider unavailable" }, + }); + expect(lines[0].level).toBe(stopReason === "error" ? 1 : 2); + }); +}); diff --git a/packages/integrations/pi-sdk/tests/eventRetention.test.ts b/packages/integrations/pi-sdk/tests/eventRetention.test.ts new file mode 100644 index 000000000..c63957618 --- /dev/null +++ b/packages/integrations/pi-sdk/tests/eventRetention.test.ts @@ -0,0 +1,82 @@ +import { afterEach, describe, expect, it, vi } from "vitest"; +import { compactPiEvent, runPiSession, summarizePiEvent, type PiEvent } from "../src/session.js"; + +afterEach(() => vi.restoreAllMocks()); + +describe("Pi retained event limits", () => { + it.each(["assistant", "tool", "unknown"])("redacts before clipping %s details", (kind) => { + const secret = `AIza${"A".repeat(35)}`; + const makeEvent = (text: string): PiEvent => + kind === "assistant" + ? { type: "message_end", message: { role: "assistant", content: [{ type: "text", text }] } } + : kind === "tool" + ? { type: "tool_execution_end", result: text } + : { type: "custom", detail: text }; + const raw = kind === "assistant" ? secret : JSON.stringify(makeEvent(secret)); + const padding = "x".repeat(20_000 - 16 - raw.indexOf(secret)) + " "; + const summary = summarizePiEvent(makeEvent(padding + secret)); + expect(summary.detail?.length).toBeLessThanOrEqual(20_000); + expect(summary.detail?.includes("AIzaAAAA")).toBe(false); + expect(summary.detail).toContain("AIza[redacted]"); + }); + + it("rejects a large image before allocating a decoded buffer", () => { + const data = "A".repeat(12 * 1024 * 1024); + const from = vi.spyOn(Buffer, "from"); + const event = imageEvent(data); + const retained = compactPiEvent(event); + expect(from.mock.calls.some(([value]) => value === data)).toBe(false); + expect(retained.result).toMatchObject({ + content: [{ type: "text", text: expect.stringContaining("Screenshot omitted") }], + }); + expect(JSON.stringify(retained)).not.toContain(data); + expect((event.result as { content: Array<{ data: string }> }).content[0].data).toBe(data); + }); + + it("shares the 64 MiB retained-image budget across the entire run", async () => { + const data = Buffer.alloc(1024 * 1024).toString("base64"); + let listener: (event: PiEvent) => void = () => {}; + const result = await runPiSession({ + prompt: "Fixture task", + model: "fixture/model", + logger: { log: () => {}, warn: () => {}, error: () => {} }, + session: {}, + sdk: { + async createSession() { + return { + agent: { state: {} }, + subscribe(fn) { + listener = fn; + return () => {}; + }, + async prompt() { + for (let i = 0; i < 65; i++) listener(imageEvent(data)); + }, + async abort() {}, + dispose() {}, + }; + }, + }, + }); + const blocks = result.events.map( + (event) => (event.result as { content: Array> }).content[0], + ); + expect(blocks.filter((block) => Buffer.isBuffer(block.bytes))).toHaveLength(64); + expect(blocks[64]).toMatchObject({ + type: "text", + text: expect.stringContaining("Screenshot omitted"), + }); + expect(blocks[64]).not.toHaveProperty("data"); + expect(blocks[64]).not.toHaveProperty("bytes"); + expect(result.status).toBe("completed"); + }); +}); + +function imageEvent(data: string): PiEvent { + return { + type: "tool_execution_end", + toolCallId: "fixture", + toolName: "screenshot", + result: { content: [{ type: "image", data, mimeType: "image/png" }] }, + }; +} diff --git a/packages/integrations/pi-sdk/tests/session.test.ts b/packages/integrations/pi-sdk/tests/session.test.ts index 03f41ca91..e6b06657e 100644 --- a/packages/integrations/pi-sdk/tests/session.test.ts +++ b/packages/integrations/pi-sdk/tests/session.test.ts @@ -288,27 +288,75 @@ describe("pi SDK session", () => { expect(result.stopReason).toBe("cancelled"); }); - it("forwards an active abort and removes the signal listener", async () => { - const fake = scriptedSdk([], { promptUntilAbort: true }); - const controller = new AbortController(); - const removeListener = vi.spyOn(controller.signal, "removeEventListener"); - const execution = runPiSession({ + it("retains screenshots once as bytes and keeps base64 out of logs", async () => { + const png = Buffer.alloc(300_000, 7); + const base64 = png.toString("base64"); + const toolResult = { + content: [ + { type: "text", text: "Screenshot captured." }, + { type: "image", data: base64, mimeType: "image/png" }, + ], + details: {}, + }; + const log = vi.fn(); + const fake = scriptedSdk([ + { + type: "message_end", + message: { + role: "assistant", + content: [{ type: "toolCall", id: "1", name: "mcp__stagehand__screenshot" }], + usage: { input: 1, output: 1 }, + stopReason: "toolUse", + }, + }, + { type: "tool_execution_end", toolCallId: "1", toolName: "shot", result: toolResult }, + { + type: "message_end", + message: { + role: "toolResult", + toolCallId: "1", + content: [ + { type: "text", text: "Screenshot captured." }, + { type: "image", data: base64, mimeType: "image/png" }, + ], + }, + }, + { type: "turn_end" }, + assistant("done", { input: 1, output: 1 }), + { type: "turn_end" }, + ]); + const result = await runPiSession({ prompt: "task", - model: "model", + model: "openai/gpt-5.4-mini", sdk: fake.sdk, - signal: controller.signal, - logger, + logger: { ...logger, log }, session: {}, }); - await vi.waitFor(() => expect(fake.createOptions).toBeDefined()); - controller.abort(new Error("active cancellation")); - const result = await execution; + const toolEnd = result.events.find((event) => event.type === "tool_execution_end"); + const image = (toolEnd?.result as { content: Array> }).content[1]; + expect(Buffer.isBuffer(image.bytes)).toBe(true); + expect((image.bytes as Buffer).equals(png)).toBe(true); + expect(image.data).toBeUndefined(); + expect(image.mimeType).toBe("image/png"); + // The original event object handed to pi is left intact. + expect(toolResult.content[1]).toMatchObject({ data: base64 }); - expect(fake.abortCount).toBe(1); - expect(result.status).toBe("sdk_error"); - expect(result.stopReason).toBe("active cancellation"); - expect(removeListener).toHaveBeenCalledWith("abort", expect.any(Function)); + const toolMessage = result.events.find( + (event) => + event.type === "message_end" && (event.message as { role: string }).role === "toolResult", + ); + const serialized = JSON.stringify(result.events, (_key, value) => + Buffer.isBuffer(value) ? "" : value, + ); + expect(serialized).not.toContain(base64.slice(0, 1_000)); + expect(JSON.stringify(toolMessage)).toContain("[image 300000 bytes]"); + + const logged = JSON.stringify(log.mock.calls); + expect(logged).not.toContain(base64.slice(0, 1_000)); + expect(logged).toContain("[image 300000 bytes]"); + expect(buildPiTranscript(result.events)).not.toContain(base64.slice(0, 1_000)); + expect(result.status).toBe("completed"); }); it("normalizes models, MCP names/results, code tools, and statuses", async () => { @@ -356,4 +404,27 @@ describe("pi SDK session", () => { "sdk_error", ); }); + + it("forwards an active abort and removes the signal listener", async () => { + const fake = scriptedSdk([], { promptUntilAbort: true }); + const controller = new AbortController(); + const removeListener = vi.spyOn(controller.signal, "removeEventListener"); + const execution = runPiSession({ + prompt: "task", + model: "model", + sdk: fake.sdk, + signal: controller.signal, + logger, + session: {}, + }); + await vi.waitFor(() => expect(fake.createOptions).toBeDefined()); + + controller.abort(new Error("active cancellation")); + const result = await execution; + + expect(fake.abortCount).toBe(1); + expect(result.status).toBe("sdk_error"); + expect(result.stopReason).toBe("active cancellation"); + expect(removeListener).toHaveBeenCalledWith("abort", expect.any(Function)); + }); });