diff --git a/packages/agent/CHANGELOG.md b/packages/agent/CHANGELOG.md index 2cb3f71dc..cc27d38e8 100644 --- a/packages/agent/CHANGELOG.md +++ b/packages/agent/CHANGELOG.md @@ -2,6 +2,10 @@ ## [Unreleased] +### Fixed + +- Fixed proxy-stream clients dropping finalized provider-only content blocks, including Anthropic native web-search history, by carrying terminal assistant content in `done` and `error` events ([#6703](https://github.com/can1357/oh-my-pi/issues/6703)). + ## [17.1.2] - 2026-07-24 ### Added diff --git a/packages/agent/src/proxy.ts b/packages/agent/src/proxy.ts index f96d0b766..b02ca54e3 100644 --- a/packages/agent/src/proxy.ts +++ b/packages/agent/src/proxy.ts @@ -56,12 +56,14 @@ export type ProxyAssistantMessageEvent = type: "done"; reason: Extract; usage: AssistantMessage["usage"]; + content: AssistantMessage["content"]; } | { type: "error"; reason: Extract; errorMessage?: string; usage: AssistantMessage["usage"]; + content: AssistantMessage["content"]; }; export interface ProxyStreamOptions extends SimpleStreamOptions { @@ -372,6 +374,7 @@ function processProxyEvent( case "done": partial.stopReason = proxyEvent.reason; partial.usage = proxyEvent.usage; + partial.content = proxyEvent.content; calculateCost(model, partial.usage); scrubPartialJson(partial); return { type: "done", reason: proxyEvent.reason, message: partial }; @@ -380,6 +383,7 @@ function processProxyEvent( partial.stopReason = proxyEvent.reason; partial.errorMessage = proxyEvent.errorMessage; partial.usage = proxyEvent.usage; + partial.content = proxyEvent.content; calculateCost(model, partial.usage); scrubPartialJson(partial); return { type: "error", reason: proxyEvent.reason, error: partial }; diff --git a/packages/agent/test/proxy-stream-disconnect.test.ts b/packages/agent/test/proxy-stream-disconnect.test.ts index 9b8861143..0d8d92d29 100644 --- a/packages/agent/test/proxy-stream-disconnect.test.ts +++ b/packages/agent/test/proxy-stream-disconnect.test.ts @@ -9,7 +9,7 @@ import { describe, expect, it } from "bun:test"; import type { ProxyAssistantMessageEvent } from "@oh-my-pi/pi-agent-core/proxy"; import { type ProxyMessageEventStream, streamProxy } from "@oh-my-pi/pi-agent-core/proxy"; -import type { AssistantMessageEvent, Context, FetchImpl, Model, ToolCall } from "@oh-my-pi/pi-ai"; +import type { AssistantMessage, AssistantMessageEvent, Context, FetchImpl, Model, ToolCall } from "@oh-my-pi/pi-ai"; import { getStreamingPartialJson } from "@oh-my-pi/pi-ai/utils/block-symbols"; import { buildModel } from "@oh-my-pi/pi-catalog/build"; @@ -178,6 +178,7 @@ describe("streamProxy — server disconnect without terminal event", () => { type: "done", reason: "stop", usage: { ...baseUsage }, + content: [{ type: "text", text: "Hello" }], }, ]; const body = buildSseBody(events); @@ -197,6 +198,54 @@ describe("streamProxy — server disconnect without terminal event", () => { expect(result.content.length).toBeGreaterThan(0); }); + it("restores terminal blocks that have no proxy stream events", async () => { + const finalizedContent: AssistantMessage["content"] = [ + { type: "thinking", thinking: "Search first.", thinkingSignature: "sig-1" }, + { + type: "anthropicServerTool", + block: { + type: "server_tool_use", + id: "srvtoolu_1", + name: "web_search", + input: { query: "current UTC date" }, + }, + }, + { + type: "anthropicServerTool", + block: { + type: "web_search_tool_result", + tool_use_id: "srvtoolu_1", + content: [{ type: "web_search_result", encrypted_content: "opaque-result" }], + }, + }, + { type: "thinking", thinking: "Use the result.", thinkingSignature: "sig-2" }, + { type: "toolCall", id: "toolu_write", name: "write", arguments: { path: "date.txt" } }, + ]; + const events: ProxyAssistantMessageEvent[] = [ + { type: "start" }, + { type: "thinking_start", contentIndex: 0 }, + { type: "thinking_delta", contentIndex: 0, delta: "Search first." }, + { type: "thinking_end", contentIndex: 0, contentSignature: "sig-1" }, + { type: "thinking_start", contentIndex: 1 }, + { type: "thinking_delta", contentIndex: 1, delta: "Use the result." }, + { type: "thinking_end", contentIndex: 1, contentSignature: "sig-2" }, + { type: "toolcall_start", contentIndex: 2, id: "toolu_write", toolName: "write" }, + { type: "toolcall_delta", contentIndex: 2, delta: '{"path":"date.txt"}' }, + { type: "toolcall_end", contentIndex: 2 }, + { type: "done", reason: "toolUse", usage: { ...baseUsage }, content: finalizedContent }, + ]; + const body = buildSseBody(events); + const fetchMock: FetchImpl = () => Promise.resolve(new Response(body, { status: 200 })); + + const result = await streamProxy(mockModel, mockContext, { + proxyUrl: "http://localhost:0", + authToken: "test", + fetch: fetchMock, + }).result(); + + expect(result.content).toEqual(finalizedContent); + }); + it("completes with error event when server sends an 'error' terminal event", async () => { const events: ProxyAssistantMessageEvent[] = [ { type: "start" }, @@ -207,6 +256,7 @@ describe("streamProxy — server disconnect without terminal event", () => { reason: "error", errorMessage: "rate_limit_exceeded", usage: { ...baseUsage }, + content: [{ type: "text", text: "Hel" }], }, ]; const body = buildSseBody(events); diff --git a/packages/agent/test/proxy-toolcall-partial-json.test.ts b/packages/agent/test/proxy-toolcall-partial-json.test.ts index 84bf0c408..0d1892e0c 100644 --- a/packages/agent/test/proxy-toolcall-partial-json.test.ts +++ b/packages/agent/test/proxy-toolcall-partial-json.test.ts @@ -87,7 +87,12 @@ describe("streamProxy — tool-call streaming and partialJson isolation", () => { type: "toolcall_delta", contentIndex: 0, delta: '{"comm' }, { type: "toolcall_delta", contentIndex: 0, delta: 'and":"ls"}' }, { type: "toolcall_end", contentIndex: 0 }, - { type: "done", reason: "toolUse", usage: { ...baseUsage } }, + { + type: "done", + reason: "toolUse", + usage: { ...baseUsage }, + content: [{ type: "toolCall", id: "call_1", name: "bash", arguments: { command: "ls" } }], + }, ]; const body = buildSseBody(events); const fetchMock: FetchImpl = () => Promise.resolve(new Response(body, { status: 200 })); @@ -119,7 +124,12 @@ describe("streamProxy — tool-call streaming and partialJson isolation", () => { type: "toolcall_delta", contentIndex: 0, delta: '{"comm' }, { type: "toolcall_delta", contentIndex: 0, delta: 'and":"ls"}' }, { type: "toolcall_end", contentIndex: 0 }, - { type: "done", reason: "toolUse", usage: { ...baseUsage } }, + { + type: "done", + reason: "toolUse", + usage: { ...baseUsage }, + content: [{ type: "toolCall", id: "call_1", name: "bash", arguments: { command: "ls" } }], + }, ]; const body = buildSseBody(events); const fetchMock: FetchImpl = () => Promise.resolve(new Response(body, { status: 200 })); @@ -173,7 +183,12 @@ describe("streamProxy — tool-call streaming and partialJson isolation", () => { type: "toolcall_delta", contentIndex: 0, delta: '{"path' }, { type: "toolcall_delta", contentIndex: 0, delta: '":"/tmp/x"}' }, { type: "toolcall_end", contentIndex: 0 }, - { type: "done", reason: "toolUse", usage: { ...baseUsage } }, + { + type: "done", + reason: "toolUse", + usage: { ...baseUsage }, + content: [{ type: "toolCall", id: "call_1", name: "read", arguments: { path: "/tmp/x" } }], + }, ]; const body = buildSseBody(events); const fetchMock: FetchImpl = () => Promise.resolve(new Response(body, { status: 200 })); @@ -202,7 +217,12 @@ describe("streamProxy — tool-call streaming and partialJson isolation", () => { type: "toolcall_delta", contentIndex: 0, delta: '{"path' }, { type: "toolcall_delta", contentIndex: 0, delta: '":"/a"}' }, // Missing toolcall_end — stream goes straight to done - { type: "done", reason: "toolUse", usage: { ...baseUsage } }, + { + type: "done", + reason: "toolUse", + usage: { ...baseUsage }, + content: [{ type: "toolCall", id: "call_1", name: "edit", arguments: { path: "/a" } }], + }, ]; const body = buildSseBody(events); const fetchMock: FetchImpl = () => Promise.resolve(new Response(body, { status: 200 })); @@ -231,7 +251,15 @@ describe("streamProxy — tool-call streaming and partialJson isolation", () => { type: "toolcall_delta", contentIndex: 1, delta: 'ls"}' }, { type: "toolcall_end", contentIndex: 0 }, { type: "toolcall_end", contentIndex: 1 }, - { type: "done", reason: "toolUse", usage: { ...baseUsage } }, + { + type: "done", + reason: "toolUse", + usage: { ...baseUsage }, + content: [ + { type: "toolCall", id: "call_1", name: "read", arguments: { path: "a" } }, + { type: "toolCall", id: "call_2", name: "bash", arguments: { command: "ls" } }, + ], + }, ]; const body = buildSseBody(events); const fetchMock: FetchImpl = () => Promise.resolve(new Response(body, { status: 200 }));