Skip to content

Commit 6fcc945

Browse files
authored
fix(agents): trim dense text delta snapshots
Trim dense plain text-delta stream snapshots for OpenAI-compatible, Responses, and Ollama providers while preserving full snapshots on stream checkpoints and terminal events. Reconstruct partial-less text deltas in the agent loop so live message updates continue to advance for immutable snapshot providers, and document the optional text_delta.partial contract. Fixes #86599.
1 parent b3c9469 commit 6fcc945

7 files changed

Lines changed: 208 additions & 19 deletions

File tree

extensions/ollama/src/stream-runtime.test.ts

Lines changed: 6 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -1630,9 +1630,9 @@ describe("createOllamaStreamFn streaming events", () => {
16301630
const textStartEvent = events.find((e) => e.type === "text_start");
16311631
expect(textStartEvent?.partial.content).toStrictEqual([]);
16321632

1633-
// text_delta partials accumulate content progressively
1634-
expect(deltas[0].partial.content).toEqual([{ type: "text", text: "Hello" }]);
1635-
expect(deltas[1].partial.content).toEqual([{ type: "text", text: "Hello world" }]);
1633+
// text_delta events stay lightweight; text_end/done carry the full snapshot.
1634+
expect(deltas[0]).not.toHaveProperty("partial");
1635+
expect(deltas[1]).not.toHaveProperty("partial");
16361636

16371637
// done event contains the final message
16381638
const doneEvent = events.at(-1);
@@ -2106,7 +2106,7 @@ describe("createOllamaStreamFn streaming events", () => {
21062106
);
21072107
});
21082108

2109-
it("sanitizes Kimi inline reasoning in text_delta, text_end, partial, and done output", async () => {
2109+
it("sanitizes Kimi inline reasoning in text_delta, text_end, and done output", async () => {
21102110
await withMockNdjsonFetch(
21112111
[
21122112
JSON.stringify({
@@ -2149,8 +2149,8 @@ describe("createOllamaStreamFn streaming events", () => {
21492149
expect(deltas).toHaveLength(2);
21502150
expect(deltas[0]?.delta).toBe("Final answer");
21512151
expect(deltas[1]?.delta).toBe(" only.");
2152-
expect(deltas[0]?.partial.content).toEqual([{ type: "text", text: "Final answer" }]);
2153-
expect(deltas[1]?.partial.content).toEqual([{ type: "text", text: "Final answer only." }]);
2152+
expect(deltas[0]).not.toHaveProperty("partial");
2153+
expect(deltas[1]).not.toHaveProperty("partial");
21542154

21552155
const textEnd = events.find((e) => e.type === "text_end");
21562156
expect(textEnd?.content).toBe("Final answer only.");

extensions/ollama/src/stream.ts

Lines changed: 0 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -1322,17 +1322,10 @@ function createRawOllamaStreamFn(
13221322
}
13231323

13241324
accumulatedVisibleContent = nextVisibleContent;
1325-
const partial = buildStreamAssistantMessage({
1326-
model: modelInfo,
1327-
content: buildCurrentContent(),
1328-
stopReason: "stop",
1329-
usage: buildUsageWithNoCost({}),
1330-
});
13311325
stream.push({
13321326
type: "text_delta",
13331327
contentIndex: textContentIndex(),
13341328
delta,
1335-
partial,
13361329
});
13371330
};
13381331

packages/agent-core/src/agent-loop.test.ts

Lines changed: 70 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,8 @@
11
// Agent Core tests cover agent loop behavior.
22
import { describe, expect, it } from "vitest";
33
import { agentLoop, agentLoopContinue } from "./agent-loop.js";
4-
import type { Message, Model } from "./llm.js";
4+
import { createAssistantMessageEventStream } from "./llm.js";
5+
import type { AssistantMessage, Message, Model } from "./llm.js";
56
import type { AgentContext, AgentEvent, AgentLoopConfig, AgentMessage, StreamFn } from "./types.js";
67

78
const model: Model = {
@@ -73,3 +74,71 @@ describe("agentLoop EventStream failures", () => {
7374
expectTerminalFailure(events, result);
7475
});
7576
});
77+
78+
describe("agentLoop streaming updates", () => {
79+
it("rebuilds assistant message snapshots for text deltas without partial snapshots", async () => {
80+
const streamFn: StreamFn = async () => {
81+
const stream = createAssistantMessageEventStream();
82+
const startMessage: AssistantMessage = {
83+
role: "assistant",
84+
content: [],
85+
api: model.api,
86+
provider: model.provider,
87+
model: model.id,
88+
usage: {
89+
input: 0,
90+
output: 0,
91+
cacheRead: 0,
92+
cacheWrite: 0,
93+
totalTokens: 0,
94+
cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 },
95+
},
96+
stopReason: "stop",
97+
timestamp: 1,
98+
};
99+
const textStartMessage: AssistantMessage = { ...startMessage, content: [] };
100+
const finalMessage: AssistantMessage = {
101+
...startMessage,
102+
content: [{ type: "text", text: "Hello world" }],
103+
};
104+
105+
queueMicrotask(() => {
106+
stream.push({ type: "start", partial: startMessage });
107+
stream.push({ type: "text_start", contentIndex: 0, partial: textStartMessage });
108+
stream.push({ type: "text_delta", contentIndex: 0, delta: "Hello" });
109+
stream.push({ type: "text_delta", contentIndex: 0, delta: " world" });
110+
stream.push({
111+
type: "text_end",
112+
contentIndex: 0,
113+
content: "Hello world",
114+
partial: finalMessage,
115+
});
116+
stream.push({ type: "done", reason: "stop", message: finalMessage });
117+
});
118+
119+
return stream;
120+
};
121+
122+
const stream = agentLoop(
123+
[{ role: "user", content: "hello", timestamp: 1 }],
124+
{ systemPrompt: "", messages: [] },
125+
config,
126+
undefined,
127+
streamFn,
128+
);
129+
const events = await collectEvents(stream);
130+
131+
const deltaUpdates = events.filter(
132+
(event): event is Extract<AgentEvent, { type: "message_update" }> =>
133+
event.type === "message_update" && event.assistantMessageEvent.type === "text_delta",
134+
);
135+
expect(deltaUpdates).toHaveLength(2);
136+
expect(deltaUpdates.map((event) => event.message)).toMatchObject([
137+
{ role: "assistant", content: [{ type: "text", text: "Hello" }] },
138+
{ role: "assistant", content: [{ type: "text", text: "Hello world" }] },
139+
]);
140+
for (const update of deltaUpdates) {
141+
expect(update.assistantMessageEvent).not.toHaveProperty("partial");
142+
}
143+
});
144+
});

packages/agent-core/src/agent-loop.ts

Lines changed: 45 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,7 @@
33
import { EventStream as LlmEventStream } from "@openclaw/llm-core";
44
import type {
55
AssistantMessage,
6+
AssistantMessageEvent,
67
Context,
78
EventStream,
89
ToolResultMessage,
@@ -35,6 +36,49 @@ const EMPTY_USAGE = {
3536

3637
const EventStreamConstructor: typeof SourceEventStream = LlmEventStream;
3738

39+
type AssistantMessageUpdateEvent = Extract<
40+
AssistantMessageEvent,
41+
{
42+
type:
43+
| "text_start"
44+
| "text_delta"
45+
| "text_end"
46+
| "thinking_start"
47+
| "thinking_delta"
48+
| "thinking_end"
49+
| "toolcall_start"
50+
| "toolcall_delta"
51+
| "toolcall_end";
52+
}
53+
>;
54+
55+
function appendTextDeltaToAssistantMessage(
56+
message: AssistantMessage,
57+
contentIndex: number,
58+
delta: string,
59+
): AssistantMessage {
60+
const content = [...message.content];
61+
const currentContent = content[contentIndex];
62+
content[contentIndex] =
63+
currentContent?.type === "text"
64+
? { ...currentContent, text: currentContent.text + delta }
65+
: { type: "text", text: delta };
66+
return { ...message, content };
67+
}
68+
69+
function resolveAssistantMessageUpdate(
70+
event: AssistantMessageUpdateEvent,
71+
currentMessage: AssistantMessage,
72+
): AssistantMessage {
73+
if ("partial" in event && event.partial) {
74+
return event.partial;
75+
}
76+
if (event.type === "text_delta") {
77+
return appendTextDeltaToAssistantMessage(currentMessage, event.contentIndex, event.delta);
78+
}
79+
return currentMessage;
80+
}
81+
3882
/**
3983
* Start an agent loop with a new prompt message.
4084
* The prompt is added to the context and events are emitted for it.
@@ -402,7 +446,7 @@ async function streamAssistantResponse(
402446
case "toolcall_delta":
403447
case "toolcall_end":
404448
if (partialMessage) {
405-
const message = event.partial;
449+
const message = resolveAssistantMessageUpdate(event, partialMessage);
406450
partialMessage = message;
407451
context.messages[context.messages.length - 1] = message;
408452
await emit({

packages/llm-core/src/types.ts

Lines changed: 6 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -366,7 +366,12 @@ export interface Context {
366366
export type AssistantMessageEvent =
367367
| { type: "start"; partial: AssistantMessage }
368368
| { type: "text_start"; contentIndex: number; partial: AssistantMessage }
369-
| { type: "text_delta"; contentIndex: number; delta: string; partial: AssistantMessage }
369+
/**
370+
* Plain text deltas may omit `partial` to avoid retaining one full assistant
371+
* snapshot per token. Consumers that need current text should replay `delta`
372+
* from the latest start/end partial checkpoint.
373+
*/
374+
| { type: "text_delta"; contentIndex: number; delta: string; partial?: AssistantMessage }
370375
| { type: "text_end"; contentIndex: number; content: string; partial: AssistantMessage }
371376
| { type: "thinking_start"; contentIndex: number; partial: AssistantMessage }
372377
| { type: "thinking_delta"; contentIndex: number; delta: string; partial: AssistantMessage }

src/agents/openai-transport-stream.test.ts

Lines changed: 81 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -25,7 +25,7 @@ import { SYSTEM_PROMPT_CACHE_BOUNDARY } from "./system-prompt-cache-boundary.js"
2525
type OpenAICompletionsOutput = Parameters<typeof testing.processOpenAICompletionsStream>[1];
2626
type OpenAIResponsesOutput = Parameters<typeof testing.processResponsesStream>[1];
2727

28-
type CapturedStreamEvent = { type?: string; delta?: string };
28+
type CapturedStreamEvent = { type?: string; delta?: string; partial?: unknown };
2929

3030
function createDeepSeekCompletionsModel(): Model<"openai-completions"> {
3131
return {
@@ -1859,6 +1859,64 @@ describe("openai transport stream", () => {
18591859
expect(stream.push.mock.calls.length).toBeLessThan(512);
18601860
});
18611861

1862+
it("omits accumulated partial snapshots from OpenAI-compatible text deltas", async () => {
1863+
const model = {
1864+
id: "dense-local",
1865+
name: "Dense Local",
1866+
api: "openai-completions",
1867+
provider: "local",
1868+
baseUrl: "http://127.0.0.1:18065/v1",
1869+
reasoning: false,
1870+
input: ["text"],
1871+
cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0 },
1872+
contextWindow: 128000,
1873+
maxTokens: 4096,
1874+
} satisfies Model<"openai-completions">;
1875+
const output = createAssistantOutput(model);
1876+
const events: CapturedStreamEvent[] = [];
1877+
1878+
await testing.processOpenAICompletionsStream(
1879+
streamChunks([
1880+
{
1881+
id: "chatcmpl-dense",
1882+
object: "chat.completion.chunk" as const,
1883+
created: 1775425651,
1884+
model: model.id,
1885+
choices: [
1886+
{
1887+
index: 0,
1888+
delta: { role: "assistant" as const, content: "a" },
1889+
logprobs: null,
1890+
finish_reason: null,
1891+
},
1892+
],
1893+
},
1894+
{
1895+
id: "chatcmpl-dense",
1896+
object: "chat.completion.chunk" as const,
1897+
created: 1775425651,
1898+
model: model.id,
1899+
choices: [
1900+
{
1901+
index: 0,
1902+
delta: { content: "b" },
1903+
logprobs: null,
1904+
finish_reason: null,
1905+
},
1906+
],
1907+
},
1908+
]),
1909+
output,
1910+
model,
1911+
{ push: (event) => events.push(event as CapturedStreamEvent) },
1912+
);
1913+
1914+
const textDeltas = events.filter((event) => event.type === "text_delta");
1915+
expect(textDeltas).toHaveLength(2);
1916+
expect(textDeltas.every((event) => !("partial" in event))).toBe(true);
1917+
expect(output.content).toEqual([{ type: "text", text: "ab" }]);
1918+
});
1919+
18621920
it("yields to aborts during bursty Responses streams", async () => {
18631921
const model = createAzureResponsesModel();
18641922
const output = createResponsesAssistantOutput(model);
@@ -1887,6 +1945,28 @@ describe("openai transport stream", () => {
18871945
expect(stream.push.mock.calls.length).toBeLessThan(512);
18881946
});
18891947

1948+
it("omits accumulated partial snapshots from Responses text deltas", async () => {
1949+
const model = createAzureResponsesModel();
1950+
const output = createResponsesAssistantOutput(model);
1951+
const events: CapturedStreamEvent[] = [];
1952+
1953+
await testing.processResponsesStream(
1954+
streamChunks([
1955+
{ type: "response.output_item.added", item: { type: "message" } },
1956+
{ type: "response.output_text.delta", delta: "a" },
1957+
{ type: "response.output_text.delta", delta: "b" },
1958+
]),
1959+
output,
1960+
{ push: (event) => events.push(event as CapturedStreamEvent) },
1961+
model,
1962+
);
1963+
1964+
const textDeltas = events.filter((event) => event.type === "text_delta");
1965+
expect(textDeltas).toHaveLength(2);
1966+
expect(textDeltas.every((event) => !("partial" in event))).toBe(true);
1967+
expect(output.content).toEqual([{ type: "text", text: "ab" }]);
1968+
});
1969+
18901970
it("skips null and non-object OpenAI-compatible stream chunks", async () => {
18911971
const model = {
18921972
id: "glm-5",

src/agents/openai-transport-stream.ts

Lines changed: 0 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1578,7 +1578,6 @@ async function processResponsesStream(
15781578
type: "text_delta",
15791579
contentIndex: blockIndex(),
15801580
delta: stringifyUnknown(event.delta),
1581-
partial: output,
15821581
});
15831582
}
15841583
} else if (type === "response.function_call_arguments.delta") {
@@ -2740,7 +2739,6 @@ async function processOpenAICompletionsStream(
27402739
type: "text_delta",
27412740
contentIndex: blockIndex(),
27422741
delta: text,
2743-
partial: output,
27442742
});
27452743
};
27462744
const flushPendingPostToolCallDeltas = () => {

0 commit comments

Comments
 (0)