From 74310256f83bf33f99d90f9f308770a029da6443 Mon Sep 17 00:00:00 2001 From: "vercel-gh-bot-3[bot]" <282332853+vercel-gh-bot-3[bot]@users.noreply.github.com> Date: Fri, 21 Aug 2026 19:54:39 +0000 Subject: [PATCH 1/5] fix(eve): preserve streamed tool inputs Co-authored-by: ruiconti <1834568+ruiconti@users.noreply.github.com> Signed-off-by: Rui Conti --- .changeset/stream-tool-input.md | 5 + docs/concepts/sessions-runs-and-streaming.md | 7 +- docs/guides/client/streaming.mdx | 33 ++--- .../src/client/message-reducer-primitives.ts | 71 +++++++++++ .../eve/src/client/message-reducer-types.ts | 1 + .../eve/src/client/message-reducer.test.ts | 109 ++++++++++++++++ packages/eve/src/client/message-reducer.ts | 117 +++++++++--------- packages/eve/src/harness/emission.test.ts | 107 ++++++++++++++++ packages/eve/src/harness/emission.ts | 60 +++++++-- .../harness/ordered-stream-emitter.test.ts | 38 ++++++ .../eve/src/harness/ordered-stream-emitter.ts | 27 +++- packages/eve/src/protocol/message.test.ts | 2 +- packages/eve/src/protocol/message.ts | 44 ++++++- packages/eve/src/public/definitions/hook.ts | 1 + 14 files changed, 532 insertions(+), 90 deletions(-) create mode 100644 .changeset/stream-tool-input.md create mode 100644 packages/eve/src/client/message-reducer-primitives.ts diff --git a/.changeset/stream-tool-input.md b/.changeset/stream-tool-input.md new file mode 100644 index 0000000000..7969307afe --- /dev/null +++ b/.changeset/stream-tool-input.md @@ -0,0 +1,5 @@ +--- +"eve": patch +--- + +Tool inputs now stream through the durable event protocol as `action.input.appended` before the matching validated `actions.requested` event. The default message reducer exposes the cumulative raw input on `dynamic-tool.inputText` while its state is `input-streaming`, so UIs can progressively render long JSON inputs. This advances the stream protocol to version 24; when assistant text precedes a tool call, `message.completed` now arrives before that call's streamed input events. diff --git a/docs/concepts/sessions-runs-and-streaming.md b/docs/concepts/sessions-runs-and-streaming.md index c35c6cd824..1fd74a8cff 100644 --- a/docs/concepts/sessions-runs-and-streaming.md +++ b/docs/concepts/sessions-runs-and-streaming.md @@ -60,6 +60,7 @@ The stream is newline-delimited JSON (NDJSON), one event per line: | `reasoning.completed` | The finalized reasoning block. | | `message.appended` | An assistant text delta (incremental, with cumulative text so far). | | `message.completed` | A finalized assistant text block. | +| `action.input.appended` | A tool-input text delta, with the cumulative raw input and tool-call identity. | | `result.completed` | The finalized structured result for a turn that requested an output schema; carries `result`. | | `compaction.requested` | Context-window compaction began; carries `modelId`, `sessionId`, `turnId`, `usageInputTokens`. | | `compaction.completed` | A compaction checkpoint was written to durable history. | @@ -76,7 +77,9 @@ The stream is newline-delimited JSON (NDJSON), one event per line: The optional `data.trace` on session and turn starts contains eve-owned W3C trace coordinates: `traceId`, `spanId`, and `traceFlags`. Use it to correlate stream consumers such as eval reporters with an observability backend. An uninstrumented target omits it. -`reasoning.appended` and `message.appended` stream incremental output as it arrives. When the durable stream writer is busy, eve may coalesce adjacent deltas of the same type; the text remains in source order, and any other event forms an ordering barrier. Each append carries both the new delta and the cumulative text for the current block. The finalized block shows up on `message.completed` and `reasoning.completed`, which is the compatibility path for clients that don't render incremental streaming. +`reasoning.appended`, `message.appended`, and `action.input.appended` stream incremental output as it arrives. When the durable stream writer is busy, eve may coalesce adjacent deltas of the same type; the text remains in source order, and any other event forms an ordering barrier. Each append carries both the new delta and the cumulative text for the current block. The finalized text and reasoning blocks show up on `message.completed` and `reasoning.completed`, which is the compatibility path for clients that don't render incremental streaming. + +`action.input.appended` arrives before the matching `actions.requested` event. It carries `callId`, `toolName`, `inputTextDelta`, and `inputTextSoFar`; the input text may be incomplete JSON. The default client reducer projects it as a `dynamic-tool` part with `state: "input-streaming"` and the cumulative text in `inputText`. `actions.requested` replaces that part with `state: "input-available"` and the validated `input`. Excluded internal actions never publish their input stream. `action.partial` carries one complete preliminary output snapshot from an authored async-generator tool. A later partial for the same `callId` replaces it, and `action.result` is the final snapshot. When the durable writer is busy, eve may keep only the newest adjacent partial for a call. Treat partials as last-write-wins: a durable step can retry and replay overlapping event runs. Provider-executed tool progress and MCP progress notifications are not projected as `action.partial` events. @@ -111,7 +114,7 @@ Alongside `type` and `data`, every event carries a `meta` envelope: `meta.id` is stable. eve mints it once, when the event is written to the durable stream, and stores it with the event. Reconnecting from a cursor, rewinding to `startIndex=0`, or replaying a finished session all return the same id for the same event. -`meta.at` has always been there; `meta.id` arrived in stream version 20. Events written by an earlier version are stored with the envelope but no id inside it, so rewinding into the part of a session that ran before you upgraded yields events whose `meta.id` is absent, even though the type says it is always a string. eve passes those events through rather than dropping them, and they cannot be deduplicated. The exposure ends when the sessions that predate your upgrade do. +`meta.at` has always been there; `meta.id` arrived in stream version 20, and `action.input.appended` arrived in version 24. Events written by an earlier version are stored with the envelope but no id inside it, so rewinding into the part of a session that ran before you upgraded yields events whose `meta.id` is absent, even though the type says it is always a string. eve passes those events through rather than dropping them, and they cannot be deduplicated. The exposure ends when the sessions that predate your upgrade do. That makes it the key for ingesting a stream into a database without duplicating rows when you re-read it: diff --git a/docs/guides/client/streaming.mdx b/docs/guides/client/streaming.mdx index 13ce0b7068..fb4839927a 100644 --- a/docs/guides/client/streaming.mdx +++ b/docs/guides/client/streaming.mdx @@ -75,7 +75,7 @@ for await (const event of response) { } ``` -`message.appended` and `reasoning.appended` are incremental delta events. eve may combine adjacent deltas of the same type while a durable stream write is in flight, but preserves their text and event ordering; any different event is a barrier. Their completed forms, `message.completed` and `reasoning.completed`, are the compatibility path for clients that don't render deltas. +`message.appended`, `reasoning.appended`, and `action.input.appended` are incremental delta events. eve may combine adjacent deltas of the same type while a durable stream write is in flight, but preserves their text and event ordering; any different event is a barrier. The completed text forms, `message.completed` and `reasoning.completed`, are the compatibility path for clients that don't render deltas. A streamed tool input is complete when the matching validated call arrives in `actions.requested`. ## Handle event types @@ -98,23 +98,26 @@ function handleEvent(event: MessageStreamEvent) { The most common UI events are: -| Event | Use | -| -------------------- | ------------------------------------------------------------------------------ | -| `message.received` | Confirm the user message landed; `data.parts` includes text and file metadata. | -| `reasoning.appended` | Render reasoning deltas when the model provides them. | -| `message.appended` | Render assistant text deltas. | -| `actions.requested` | Show tool calls as the model requests them, before execution. | -| `action.partial` | Update a generator tool's provisional output snapshot. | -| `action.result` | Show tool call results. | -| `input.requested` | Pause the UI for approval or a question answer. | -| `input.resolved` | Record the server-accepted outcome and response for each human-input request. | -| `result.completed` | Read structured output from an [output schema](./output-schema). | -| `session.waiting` | Enable the composer; the same fixed session handle accepts the next message. | -| `session.completed` | Mark the conversation terminal. | -| `session.failed` | Mark the conversation failed. | +| Event | Use | +| ----------------------- | ------------------------------------------------------------------------------ | +| `message.received` | Confirm the user message landed; `data.parts` includes text and file metadata. | +| `reasoning.appended` | Render reasoning deltas when the model provides them. | +| `message.appended` | Render assistant text deltas. | +| `action.input.appended` | Render cumulative raw tool input before validation completes. | +| `actions.requested` | Show tool calls as the model requests them, before execution. | +| `action.partial` | Update a generator tool's provisional output snapshot. | +| `action.result` | Show tool call results. | +| `input.requested` | Pause the UI for approval or a question answer. | +| `input.resolved` | Record the server-accepted outcome and response for each human-input request. | +| `result.completed` | Read structured output from an [output schema](./output-schema). | +| `session.waiting` | Enable the composer; the same fixed session handle accepts the next message. | +| `session.completed` | Mark the conversation terminal. | +| `session.failed` | Mark the conversation failed. | For the complete event table, see [Sessions, runs & streaming](../../concepts/sessions-runs-and-streaming). +The default message reducer turns `action.input.appended` into a `dynamic-tool` part with `state: "input-streaming"`. Its `inputText` field contains cumulative raw text that may be incomplete JSON. The matching `actions.requested` event upgrades the same `toolCallId` to `state: "input-available"` and puts the validated value in `input`. + When a submitted message includes attachments, `message.received.data.message` stays the flattened compatibility summary, while `message.received.data.parts` carries renderable text and file metadata. File parts never include raw bytes or internal sandbox paths; `url` appears only for diff --git a/packages/eve/src/client/message-reducer-primitives.ts b/packages/eve/src/client/message-reducer-primitives.ts new file mode 100644 index 0000000000..d221bb8f7a --- /dev/null +++ b/packages/eve/src/client/message-reducer-primitives.ts @@ -0,0 +1,71 @@ +import type { EveMessage, EveMessageData, EveMessagePart } from "#client/message-reducer-types.js"; +import type { MessageReceivedPart } from "#protocol/message.js"; + +export function projectReceivedParts( + parts: readonly MessageReceivedPart[] | undefined, + message: string, +): readonly EveMessagePart[] { + return ( + parts?.map((part) => + part.type === "text" + ? { state: "done", text: part.text, type: "text" } + : { + filename: part.filename, + mediaType: part.mediaType, + size: part.size, + type: "file", + url: part.url, + }, + ) ?? [{ state: "done", text: message, type: "text" }] + ); +} + +export function partKey(part: EveMessagePart): string { + switch (part.type) { + case "text": + return `text:${part.stepIndex ?? 0}`; + case "reasoning": + return `reasoning:${part.stepIndex ?? 0}`; + case "file": + return `file:${part.stepIndex ?? 0}:${part.filename ?? part.url ?? part.mediaType}`; + case "step-start": + return "step-start"; + case "authorization": + return `authorization:${part.turnId}:${part.stepIndex}:${part.name}`; + case "dynamic-tool": + return `dynamic-tool:${part.toolCallId}`; + } +} + +export function upsertMessage(data: EveMessageData, next: EveMessage): EveMessageData { + const index = data.messages.findIndex((message) => message.id === next.id); + if (index === -1) { + return { messages: [...data.messages, next] }; + } + + return { + messages: [...data.messages.slice(0, index), next, ...data.messages.slice(index + 1)], + }; +} + +export function removeStreamingToolPartsForTurn( + data: EveMessageData, + turnId: string, +): EveMessageData { + const index = data.messages.findIndex( + (message) => message.role === "assistant" && message.metadata?.turnId === turnId, + ); + const message = data.messages[index]; + if (message === undefined) return data; + + return upsertMessage(data, { + ...message, + parts: message.parts.filter( + (part) => part.type !== "dynamic-tool" || part.state !== "input-streaming", + ), + }); +} + +export function optimisticUserMessageId(submissionId: string): string { + return `optimistic:${submissionId}:user`; +} diff --git a/packages/eve/src/client/message-reducer-types.ts b/packages/eve/src/client/message-reducer-types.ts index 2fa6e34a8c..5bae08f205 100644 --- a/packages/eve/src/client/message-reducer-types.ts +++ b/packages/eve/src/client/message-reducer-types.ts @@ -140,6 +140,7 @@ export type EveDynamicToolPart = { readonly approval?: never; readonly errorText?: never; readonly input: unknown | undefined; + readonly inputText: string; readonly output?: never; readonly state: "input-streaming"; } diff --git a/packages/eve/src/client/message-reducer.test.ts b/packages/eve/src/client/message-reducer.test.ts index 3a8981cae6..02ff0d5136 100644 --- a/packages/eve/src/client/message-reducer.test.ts +++ b/packages/eve/src/client/message-reducer.test.ts @@ -3,6 +3,7 @@ import { describe, expect, it } from "vitest"; import { defaultMessageReducer } from "#client/message-reducer.js"; import { stampTestEvents } from "#internal/testing/events.js"; import { + createActionInputAppendedEvent, createActionPartialEvent, createActionResultEvent, createActionsRequestedEvent, @@ -17,6 +18,7 @@ import { createResultCompletedEvent, createStepStartedEvent, createTurnCancelledEvent, + createTurnFailedEvent, type UnstampedMessageStreamEvent, } from "#protocol/message.js"; @@ -33,6 +35,113 @@ function reduceServerEvents( } describe("defaultMessageReducer", () => { + it("projects streamed tool input and upgrades it to the validated request", () => { + const reducer = defaultMessageReducer(); + let data = reduceServerEvents(reducer, reducer.initial(), [ + createActionInputAppendedEvent({ + callId: "call_render", + inputTextDelta: "", + inputTextSoFar: "", + sequence: 1, + stepIndex: 0, + toolName: "render", + turnId: "turn_1", + }), + createActionInputAppendedEvent({ + callId: "call_render", + inputTextDelta: '{"title":"Hel', + inputTextSoFar: '{"title":"Hel', + sequence: 1, + stepIndex: 0, + toolName: "render", + turnId: "turn_1", + }), + ]); + + expect(data.messages[0]?.parts).toContainEqual({ + input: undefined, + inputText: '{"title":"Hel', + state: "input-streaming", + stepIndex: 0, + toolCallId: "call_render", + toolMetadata: { eve: { kind: "unknown", name: "render" } }, + toolName: "render", + type: "dynamic-tool", + }); + + data = reduceServerEvents(reducer, data, [ + createActionsRequestedEvent({ + actions: [ + { + callId: "call_render", + input: { title: "Hello" }, + kind: "tool-call", + toolName: "render", + }, + ], + sequence: 1, + stepIndex: 0, + turnId: "turn_1", + }), + ]); + + expect(data.messages[0]?.parts).toContainEqual({ + input: { title: "Hello" }, + state: "input-available", + stepIndex: 0, + toolCallId: "call_render", + toolMetadata: { eve: { inputRequest: undefined, kind: "tool-call", name: "render" } }, + toolName: "render", + type: "dynamic-tool", + }); + + const settled = data; + data = reduceServerEvents(reducer, data, [ + createActionInputAppendedEvent({ + callId: "call_render", + inputTextDelta: "late", + inputTextSoFar: "late", + sequence: 1, + stepIndex: 0, + toolName: "render", + turnId: "turn_1", + }), + ]); + expect(data).toBe(settled); + }); + + it("removes an unfinished streamed tool input when the turn is cancelled", () => { + const reducer = defaultMessageReducer(); + const data = reduceServerEvents(reducer, reducer.initial(), [ + createActionInputAppendedEvent({ + callId: "call_render", + inputTextDelta: "{", + inputTextSoFar: "{", + sequence: 1, + stepIndex: 0, + toolName: "render", + turnId: "turn_1", + }), + createTurnCancelledEvent({ sequence: 1, turnId: "turn_1" }), + ]); + + expect(data.messages[0]?.parts).toEqual([{ type: "step-start" }]); + }); + + it("does not create an assistant message when a turn fails before streaming", () => { + const reducer = defaultMessageReducer(); + const data = reduceServerEvents(reducer, reducer.initial(), [ + createTurnFailedEvent({ + code: "MODEL_FAILED", + message: "model failed", + sequence: 1, + turnId: "turn_1", + }), + ]); + + expect(data.messages).toEqual([]); + }); + it("replaces tool-generator snapshots and ignores a late partial after the terminal result", () => { const reducer = defaultMessageReducer(); let data = reduceServerEvents(reducer, reducer.initial(), [ diff --git a/packages/eve/src/client/message-reducer.ts b/packages/eve/src/client/message-reducer.ts index d3ab552547..39c9691d95 100644 --- a/packages/eve/src/client/message-reducer.ts +++ b/packages/eve/src/client/message-reducer.ts @@ -20,12 +20,15 @@ import { stringifyUnknown, toMessageInputRequest, } from "#client/message-action-parts.js"; +import { + optimisticUserMessageId, + partKey, + projectReceivedParts, + removeStreamingToolPartsForTurn, + upsertMessage, +} from "#client/message-reducer-primitives.js"; import type { InputResponse } from "#shared/input.js"; -import type { - AuthorizationCompletedStreamEvent, - InputResolution, - MessageReceivedPart, -} from "#protocol/message.js"; +import type { AuthorizationCompletedStreamEvent, InputResolution } from "#protocol/message.js"; export type { EveAuthorizationChallenge, @@ -137,6 +140,37 @@ function reduceMessageData(data: EveMessageData, event: EveAgentReducerEvent): E }), ); + case "action.input.appended": { + const existing = findToolPart(data, event.data.callId); + if (existing !== undefined && existing.state !== "input-streaming") { + return data; + } + + const nextPart: EveDynamicToolPart = { + input: undefined, + inputText: event.data.inputTextSoFar, + state: "input-streaming", + stepIndex: event.data.stepIndex, + toolCallId: event.data.callId, + toolMetadata: existing?.toolMetadata ?? { + eve: { + kind: "unknown", + name: event.data.toolName, + }, + }, + toolName: event.data.toolName, + type: "dynamic-tool", + }; + + if (existing !== undefined) { + return updateToolPart(data, event.data.callId, nextPart); + } + + return updateAssistantMessage(data, event.data.turnId, (message) => + upsertPart(ensureStepStartPart(message, event.data.stepIndex), nextPart), + ); + } + case "actions.requested": { let next = data; for (const action of event.data.actions) { @@ -341,7 +375,11 @@ function reduceMessageData(data: EveMessageData, event: EveAgentReducerEvent): E return updateAssistantMetadata(data, event.data.turnId, { result: event.data.result }); case "turn.completed": - return updateAssistantMetadata(data, event.data.turnId, { status: "complete" }); + return updateAssistantMessage(data, event.data.turnId, (message) => ({ + ...message, + metadata: { ...message.metadata, status: "complete" }, + parts: removeStreamingToolParts(message.parts), + })); case "turn.cancelled": // Finalize whatever the cancelled turn streamed: no message.completed @@ -349,14 +387,18 @@ function reduceMessageData(data: EveMessageData, event: EveAgentReducerEvent): E return updateAssistantMessage(data, event.data.turnId, (message) => ({ ...message, metadata: { ...message.metadata, status: "complete" }, - parts: message.parts.map((part) => - (part.type === "text" || part.type === "reasoning") && part.state === "streaming" - ? { ...part, state: "done" } - : part, + parts: removeStreamingToolParts( + message.parts.map((part) => + (part.type === "text" || part.type === "reasoning") && part.state === "streaming" + ? { ...part, state: "done" } + : part, + ), ), })); case "turn.failed": + return removeStreamingToolPartsForTurn(data, event.data.turnId); + case "session.failed": return data; @@ -365,6 +407,10 @@ function reduceMessageData(data: EveMessageData, event: EveAgentReducerEvent): E } } +function removeStreamingToolParts(parts: readonly EveMessagePart[]): readonly EveMessagePart[] { + return parts.filter((part) => part.type !== "dynamic-tool" || part.state !== "input-streaming"); +} + function respondToInputRequest(data: EveMessageData, response: InputResponse): EveMessageData { const existing = findToolPartByApprovalId(data, response.requestId); if (!existing) return data; @@ -646,54 +692,3 @@ function findToolPartByApprovalId( } return undefined; } - -function projectReceivedParts( - parts: readonly MessageReceivedPart[] | undefined, - message: string, -): readonly EveMessagePart[] { - return ( - parts?.map((part) => - part.type === "text" - ? { state: "done", text: part.text, type: "text" } - : { - filename: part.filename, - mediaType: part.mediaType, - size: part.size, - type: "file", - url: part.url, - }, - ) ?? [{ state: "done", text: message, type: "text" }] - ); -} - -function partKey(part: EveMessagePart): string { - switch (part.type) { - case "text": - return `text:${part.stepIndex ?? 0}`; - case "reasoning": - return `reasoning:${part.stepIndex ?? 0}`; - case "file": - return `file:${part.stepIndex ?? 0}:${part.filename ?? part.url ?? part.mediaType}`; - case "step-start": - return "step-start"; - case "authorization": - return `authorization:${part.turnId}:${part.stepIndex}:${part.name}`; - case "dynamic-tool": - return `dynamic-tool:${part.toolCallId}`; - } -} - -function upsertMessage(data: EveMessageData, next: EveMessage): EveMessageData { - const index = data.messages.findIndex((message) => message.id === next.id); - if (index === -1) { - return { messages: [...data.messages, next] }; - } - - return { - messages: [...data.messages.slice(0, index), next, ...data.messages.slice(index + 1)], - }; -} - -function optimisticUserMessageId(submissionId: string): string { - return `optimistic:${submissionId}:user`; -} diff --git a/packages/eve/src/harness/emission.test.ts b/packages/eve/src/harness/emission.test.ts index 7898cf1a84..79d1c95637 100644 --- a/packages/eve/src/harness/emission.test.ts +++ b/packages/eve/src/harness/emission.test.ts @@ -277,6 +277,108 @@ describe("emitStreamContent empty delivery", () => { }); describe("emitStreamContent action requests", () => { + it("streams a visible action input lifecycle before the completed request", async () => { + const emit = createEmitStub(); + const tools = new Map([ + [ + "render", + { + description: "Render a JSON document.", + inputSchema: jsonSchema({ type: "object" }), + name: "render", + }, + ], + ]); + + await emitStreamContent( + emit, + EMISSION_STATE, + streamOf([ + { id: "call-render", toolName: "render", type: "tool-input-start" }, + { delta: '{"title":"Hel', id: "call-render", type: "tool-input-delta" }, + { delta: 'lo"}', id: "call-render", type: "tool-input-delta" }, + { id: "call-render", type: "tool-input-end" }, + { + input: { title: "Hello" }, + toolCallId: "call-render", + toolName: "render", + type: "tool-call", + }, + { finishReason: "tool-calls", type: "finish-step" }, + ] as TextStreamPart[]), + { + excludedActionToolNames: new Set(), + tools, + }, + ); + + const events = vi.mocked(emit).mock.calls.map(([event]) => event); + expect(events.map((event) => event.type)).toEqual([ + "action.input.appended", + "action.input.appended", + "action.input.appended", + "actions.requested", + ]); + expect(events.slice(1, 3)).toMatchObject([ + { + data: { + callId: "call-render", + inputTextDelta: '{"title":"Hel', + inputTextSoFar: '{"title":"Hel', + }, + }, + { + data: { + callId: "call-render", + inputTextDelta: 'lo"}', + inputTextSoFar: '{"title":"Hello"}', + }, + }, + ]); + }); + + it("does not expose streamed input for excluded actions", async () => { + const emit = createEmitStub(); + + await emitStreamContent( + emit, + EMISSION_STATE, + streamOf([ + { id: "call-hidden", toolName: "hidden", type: "tool-input-start" }, + { delta: '{"secret":true}', id: "call-hidden", type: "tool-input-delta" }, + { id: "call-hidden", type: "tool-input-end" }, + ] as TextStreamPart[]), + { + excludedActionToolNames: new Set(["hidden"]), + tools: new Map(), + }, + ); + + expect(emit).not.toHaveBeenCalled(); + }); + + it("does not expose streamed input for provider-executed actions", async () => { + const emit = createEmitStub(); + + await emitStreamContent( + emit, + EMISSION_STATE, + streamOf([ + { + id: "call-provider", + providerExecuted: true, + toolName: "web_search", + type: "tool-input-start", + }, + { delta: '{"query":"eve"}', id: "call-provider", type: "tool-input-delta" }, + { id: "call-provider", type: "tool-input-end" }, + ] as TextStreamPart[]), + { excludedActionToolNames: new Set(), tools: new Map() }, + ); + + expect(emit).not.toHaveBeenCalled(); + }); + it("cancels a pending provider action batch when the stream aborts", async () => { vi.useFakeTimers(); const emit = createEmitStub(); @@ -382,6 +484,9 @@ describe("emitStreamContent action requests", () => { EMISSION_STATE, streamOf([ { id: "message-1", text: "Checking the release notes.", type: "text-delta" }, + { id: "call-delegate", toolName: "delegate", type: "tool-input-start" }, + { delta: '{"task":"research the release"}', id: "call-delegate", type: "tool-input-delta" }, + { id: "call-delegate", type: "tool-input-end" }, { input: { task: "research the release" }, toolCallId: "call-delegate", @@ -400,6 +505,8 @@ describe("emitStreamContent action requests", () => { expect(events.map((event) => event.type)).toEqual([ "message.appended", "message.completed", + "action.input.appended", + "action.input.appended", "actions.requested", ]); expect(events[1]).toMatchObject({ diff --git a/packages/eve/src/harness/emission.ts b/packages/eve/src/harness/emission.ts index 21337b9020..92b1bdc21e 100644 --- a/packages/eve/src/harness/emission.ts +++ b/packages/eve/src/harness/emission.ts @@ -17,6 +17,7 @@ import type { } from "#protocol/message.js"; import { createActionsRequestedEvent, + createActionInputAppendedEvent, createActionPartialEvent, createActionResultEvent, createMessageAppendedEvent, @@ -284,14 +285,7 @@ interface StreamActionEmissionOptions { readonly tools: HarnessToolMap; } -/** - * Consumes the AI SDK `fullStream` and emits real-time text and reasoning - * events. - * - * Emits local tool events in source order. Provider calls that arrive in one - * stream batch into one request event before their first result. A result - * without a streamed call resumes a call from an earlier step. - */ +/** Consumes `fullStream` in source order, batching provider calls before their first result. */ export async function emitStreamContent( emitFn: HarnessEmitFn, state: HarnessEmissionState, @@ -339,6 +333,7 @@ async function consumeStreamContent( const invalidInputToolCallIds = new Set(); const inlineAuthorizationResults: TypedToolResult[] = []; const trailingInlineToolResultParts: InlineToolResultPart[] = []; + const streamingActionInputs = new Map(); const flushCurrentMessage = async (): Promise => { if (currentMessage.length === 0) { @@ -356,6 +351,24 @@ async function consumeStreamContent( currentMessage = ""; }; + const emitActionInput = async ( + callId: string, + toolName: string, + inputTextDelta: string, + inputTextSoFar: string, + ): Promise => + emitFn( + createActionInputAppendedEvent({ + callId, + inputTextDelta, + inputTextSoFar, + sequence: state.sequence, + stepIndex: state.stepIndex, + toolName, + turnId: state.turnId, + }), + ); + const emitActionRequest = async (action: RuntimeActionRequest): Promise => { if (emittedActionCallIds.has(action.callId)) { return; @@ -510,8 +523,39 @@ async function consumeStreamContent( }), ); break; + case "tool-input-start": { + if ( + options === undefined || + part.providerExecuted === true || + options.excludedActionToolNames.has(part.toolName) + ) { + streamingActionInputs.delete(part.id); + break; + } + await providerActionBatch.flush(); + if (currentMessage.trim().length > 0) { + await flushCurrentMessage(); + } + streamingActionInputs.set(part.id, { text: "", toolName: part.toolName }); + await emitActionInput(part.id, part.toolName, "", ""); + break; + } + case "tool-input-delta": { + const input = streamingActionInputs.get(part.id); + if (input === undefined) { + break; + } + await providerActionBatch.flush(); + input.text += part.delta; + await emitActionInput(part.id, input.toolName, part.delta, input.text); + break; + } + case "tool-input-end": + streamingActionInputs.delete(part.id); + break; case "tool-call": { const toolCall = part as TypedToolCall; + streamingActionInputs.delete(toolCall.toolCallId); toolCallIdsSeenInStream.add(toolCall.toolCallId); if (toolCall.providerExecuted === true) { await collectProviderToolCall(toolCall); diff --git a/packages/eve/src/harness/ordered-stream-emitter.test.ts b/packages/eve/src/harness/ordered-stream-emitter.test.ts index 0351417fe0..75fcf1e9f6 100644 --- a/packages/eve/src/harness/ordered-stream-emitter.test.ts +++ b/packages/eve/src/harness/ordered-stream-emitter.test.ts @@ -2,6 +2,7 @@ import { describe, expect, it, vi } from "vitest"; import { createOrderedStreamEmitter } from "#harness/ordered-stream-emitter.js"; import { + createActionInputAppendedEvent, createActionPartialEvent, createActionResultEvent, createMessageAppendedEvent, @@ -47,6 +48,18 @@ function partial(callId: string, output: string) { }); } +function input(callId: string, delta: string, soFar: string) { + return createActionInputAppendedEvent({ + callId, + inputTextDelta: delta, + inputTextSoFar: soFar, + sequence: 1, + stepIndex: 0, + toolName: "render", + turnId: "turn_1", + }); +} + describe("createOrderedStreamEmitter", () => { it("keeps consuming while a write is active and preserves the latest event payload", async () => { const firstWrite = deferred(); @@ -68,6 +81,31 @@ describe("createOrderedStreamEmitter", () => { expect(events).toEqual([message("A", "A"), message("BC", "ABC")]); }); + it("coalesces adjacent input deltas for the same tool call", async () => { + const firstWrite = deferred(); + const events: UnstampedMessageStreamEvent[] = []; + const emitFn = vi.fn(async (event: UnstampedMessageStreamEvent) => { + events.push(event); + if (events.length === 1) await firstWrite.promise; + }); + const emitter = createOrderedStreamEmitter(emitFn); + + await emitter.emit(message("A", "A")); + await emitter.emit(input("call_1", "{", "{")); + await emitter.emit(input("call_1", '"title":', '{"title":')); + await emitter.emit(input("call_1", '"Hello"}', '{"title":"Hello"}')); + await emitter.emit(input("call_2", "{}", "{}")); + + firstWrite.resolve(); + await emitter.closeAndDrain(); + + expect(events).toEqual([ + message("A", "A"), + input("call_1", '{"title":"Hello"}', '{"title":"Hello"}'), + input("call_2", "{}", "{}"), + ]); + }); + it("treats other event types and stream coordinates as ordering barriers", async () => { const firstWrite = deferred(); const events: UnstampedMessageStreamEvent[] = []; diff --git a/packages/eve/src/harness/ordered-stream-emitter.ts b/packages/eve/src/harness/ordered-stream-emitter.ts index 28d3635b9d..86c81cee1f 100644 --- a/packages/eve/src/harness/ordered-stream-emitter.ts +++ b/packages/eve/src/harness/ordered-stream-emitter.ts @@ -1,4 +1,5 @@ import type { + ActionInputAppendedStreamEvent, UnstampedMessageStreamEvent, MessageAppendedStreamEvent, ReasoningAppendedStreamEvent, @@ -185,6 +186,17 @@ function mergeAdjacentEmissions( return true; } + if (left.event.type === "action.input.appended" && right.type === "action.input.appended") { + if (left.event.data.callId !== right.data.callId || !sameCoordinates(left.event, right)) { + return false; + } + left.deltaParts ??= [left.event.data.inputTextDelta]; + left.deltaParts.push(right.data.inputTextDelta); + left.event = right; + left.messages = messages; + return true; + } + if (left.event.type === "action.partial" && right.type === "action.partial") { if (left.event.data.result.callId !== right.data.result.callId) return false; left.event = right; @@ -198,6 +210,7 @@ function mergeAdjacentEmissions( function appendDelta(event: UnstampedMessageStreamEvent): string | undefined { if (event.type === "message.appended") return event.data.messageDelta; if (event.type === "reasoning.appended") return event.data.reasoningDelta; + if (event.type === "action.input.appended") return event.data.inputTextDelta; return undefined; } @@ -224,12 +237,22 @@ function materializeEvent(emission: PendingEmission): UnstampedMessageStreamEven }; } + if (emission.event.type === "action.input.appended") { + return { + ...emission.event, + data: { + ...emission.event.data, + inputTextDelta: emission.deltaParts.join(""), + }, + }; + } + return emission.event; } function sameCoordinates( - left: MessageAppendedStreamEvent | ReasoningAppendedStreamEvent, - right: MessageAppendedStreamEvent | ReasoningAppendedStreamEvent, + left: ActionInputAppendedStreamEvent | MessageAppendedStreamEvent | ReasoningAppendedStreamEvent, + right: ActionInputAppendedStreamEvent | MessageAppendedStreamEvent | ReasoningAppendedStreamEvent, ): boolean { return ( left.data.sequence === right.data.sequence && diff --git a/packages/eve/src/protocol/message.test.ts b/packages/eve/src/protocol/message.test.ts index 513b225e9b..4a3698579b 100644 --- a/packages/eve/src/protocol/message.test.ts +++ b/packages/eve/src/protocol/message.test.ts @@ -25,7 +25,7 @@ import { createEveConnectionCallbackRoutePath } from "#protocol/routes.js"; describe("message stream protocol", () => { it("pins the stream version for timed session events", () => { - expect(EVE_MESSAGE_STREAM_VERSION).toBe("23"); + expect(EVE_MESSAGE_STREAM_VERSION).toBe("24"); }); it("creates authoritative input resolution batches", () => { diff --git a/packages/eve/src/protocol/message.ts b/packages/eve/src/protocol/message.ts index 517ad6bf03..5a5f2c3e6f 100644 --- a/packages/eve/src/protocol/message.ts +++ b/packages/eve/src/protocol/message.ts @@ -27,7 +27,7 @@ export const EVE_STREAM_TAIL_INDEX_HEADER = "x-eve-stream-tail-index"; export const EVE_STREAM_VERSION_HEADER = "x-eve-stream-version"; export const EVE_MESSAGE_STREAM_CONTENT_TYPE = "application/x-ndjson; charset=utf-8"; export const EVE_MESSAGE_STREAM_FORMAT = "ndjson"; -export const EVE_MESSAGE_STREAM_VERSION = "23"; +export const EVE_MESSAGE_STREAM_VERSION = "24"; /** * eve-owned finish reason for one completed assistant step. @@ -423,6 +423,23 @@ export interface MessageAppendedStreamEvent { type: "message.appended"; } +/** + * Stream event emitted while the model is generating the input for one tool + * call, before the validated call is announced via `actions.requested`. + */ +export interface ActionInputAppendedStreamEvent { + data: { + callId: string; + inputTextDelta: string; + inputTextSoFar: string; + sequence: number; + stepIndex: number; + toolName: string; + turnId: string; + }; + type: "action.input.appended"; +} + /** * Stream event emitted when one reasoning delta is appended to the current * reasoning block for the current step. @@ -715,6 +732,7 @@ export interface SessionCompletedStreamEvent { * consumers receive {@link MessageStreamEvent}. */ export type UnstampedMessageStreamEvent = + | ActionInputAppendedStreamEvent | ApprovalCandidateStreamEvent | ApprovalSettledStreamEvent | ContextClearedStreamEvent @@ -1081,6 +1099,30 @@ export function createActionsRequestedEvent(input: { }; } +/** Creates an `action.input.appended` event for streamed tool input text. */ +export function createActionInputAppendedEvent(input: { + readonly callId: string; + readonly inputTextDelta: string; + readonly inputTextSoFar: string; + readonly sequence: number; + readonly stepIndex: number; + readonly toolName: string; + readonly turnId: string; +}): ActionInputAppendedStreamEvent { + return { + data: { + callId: input.callId, + inputTextDelta: input.inputTextDelta, + inputTextSoFar: input.inputTextSoFar, + sequence: input.sequence, + stepIndex: input.stepIndex, + toolName: input.toolName, + turnId: input.turnId, + }, + type: "action.input.appended", + }; +} + /** * Creates the `authorization.required` event for one authorization source * that needs user authorization before it can continue. diff --git a/packages/eve/src/public/definitions/hook.ts b/packages/eve/src/public/definitions/hook.ts index d2520f573b..48e2bc3b5a 100644 --- a/packages/eve/src/public/definitions/hook.ts +++ b/packages/eve/src/public/definitions/hook.ts @@ -15,6 +15,7 @@ type ProtocolEvent = Extract< * until eve exposes them here. */ export interface HookEventMap { + readonly "action.input.appended": ProtocolEvent<"action.input.appended">; readonly "action.partial": ProtocolEvent<"action.partial">; readonly "action.result": ProtocolEvent<"action.result">; readonly "approval.candidate": ProtocolEvent<"approval.candidate">; From cbba52cfbd380213fda8b87a0a5016b5c03290a0 Mon Sep 17 00:00:00 2001 From: "vercel-gh-bot-5[bot]" <312521305+vercel-gh-bot-5[bot]@users.noreply.github.com> Date: Fri, 21 Aug 2026 20:40:03 +0000 Subject: [PATCH 2/5] refactor(eve): contain streamed tool input state Co-authored-by: ruiconti <1834568+ruiconti@users.noreply.github.com> Signed-off-by: Rui Conti --- .../compatibility/channel/v9.ts | 3 + .../compatibility/dynamicInstructions/v13.ts | 10 + .../compatibility/dynamicSkill/v12.ts | 11 + .../compatibility/dynamicTool/v19.ts | 19 + .../compatibility/hook/v14.ts | 13 + .../compatibility/schedule/v3.ts | 18 + .../compatibility/subagent/v4.ts | 7 + .../reports/channel/v10.json | 21 + .../reports/dynamicInstructions/v14.json | 7 + .../reports/dynamicSkill/v13.json | 7 + .../reports/dynamicTool/v20.json | 13 + .../extension-contracts/reports/hook/v15.json | 7 + .../reports/schedule/v4.json | 14 + .../reports/subagent/v5.json | 20 + .../src/client/message-reducer-primitives.ts | 71 --- .../eve/src/client/message-reducer-state.ts | 286 +++++++++ packages/eve/src/client/message-reducer.ts | 246 +------- .../src/compiler/extension-compatibility.ts | 26 +- packages/eve/src/harness/emission.ts | 545 +----------------- .../eve/src/harness/ordered-stream-emitter.ts | 76 ++- .../src/harness/stream-content-emission.ts | 495 ++++++++++++++++ 21 files changed, 1043 insertions(+), 872 deletions(-) create mode 100644 packages/eve/extension-contracts/compatibility/channel/v9.ts create mode 100644 packages/eve/extension-contracts/compatibility/dynamicInstructions/v13.ts create mode 100644 packages/eve/extension-contracts/compatibility/dynamicSkill/v12.ts create mode 100644 packages/eve/extension-contracts/compatibility/dynamicTool/v19.ts create mode 100644 packages/eve/extension-contracts/compatibility/hook/v14.ts create mode 100644 packages/eve/extension-contracts/compatibility/schedule/v3.ts create mode 100644 packages/eve/extension-contracts/compatibility/subagent/v4.ts create mode 100644 packages/eve/extension-contracts/reports/channel/v10.json create mode 100644 packages/eve/extension-contracts/reports/dynamicInstructions/v14.json create mode 100644 packages/eve/extension-contracts/reports/dynamicSkill/v13.json create mode 100644 packages/eve/extension-contracts/reports/dynamicTool/v20.json create mode 100644 packages/eve/extension-contracts/reports/hook/v15.json create mode 100644 packages/eve/extension-contracts/reports/schedule/v4.json create mode 100644 packages/eve/extension-contracts/reports/subagent/v5.json delete mode 100644 packages/eve/src/client/message-reducer-primitives.ts create mode 100644 packages/eve/src/client/message-reducer-state.ts create mode 100644 packages/eve/src/harness/stream-content-emission.ts diff --git a/packages/eve/extension-contracts/compatibility/channel/v9.ts b/packages/eve/extension-contracts/compatibility/channel/v9.ts new file mode 100644 index 0000000000..67e58f0bbe --- /dev/null +++ b/packages/eve/extension-contracts/compatibility/channel/v9.ts @@ -0,0 +1,3 @@ +import { disableRoute } from "#public/channels/index.js"; + +export default disableRoute(); diff --git a/packages/eve/extension-contracts/compatibility/dynamicInstructions/v13.ts b/packages/eve/extension-contracts/compatibility/dynamicInstructions/v13.ts new file mode 100644 index 0000000000..02a61cf705 --- /dev/null +++ b/packages/eve/extension-contracts/compatibility/dynamicInstructions/v13.ts @@ -0,0 +1,10 @@ +import { defineDynamic, defineInstructions } from "#public/instructions/index.js"; + +export default defineDynamic({ + events: { + "session.started": (_event, ctx) => + defineInstructions({ + markdown: `Use session ${ctx.session.id} when correlating evidence.`, + }), + }, +}); diff --git a/packages/eve/extension-contracts/compatibility/dynamicSkill/v12.ts b/packages/eve/extension-contracts/compatibility/dynamicSkill/v12.ts new file mode 100644 index 0000000000..3c6e1498e9 --- /dev/null +++ b/packages/eve/extension-contracts/compatibility/dynamicSkill/v12.ts @@ -0,0 +1,11 @@ +import { defineDynamic, defineSkill } from "#public/skills/index.js"; + +export default defineDynamic({ + events: { + "turn.started": (_event, ctx) => + defineSkill({ + description: `Review evidence for session ${ctx.session.id}.`, + markdown: "# Evidence review\n\nCheck every claim against its source.", + }), + }, +}); diff --git a/packages/eve/extension-contracts/compatibility/dynamicTool/v19.ts b/packages/eve/extension-contracts/compatibility/dynamicTool/v19.ts new file mode 100644 index 0000000000..6252c112c4 --- /dev/null +++ b/packages/eve/extension-contracts/compatibility/dynamicTool/v19.ts @@ -0,0 +1,19 @@ +import { z } from "zod"; + +import { defineDynamic, defineTool } from "#public/tools/index.js"; + +export default defineDynamic({ + events: { + "turn.started": (_event, ctx) => + defineTool({ + description: "Echo the tool call identifiers.", + inputSchema: z.object({ note: z.string() }), + execute: ({ note }, toolCtx) => ({ + callId: toolCtx.callId, + note, + sessionId: ctx.session.id, + toolName: toolCtx.toolName, + }), + }), + }, +}); diff --git a/packages/eve/extension-contracts/compatibility/hook/v14.ts b/packages/eve/extension-contracts/compatibility/hook/v14.ts new file mode 100644 index 0000000000..38ecc7e065 --- /dev/null +++ b/packages/eve/extension-contracts/compatibility/hook/v14.ts @@ -0,0 +1,13 @@ +import { defineHook } from "#public/hooks/index.js"; + +export default defineHook({ + events: { + "subagent.completed"(event, ctx) { + console.info("subagent completed", { + output: event.data.output, + sessionId: ctx.session.id, + subagentName: event.data.subagentName, + }); + }, + }, +}); diff --git a/packages/eve/extension-contracts/compatibility/schedule/v3.ts b/packages/eve/extension-contracts/compatibility/schedule/v3.ts new file mode 100644 index 0000000000..15ba030ed8 --- /dev/null +++ b/packages/eve/extension-contracts/compatibility/schedule/v3.ts @@ -0,0 +1,18 @@ +import channel from "../channel/v7.js"; +import { defineSchedule } from "#public/schedules/index.js"; + +export default defineSchedule({ + cron: "0 0 * * *", + async run({ appAuth, to, waitUntil }) { + waitUntil( + (async () => { + const session = await to(channel, { sessionRef: "daily" }).send("Start review", { + auth: appAuth, + }); + await session.respond([{ optionId: "approve", requestId: "approval-1" }], { + auth: appAuth, + }); + })(), + ); + }, +}); diff --git a/packages/eve/extension-contracts/compatibility/subagent/v4.ts b/packages/eve/extension-contracts/compatibility/subagent/v4.ts new file mode 100644 index 0000000000..e3bdcaa242 --- /dev/null +++ b/packages/eve/extension-contracts/compatibility/subagent/v4.ts @@ -0,0 +1,7 @@ +import { defineAgent } from "#public/index.js"; + +export default defineAgent({ + compaction: { thresholdPercent: 0.8 }, + description: "Delegate research tasks.", + model: "anthropic/claude-sonnet-5", +}); diff --git a/packages/eve/extension-contracts/reports/channel/v10.json b/packages/eve/extension-contracts/reports/channel/v10.json new file mode 100644 index 0000000000..bfb50159bd --- /dev/null +++ b/packages/eve/extension-contracts/reports/channel/v10.json @@ -0,0 +1,21 @@ +{ + "kind": "eve-extension-capability-contract", + "capability": "channel", + "epoch": 10, + "sha256": "5d17981853027789da610ed4ae859c00b3ffb5dec24d12bac8bdfd11370a8df6", + "exports": [ + "DELETE", + "GET", + "HEAD", + "OPTIONS", + "PATCH", + "POST", + "PUT", + "WS", + "createWebSocketUpgradeServer", + "defineChannel", + "disableRoute", + "isChannel", + "isDisabledRouteSentinel" + ] +} diff --git a/packages/eve/extension-contracts/reports/dynamicInstructions/v14.json b/packages/eve/extension-contracts/reports/dynamicInstructions/v14.json new file mode 100644 index 0000000000..20f5f395a0 --- /dev/null +++ b/packages/eve/extension-contracts/reports/dynamicInstructions/v14.json @@ -0,0 +1,7 @@ +{ + "kind": "eve-extension-capability-contract", + "capability": "dynamicInstructions", + "epoch": 14, + "sha256": "4734df4dbc68c6082625975f794c19a0d39004c357d820e445a093d18b6c339b", + "exports": ["defineDynamic"] +} diff --git a/packages/eve/extension-contracts/reports/dynamicSkill/v13.json b/packages/eve/extension-contracts/reports/dynamicSkill/v13.json new file mode 100644 index 0000000000..d1556bef4f --- /dev/null +++ b/packages/eve/extension-contracts/reports/dynamicSkill/v13.json @@ -0,0 +1,7 @@ +{ + "kind": "eve-extension-capability-contract", + "capability": "dynamicSkill", + "epoch": 13, + "sha256": "d0194e3d413ab385e1f2fac57c8ea9bcba9da6a3837a82e5fa0b68b6896833cc", + "exports": ["defineDynamic"] +} diff --git a/packages/eve/extension-contracts/reports/dynamicTool/v20.json b/packages/eve/extension-contracts/reports/dynamicTool/v20.json new file mode 100644 index 0000000000..96dbef7d12 --- /dev/null +++ b/packages/eve/extension-contracts/reports/dynamicTool/v20.json @@ -0,0 +1,13 @@ +{ + "kind": "eve-extension-capability-contract", + "capability": "dynamicTool", + "epoch": 20, + "sha256": "202c8444839e924ac223a4654c6802bbe91e63e4949adf2d044de30d502978d7", + "exports": [ + "DynamicToolEntry", + "DynamicToolEvents", + "DynamicToolResult", + "DynamicToolSet", + "defineDynamic" + ] +} diff --git a/packages/eve/extension-contracts/reports/hook/v15.json b/packages/eve/extension-contracts/reports/hook/v15.json new file mode 100644 index 0000000000..def9c36f26 --- /dev/null +++ b/packages/eve/extension-contracts/reports/hook/v15.json @@ -0,0 +1,7 @@ +{ + "kind": "eve-extension-capability-contract", + "capability": "hook", + "epoch": 15, + "sha256": "bb84054bca066e8f1baa191e6e17b7274792bea39664e810a0c8ab789838b498", + "exports": ["defineHook"] +} diff --git a/packages/eve/extension-contracts/reports/schedule/v4.json b/packages/eve/extension-contracts/reports/schedule/v4.json new file mode 100644 index 0000000000..4d360bf02d --- /dev/null +++ b/packages/eve/extension-contracts/reports/schedule/v4.json @@ -0,0 +1,14 @@ +{ + "kind": "eve-extension-capability-contract", + "capability": "schedule", + "epoch": 4, + "sha256": "19917f781b201d389a25b49c33176301be1298facda8e0e9e3b31db514a25670", + "exports": [ + "ScheduleDefinition", + "ScheduleHandlerArgs", + "ScheduleRunHandler", + "ScheduleToFn", + "TypedReceiveTarget", + "defineSchedule" + ] +} diff --git a/packages/eve/extension-contracts/reports/subagent/v5.json b/packages/eve/extension-contracts/reports/subagent/v5.json new file mode 100644 index 0000000000..95961778e2 --- /dev/null +++ b/packages/eve/extension-contracts/reports/subagent/v5.json @@ -0,0 +1,20 @@ +{ + "kind": "eve-extension-capability-contract", + "capability": "subagent", + "epoch": 5, + "sha256": "b1ce511edad297f42a47572994c0a5af126fe53c63b1b4140ffc6af7aa6c6734", + "exports": [ + "AgentCompactionDefinition", + "AgentDefinition", + "AgentModelDefinition", + "AgentStaticModelDefinition", + "DefinedAgent", + "DynamicLocalSubagentDefinition", + "DynamicSubagentDefinition", + "RemoteAgentDefinition", + "RemoteAgentDefinitionInput", + "defineAgent", + "defineDynamic", + "defineRemoteAgent" + ] +} diff --git a/packages/eve/src/client/message-reducer-primitives.ts b/packages/eve/src/client/message-reducer-primitives.ts deleted file mode 100644 index d221bb8f7a..0000000000 --- a/packages/eve/src/client/message-reducer-primitives.ts +++ /dev/null @@ -1,71 +0,0 @@ -import type { EveMessage, EveMessageData, EveMessagePart } from "#client/message-reducer-types.js"; -import type { MessageReceivedPart } from "#protocol/message.js"; - -export function projectReceivedParts( - parts: readonly MessageReceivedPart[] | undefined, - message: string, -): readonly EveMessagePart[] { - return ( - parts?.map((part) => - part.type === "text" - ? { state: "done", text: part.text, type: "text" } - : { - filename: part.filename, - mediaType: part.mediaType, - size: part.size, - type: "file", - url: part.url, - }, - ) ?? [{ state: "done", text: message, type: "text" }] - ); -} - -export function partKey(part: EveMessagePart): string { - switch (part.type) { - case "text": - return `text:${part.stepIndex ?? 0}`; - case "reasoning": - return `reasoning:${part.stepIndex ?? 0}`; - case "file": - return `file:${part.stepIndex ?? 0}:${part.filename ?? part.url ?? part.mediaType}`; - case "step-start": - return "step-start"; - case "authorization": - return `authorization:${part.turnId}:${part.stepIndex}:${part.name}`; - case "dynamic-tool": - return `dynamic-tool:${part.toolCallId}`; - } -} - -export function upsertMessage(data: EveMessageData, next: EveMessage): EveMessageData { - const index = data.messages.findIndex((message) => message.id === next.id); - if (index === -1) { - return { messages: [...data.messages, next] }; - } - - return { - messages: [...data.messages.slice(0, index), next, ...data.messages.slice(index + 1)], - }; -} - -export function removeStreamingToolPartsForTurn( - data: EveMessageData, - turnId: string, -): EveMessageData { - const index = data.messages.findIndex( - (message) => message.role === "assistant" && message.metadata?.turnId === turnId, - ); - const message = data.messages[index]; - if (message === undefined) return data; - - return upsertMessage(data, { - ...message, - parts: message.parts.filter( - (part) => part.type !== "dynamic-tool" || part.state !== "input-streaming", - ), - }); -} - -export function optimisticUserMessageId(submissionId: string): string { - return `optimistic:${submissionId}:user`; -} diff --git a/packages/eve/src/client/message-reducer-state.ts b/packages/eve/src/client/message-reducer-state.ts new file mode 100644 index 0000000000..ebba6cd9e1 --- /dev/null +++ b/packages/eve/src/client/message-reducer-state.ts @@ -0,0 +1,286 @@ +import type { + EveAuthorizationPart, + EveDynamicToolPart, + EveMessage, + EveMessageData, + EveMessageMetadata, + EveMessagePart, +} from "#client/message-reducer-types.js"; +import type { MessageReceivedPart } from "#protocol/message.js"; + +export type EveAssistantMessage = EveMessage & { readonly role: "assistant" }; + +export function projectReceivedParts( + parts: readonly MessageReceivedPart[] | undefined, + message: string, +): readonly EveMessagePart[] { + return ( + parts?.map((part) => + "text" in part + ? { state: "done", text: part.text, type: "text" } + : { + filename: part.filename, + mediaType: part.mediaType, + size: part.size, + type: "file", + url: part.url, + }, + ) ?? [{ state: "done", text: message, type: "text" }] + ); +} + +function partKey(part: EveMessagePart): string { + switch (part.type) { + case "text": + return `text:${part.stepIndex ?? 0}`; + case "reasoning": + return `reasoning:${part.stepIndex ?? 0}`; + case "file": + return `file:${part.stepIndex ?? 0}:${part.filename ?? part.url ?? part.mediaType}`; + case "step-start": + return "step-start"; + case "authorization": + return `authorization:${part.turnId}:${part.stepIndex}:${part.name}`; + case "dynamic-tool": + return `dynamic-tool:${part.toolCallId}`; + } +} + +export function upsertMessage(data: EveMessageData, next: EveMessage): EveMessageData { + const index = data.messages.findIndex((message) => message.id === next.id); + if (index === -1) { + return { messages: [...data.messages, next] }; + } + if (data.messages[index] === next) return data; + + return { + messages: [...data.messages.slice(0, index), next, ...data.messages.slice(index + 1)], + }; +} + +export function optimisticUserMessageId(submissionId: string): string { + return `optimistic:${submissionId}:user`; +} + +export function updateAssistantMessage( + data: EveMessageData, + turnId: string, + update: (message: EveAssistantMessage) => EveAssistantMessage, +): EveMessageData { + const message = findAssistantMessage(data, turnId) ?? createAssistantMessage(turnId); + return upsertMessage(data, update(message)); +} + +export function updateExistingAssistantMessage( + data: EveMessageData, + turnId: string, + update: (message: EveAssistantMessage) => EveAssistantMessage, +): EveMessageData { + const message = findAssistantMessage(data, turnId); + return message === undefined ? data : upsertMessage(data, update(message)); +} + +export function updateAssistantMetadata( + data: EveMessageData, + turnId: string, + metadata: EveMessageMetadata, +): EveMessageData { + return updateAssistantMessage(data, turnId, (message) => ({ + ...message, + metadata: { + ...message.metadata, + ...metadata, + }, + })); +} + +export function ensureStepStartPart( + message: EveAssistantMessage, + stepIndex: number, +): EveAssistantMessage { + const stepStartCount = message.parts.filter((part) => part.type === "step-start").length; + if (stepStartCount > stepIndex) return message; + + const missingCount = stepIndex - stepStartCount + 1; + return { + ...message, + parts: [ + ...message.parts, + ...Array.from({ length: missingCount }, () => ({ type: "step-start" as const })), + ], + }; +} + +export function upsertPart( + message: EveAssistantMessage, + next: EveMessagePart, +): EveAssistantMessage { + const index = message.parts.findIndex((part) => partKey(part) === partKey(next)); + const parts = + index === -1 + ? [...message.parts, next] + : [...message.parts.slice(0, index), next, ...message.parts.slice(index + 1)]; + + return { + ...message, + metadata: { + ...message.metadata, + status: next.type === "text" && next.state === "done" ? "complete" : "streaming", + }, + parts, + }; +} + +type EveRunPart = Extract; + +// Upserts a text/reasoning part, keeping multiple runs per step distinct: one +// step can produce text, call tools, then produce more text (see +// `MessageCompletedStreamEvent`), so a step-only key would collapse them. +// +// We find the latest same-step run of this type: while it is still streaming, +// its snapshots replace it in place; once it is done (or there is none), `next` +// begins a new run appended in arrival order. +export function upsertRun(message: EveAssistantMessage, next: EveRunPart): EveAssistantMessage { + let lastIndex = -1; + for (let index = message.parts.length - 1; index >= 0; index -= 1) { + const part = message.parts[index]; + if (part?.type === next.type && part.stepIndex === next.stepIndex) { + lastIndex = index; + break; + } + } + + const openRun = + lastIndex !== -1 && (message.parts[lastIndex] as EveRunPart).state === "streaming"; + const parts = openRun + ? [...message.parts.slice(0, lastIndex), next, ...message.parts.slice(lastIndex + 1)] + : [...message.parts, next]; + + return { + ...message, + metadata: { + ...message.metadata, + status: next.type === "text" && next.state === "done" ? "complete" : "streaming", + }, + parts, + }; +} + +export function removeTextPart( + message: EveAssistantMessage, + stepIndex: number, +): EveAssistantMessage { + const parts = message.parts.filter( + (part) => part.type !== "text" || part.stepIndex !== stepIndex, + ); + if (parts.length === message.parts.length) return message; + + return { + ...message, + metadata: { + ...message.metadata, + status: "complete", + }, + parts, + }; +} + +export function removeStreamingToolParts( + parts: readonly EveMessagePart[], +): readonly EveMessagePart[] { + const next = parts.filter( + (part) => part.type !== "dynamic-tool" || part.state !== "input-streaming", + ); + return next.length === parts.length ? parts : next; +} + +export function updateToolPart( + data: EveMessageData, + toolCallId: string, + next: EveDynamicToolPart, +): EveMessageData { + const message = data.messages.find( + (candidate): candidate is EveAssistantMessage => + candidate.role === "assistant" && + candidate.parts.some( + (part) => part.type === "dynamic-tool" && part.toolCallId === toolCallId, + ), + ); + return message === undefined ? data : upsertMessage(data, upsertPart(message, next)); +} + +export function updateAuthorizationPart( + data: EveMessageData, + existing: EveAuthorizationPart, + next: EveAuthorizationPart, +): EveMessageData { + const message = data.messages.find( + (candidate): candidate is EveAssistantMessage => + candidate.role === "assistant" && candidate.parts.some((part) => part === existing), + ); + return message === undefined ? data : upsertMessage(data, upsertPart(message, next)); +} + +export function findToolPart( + data: EveMessageData, + toolCallId: string, +): EveDynamicToolPart | undefined { + for (const message of data.messages) { + for (const part of message.parts) { + if (part.type === "dynamic-tool" && part.toolCallId === toolCallId) return part; + } + } + return undefined; +} + +export function findLatestPendingAuthorizationPart( + data: EveMessageData, + name: string, +): EveAuthorizationPart | undefined { + for (let messageIndex = data.messages.length - 1; messageIndex >= 0; messageIndex -= 1) { + const message = data.messages[messageIndex]; + if (message?.role !== "assistant") continue; + + for (let partIndex = message.parts.length - 1; partIndex >= 0; partIndex -= 1) { + const part = message.parts[partIndex]; + if (part?.type === "authorization" && part.state === "required" && part.name === name) { + return part; + } + } + } + return undefined; +} + +export function findToolPartByApprovalId( + data: EveMessageData, + approvalId: string, +): EveDynamicToolPart | undefined { + for (const message of data.messages) { + for (const part of message.parts) { + if (part.type === "dynamic-tool" && part.approval?.id === approvalId) return part; + } + } + return undefined; +} + +function findAssistantMessage( + data: EveMessageData, + turnId: string, +): EveAssistantMessage | undefined { + return data.messages.find( + (message): message is EveAssistantMessage => + message.role === "assistant" && message.metadata?.turnId === turnId, + ); +} + +function createAssistantMessage(turnId: string): EveAssistantMessage { + return { + id: `${turnId}:assistant`, + metadata: { + status: "streaming", + turnId, + }, + parts: [], + role: "assistant", + }; +} diff --git a/packages/eve/src/client/message-reducer.ts b/packages/eve/src/client/message-reducer.ts index 39c9691d95..0a8bb1c578 100644 --- a/packages/eve/src/client/message-reducer.ts +++ b/packages/eve/src/client/message-reducer.ts @@ -3,14 +3,7 @@ import { createAuthorizationCompletedPart, createAuthorizationRequiredPart, } from "#client/authorization-message-parts.js"; -import type { - EveAuthorizationPart, - EveMessageData, - EveDynamicToolPart, - EveMessage, - EveMessageMetadata, - EveMessagePart, -} from "#client/message-reducer-types.js"; +import type { EveDynamicToolPart, EveMessageData } from "#client/message-reducer-types.js"; import { approvedApproval, createToolMetadata, @@ -21,12 +14,23 @@ import { toMessageInputRequest, } from "#client/message-action-parts.js"; import { + ensureStepStartPart, + findLatestPendingAuthorizationPart, + findToolPart, + findToolPartByApprovalId, optimisticUserMessageId, - partKey, projectReceivedParts, - removeStreamingToolPartsForTurn, + removeStreamingToolParts, + removeTextPart, + updateAssistantMessage, + updateAssistantMetadata, + updateAuthorizationPart, + updateExistingAssistantMessage, + updateToolPart, upsertMessage, -} from "#client/message-reducer-primitives.js"; + upsertPart, + upsertRun, +} from "#client/message-reducer-state.js"; import type { InputResponse } from "#shared/input.js"; import type { AuthorizationCompletedStreamEvent, InputResolution } from "#protocol/message.js"; @@ -43,8 +47,6 @@ export type { EveMessageToolMetadata, } from "#client/message-reducer-types.js"; -type EveAssistantMessage = EveMessage & { readonly role: "assistant" }; - /** * Creates a UIMessage-compatible eve reducer for chat and agent UIs. * @@ -397,7 +399,10 @@ function reduceMessageData(data: EveMessageData, event: EveAgentReducerEvent): E })); case "turn.failed": - return removeStreamingToolPartsForTurn(data, event.data.turnId); + return updateExistingAssistantMessage(data, event.data.turnId, (message) => { + const parts = removeStreamingToolParts(message.parts); + return parts === message.parts ? message : { ...message, parts }; + }); case "session.failed": return data; @@ -407,10 +412,6 @@ function reduceMessageData(data: EveMessageData, event: EveAgentReducerEvent): E } } -function removeStreamingToolParts(parts: readonly EveMessagePart[]): readonly EveMessagePart[] { - return parts.filter((part) => part.type !== "dynamic-tool" || part.state !== "input-streaming"); -} - function respondToInputRequest(data: EveMessageData, response: InputResponse): EveMessageData { const existing = findToolPartByApprovalId(data, response.requestId); if (!existing) return data; @@ -460,152 +461,6 @@ function resolveInputRequest(data: EveMessageData, resolution: InputResolution): }); } -function updateAssistantMessage( - data: EveMessageData, - turnId: string, - update: (message: EveAssistantMessage) => EveAssistantMessage, -): EveMessageData { - const existing = data.messages.find( - (message): message is EveAssistantMessage => - message.role === "assistant" && message.metadata?.turnId === turnId, - ); - - const message = existing ?? createAssistantMessage(turnId); - return upsertMessage(data, update(message)); -} - -function updateAssistantMetadata( - data: EveMessageData, - turnId: string, - metadata: EveMessageMetadata, -): EveMessageData { - return updateAssistantMessage(data, turnId, (message) => ({ - ...message, - metadata: { - ...message.metadata, - ...metadata, - }, - })); -} - -function createAssistantMessage(turnId: string): EveAssistantMessage { - return { - id: `${turnId}:assistant`, - metadata: { - status: "streaming", - turnId, - }, - parts: [], - role: "assistant", - }; -} - -function ensureStepStartPart(message: EveAssistantMessage, stepIndex: number): EveAssistantMessage { - const stepStartCount = message.parts.filter((part) => part.type === "step-start").length; - if (stepStartCount > stepIndex) { - return message; - } - - const missingCount = stepIndex - stepStartCount + 1; - return { - ...message, - parts: [ - ...message.parts, - ...Array.from({ length: missingCount }, () => ({ type: "step-start" as const })), - ], - }; -} - -function upsertPart(message: EveAssistantMessage, next: EveMessagePart): EveAssistantMessage { - const index = message.parts.findIndex((part) => partKey(part) === partKey(next)); - const parts = - index === -1 - ? [...message.parts, next] - : [...message.parts.slice(0, index), next, ...message.parts.slice(index + 1)]; - - return { - ...message, - metadata: { - ...message.metadata, - status: next.type === "text" && next.state === "done" ? "complete" : "streaming", - }, - parts, - }; -} - -type EveRunPart = Extract; - -// Upserts a text/reasoning part, keeping multiple runs per step distinct: one -// step can produce text, call tools, then produce more text (see -// `MessageCompletedStreamEvent`), so a step-only key would collapse them. -// -// We find the latest same-step run of this type: while it is still streaming, -// its snapshots replace it in place; once it is done (or there is none), `next` -// begins a new run appended in arrival order. -function upsertRun(message: EveAssistantMessage, next: EveRunPart): EveAssistantMessage { - let lastIndex = -1; - for (let index = message.parts.length - 1; index >= 0; index -= 1) { - const part = message.parts[index]; - if (part?.type === next.type && part.stepIndex === next.stepIndex) { - lastIndex = index; - break; - } - } - - const openRun = - lastIndex !== -1 && (message.parts[lastIndex] as EveRunPart).state === "streaming"; - const parts = openRun - ? [...message.parts.slice(0, lastIndex), next, ...message.parts.slice(lastIndex + 1)] - : [...message.parts, next]; - - return { - ...message, - metadata: { - ...message.metadata, - status: next.type === "text" && next.state === "done" ? "complete" : "streaming", - }, - parts, - }; -} - -function removeTextPart(message: EveAssistantMessage, stepIndex: number): EveAssistantMessage { - const parts = message.parts.filter( - (part) => part.type !== "text" || part.stepIndex !== stepIndex, - ); - if (parts.length === message.parts.length) { - return message; - } - - return { - ...message, - metadata: { - ...message.metadata, - status: "complete", - }, - parts, - }; -} - -function updateToolPart( - data: EveMessageData, - toolCallId: string, - next: EveDynamicToolPart, -): EveMessageData { - const message = data.messages.find( - (candidate): candidate is EveAssistantMessage => - candidate.role === "assistant" && - candidate.parts.some( - (part) => part.type === "dynamic-tool" && part.toolCallId === toolCallId, - ), - ); - - if (!message) { - return data; - } - - return upsertMessage(data, upsertPart(message, next)); -} - function completeAuthorization( data: EveMessageData, event: AuthorizationCompletedStreamEvent, @@ -622,34 +477,6 @@ function completeAuthorization( ); } -function updateAuthorizationPart( - data: EveMessageData, - existing: EveAuthorizationPart, - next: EveAuthorizationPart, -): EveMessageData { - const message = data.messages.find( - (candidate): candidate is EveAssistantMessage => - candidate.role === "assistant" && candidate.parts.some((part) => part === existing), - ); - - if (!message) { - return data; - } - - return upsertMessage(data, upsertPart(message, next)); -} - -function findToolPart(data: EveMessageData, toolCallId: string): EveDynamicToolPart | undefined { - for (const message of data.messages) { - for (const part of message.parts) { - if (part.type === "dynamic-tool" && part.toolCallId === toolCallId) { - return part; - } - } - } - return undefined; -} - function isSettledToolPart(part: EveDynamicToolPart): boolean { return ( part.state === "output-denied" || @@ -657,38 +484,3 @@ function isSettledToolPart(part: EveDynamicToolPart): boolean { (part.state === "output-available" && part.partial !== true) ); } - -function findLatestPendingAuthorizationPart( - data: EveMessageData, - name: string, -): EveAuthorizationPart | undefined { - for (let messageIndex = data.messages.length - 1; messageIndex >= 0; messageIndex -= 1) { - const message = data.messages[messageIndex]; - if (message?.role !== "assistant") { - continue; - } - - for (let partIndex = message.parts.length - 1; partIndex >= 0; partIndex -= 1) { - const part = message.parts[partIndex]; - if (part?.type === "authorization" && part.state === "required" && part.name === name) { - return part; - } - } - } - - return undefined; -} - -function findToolPartByApprovalId( - data: EveMessageData, - approvalId: string, -): EveDynamicToolPart | undefined { - for (const message of data.messages) { - for (const part of message.parts) { - if (part.type === "dynamic-tool" && part.approval?.id === approvalId) { - return part; - } - } - } - return undefined; -} diff --git a/packages/eve/src/compiler/extension-compatibility.ts b/packages/eve/src/compiler/extension-compatibility.ts index fe00e407b6..44869506d5 100644 --- a/packages/eve/src/compiler/extension-compatibility.ts +++ b/packages/eve/src/compiler/extension-compatibility.ts @@ -27,15 +27,15 @@ const EXTENSION_CAPABILITY_CONTRACTS = { dropped: { 15: "TaskExec replaces stageEffect with send" }, }, dynamicTool: { - current: 19, - supported: [1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 11, 12, 13, 14, 15, 16, 17, 18, 19], + current: 20, + supported: [1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 11, 12, 13, 14, 15, 16, 17, 18, 19, 20], dropped: {}, }, - channel: { current: 9, supported: [1, 2, 3, 4, 5, 6, 7, 8, 9], dropped: {} }, - schedule: { current: 3, supported: [1, 2, 3], dropped: {} }, + channel: { current: 10, supported: [1, 2, 3, 4, 5, 6, 7, 8, 9, 10], dropped: {} }, + schedule: { current: 4, supported: [1, 2, 3, 4], dropped: {} }, subagent: { - current: 4, - supported: [3, 4], + current: 5, + supported: [3, 4, 5], dropped: { 1: "Persistent subagent sessions are now the default and the experimental opt-in was removed", 2: "Persistent subagent sessions are now the default and the experimental opt-in was removed", @@ -47,8 +47,8 @@ const EXTENSION_CAPABILITY_CONTRACTS = { dropped: {}, }, hook: { - current: 14, - supported: [10, 11, 12, 13, 14], + current: 15, + supported: [10, 11, 12, 13, 14, 15], dropped: { 1: "Model identity moved from session.started runtime metadata to step.started call attribution.", 2: "Model identity moved from session.started runtime metadata to step.started call attribution.", @@ -62,13 +62,17 @@ const EXTENSION_CAPABILITY_CONTRACTS = { }, }, skill: { current: 1, supported: [1], dropped: {} }, - dynamicSkill: { current: 12, supported: [1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 11, 12], dropped: {} }, - instructions: { current: 2, supported: [1, 2], dropped: {} }, - dynamicInstructions: { + dynamicSkill: { current: 13, supported: [1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 11, 12, 13], dropped: {}, }, + instructions: { current: 2, supported: [1, 2], dropped: {} }, + dynamicInstructions: { + current: 14, + supported: [1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 11, 12, 13, 14], + dropped: {}, + }, config: { current: 1, supported: [1], dropped: {} }, state: { current: 3, supported: [1, 2, 3], dropped: {} }, } as const satisfies Record; diff --git a/packages/eve/src/harness/emission.ts b/packages/eve/src/harness/emission.ts index 92b1bdc21e..1635ec555f 100644 --- a/packages/eve/src/harness/emission.ts +++ b/packages/eve/src/harness/emission.ts @@ -1,30 +1,6 @@ -import type { - ModelMessage, - TextStreamPart, - ToolSet, - TypedToolCall, - TypedToolError, - TypedToolResult, -} from "ai"; - -type ToolResponsePart = Extract["content"][number]; -type InlineToolResultPart = Extract; - -import type { - AssistantStepFinishReason, - RuntimeIdentity, - RuntimeTraceContext, -} from "#protocol/message.js"; +import type { RuntimeIdentity, RuntimeTraceContext } from "#protocol/message.js"; import { - createActionsRequestedEvent, - createActionInputAppendedEvent, - createActionPartialEvent, - createActionResultEvent, - createMessageAppendedEvent, - createMessageCompletedEvent, createMessageReceivedEvent, - createReasoningAppendedEvent, - createReasoningCompletedEvent, createSessionCompletedEvent, createSessionFailedEvent, createSessionStartedEvent, @@ -36,27 +12,9 @@ import { createTurnStartedEvent, } from "#protocol/message.js"; import type { RunMode } from "#shared/run-mode.js"; -import { hasEmptyDeliverySentinel } from "#shared/empty-delivery.js"; import type { JsonObject } from "#shared/json.js"; -import { - createRuntimeToolResultFromStepResult, - createRuntimeToolResultFromToolError, - createToolResultMessagePartFromToolError, -} from "#harness/action-result-helpers.js"; -import { createRuntimeActionRequestFromToolCall } from "#harness/runtime-actions.js"; -import { - createInvalidToolCallInputError, - isInvalidToolCall, - resolveProviderToolCallRequest, -} from "#harness/tool-call-input-errors.js"; -import type { RuntimeActionRequest, RuntimeToolResultActionResult } from "#shared/action-types.js"; -import { createProviderStreamActionBatch } from "#harness/stream-actions.js"; -import { normalizeModelStreamError } from "#harness/model-call-error.js"; -import { createOrderedStreamEmitter } from "#harness/ordered-stream-emitter.js"; -import { interruptStreamOnFailure } from "#harness/interruptible-stream.js"; -import { isInlineAuthorizationToolResult } from "#harness/inline-tool-authorization.js"; import type { HarnessEmissionState } from "#harness/emission-state.js"; -import type { HarnessEmitFn, HarnessToolMap, StepInput } from "#harness/types.js"; +import type { HarnessEmitFn, StepInput } from "#harness/types.js"; export { getHarnessEmissionState, @@ -64,10 +22,10 @@ export { setHarnessEmissionState, } from "#harness/emission-state.js"; export type { HarnessEmissionState } from "#harness/emission-state.js"; - -// --------------------------------------------------------------------------- -// Turn lifecycle helpers -// --------------------------------------------------------------------------- +export { + emitStreamContent, + normalizeAssistantStepFinishReason, +} from "#harness/stream-content-emission.js"; /** * Emits `session.started` (once), `turn.started`, and `message.received` at the @@ -106,9 +64,7 @@ export async function emitTurnPreamble( }; } -/** - * Emits `step.started` for one model call. - */ +/** Emits `step.started` for one model call. */ export async function emitStepStarted( emitFn: HarnessEmitFn, state: HarnessEmissionState, @@ -132,11 +88,7 @@ interface FailedStepPayload { readonly message: string; } -/** - * Emits the shared head of both failure cascades: `step.failed` → - * `turn.failed`. Both terminal and recoverable paths diverge only on - * the third event (`session.failed` vs. `session.waiting`). - */ +/** Emits the shared `step.failed` → `turn.failed` failure prefix. */ async function emitStepAndTurnFailed( emitFn: HarnessEmitFn, state: HarnessEmissionState, @@ -159,15 +111,7 @@ async function emitStepAndTurnFailed( ); } -/** - * Emits the full terminal failure cascade: `step.failed` → - * `turn.failed` → `session.failed`. - * - * Use this when the session cannot be salvaged (structural config - * error, auth misconfig, non-recoverable provider response). The - * `session.failed` tail tells adapters the session is dead and no - * further follow-up is possible on the same continuation token. - */ +/** Emits the terminal `step.failed` → `turn.failed` → `session.failed` cascade. */ export async function emitFailedStep( emitFn: HarnessEmitFn, state: HarnessEmissionState, @@ -177,10 +121,7 @@ export async function emitFailedStep( await emitFn(createSessionFailedEvent(input)); } -/** - * Emits the recoverable failure cascade: `step.failed` → - * `turn.failed` → `session.waiting`. - */ +/** Emits the recoverable `step.failed` → `turn.failed` → `session.waiting` cascade. */ export async function emitRecoverableFailedTurn( emitFn: HarnessEmitFn, state: HarnessEmissionState, @@ -197,9 +138,7 @@ export async function emitRecoverableFailedTurn( }; } -/** - * Returns updated emission state for the next step in the current turn. - */ +/** Returns updated emission state for the next step in the current turn. */ export function advanceStep(state: HarnessEmissionState): HarnessEmissionState { return { ...state, @@ -207,10 +146,7 @@ export function advanceStep(state: HarnessEmissionState): HarnessEmissionState { }; } -/** - * Emits `turn.completed` and either `session.waiting` or `session.completed`. - * Returns updated emission state with an incremented sequence. - */ +/** Emits the turn terminal events and advances emission state to the next turn. */ export async function emitTurnEpilogue( emitFn: HarnessEmitFn, state: HarnessEmissionState, @@ -238,460 +174,3 @@ export async function emitTurnEpilogue( turnId: "", }; } - -// --------------------------------------------------------------------------- -// Shared helpers -// --------------------------------------------------------------------------- - -/** - * Maps an AI SDK finish reason string to the eve-owned - * {@link AssistantStepFinishReason} union. Unknown values become `"other"`. - */ -export function normalizeAssistantStepFinishReason( - value: string | undefined, -): AssistantStepFinishReason { - switch (value) { - case "content-filter": - case "error": - case "length": - case "stop": - case "tool-calls": - return value; - default: - return "other"; - } -} - -// --------------------------------------------------------------------------- -// Stream content emission -// --------------------------------------------------------------------------- - -/** - * Result of consuming one step's `fullStream`. - * - * Inline results avoid duplicate post-step events. Approval-resume - * authorization results also route back to the park detector. - */ -interface EmittedStreamContent { - readonly emittedActionCallIds: ReadonlySet; - readonly handledInlineToolResultCallIds: ReadonlySet; - readonly invalidInputToolCallIds: ReadonlySet; - readonly inlineAuthorizationResults: readonly TypedToolResult[]; - readonly trailingInlineToolResultParts: readonly InlineToolResultPart[]; -} - -interface StreamActionEmissionOptions { - readonly excludedActionToolNames: ReadonlySet; - readonly tools: HarnessToolMap; -} - -/** Consumes `fullStream` in source order, batching provider calls before their first result. */ -export async function emitStreamContent( - emitFn: HarnessEmitFn, - state: HarnessEmissionState, - fullStream: AsyncIterable>, - options?: StreamActionEmissionOptions, -): Promise { - const orderedEmitter = createOrderedStreamEmitter(emitFn); - const providerActionBatch = createProviderStreamActionBatch({ - emitFn: orderedEmitter.emit, - state, - }); - try { - return await consumeStreamContent( - orderedEmitter.emit, - state, - interruptStreamOnFailure(fullStream, orderedEmitter.failureSignal), - providerActionBatch, - options, - ); - } finally { - try { - await providerActionBatch.cancel(); - } finally { - await orderedEmitter.closeAndDrain(); - } - } -} - -async function consumeStreamContent( - emitFn: HarnessEmitFn, - state: HarnessEmissionState, - fullStream: AsyncIterable>, - providerActionBatch: ReturnType, - options?: StreamActionEmissionOptions, -): Promise { - let currentReasoning = ""; - let currentMessage = ""; - let finishReason: AssistantStepFinishReason = "stop"; - let streamError: Error | undefined; - const toolCallIdsSeenInStream = new Set(); - const emittedActionCallIds = new Set(); - const emittedActionResultCallIds = new Set(); - const providerToolCallIdsSeen = new Set(); - const handledInlineToolResultCallIds = new Set(); - const invalidInputToolCallIds = new Set(); - const inlineAuthorizationResults: TypedToolResult[] = []; - const trailingInlineToolResultParts: InlineToolResultPart[] = []; - const streamingActionInputs = new Map(); - - const flushCurrentMessage = async (): Promise => { - if (currentMessage.length === 0) { - return; - } - await emitFn( - createMessageCompletedEvent({ - finishReason: "tool-calls", - message: currentMessage, - sequence: state.sequence, - stepIndex: state.stepIndex, - turnId: state.turnId, - }), - ); - currentMessage = ""; - }; - - const emitActionInput = async ( - callId: string, - toolName: string, - inputTextDelta: string, - inputTextSoFar: string, - ): Promise => - emitFn( - createActionInputAppendedEvent({ - callId, - inputTextDelta, - inputTextSoFar, - sequence: state.sequence, - stepIndex: state.stepIndex, - toolName, - turnId: state.turnId, - }), - ); - - const emitActionRequest = async (action: RuntimeActionRequest): Promise => { - if (emittedActionCallIds.has(action.callId)) { - return; - } - - if (currentMessage.trim().length > 0) { - await flushCurrentMessage(); - } - - emittedActionCallIds.add(action.callId); - await emitFn( - createActionsRequestedEvent({ - actions: [action], - sequence: state.sequence, - stepIndex: state.stepIndex, - turnId: state.turnId, - }), - ); - }; - - const collectProviderToolCall = async (toolCall: { - readonly input?: unknown; - readonly toolCallId: string; - readonly toolName: string; - }): Promise => { - if (providerToolCallIdsSeen.has(toolCall.toolCallId)) { - return; - } - providerToolCallIdsSeen.add(toolCall.toolCallId); - if (emittedActionCallIds.has(toolCall.toolCallId)) { - return; - } - emittedActionCallIds.add(toolCall.toolCallId); - - if (currentMessage.trim().length > 0) { - await flushCurrentMessage(); - } - - const resolved = resolveProviderToolCallRequest(toolCall); - if (resolved.toolError !== undefined) { - invalidInputToolCallIds.add(toolCall.toolCallId); - await emitActionResult(createRuntimeToolResultFromToolError(resolved.toolError)); - handledInlineToolResultCallIds.add(toolCall.toolCallId); - trailingInlineToolResultParts.push( - createToolResultMessagePartFromToolError(resolved.toolError), - ); - return; - } - - providerActionBatch.observe(resolved.request); - }; - - const emitActionResult = async (result: RuntimeToolResultActionResult): Promise => { - if (emittedActionResultCallIds.has(result.callId)) { - return; - } - emittedActionResultCallIds.add(result.callId); - await emitFn( - createActionResultEvent({ - result, - sequence: state.sequence, - stepIndex: state.stepIndex, - turnId: state.turnId, - }), - ); - }; - - const emitActionPartial = async (result: RuntimeToolResultActionResult): Promise => { - await emitFn( - createActionPartialEvent({ - result, - sequence: state.sequence, - stepIndex: state.stepIndex, - turnId: state.turnId, - }), - ); - }; - - const emitToolCall = async (toolCall: TypedToolCall): Promise => { - if (isInvalidToolCall(toolCall)) { - invalidInputToolCallIds.add(toolCall.toolCallId); - return; - } - if (options === undefined || options.excludedActionToolNames.has(toolCall.toolName)) { - return; - } - - try { - await emitActionRequest( - createRuntimeActionRequestFromToolCall({ - toolCall, - tools: options.tools, - }), - ); - } catch (error) { - if (error instanceof TypeError) { - const toolError = createInvalidToolCallInputError({ error, toolCall }); - invalidInputToolCallIds.add(toolCall.toolCallId); - if (currentMessage.trim().length > 0) { - await flushCurrentMessage(); - } - await emitActionResult(createRuntimeToolResultFromToolError(toolError)); - handledInlineToolResultCallIds.add(toolCall.toolCallId); - trailingInlineToolResultParts.push(createToolResultMessagePartFromToolError(toolError)); - return; - } - throw error; - } - }; - - for await (const part of fullStream) { - if (streamError !== undefined) { - continue; - } - - switch (part.type) { - case "reasoning-delta": - await providerActionBatch.flush(); - currentReasoning += part.text; - await emitFn( - createReasoningAppendedEvent({ - reasoningDelta: part.text, - reasoningSoFar: currentReasoning, - sequence: state.sequence, - stepIndex: state.stepIndex, - turnId: state.turnId, - }), - ); - break; - case "text-delta": - await providerActionBatch.flush(); - // Flush accumulated reasoning before text begins. - if (currentReasoning.trim().length > 0) { - await emitFn( - createReasoningCompletedEvent({ - reasoning: currentReasoning, - sequence: state.sequence, - stepIndex: state.stepIndex, - turnId: state.turnId, - }), - ); - currentReasoning = ""; - } - currentMessage += part.text; - await emitFn( - createMessageAppendedEvent({ - messageDelta: part.text, - messageSoFar: currentMessage, - sequence: state.sequence, - stepIndex: state.stepIndex, - turnId: state.turnId, - }), - ); - break; - case "tool-input-start": { - if ( - options === undefined || - part.providerExecuted === true || - options.excludedActionToolNames.has(part.toolName) - ) { - streamingActionInputs.delete(part.id); - break; - } - await providerActionBatch.flush(); - if (currentMessage.trim().length > 0) { - await flushCurrentMessage(); - } - streamingActionInputs.set(part.id, { text: "", toolName: part.toolName }); - await emitActionInput(part.id, part.toolName, "", ""); - break; - } - case "tool-input-delta": { - const input = streamingActionInputs.get(part.id); - if (input === undefined) { - break; - } - await providerActionBatch.flush(); - input.text += part.delta; - await emitActionInput(part.id, input.toolName, part.delta, input.text); - break; - } - case "tool-input-end": - streamingActionInputs.delete(part.id); - break; - case "tool-call": { - const toolCall = part as TypedToolCall; - streamingActionInputs.delete(toolCall.toolCallId); - toolCallIdsSeenInStream.add(toolCall.toolCallId); - if (toolCall.providerExecuted === true) { - await collectProviderToolCall(toolCall); - } else { - await providerActionBatch.flush(); - await emitToolCall(toolCall); - } - break; - } - case "tool-result": { - const inlineToolResult = part as TypedToolResult; - if (inlineToolResult.preliminary === true) { - if (inlineToolResult.providerExecuted !== true) { - await emitActionPartial(createRuntimeToolResultFromStepResult(inlineToolResult)); - } - break; - } - if (inlineToolResult.providerExecuted === true) { - await collectProviderToolCall({ - input: "input" in inlineToolResult ? inlineToolResult.input : undefined, - toolCallId: inlineToolResult.toolCallId, - toolName: inlineToolResult.toolName, - }); - await providerActionBatch.flush(); - await emitActionResult(createRuntimeToolResultFromStepResult(inlineToolResult)); - // Provider results already live in the assistant response. Do not - // add a local tool message. - break; - } - - if (toolCallIdsSeenInStream.has(part.toolCallId)) { - if (isInlineAuthorizationToolResult(inlineToolResult)) { - break; - } - if (emittedActionCallIds.has(part.toolCallId)) { - await emitActionResult(createRuntimeToolResultFromStepResult(inlineToolResult)); - handledInlineToolResultCallIds.add(part.toolCallId); - } - break; - } - - // An approved tool can resume with its result but no matching call in - // this step. Emit it before the message that consumes it. - await providerActionBatch.flush(); - await flushCurrentMessage(); - if (isInlineAuthorizationToolResult(inlineToolResult)) { - // Keep authorization output for the park detector instead of - // emitting a normal tool result. - handledInlineToolResultCallIds.add(part.toolCallId); - inlineAuthorizationResults.push(inlineToolResult); - break; - } - await emitActionResult(createRuntimeToolResultFromStepResult(inlineToolResult)); - handledInlineToolResultCallIds.add(part.toolCallId); - break; - } - case "tool-error": { - const toolError = part as TypedToolError; - if (toolError.providerExecuted === true) { - await collectProviderToolCall(toolError); - await providerActionBatch.flush(); - await emitActionResult(createRuntimeToolResultFromToolError(toolError)); - } else if (emittedActionCallIds.has(toolError.toolCallId)) { - await emitActionResult(createRuntimeToolResultFromToolError(toolError)); - handledInlineToolResultCallIds.add(toolError.toolCallId); - trailingInlineToolResultParts.push(createToolResultMessagePartFromToolError(toolError)); - } - break; - } - case "finish-step": - finishReason = normalizeAssistantStepFinishReason(part.finishReason); - await providerActionBatch.flush(); - break; - case "error": - // `part.error` is typed as `unknown` — AI SDK providers emit - // whatever the upstream service threw. Coerce through `toError` - // so plain-object shapes (structured-clone survivors, typed - // gateway payloads) keep their `message`, `name`, `stack`, and - // `cause` instead of degrading to `new Error("[object Object]")`. - streamError = normalizeModelStreamError(part.error); - break; - case "abort": - // The SDK does not resolve step results for aborted in-flight steps. - throw new DOMException(part.reason ?? "The model stream was aborted.", "AbortError"); - default: - break; - } - } - - await providerActionBatch.flush(); - - if (streamError !== undefined) { - throw streamError; - } - - // Flush remaining reasoning. - if (currentReasoning.trim().length > 0) { - await emitFn( - createReasoningCompletedEvent({ - reasoning: currentReasoning, - sequence: state.sequence, - stepIndex: state.stepIndex, - turnId: state.turnId, - }), - ); - } - - // Channel adapters deliver terminal completions, so the reserved marker - // becomes a null completion without delaying normal streaming deltas. - if (finishReason !== "tool-calls" && hasEmptyDeliverySentinel(currentMessage)) { - await emitFn( - createMessageCompletedEvent({ - finishReason, - message: null, - sequence: state.sequence, - stepIndex: state.stepIndex, - turnId: state.turnId, - }), - ); - } else if (currentMessage.trim().length > 0) { - await emitFn( - createMessageCompletedEvent({ - finishReason, - message: currentMessage, - sequence: state.sequence, - stepIndex: state.stepIndex, - turnId: state.turnId, - }), - ); - } - - return { - emittedActionCallIds, - handledInlineToolResultCallIds, - invalidInputToolCallIds, - inlineAuthorizationResults, - trailingInlineToolResultParts, - }; -} diff --git a/packages/eve/src/harness/ordered-stream-emitter.ts b/packages/eve/src/harness/ordered-stream-emitter.ts index 86c81cee1f..b022f9afe5 100644 --- a/packages/eve/src/harness/ordered-stream-emitter.ts +++ b/packages/eve/src/harness/ordered-stream-emitter.ts @@ -6,6 +6,11 @@ import type { } from "#protocol/message.js"; import type { HarnessEmitFn } from "#harness/types.js"; +type AppendStreamEvent = + | ActionInputAppendedStreamEvent + | MessageAppendedStreamEvent + | ReasoningAppendedStreamEvent; + const MAX_PENDING_EVENTS = 64; const MAX_PENDING_DELTA_CHARACTERS = 64 * 1024; @@ -168,30 +173,20 @@ function mergeAdjacentEmissions( right: UnstampedMessageStreamEvent, messages: readonly import("ai").ModelMessage[] | undefined, ): boolean { - if (left.event.type === "message.appended" && right.type === "message.appended") { - if (!sameCoordinates(left.event, right)) return false; - left.deltaParts ??= [left.event.data.messageDelta]; - left.deltaParts.push(right.data.messageDelta); - left.event = right; - left.messages = messages; - return true; - } - - if (left.event.type === "reasoning.appended" && right.type === "reasoning.appended") { - if (!sameCoordinates(left.event, right)) return false; - left.deltaParts ??= [left.event.data.reasoningDelta]; - left.deltaParts.push(right.data.reasoningDelta); - left.event = right; - left.messages = messages; - return true; - } - - if (left.event.type === "action.input.appended" && right.type === "action.input.appended") { - if (left.event.data.callId !== right.data.callId || !sameCoordinates(left.event, right)) { + const leftAppendKey = appendKey(left.event); + const rightAppendKey = appendKey(right); + if (leftAppendKey !== undefined || rightAppendKey !== undefined) { + if ( + leftAppendKey === undefined || + leftAppendKey !== rightAppendKey || + !isAppendEvent(left.event) || + !isAppendEvent(right) || + !sameCoordinates(left.event, right) + ) { return false; } - left.deltaParts ??= [left.event.data.inputTextDelta]; - left.deltaParts.push(right.data.inputTextDelta); + left.deltaParts ??= [appendDelta(left.event)]; + left.deltaParts.push(appendDelta(right)); left.event = right; left.messages = messages; return true; @@ -207,11 +202,35 @@ function mergeAdjacentEmissions( return false; } +function appendKey(event: UnstampedMessageStreamEvent): string | undefined { + switch (event.type) { + case "message.appended": + case "reasoning.appended": + return event.type; + case "action.input.appended": + return `${event.type}:${event.data.callId}`; + default: + return undefined; + } +} + +function isAppendEvent(event: UnstampedMessageStreamEvent): event is AppendStreamEvent { + return appendKey(event) !== undefined; +} + +function appendDelta(event: AppendStreamEvent): string; +function appendDelta(event: UnstampedMessageStreamEvent): string | undefined; function appendDelta(event: UnstampedMessageStreamEvent): string | undefined { - if (event.type === "message.appended") return event.data.messageDelta; - if (event.type === "reasoning.appended") return event.data.reasoningDelta; - if (event.type === "action.input.appended") return event.data.inputTextDelta; - return undefined; + switch (event.type) { + case "message.appended": + return event.data.messageDelta; + case "reasoning.appended": + return event.data.reasoningDelta; + case "action.input.appended": + return event.data.inputTextDelta; + default: + return undefined; + } } function materializeEvent(emission: PendingEmission): UnstampedMessageStreamEvent { @@ -250,10 +269,7 @@ function materializeEvent(emission: PendingEmission): UnstampedMessageStreamEven return emission.event; } -function sameCoordinates( - left: ActionInputAppendedStreamEvent | MessageAppendedStreamEvent | ReasoningAppendedStreamEvent, - right: ActionInputAppendedStreamEvent | MessageAppendedStreamEvent | ReasoningAppendedStreamEvent, -): boolean { +function sameCoordinates(left: AppendStreamEvent, right: AppendStreamEvent): boolean { return ( left.data.sequence === right.data.sequence && left.data.stepIndex === right.data.stepIndex && diff --git a/packages/eve/src/harness/stream-content-emission.ts b/packages/eve/src/harness/stream-content-emission.ts new file mode 100644 index 0000000000..89f0ea5abe --- /dev/null +++ b/packages/eve/src/harness/stream-content-emission.ts @@ -0,0 +1,495 @@ +import type { + ModelMessage, + TextStreamPart, + ToolSet, + TypedToolCall, + TypedToolError, + TypedToolResult, +} from "ai"; + +type ToolResponsePart = Extract["content"][number]; +type InlineToolResultPart = Extract; + +import type { AssistantStepFinishReason } from "#protocol/message.js"; +import { + createActionsRequestedEvent, + createActionInputAppendedEvent, + createActionPartialEvent, + createActionResultEvent, + createMessageAppendedEvent, + createMessageCompletedEvent, + createReasoningAppendedEvent, + createReasoningCompletedEvent, +} from "#protocol/message.js"; +import { hasEmptyDeliverySentinel } from "#shared/empty-delivery.js"; +import { + createRuntimeToolResultFromStepResult, + createRuntimeToolResultFromToolError, + createToolResultMessagePartFromToolError, +} from "#harness/action-result-helpers.js"; +import { createRuntimeActionRequestFromToolCall } from "#harness/runtime-actions.js"; +import { + createInvalidToolCallInputError, + isInvalidToolCall, + resolveProviderToolCallRequest, +} from "#harness/tool-call-input-errors.js"; +import type { + RuntimeActionRequest, + RuntimeToolResultActionResult, +} from "#shared/action-types.js"; +import { createProviderStreamActionBatch } from "#harness/stream-actions.js"; +import { normalizeModelStreamError } from "#harness/model-call-error.js"; +import { createOrderedStreamEmitter } from "#harness/ordered-stream-emitter.js"; +import { interruptStreamOnFailure } from "#harness/interruptible-stream.js"; +import { isInlineAuthorizationToolResult } from "#harness/inline-tool-authorization.js"; +import type { HarnessEmissionState } from "#harness/emission-state.js"; +import type { HarnessEmitFn, HarnessToolMap } from "#harness/types.js"; + +/** + * Maps an AI SDK finish reason string to the eve-owned + * {@link AssistantStepFinishReason} union. Unknown values become `"other"`. + */ +export function normalizeAssistantStepFinishReason( + value: string | undefined, +): AssistantStepFinishReason { + switch (value) { + case "content-filter": + case "error": + case "length": + case "stop": + case "tool-calls": + return value; + default: + return "other"; + } +} + +/** + * Result of consuming one step's `fullStream`. + * + * Inline results avoid duplicate post-step events. Approval-resume + * authorization results also route back to the park detector. + */ +interface EmittedStreamContent { + readonly emittedActionCallIds: ReadonlySet; + readonly handledInlineToolResultCallIds: ReadonlySet; + readonly invalidInputToolCallIds: ReadonlySet; + readonly inlineAuthorizationResults: readonly TypedToolResult[]; + readonly trailingInlineToolResultParts: readonly InlineToolResultPart[]; +} + +interface StreamActionEmissionOptions { + readonly excludedActionToolNames: ReadonlySet; + readonly tools: HarnessToolMap; +} + +/** Consumes `fullStream` in source order, batching provider calls before their first result. */ +export async function emitStreamContent( + emitFn: HarnessEmitFn, + state: HarnessEmissionState, + fullStream: AsyncIterable>, + options?: StreamActionEmissionOptions, +): Promise { + const orderedEmitter = createOrderedStreamEmitter(emitFn); + const providerActionBatch = createProviderStreamActionBatch({ + emitFn: orderedEmitter.emit, + state, + }); + try { + return await consumeStreamContent( + orderedEmitter.emit, + state, + interruptStreamOnFailure(fullStream, orderedEmitter.failureSignal), + providerActionBatch, + options, + ); + } finally { + try { + await providerActionBatch.cancel(); + } finally { + await orderedEmitter.closeAndDrain(); + } + } +} + +async function consumeStreamContent( + emitFn: HarnessEmitFn, + state: HarnessEmissionState, + fullStream: AsyncIterable>, + providerActionBatch: ReturnType, + options?: StreamActionEmissionOptions, +): Promise { + let currentReasoning = ""; + let currentMessage = ""; + let finishReason: AssistantStepFinishReason = "stop"; + let streamError: Error | undefined; + const toolCallIdsSeenInStream = new Set(); + const emittedActionCallIds = new Set(); + const emittedActionResultCallIds = new Set(); + const providerToolCallIdsSeen = new Set(); + const handledInlineToolResultCallIds = new Set(); + const invalidInputToolCallIds = new Set(); + const inlineAuthorizationResults: TypedToolResult[] = []; + const trailingInlineToolResultParts: InlineToolResultPart[] = []; + const streamingActionInputs = new Map(); + + const flushCurrentMessage = async (): Promise => { + if (currentMessage.length === 0) { + return; + } + await emitFn( + createMessageCompletedEvent({ + finishReason: "tool-calls", + message: currentMessage, + sequence: state.sequence, + stepIndex: state.stepIndex, + turnId: state.turnId, + }), + ); + currentMessage = ""; + }; + + const emitActionInput = async ( + callId: string, + toolName: string, + inputTextDelta: string, + inputTextSoFar: string, + ): Promise => + emitFn( + createActionInputAppendedEvent({ + callId, + inputTextDelta, + inputTextSoFar, + sequence: state.sequence, + stepIndex: state.stepIndex, + toolName, + turnId: state.turnId, + }), + ); + + const emitActionRequest = async (action: RuntimeActionRequest): Promise => { + if (emittedActionCallIds.has(action.callId)) { + return; + } + + if (currentMessage.trim().length > 0) { + await flushCurrentMessage(); + } + + emittedActionCallIds.add(action.callId); + await emitFn( + createActionsRequestedEvent({ + actions: [action], + sequence: state.sequence, + stepIndex: state.stepIndex, + turnId: state.turnId, + }), + ); + }; + + const collectProviderToolCall = async (toolCall: { + readonly input?: unknown; + readonly toolCallId: string; + readonly toolName: string; + }): Promise => { + if (providerToolCallIdsSeen.has(toolCall.toolCallId)) { + return; + } + providerToolCallIdsSeen.add(toolCall.toolCallId); + if (emittedActionCallIds.has(toolCall.toolCallId)) { + return; + } + emittedActionCallIds.add(toolCall.toolCallId); + + if (currentMessage.trim().length > 0) { + await flushCurrentMessage(); + } + + const resolved = resolveProviderToolCallRequest(toolCall); + if (resolved.toolError !== undefined) { + invalidInputToolCallIds.add(toolCall.toolCallId); + await emitActionResult(createRuntimeToolResultFromToolError(resolved.toolError)); + handledInlineToolResultCallIds.add(toolCall.toolCallId); + trailingInlineToolResultParts.push( + createToolResultMessagePartFromToolError(resolved.toolError), + ); + return; + } + + providerActionBatch.observe(resolved.request); + }; + + const emitActionResult = async (result: RuntimeToolResultActionResult): Promise => { + if (emittedActionResultCallIds.has(result.callId)) { + return; + } + emittedActionResultCallIds.add(result.callId); + await emitFn( + createActionResultEvent({ + result, + sequence: state.sequence, + stepIndex: state.stepIndex, + turnId: state.turnId, + }), + ); + }; + + const emitActionPartial = async (result: RuntimeToolResultActionResult): Promise => { + await emitFn( + createActionPartialEvent({ + result, + sequence: state.sequence, + stepIndex: state.stepIndex, + turnId: state.turnId, + }), + ); + }; + + const emitToolCall = async (toolCall: TypedToolCall): Promise => { + if (isInvalidToolCall(toolCall)) { + invalidInputToolCallIds.add(toolCall.toolCallId); + return; + } + if (options === undefined || options.excludedActionToolNames.has(toolCall.toolName)) { + return; + } + + try { + await emitActionRequest( + createRuntimeActionRequestFromToolCall({ + toolCall, + tools: options.tools, + }), + ); + } catch (error) { + if (error instanceof TypeError) { + const toolError = createInvalidToolCallInputError({ error, toolCall }); + invalidInputToolCallIds.add(toolCall.toolCallId); + if (currentMessage.trim().length > 0) { + await flushCurrentMessage(); + } + await emitActionResult(createRuntimeToolResultFromToolError(toolError)); + handledInlineToolResultCallIds.add(toolCall.toolCallId); + trailingInlineToolResultParts.push(createToolResultMessagePartFromToolError(toolError)); + return; + } + throw error; + } + }; + + for await (const part of fullStream) { + if (streamError !== undefined) { + continue; + } + + switch (part.type) { + case "reasoning-delta": + await providerActionBatch.flush(); + currentReasoning += part.text; + await emitFn( + createReasoningAppendedEvent({ + reasoningDelta: part.text, + reasoningSoFar: currentReasoning, + sequence: state.sequence, + stepIndex: state.stepIndex, + turnId: state.turnId, + }), + ); + break; + case "text-delta": + await providerActionBatch.flush(); + // Flush accumulated reasoning before text begins. + if (currentReasoning.trim().length > 0) { + await emitFn( + createReasoningCompletedEvent({ + reasoning: currentReasoning, + sequence: state.sequence, + stepIndex: state.stepIndex, + turnId: state.turnId, + }), + ); + currentReasoning = ""; + } + currentMessage += part.text; + await emitFn( + createMessageAppendedEvent({ + messageDelta: part.text, + messageSoFar: currentMessage, + sequence: state.sequence, + stepIndex: state.stepIndex, + turnId: state.turnId, + }), + ); + break; + case "tool-input-start": { + if ( + options === undefined || + part.providerExecuted === true || + options.excludedActionToolNames.has(part.toolName) + ) { + streamingActionInputs.delete(part.id); + break; + } + await providerActionBatch.flush(); + if (currentMessage.trim().length > 0) { + await flushCurrentMessage(); + } + streamingActionInputs.set(part.id, { text: "", toolName: part.toolName }); + await emitActionInput(part.id, part.toolName, "", ""); + break; + } + case "tool-input-delta": { + const input = streamingActionInputs.get(part.id); + if (input === undefined) { + break; + } + await providerActionBatch.flush(); + input.text += part.delta; + await emitActionInput(part.id, input.toolName, part.delta, input.text); + break; + } + case "tool-input-end": + streamingActionInputs.delete(part.id); + break; + case "tool-call": { + const toolCall = part as TypedToolCall; + streamingActionInputs.delete(toolCall.toolCallId); + toolCallIdsSeenInStream.add(toolCall.toolCallId); + if (toolCall.providerExecuted === true) { + await collectProviderToolCall(toolCall); + } else { + await providerActionBatch.flush(); + await emitToolCall(toolCall); + } + break; + } + case "tool-result": { + const inlineToolResult = part as TypedToolResult; + if (inlineToolResult.preliminary === true) { + if (inlineToolResult.providerExecuted !== true) { + await emitActionPartial(createRuntimeToolResultFromStepResult(inlineToolResult)); + } + break; + } + if (inlineToolResult.providerExecuted === true) { + await collectProviderToolCall({ + input: "input" in inlineToolResult ? inlineToolResult.input : undefined, + toolCallId: inlineToolResult.toolCallId, + toolName: inlineToolResult.toolName, + }); + await providerActionBatch.flush(); + await emitActionResult(createRuntimeToolResultFromStepResult(inlineToolResult)); + // Provider results already live in the assistant response. Do not + // add a local tool message. + break; + } + + if (toolCallIdsSeenInStream.has(part.toolCallId)) { + if (isInlineAuthorizationToolResult(inlineToolResult)) { + break; + } + if (emittedActionCallIds.has(part.toolCallId)) { + await emitActionResult(createRuntimeToolResultFromStepResult(inlineToolResult)); + handledInlineToolResultCallIds.add(part.toolCallId); + } + break; + } + + // An approved tool can resume with its result but no matching call in + // this step. Emit it before the message that consumes it. + await providerActionBatch.flush(); + await flushCurrentMessage(); + if (isInlineAuthorizationToolResult(inlineToolResult)) { + // Keep authorization output for the park detector instead of + // emitting a normal tool result. + handledInlineToolResultCallIds.add(part.toolCallId); + inlineAuthorizationResults.push(inlineToolResult); + break; + } + await emitActionResult(createRuntimeToolResultFromStepResult(inlineToolResult)); + handledInlineToolResultCallIds.add(part.toolCallId); + break; + } + case "tool-error": { + const toolError = part as TypedToolError; + if (toolError.providerExecuted === true) { + await collectProviderToolCall(toolError); + await providerActionBatch.flush(); + await emitActionResult(createRuntimeToolResultFromToolError(toolError)); + } else if (emittedActionCallIds.has(toolError.toolCallId)) { + await emitActionResult(createRuntimeToolResultFromToolError(toolError)); + handledInlineToolResultCallIds.add(toolError.toolCallId); + trailingInlineToolResultParts.push(createToolResultMessagePartFromToolError(toolError)); + } + break; + } + case "finish-step": + finishReason = normalizeAssistantStepFinishReason(part.finishReason); + await providerActionBatch.flush(); + break; + case "error": + // `part.error` is typed as `unknown` — AI SDK providers emit + // whatever the upstream service threw. Coerce through `toError` + // so plain-object shapes (structured-clone survivors, typed + // gateway payloads) keep their `message`, `name`, `stack`, and + // `cause` instead of degrading to `new Error("[object Object]")`. + streamError = normalizeModelStreamError(part.error); + break; + case "abort": + // The SDK does not resolve step results for aborted in-flight steps. + throw new DOMException(part.reason ?? "The model stream was aborted.", "AbortError"); + default: + break; + } + } + + await providerActionBatch.flush(); + + if (streamError !== undefined) { + throw streamError; + } + + // Flush remaining reasoning. + if (currentReasoning.trim().length > 0) { + await emitFn( + createReasoningCompletedEvent({ + reasoning: currentReasoning, + sequence: state.sequence, + stepIndex: state.stepIndex, + turnId: state.turnId, + }), + ); + } + + // Channel adapters deliver terminal completions, so the reserved marker + // becomes a null completion without delaying normal streaming deltas. + if (finishReason !== "tool-calls" && hasEmptyDeliverySentinel(currentMessage)) { + await emitFn( + createMessageCompletedEvent({ + finishReason, + message: null, + sequence: state.sequence, + stepIndex: state.stepIndex, + turnId: state.turnId, + }), + ); + } else if (currentMessage.trim().length > 0) { + await emitFn( + createMessageCompletedEvent({ + finishReason, + message: currentMessage, + sequence: state.sequence, + stepIndex: state.stepIndex, + turnId: state.turnId, + }), + ); + } + + return { + emittedActionCallIds, + handledInlineToolResultCallIds, + invalidInputToolCallIds, + inlineAuthorizationResults, + trailingInlineToolResultParts, + }; +} From 0cce862eb6c8c0a0d5e4c0611fd58906896f2d25 Mon Sep 17 00:00:00 2001 From: "vercel-gh-bot-4[bot]" <312518292+vercel-gh-bot-4[bot]@users.noreply.github.com> Date: Sat, 22 Aug 2026 00:13:38 +0000 Subject: [PATCH 3/5] refactor(eve): remove reducer state extraction Co-authored-by: ruiconti <1834568+ruiconti@users.noreply.github.com> Signed-off-by: Rui Conti --- .../src/client/message-reducer-primitives.ts | 71 +++++ .../eve/src/client/message-reducer-state.ts | 286 ------------------ packages/eve/src/client/message-reducer.ts | 246 +++++++++++++-- 3 files changed, 298 insertions(+), 305 deletions(-) create mode 100644 packages/eve/src/client/message-reducer-primitives.ts delete mode 100644 packages/eve/src/client/message-reducer-state.ts diff --git a/packages/eve/src/client/message-reducer-primitives.ts b/packages/eve/src/client/message-reducer-primitives.ts new file mode 100644 index 0000000000..d221bb8f7a --- /dev/null +++ b/packages/eve/src/client/message-reducer-primitives.ts @@ -0,0 +1,71 @@ +import type { EveMessage, EveMessageData, EveMessagePart } from "#client/message-reducer-types.js"; +import type { MessageReceivedPart } from "#protocol/message.js"; + +export function projectReceivedParts( + parts: readonly MessageReceivedPart[] | undefined, + message: string, +): readonly EveMessagePart[] { + return ( + parts?.map((part) => + part.type === "text" + ? { state: "done", text: part.text, type: "text" } + : { + filename: part.filename, + mediaType: part.mediaType, + size: part.size, + type: "file", + url: part.url, + }, + ) ?? [{ state: "done", text: message, type: "text" }] + ); +} + +export function partKey(part: EveMessagePart): string { + switch (part.type) { + case "text": + return `text:${part.stepIndex ?? 0}`; + case "reasoning": + return `reasoning:${part.stepIndex ?? 0}`; + case "file": + return `file:${part.stepIndex ?? 0}:${part.filename ?? part.url ?? part.mediaType}`; + case "step-start": + return "step-start"; + case "authorization": + return `authorization:${part.turnId}:${part.stepIndex}:${part.name}`; + case "dynamic-tool": + return `dynamic-tool:${part.toolCallId}`; + } +} + +export function upsertMessage(data: EveMessageData, next: EveMessage): EveMessageData { + const index = data.messages.findIndex((message) => message.id === next.id); + if (index === -1) { + return { messages: [...data.messages, next] }; + } + + return { + messages: [...data.messages.slice(0, index), next, ...data.messages.slice(index + 1)], + }; +} + +export function removeStreamingToolPartsForTurn( + data: EveMessageData, + turnId: string, +): EveMessageData { + const index = data.messages.findIndex( + (message) => message.role === "assistant" && message.metadata?.turnId === turnId, + ); + const message = data.messages[index]; + if (message === undefined) return data; + + return upsertMessage(data, { + ...message, + parts: message.parts.filter( + (part) => part.type !== "dynamic-tool" || part.state !== "input-streaming", + ), + }); +} + +export function optimisticUserMessageId(submissionId: string): string { + return `optimistic:${submissionId}:user`; +} diff --git a/packages/eve/src/client/message-reducer-state.ts b/packages/eve/src/client/message-reducer-state.ts deleted file mode 100644 index ebba6cd9e1..0000000000 --- a/packages/eve/src/client/message-reducer-state.ts +++ /dev/null @@ -1,286 +0,0 @@ -import type { - EveAuthorizationPart, - EveDynamicToolPart, - EveMessage, - EveMessageData, - EveMessageMetadata, - EveMessagePart, -} from "#client/message-reducer-types.js"; -import type { MessageReceivedPart } from "#protocol/message.js"; - -export type EveAssistantMessage = EveMessage & { readonly role: "assistant" }; - -export function projectReceivedParts( - parts: readonly MessageReceivedPart[] | undefined, - message: string, -): readonly EveMessagePart[] { - return ( - parts?.map((part) => - "text" in part - ? { state: "done", text: part.text, type: "text" } - : { - filename: part.filename, - mediaType: part.mediaType, - size: part.size, - type: "file", - url: part.url, - }, - ) ?? [{ state: "done", text: message, type: "text" }] - ); -} - -function partKey(part: EveMessagePart): string { - switch (part.type) { - case "text": - return `text:${part.stepIndex ?? 0}`; - case "reasoning": - return `reasoning:${part.stepIndex ?? 0}`; - case "file": - return `file:${part.stepIndex ?? 0}:${part.filename ?? part.url ?? part.mediaType}`; - case "step-start": - return "step-start"; - case "authorization": - return `authorization:${part.turnId}:${part.stepIndex}:${part.name}`; - case "dynamic-tool": - return `dynamic-tool:${part.toolCallId}`; - } -} - -export function upsertMessage(data: EveMessageData, next: EveMessage): EveMessageData { - const index = data.messages.findIndex((message) => message.id === next.id); - if (index === -1) { - return { messages: [...data.messages, next] }; - } - if (data.messages[index] === next) return data; - - return { - messages: [...data.messages.slice(0, index), next, ...data.messages.slice(index + 1)], - }; -} - -export function optimisticUserMessageId(submissionId: string): string { - return `optimistic:${submissionId}:user`; -} - -export function updateAssistantMessage( - data: EveMessageData, - turnId: string, - update: (message: EveAssistantMessage) => EveAssistantMessage, -): EveMessageData { - const message = findAssistantMessage(data, turnId) ?? createAssistantMessage(turnId); - return upsertMessage(data, update(message)); -} - -export function updateExistingAssistantMessage( - data: EveMessageData, - turnId: string, - update: (message: EveAssistantMessage) => EveAssistantMessage, -): EveMessageData { - const message = findAssistantMessage(data, turnId); - return message === undefined ? data : upsertMessage(data, update(message)); -} - -export function updateAssistantMetadata( - data: EveMessageData, - turnId: string, - metadata: EveMessageMetadata, -): EveMessageData { - return updateAssistantMessage(data, turnId, (message) => ({ - ...message, - metadata: { - ...message.metadata, - ...metadata, - }, - })); -} - -export function ensureStepStartPart( - message: EveAssistantMessage, - stepIndex: number, -): EveAssistantMessage { - const stepStartCount = message.parts.filter((part) => part.type === "step-start").length; - if (stepStartCount > stepIndex) return message; - - const missingCount = stepIndex - stepStartCount + 1; - return { - ...message, - parts: [ - ...message.parts, - ...Array.from({ length: missingCount }, () => ({ type: "step-start" as const })), - ], - }; -} - -export function upsertPart( - message: EveAssistantMessage, - next: EveMessagePart, -): EveAssistantMessage { - const index = message.parts.findIndex((part) => partKey(part) === partKey(next)); - const parts = - index === -1 - ? [...message.parts, next] - : [...message.parts.slice(0, index), next, ...message.parts.slice(index + 1)]; - - return { - ...message, - metadata: { - ...message.metadata, - status: next.type === "text" && next.state === "done" ? "complete" : "streaming", - }, - parts, - }; -} - -type EveRunPart = Extract; - -// Upserts a text/reasoning part, keeping multiple runs per step distinct: one -// step can produce text, call tools, then produce more text (see -// `MessageCompletedStreamEvent`), so a step-only key would collapse them. -// -// We find the latest same-step run of this type: while it is still streaming, -// its snapshots replace it in place; once it is done (or there is none), `next` -// begins a new run appended in arrival order. -export function upsertRun(message: EveAssistantMessage, next: EveRunPart): EveAssistantMessage { - let lastIndex = -1; - for (let index = message.parts.length - 1; index >= 0; index -= 1) { - const part = message.parts[index]; - if (part?.type === next.type && part.stepIndex === next.stepIndex) { - lastIndex = index; - break; - } - } - - const openRun = - lastIndex !== -1 && (message.parts[lastIndex] as EveRunPart).state === "streaming"; - const parts = openRun - ? [...message.parts.slice(0, lastIndex), next, ...message.parts.slice(lastIndex + 1)] - : [...message.parts, next]; - - return { - ...message, - metadata: { - ...message.metadata, - status: next.type === "text" && next.state === "done" ? "complete" : "streaming", - }, - parts, - }; -} - -export function removeTextPart( - message: EveAssistantMessage, - stepIndex: number, -): EveAssistantMessage { - const parts = message.parts.filter( - (part) => part.type !== "text" || part.stepIndex !== stepIndex, - ); - if (parts.length === message.parts.length) return message; - - return { - ...message, - metadata: { - ...message.metadata, - status: "complete", - }, - parts, - }; -} - -export function removeStreamingToolParts( - parts: readonly EveMessagePart[], -): readonly EveMessagePart[] { - const next = parts.filter( - (part) => part.type !== "dynamic-tool" || part.state !== "input-streaming", - ); - return next.length === parts.length ? parts : next; -} - -export function updateToolPart( - data: EveMessageData, - toolCallId: string, - next: EveDynamicToolPart, -): EveMessageData { - const message = data.messages.find( - (candidate): candidate is EveAssistantMessage => - candidate.role === "assistant" && - candidate.parts.some( - (part) => part.type === "dynamic-tool" && part.toolCallId === toolCallId, - ), - ); - return message === undefined ? data : upsertMessage(data, upsertPart(message, next)); -} - -export function updateAuthorizationPart( - data: EveMessageData, - existing: EveAuthorizationPart, - next: EveAuthorizationPart, -): EveMessageData { - const message = data.messages.find( - (candidate): candidate is EveAssistantMessage => - candidate.role === "assistant" && candidate.parts.some((part) => part === existing), - ); - return message === undefined ? data : upsertMessage(data, upsertPart(message, next)); -} - -export function findToolPart( - data: EveMessageData, - toolCallId: string, -): EveDynamicToolPart | undefined { - for (const message of data.messages) { - for (const part of message.parts) { - if (part.type === "dynamic-tool" && part.toolCallId === toolCallId) return part; - } - } - return undefined; -} - -export function findLatestPendingAuthorizationPart( - data: EveMessageData, - name: string, -): EveAuthorizationPart | undefined { - for (let messageIndex = data.messages.length - 1; messageIndex >= 0; messageIndex -= 1) { - const message = data.messages[messageIndex]; - if (message?.role !== "assistant") continue; - - for (let partIndex = message.parts.length - 1; partIndex >= 0; partIndex -= 1) { - const part = message.parts[partIndex]; - if (part?.type === "authorization" && part.state === "required" && part.name === name) { - return part; - } - } - } - return undefined; -} - -export function findToolPartByApprovalId( - data: EveMessageData, - approvalId: string, -): EveDynamicToolPart | undefined { - for (const message of data.messages) { - for (const part of message.parts) { - if (part.type === "dynamic-tool" && part.approval?.id === approvalId) return part; - } - } - return undefined; -} - -function findAssistantMessage( - data: EveMessageData, - turnId: string, -): EveAssistantMessage | undefined { - return data.messages.find( - (message): message is EveAssistantMessage => - message.role === "assistant" && message.metadata?.turnId === turnId, - ); -} - -function createAssistantMessage(turnId: string): EveAssistantMessage { - return { - id: `${turnId}:assistant`, - metadata: { - status: "streaming", - turnId, - }, - parts: [], - role: "assistant", - }; -} diff --git a/packages/eve/src/client/message-reducer.ts b/packages/eve/src/client/message-reducer.ts index 0a8bb1c578..39c9691d95 100644 --- a/packages/eve/src/client/message-reducer.ts +++ b/packages/eve/src/client/message-reducer.ts @@ -3,7 +3,14 @@ import { createAuthorizationCompletedPart, createAuthorizationRequiredPart, } from "#client/authorization-message-parts.js"; -import type { EveDynamicToolPart, EveMessageData } from "#client/message-reducer-types.js"; +import type { + EveAuthorizationPart, + EveMessageData, + EveDynamicToolPart, + EveMessage, + EveMessageMetadata, + EveMessagePart, +} from "#client/message-reducer-types.js"; import { approvedApproval, createToolMetadata, @@ -14,23 +21,12 @@ import { toMessageInputRequest, } from "#client/message-action-parts.js"; import { - ensureStepStartPart, - findLatestPendingAuthorizationPart, - findToolPart, - findToolPartByApprovalId, optimisticUserMessageId, + partKey, projectReceivedParts, - removeStreamingToolParts, - removeTextPart, - updateAssistantMessage, - updateAssistantMetadata, - updateAuthorizationPart, - updateExistingAssistantMessage, - updateToolPart, + removeStreamingToolPartsForTurn, upsertMessage, - upsertPart, - upsertRun, -} from "#client/message-reducer-state.js"; +} from "#client/message-reducer-primitives.js"; import type { InputResponse } from "#shared/input.js"; import type { AuthorizationCompletedStreamEvent, InputResolution } from "#protocol/message.js"; @@ -47,6 +43,8 @@ export type { EveMessageToolMetadata, } from "#client/message-reducer-types.js"; +type EveAssistantMessage = EveMessage & { readonly role: "assistant" }; + /** * Creates a UIMessage-compatible eve reducer for chat and agent UIs. * @@ -399,10 +397,7 @@ function reduceMessageData(data: EveMessageData, event: EveAgentReducerEvent): E })); case "turn.failed": - return updateExistingAssistantMessage(data, event.data.turnId, (message) => { - const parts = removeStreamingToolParts(message.parts); - return parts === message.parts ? message : { ...message, parts }; - }); + return removeStreamingToolPartsForTurn(data, event.data.turnId); case "session.failed": return data; @@ -412,6 +407,10 @@ function reduceMessageData(data: EveMessageData, event: EveAgentReducerEvent): E } } +function removeStreamingToolParts(parts: readonly EveMessagePart[]): readonly EveMessagePart[] { + return parts.filter((part) => part.type !== "dynamic-tool" || part.state !== "input-streaming"); +} + function respondToInputRequest(data: EveMessageData, response: InputResponse): EveMessageData { const existing = findToolPartByApprovalId(data, response.requestId); if (!existing) return data; @@ -461,6 +460,152 @@ function resolveInputRequest(data: EveMessageData, resolution: InputResolution): }); } +function updateAssistantMessage( + data: EveMessageData, + turnId: string, + update: (message: EveAssistantMessage) => EveAssistantMessage, +): EveMessageData { + const existing = data.messages.find( + (message): message is EveAssistantMessage => + message.role === "assistant" && message.metadata?.turnId === turnId, + ); + + const message = existing ?? createAssistantMessage(turnId); + return upsertMessage(data, update(message)); +} + +function updateAssistantMetadata( + data: EveMessageData, + turnId: string, + metadata: EveMessageMetadata, +): EveMessageData { + return updateAssistantMessage(data, turnId, (message) => ({ + ...message, + metadata: { + ...message.metadata, + ...metadata, + }, + })); +} + +function createAssistantMessage(turnId: string): EveAssistantMessage { + return { + id: `${turnId}:assistant`, + metadata: { + status: "streaming", + turnId, + }, + parts: [], + role: "assistant", + }; +} + +function ensureStepStartPart(message: EveAssistantMessage, stepIndex: number): EveAssistantMessage { + const stepStartCount = message.parts.filter((part) => part.type === "step-start").length; + if (stepStartCount > stepIndex) { + return message; + } + + const missingCount = stepIndex - stepStartCount + 1; + return { + ...message, + parts: [ + ...message.parts, + ...Array.from({ length: missingCount }, () => ({ type: "step-start" as const })), + ], + }; +} + +function upsertPart(message: EveAssistantMessage, next: EveMessagePart): EveAssistantMessage { + const index = message.parts.findIndex((part) => partKey(part) === partKey(next)); + const parts = + index === -1 + ? [...message.parts, next] + : [...message.parts.slice(0, index), next, ...message.parts.slice(index + 1)]; + + return { + ...message, + metadata: { + ...message.metadata, + status: next.type === "text" && next.state === "done" ? "complete" : "streaming", + }, + parts, + }; +} + +type EveRunPart = Extract; + +// Upserts a text/reasoning part, keeping multiple runs per step distinct: one +// step can produce text, call tools, then produce more text (see +// `MessageCompletedStreamEvent`), so a step-only key would collapse them. +// +// We find the latest same-step run of this type: while it is still streaming, +// its snapshots replace it in place; once it is done (or there is none), `next` +// begins a new run appended in arrival order. +function upsertRun(message: EveAssistantMessage, next: EveRunPart): EveAssistantMessage { + let lastIndex = -1; + for (let index = message.parts.length - 1; index >= 0; index -= 1) { + const part = message.parts[index]; + if (part?.type === next.type && part.stepIndex === next.stepIndex) { + lastIndex = index; + break; + } + } + + const openRun = + lastIndex !== -1 && (message.parts[lastIndex] as EveRunPart).state === "streaming"; + const parts = openRun + ? [...message.parts.slice(0, lastIndex), next, ...message.parts.slice(lastIndex + 1)] + : [...message.parts, next]; + + return { + ...message, + metadata: { + ...message.metadata, + status: next.type === "text" && next.state === "done" ? "complete" : "streaming", + }, + parts, + }; +} + +function removeTextPart(message: EveAssistantMessage, stepIndex: number): EveAssistantMessage { + const parts = message.parts.filter( + (part) => part.type !== "text" || part.stepIndex !== stepIndex, + ); + if (parts.length === message.parts.length) { + return message; + } + + return { + ...message, + metadata: { + ...message.metadata, + status: "complete", + }, + parts, + }; +} + +function updateToolPart( + data: EveMessageData, + toolCallId: string, + next: EveDynamicToolPart, +): EveMessageData { + const message = data.messages.find( + (candidate): candidate is EveAssistantMessage => + candidate.role === "assistant" && + candidate.parts.some( + (part) => part.type === "dynamic-tool" && part.toolCallId === toolCallId, + ), + ); + + if (!message) { + return data; + } + + return upsertMessage(data, upsertPart(message, next)); +} + function completeAuthorization( data: EveMessageData, event: AuthorizationCompletedStreamEvent, @@ -477,6 +622,34 @@ function completeAuthorization( ); } +function updateAuthorizationPart( + data: EveMessageData, + existing: EveAuthorizationPart, + next: EveAuthorizationPart, +): EveMessageData { + const message = data.messages.find( + (candidate): candidate is EveAssistantMessage => + candidate.role === "assistant" && candidate.parts.some((part) => part === existing), + ); + + if (!message) { + return data; + } + + return upsertMessage(data, upsertPart(message, next)); +} + +function findToolPart(data: EveMessageData, toolCallId: string): EveDynamicToolPart | undefined { + for (const message of data.messages) { + for (const part of message.parts) { + if (part.type === "dynamic-tool" && part.toolCallId === toolCallId) { + return part; + } + } + } + return undefined; +} + function isSettledToolPart(part: EveDynamicToolPart): boolean { return ( part.state === "output-denied" || @@ -484,3 +657,38 @@ function isSettledToolPart(part: EveDynamicToolPart): boolean { (part.state === "output-available" && part.partial !== true) ); } + +function findLatestPendingAuthorizationPart( + data: EveMessageData, + name: string, +): EveAuthorizationPart | undefined { + for (let messageIndex = data.messages.length - 1; messageIndex >= 0; messageIndex -= 1) { + const message = data.messages[messageIndex]; + if (message?.role !== "assistant") { + continue; + } + + for (let partIndex = message.parts.length - 1; partIndex >= 0; partIndex -= 1) { + const part = message.parts[partIndex]; + if (part?.type === "authorization" && part.state === "required" && part.name === name) { + return part; + } + } + } + + return undefined; +} + +function findToolPartByApprovalId( + data: EveMessageData, + approvalId: string, +): EveDynamicToolPart | undefined { + for (const message of data.messages) { + for (const part of message.parts) { + if (part.type === "dynamic-tool" && part.approval?.id === approvalId) { + return part; + } + } + } + return undefined; +} From b6585a3c764bdce6aca7131fccad5b73933484a0 Mon Sep 17 00:00:00 2001 From: Rui Conti Date: Wed, 26 Aug 2026 23:57:56 -0400 Subject: [PATCH 4/5] fix(eve): store streamed tool input deltas Signed-off-by: Rui Conti --- .changeset/stream-tool-input.md | 2 +- docs/concepts/sessions-runs-and-streaming.md | 6 ++-- docs/guides/client/streaming.mdx | 6 ++-- .../reports/channel/v10.json | 2 +- .../reports/dynamicInstructions/v14.json | 2 +- .../reports/dynamicSkill/v13.json | 2 +- .../reports/dynamicTool/v20.json | 2 +- .../extension-contracts/reports/hook/v15.json | 2 +- .../reports/schedule/v4.json | 2 +- .../reports/subagent/v5.json | 2 +- .../src/client/message-reducer-primitives.ts | 10 ++++++ .../eve/src/client/message-reducer-types.ts | 1 + .../eve/src/client/message-reducer.test.ts | 22 +++++++++--- packages/eve/src/client/message-reducer.ts | 14 +++++--- packages/eve/src/harness/emission.test.ts | 36 +++++++++++++------ .../harness/ordered-stream-emitter.test.ts | 16 ++++----- .../eve/src/harness/ordered-stream-emitter.ts | 11 +++++- .../src/harness/stream-content-emission.ts | 20 +++++------ packages/eve/src/protocol/message.ts | 7 ++-- 19 files changed, 109 insertions(+), 56 deletions(-) diff --git a/.changeset/stream-tool-input.md b/.changeset/stream-tool-input.md index 7969307afe..8e792f4eac 100644 --- a/.changeset/stream-tool-input.md +++ b/.changeset/stream-tool-input.md @@ -2,4 +2,4 @@ "eve": patch --- -Tool inputs now stream through the durable event protocol as `action.input.appended` before the matching validated `actions.requested` event. The default message reducer exposes the cumulative raw input on `dynamic-tool.inputText` while its state is `input-streaming`, so UIs can progressively render long JSON inputs. This advances the stream protocol to version 24; when assistant text precedes a tool call, `message.completed` now arrives before that call's streamed input events. +Tool inputs now stream through the durable event protocol as `action.input.appended` before the matching validated `actions.requested` event. Each event stores only its raw delta and UTF-16 offset, while the default message reducer exposes cumulative raw input on `dynamic-tool.inputText` in the `input-streaming` state. This advances the stream protocol to version 24; when assistant text precedes a tool call, `message.completed` now arrives before that call's streamed input events. diff --git a/docs/concepts/sessions-runs-and-streaming.md b/docs/concepts/sessions-runs-and-streaming.md index 1fd74a8cff..0249b64579 100644 --- a/docs/concepts/sessions-runs-and-streaming.md +++ b/docs/concepts/sessions-runs-and-streaming.md @@ -49,6 +49,7 @@ The stream is newline-delimited JSON (NDJSON), one event per line: | `turn.started` | A new turn began; carries the active `trace` when the runtime is traced. | | `message.received` | An inbound user message was accepted; carries flattened text plus structured text/file parts. | | `step.started` | A model step began. | +| `action.input.appended` | A raw tool-input text delta, its character offset, and tool-call identity. | | `actions.requested` | The model requested one or more actions, including tool calls; calls stream before execution. | | `action.partial` | A locally executed tool generator yielded a preliminary output snapshot. | | `action.result` | A tool call returned. | @@ -60,7 +61,6 @@ The stream is newline-delimited JSON (NDJSON), one event per line: | `reasoning.completed` | The finalized reasoning block. | | `message.appended` | An assistant text delta (incremental, with cumulative text so far). | | `message.completed` | A finalized assistant text block. | -| `action.input.appended` | A tool-input text delta, with the cumulative raw input and tool-call identity. | | `result.completed` | The finalized structured result for a turn that requested an output schema; carries `result`. | | `compaction.requested` | Context-window compaction began; carries `modelId`, `sessionId`, `turnId`, `usageInputTokens`. | | `compaction.completed` | A compaction checkpoint was written to durable history. | @@ -77,9 +77,9 @@ The stream is newline-delimited JSON (NDJSON), one event per line: The optional `data.trace` on session and turn starts contains eve-owned W3C trace coordinates: `traceId`, `spanId`, and `traceFlags`. Use it to correlate stream consumers such as eval reporters with an observability backend. An uninstrumented target omits it. -`reasoning.appended`, `message.appended`, and `action.input.appended` stream incremental output as it arrives. When the durable stream writer is busy, eve may coalesce adjacent deltas of the same type; the text remains in source order, and any other event forms an ordering barrier. Each append carries both the new delta and the cumulative text for the current block. The finalized text and reasoning blocks show up on `message.completed` and `reasoning.completed`, which is the compatibility path for clients that don't render incremental streaming. +`reasoning.appended`, `message.appended`, and `action.input.appended` stream incremental output as it arrives. When the durable stream writer is busy, eve may coalesce adjacent deltas for the same text block or tool call; the text remains in source order, and a different event type, tool `callId`, or stream coordinate forms an ordering barrier. Text and reasoning appends carry both the new delta and the cumulative text for the current block. The finalized blocks show up on `message.completed` and `reasoning.completed`, which is the compatibility path for clients that don't render incremental streaming. -`action.input.appended` arrives before the matching `actions.requested` event. It carries `callId`, `toolName`, `inputTextDelta`, and `inputTextSoFar`; the input text may be incomplete JSON. The default client reducer projects it as a `dynamic-tool` part with `state: "input-streaming"` and the cumulative text in `inputText`. `actions.requested` replaces that part with `state: "input-available"` and the validated `input`. Excluded internal actions never publish their input stream. +When a streamed tool input becomes a validated call, its `action.input.appended` events precede the matching `actions.requested` event. Each append carries `callId`, `toolName`, `inputTextDelta`, and `inputTextOffset`; the offset is the zero-based UTF-16 code-unit position where the delta begins. Storing only the delta and offset avoids repeating the cumulative input in every durable event. The default client reducer starts or restarts accumulation at offset `0`, ignores a nonzero offset that is not contiguous, and projects the potentially incomplete JSON as a `dynamic-tool` part with `state: "input-streaming"` and cumulative text in `inputText`. `actions.requested` replaces that part with `state: "input-available"` and the validated `input`. Excluded internal actions never publish their input stream. `action.partial` carries one complete preliminary output snapshot from an authored async-generator tool. A later partial for the same `callId` replaces it, and `action.result` is the final snapshot. When the durable writer is busy, eve may keep only the newest adjacent partial for a call. Treat partials as last-write-wins: a durable step can retry and replay overlapping event runs. Provider-executed tool progress and MCP progress notifications are not projected as `action.partial` events. diff --git a/docs/guides/client/streaming.mdx b/docs/guides/client/streaming.mdx index fb4839927a..edcd8623e1 100644 --- a/docs/guides/client/streaming.mdx +++ b/docs/guides/client/streaming.mdx @@ -75,7 +75,7 @@ for await (const event of response) { } ``` -`message.appended`, `reasoning.appended`, and `action.input.appended` are incremental delta events. eve may combine adjacent deltas of the same type while a durable stream write is in flight, but preserves their text and event ordering; any different event is a barrier. The completed text forms, `message.completed` and `reasoning.completed`, are the compatibility path for clients that don't render deltas. A streamed tool input is complete when the matching validated call arrives in `actions.requested`. +`message.appended`, `reasoning.appended`, and `action.input.appended` are incremental delta events. eve may combine adjacent deltas for the same text block or tool call while a durable stream write is in flight, but preserves their text and event ordering. A different event type, tool `callId`, or stream coordinate is a barrier. The completed text forms, `message.completed` and `reasoning.completed`, are the compatibility path for clients that don't render deltas. A streamed tool input is complete when the matching validated call arrives in `actions.requested`. ## Handle event types @@ -103,7 +103,7 @@ The most common UI events are: | `message.received` | Confirm the user message landed; `data.parts` includes text and file metadata. | | `reasoning.appended` | Render reasoning deltas when the model provides them. | | `message.appended` | Render assistant text deltas. | -| `action.input.appended` | Render cumulative raw tool input before validation completes. | +| `action.input.appended` | Accumulate raw tool-input deltas before validation completes. | | `actions.requested` | Show tool calls as the model requests them, before execution. | | `action.partial` | Update a generator tool's provisional output snapshot. | | `action.result` | Show tool call results. | @@ -116,7 +116,7 @@ The most common UI events are: For the complete event table, see [Sessions, runs & streaming](../../concepts/sessions-runs-and-streaming). -The default message reducer turns `action.input.appended` into a `dynamic-tool` part with `state: "input-streaming"`. Its `inputText` field contains cumulative raw text that may be incomplete JSON. The matching `actions.requested` event upgrades the same `toolCallId` to `state: "input-available"` and puts the validated value in `input`. +Each `action.input.appended` event stores `inputTextDelta` and its zero-based UTF-16 code-unit `inputTextOffset`, so consumers can reject a gap instead of concatenating corrupt input. The default message reducer starts or restarts at offset `0` and ignores a nonzero offset that is not contiguous. It accumulates accepted deltas into a `dynamic-tool` part with `state: "input-streaming"`; its `inputText` field contains cumulative raw text that may be incomplete JSON. The matching `actions.requested` event upgrades the same `toolCallId` to `state: "input-available"` and puts the validated value in `input`. When a submitted message includes attachments, `message.received.data.message` stays the flattened compatibility summary, while `message.received.data.parts` carries renderable text and diff --git a/packages/eve/extension-contracts/reports/channel/v10.json b/packages/eve/extension-contracts/reports/channel/v10.json index bfb50159bd..95ded434a5 100644 --- a/packages/eve/extension-contracts/reports/channel/v10.json +++ b/packages/eve/extension-contracts/reports/channel/v10.json @@ -2,7 +2,7 @@ "kind": "eve-extension-capability-contract", "capability": "channel", "epoch": 10, - "sha256": "5d17981853027789da610ed4ae859c00b3ffb5dec24d12bac8bdfd11370a8df6", + "sha256": "0b7359bd92a73d6d4c9c51f63d7e7af5b169338cab6f0e1429ce17f44e322c1a", "exports": [ "DELETE", "GET", diff --git a/packages/eve/extension-contracts/reports/dynamicInstructions/v14.json b/packages/eve/extension-contracts/reports/dynamicInstructions/v14.json index 20f5f395a0..ffe62d2bfc 100644 --- a/packages/eve/extension-contracts/reports/dynamicInstructions/v14.json +++ b/packages/eve/extension-contracts/reports/dynamicInstructions/v14.json @@ -2,6 +2,6 @@ "kind": "eve-extension-capability-contract", "capability": "dynamicInstructions", "epoch": 14, - "sha256": "4734df4dbc68c6082625975f794c19a0d39004c357d820e445a093d18b6c339b", + "sha256": "9dea3c7c9cecd8c6136e6e6c38d1d21e74e7ee13427e04cd4aedb604d5ffebcd", "exports": ["defineDynamic"] } diff --git a/packages/eve/extension-contracts/reports/dynamicSkill/v13.json b/packages/eve/extension-contracts/reports/dynamicSkill/v13.json index d1556bef4f..74ea0cc407 100644 --- a/packages/eve/extension-contracts/reports/dynamicSkill/v13.json +++ b/packages/eve/extension-contracts/reports/dynamicSkill/v13.json @@ -2,6 +2,6 @@ "kind": "eve-extension-capability-contract", "capability": "dynamicSkill", "epoch": 13, - "sha256": "d0194e3d413ab385e1f2fac57c8ea9bcba9da6a3837a82e5fa0b68b6896833cc", + "sha256": "e314c50148ec36066b8708c6407c3825afe6f141a2c86125c15c6cc63699d1dc", "exports": ["defineDynamic"] } diff --git a/packages/eve/extension-contracts/reports/dynamicTool/v20.json b/packages/eve/extension-contracts/reports/dynamicTool/v20.json index 96dbef7d12..989c199850 100644 --- a/packages/eve/extension-contracts/reports/dynamicTool/v20.json +++ b/packages/eve/extension-contracts/reports/dynamicTool/v20.json @@ -2,7 +2,7 @@ "kind": "eve-extension-capability-contract", "capability": "dynamicTool", "epoch": 20, - "sha256": "202c8444839e924ac223a4654c6802bbe91e63e4949adf2d044de30d502978d7", + "sha256": "8767b4843b08e49833245ec0004d63cbe775263f349888dd075fa6f71ea6959d", "exports": [ "DynamicToolEntry", "DynamicToolEvents", diff --git a/packages/eve/extension-contracts/reports/hook/v15.json b/packages/eve/extension-contracts/reports/hook/v15.json index def9c36f26..0d884f486d 100644 --- a/packages/eve/extension-contracts/reports/hook/v15.json +++ b/packages/eve/extension-contracts/reports/hook/v15.json @@ -2,6 +2,6 @@ "kind": "eve-extension-capability-contract", "capability": "hook", "epoch": 15, - "sha256": "bb84054bca066e8f1baa191e6e17b7274792bea39664e810a0c8ab789838b498", + "sha256": "9b4ffcb120e78affa01a862934dc641327c812958fb7da37d254b974c7d7ef9d", "exports": ["defineHook"] } diff --git a/packages/eve/extension-contracts/reports/schedule/v4.json b/packages/eve/extension-contracts/reports/schedule/v4.json index 4d360bf02d..27f3214f6f 100644 --- a/packages/eve/extension-contracts/reports/schedule/v4.json +++ b/packages/eve/extension-contracts/reports/schedule/v4.json @@ -2,7 +2,7 @@ "kind": "eve-extension-capability-contract", "capability": "schedule", "epoch": 4, - "sha256": "19917f781b201d389a25b49c33176301be1298facda8e0e9e3b31db514a25670", + "sha256": "cc792aebecfd572ff328f0e5c05cd50ca662fcbaa61b59adec30aa023f5cfc67", "exports": [ "ScheduleDefinition", "ScheduleHandlerArgs", diff --git a/packages/eve/extension-contracts/reports/subagent/v5.json b/packages/eve/extension-contracts/reports/subagent/v5.json index 95961778e2..fa955366e0 100644 --- a/packages/eve/extension-contracts/reports/subagent/v5.json +++ b/packages/eve/extension-contracts/reports/subagent/v5.json @@ -2,7 +2,7 @@ "kind": "eve-extension-capability-contract", "capability": "subagent", "epoch": 5, - "sha256": "b1ce511edad297f42a47572994c0a5af126fe53c63b1b4140ffc6af7aa6c6734", + "sha256": "b940181d78346c850ffb480856db39d21fecc4ebd9b8311c4f40dc4dd3366e2e", "exports": [ "AgentCompactionDefinition", "AgentDefinition", diff --git a/packages/eve/src/client/message-reducer-primitives.ts b/packages/eve/src/client/message-reducer-primitives.ts index d221bb8f7a..72f3124958 100644 --- a/packages/eve/src/client/message-reducer-primitives.ts +++ b/packages/eve/src/client/message-reducer-primitives.ts @@ -69,3 +69,13 @@ export function removeStreamingToolPartsForTurn( export function optimisticUserMessageId(submissionId: string): string { return `optimistic:${submissionId}:user`; } + +export function appendToolInputDelta( + inputText: string | undefined, + offset: number, + delta: string, +): string | undefined { + if (offset === 0) return delta; + if (inputText?.length !== offset) return undefined; + return inputText + delta; +} diff --git a/packages/eve/src/client/message-reducer-types.ts b/packages/eve/src/client/message-reducer-types.ts index 5bae08f205..b742698502 100644 --- a/packages/eve/src/client/message-reducer-types.ts +++ b/packages/eve/src/client/message-reducer-types.ts @@ -140,6 +140,7 @@ export type EveDynamicToolPart = { readonly approval?: never; readonly errorText?: never; readonly input: unknown | undefined; + /** Accumulated raw tool input, which may be incomplete JSON. */ readonly inputText: string; readonly output?: never; readonly state: "input-streaming"; diff --git a/packages/eve/src/client/message-reducer.test.ts b/packages/eve/src/client/message-reducer.test.ts index 02ff0d5136..a4a4de8c6c 100644 --- a/packages/eve/src/client/message-reducer.test.ts +++ b/packages/eve/src/client/message-reducer.test.ts @@ -41,7 +41,7 @@ describe("defaultMessageReducer", () => { createActionInputAppendedEvent({ callId: "call_render", inputTextDelta: "", - inputTextSoFar: "", + inputTextOffset: 0, sequence: 1, stepIndex: 0, toolName: "render", @@ -50,7 +50,7 @@ describe("defaultMessageReducer", () => { createActionInputAppendedEvent({ callId: "call_render", inputTextDelta: '{"title":"Hel', - inputTextSoFar: '{"title":"Hel', + inputTextOffset: 0, sequence: 1, stepIndex: 0, toolName: "render", @@ -69,6 +69,20 @@ describe("defaultMessageReducer", () => { type: "dynamic-tool", }); + const contiguous = data; + data = reduceServerEvents(reducer, data, [ + createActionInputAppendedEvent({ + callId: "call_render", + inputTextDelta: "gap", + inputTextOffset: 99, + sequence: 1, + stepIndex: 0, + toolName: "render", + turnId: "turn_1", + }), + ]); + expect(data).toBe(contiguous); + data = reduceServerEvents(reducer, data, [ createActionsRequestedEvent({ actions: [ @@ -100,7 +114,7 @@ describe("defaultMessageReducer", () => { createActionInputAppendedEvent({ callId: "call_render", inputTextDelta: "late", - inputTextSoFar: "late", + inputTextOffset: 0, sequence: 1, stepIndex: 0, toolName: "render", @@ -116,7 +130,7 @@ describe("defaultMessageReducer", () => { createActionInputAppendedEvent({ callId: "call_render", inputTextDelta: "{", - inputTextSoFar: "{", + inputTextOffset: 0, sequence: 1, stepIndex: 0, toolName: "render", diff --git a/packages/eve/src/client/message-reducer.ts b/packages/eve/src/client/message-reducer.ts index 39c9691d95..9817d56e85 100644 --- a/packages/eve/src/client/message-reducer.ts +++ b/packages/eve/src/client/message-reducer.ts @@ -21,6 +21,7 @@ import { toMessageInputRequest, } from "#client/message-action-parts.js"; import { + appendToolInputDelta, optimisticUserMessageId, partKey, projectReceivedParts, @@ -142,13 +143,18 @@ function reduceMessageData(data: EveMessageData, event: EveAgentReducerEvent): E case "action.input.appended": { const existing = findToolPart(data, event.data.callId); - if (existing !== undefined && existing.state !== "input-streaming") { - return data; - } + if (existing !== undefined && existing.state !== "input-streaming") return data; + + const inputText = appendToolInputDelta( + existing?.state === "input-streaming" ? existing.inputText : undefined, + event.data.inputTextOffset, + event.data.inputTextDelta, + ); + if (inputText === undefined) return data; const nextPart: EveDynamicToolPart = { input: undefined, - inputText: event.data.inputTextSoFar, + inputText, state: "input-streaming", stepIndex: event.data.stepIndex, toolCallId: event.data.callId, diff --git a/packages/eve/src/harness/emission.test.ts b/packages/eve/src/harness/emission.test.ts index 79d1c95637..3650b3dbf0 100644 --- a/packages/eve/src/harness/emission.test.ts +++ b/packages/eve/src/harness/emission.test.ts @@ -319,20 +319,34 @@ describe("emitStreamContent action requests", () => { "action.input.appended", "actions.requested", ]); - expect(events.slice(1, 3)).toMatchObject([ + const inputEvents = events.filter((event) => event.type === "action.input.appended"); + expect(inputEvents.map((event) => event.data)).toEqual([ { - data: { - callId: "call-render", - inputTextDelta: '{"title":"Hel', - inputTextSoFar: '{"title":"Hel', - }, + callId: "call-render", + inputTextDelta: "", + inputTextOffset: 0, + sequence: 0, + stepIndex: 0, + toolName: "render", + turnId: "turn_0", }, { - data: { - callId: "call-render", - inputTextDelta: 'lo"}', - inputTextSoFar: '{"title":"Hello"}', - }, + callId: "call-render", + inputTextDelta: '{"title":"Hel', + inputTextOffset: 0, + sequence: 0, + stepIndex: 0, + toolName: "render", + turnId: "turn_0", + }, + { + callId: "call-render", + inputTextDelta: 'lo"}', + inputTextOffset: 13, + sequence: 0, + stepIndex: 0, + toolName: "render", + turnId: "turn_0", }, ]); }); diff --git a/packages/eve/src/harness/ordered-stream-emitter.test.ts b/packages/eve/src/harness/ordered-stream-emitter.test.ts index 75fcf1e9f6..0c6e457bf2 100644 --- a/packages/eve/src/harness/ordered-stream-emitter.test.ts +++ b/packages/eve/src/harness/ordered-stream-emitter.test.ts @@ -48,11 +48,11 @@ function partial(callId: string, output: string) { }); } -function input(callId: string, delta: string, soFar: string) { +function input(callId: string, delta: string, offset: number) { return createActionInputAppendedEvent({ callId, inputTextDelta: delta, - inputTextSoFar: soFar, + inputTextOffset: offset, sequence: 1, stepIndex: 0, toolName: "render", @@ -91,18 +91,18 @@ describe("createOrderedStreamEmitter", () => { const emitter = createOrderedStreamEmitter(emitFn); await emitter.emit(message("A", "A")); - await emitter.emit(input("call_1", "{", "{")); - await emitter.emit(input("call_1", '"title":', '{"title":')); - await emitter.emit(input("call_1", '"Hello"}', '{"title":"Hello"}')); - await emitter.emit(input("call_2", "{}", "{}")); + await emitter.emit(input("call_1", "{", 0)); + await emitter.emit(input("call_1", '"title":', 1)); + await emitter.emit(input("call_1", '"Hello"}', 9)); + await emitter.emit(input("call_2", "{}", 0)); firstWrite.resolve(); await emitter.closeAndDrain(); expect(events).toEqual([ message("A", "A"), - input("call_1", '{"title":"Hello"}', '{"title":"Hello"}'), - input("call_2", "{}", "{}"), + input("call_1", '{"title":"Hello"}', 0), + input("call_2", "{}", 0), ]); }); diff --git a/packages/eve/src/harness/ordered-stream-emitter.ts b/packages/eve/src/harness/ordered-stream-emitter.ts index b022f9afe5..caf5c2271d 100644 --- a/packages/eve/src/harness/ordered-stream-emitter.ts +++ b/packages/eve/src/harness/ordered-stream-emitter.ts @@ -187,7 +187,16 @@ function mergeAdjacentEmissions( } left.deltaParts ??= [appendDelta(left.event)]; left.deltaParts.push(appendDelta(right)); - left.event = right; + left.event = + left.event.type === "action.input.appended" && right.type === "action.input.appended" + ? { + ...right, + data: { + ...right.data, + inputTextOffset: left.event.data.inputTextOffset, + }, + } + : right; left.messages = messages; return true; } diff --git a/packages/eve/src/harness/stream-content-emission.ts b/packages/eve/src/harness/stream-content-emission.ts index 89f0ea5abe..06e51fa350 100644 --- a/packages/eve/src/harness/stream-content-emission.ts +++ b/packages/eve/src/harness/stream-content-emission.ts @@ -33,10 +33,7 @@ import { isInvalidToolCall, resolveProviderToolCallRequest, } from "#harness/tool-call-input-errors.js"; -import type { - RuntimeActionRequest, - RuntimeToolResultActionResult, -} from "#shared/action-types.js"; +import type { RuntimeActionRequest, RuntimeToolResultActionResult } from "#shared/action-types.js"; import { createProviderStreamActionBatch } from "#harness/stream-actions.js"; import { normalizeModelStreamError } from "#harness/model-call-error.js"; import { createOrderedStreamEmitter } from "#harness/ordered-stream-emitter.js"; @@ -131,7 +128,7 @@ async function consumeStreamContent( const invalidInputToolCallIds = new Set(); const inlineAuthorizationResults: TypedToolResult[] = []; const trailingInlineToolResultParts: InlineToolResultPart[] = []; - const streamingActionInputs = new Map(); + const streamingActionInputs = new Map(); const flushCurrentMessage = async (): Promise => { if (currentMessage.length === 0) { @@ -153,13 +150,13 @@ async function consumeStreamContent( callId: string, toolName: string, inputTextDelta: string, - inputTextSoFar: string, + inputTextOffset: number, ): Promise => emitFn( createActionInputAppendedEvent({ callId, inputTextDelta, - inputTextSoFar, + inputTextOffset, sequence: state.sequence, stepIndex: state.stepIndex, toolName, @@ -334,8 +331,8 @@ async function consumeStreamContent( if (currentMessage.trim().length > 0) { await flushCurrentMessage(); } - streamingActionInputs.set(part.id, { text: "", toolName: part.toolName }); - await emitActionInput(part.id, part.toolName, "", ""); + streamingActionInputs.set(part.id, { offset: 0, toolName: part.toolName }); + await emitActionInput(part.id, part.toolName, "", 0); break; } case "tool-input-delta": { @@ -344,8 +341,9 @@ async function consumeStreamContent( break; } await providerActionBatch.flush(); - input.text += part.delta; - await emitActionInput(part.id, input.toolName, part.delta, input.text); + const inputTextOffset = input.offset; + input.offset += part.delta.length; + await emitActionInput(part.id, input.toolName, part.delta, inputTextOffset); break; } case "tool-input-end": diff --git a/packages/eve/src/protocol/message.ts b/packages/eve/src/protocol/message.ts index 5a5f2c3e6f..4a9ab24eda 100644 --- a/packages/eve/src/protocol/message.ts +++ b/packages/eve/src/protocol/message.ts @@ -431,7 +431,8 @@ export interface ActionInputAppendedStreamEvent { data: { callId: string; inputTextDelta: string; - inputTextSoFar: string; + /** Zero-based UTF-16 code-unit offset where `inputTextDelta` begins. */ + inputTextOffset: number; sequence: number; stepIndex: number; toolName: string; @@ -1103,7 +1104,7 @@ export function createActionsRequestedEvent(input: { export function createActionInputAppendedEvent(input: { readonly callId: string; readonly inputTextDelta: string; - readonly inputTextSoFar: string; + readonly inputTextOffset: number; readonly sequence: number; readonly stepIndex: number; readonly toolName: string; @@ -1113,7 +1114,7 @@ export function createActionInputAppendedEvent(input: { data: { callId: input.callId, inputTextDelta: input.inputTextDelta, - inputTextSoFar: input.inputTextSoFar, + inputTextOffset: input.inputTextOffset, sequence: input.sequence, stepIndex: input.stepIndex, toolName: input.toolName, From 8df42a63b97f112f4da38a6f01fd54509b235d02 Mon Sep 17 00:00:00 2001 From: Rui Conti Date: Thu, 27 Aug 2026 00:06:56 -0400 Subject: [PATCH 5/5] test(eve): update installed extension epochs Signed-off-by: Rui Conti --- .../test/scenarios/mounted-extension-installed.scenario.test.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/packages/eve/test/scenarios/mounted-extension-installed.scenario.test.ts b/packages/eve/test/scenarios/mounted-extension-installed.scenario.test.ts index a1ac986833..eb8f6cf9cc 100644 --- a/packages/eve/test/scenarios/mounted-extension-installed.scenario.test.ts +++ b/packages/eve/test/scenarios/mounted-extension-installed.scenario.test.ts @@ -249,7 +249,7 @@ describe("mounted extension installed under node_modules", () => { ); expect( JSON.parse(extensionFiles[`node_modules/${PACKAGE_NAME}/dist/extension/_manifest.json`]!), - ).toMatchObject({ requires: { channel: 9, schedule: 3, subagent: 4 } }); + ).toMatchObject({ requires: { channel: 10, schedule: 4, subagent: 5 } }); const app = await scenarioApp({ name: "mounted-extension-installed", installDependencies: true,