fix(agent): carried terminal content through proxy streams
Made proxy done and error events carry finalized assistant content so provider-only blocks survive event reconstruction. Covered Anthropic native web-search history restored after live events omitted the opaque blocks. Fixes #6703
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -56,12 +56,14 @@ export type ProxyAssistantMessageEvent =
|
||||
type: "done";
|
||||
reason: Extract<StopReason, "stop" | "length" | "toolUse">;
|
||||
usage: AssistantMessage["usage"];
|
||||
content: AssistantMessage["content"];
|
||||
}
|
||||
| {
|
||||
type: "error";
|
||||
reason: Extract<StopReason, "aborted" | "error">;
|
||||
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 };
|
||||
|
||||
@@ -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);
|
||||
|
||||
@@ -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 }));
|
||||
|
||||
Reference in New Issue
Block a user