diff --git a/.changeset/stream-tool-input.md b/.changeset/stream-tool-input.md new file mode 100644 index 0000000000..8e792f4eac --- /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. 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 c35c6cd824..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. | @@ -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 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. + +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. @@ -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..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` 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 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 @@ -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` | 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. | +| `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). +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 file metadata. File parts never include raw bytes or internal sandbox paths; `url` appears only for 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..95ded434a5 --- /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": "0b7359bd92a73d6d4c9c51f63d7e7af5b169338cab6f0e1429ce17f44e322c1a", + "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..ffe62d2bfc --- /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": "9dea3c7c9cecd8c6136e6e6c38d1d21e74e7ee13427e04cd4aedb604d5ffebcd", + "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..74ea0cc407 --- /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": "e314c50148ec36066b8708c6407c3825afe6f141a2c86125c15c6cc63699d1dc", + "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..989c199850 --- /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": "8767b4843b08e49833245ec0004d63cbe775263f349888dd075fa6f71ea6959d", + "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..0d884f486d --- /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": "9b4ffcb120e78affa01a862934dc641327c812958fb7da37d254b974c7d7ef9d", + "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..27f3214f6f --- /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": "cc792aebecfd572ff328f0e5c05cd50ca662fcbaa61b59adec30aa023f5cfc67", + "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..fa955366e0 --- /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": "b940181d78346c850ffb480856db39d21fecc4ebd9b8311c4f40dc4dd3366e2e", + "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 new file mode 100644 index 0000000000..72f3124958 --- /dev/null +++ b/packages/eve/src/client/message-reducer-primitives.ts @@ -0,0 +1,81 @@ +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`; +} + +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 2fa6e34a8c..b742698502 100644 --- a/packages/eve/src/client/message-reducer-types.ts +++ b/packages/eve/src/client/message-reducer-types.ts @@ -140,6 +140,8 @@ 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 3a8981cae6..a4a4de8c6c 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,127 @@ 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: "", + inputTextOffset: 0, + sequence: 1, + stepIndex: 0, + toolName: "render", + turnId: "turn_1", + }), + createActionInputAppendedEvent({ + callId: "call_render", + inputTextDelta: '{"title":"Hel', + inputTextOffset: 0, + 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", + }); + + 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: [ + { + 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", + inputTextOffset: 0, + 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: "{", + inputTextOffset: 0, + 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..9817d56e85 100644 --- a/packages/eve/src/client/message-reducer.ts +++ b/packages/eve/src/client/message-reducer.ts @@ -20,12 +20,16 @@ import { stringifyUnknown, toMessageInputRequest, } from "#client/message-action-parts.js"; +import { + appendToolInputDelta, + 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 +141,42 @@ 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 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, + 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 +381,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 +393,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 +413,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 +698,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/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.test.ts b/packages/eve/src/harness/emission.test.ts index 7898cf1a84..3650b3dbf0 100644 --- a/packages/eve/src/harness/emission.test.ts +++ b/packages/eve/src/harness/emission.test.ts @@ -277,6 +277,122 @@ 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", + ]); + const inputEvents = events.filter((event) => event.type === "action.input.appended"); + expect(inputEvents.map((event) => event.data)).toEqual([ + { + callId: "call-render", + inputTextDelta: "", + inputTextOffset: 0, + sequence: 0, + stepIndex: 0, + toolName: "render", + turnId: "turn_0", + }, + { + 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", + }, + ]); + }); + + 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 +498,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 +519,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..aa778bfab4 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, + inputTextOffset: number, + ): Promise => + emitFn( + createActionInputAppendedEvent({ + callId, + inputTextDelta, + inputTextOffset, + 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,40 @@ 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, { offset: 0, toolName: part.toolName }); + await emitActionInput(part.id, part.toolName, "", 0); + break; + } + case "tool-input-delta": { + const input = streamingActionInputs.get(part.id); + if (input === undefined) { + break; + } + await providerActionBatch.flush(); + const inputTextOffset = input.offset; + input.offset += part.delta.length; + await emitActionInput(part.id, input.toolName, part.delta, inputTextOffset); + 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..0c6e457bf2 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, offset: number) { + return createActionInputAppendedEvent({ + callId, + inputTextDelta: delta, + inputTextOffset: offset, + 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", "{", 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"}', 0), + input("call_2", "{}", 0), + ]); + }); + 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..caf5c2271d 100644 --- a/packages/eve/src/harness/ordered-stream-emitter.ts +++ b/packages/eve/src/harness/ordered-stream-emitter.ts @@ -1,10 +1,16 @@ import type { + ActionInputAppendedStreamEvent, UnstampedMessageStreamEvent, MessageAppendedStreamEvent, ReasoningAppendedStreamEvent, } 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; @@ -167,20 +173,30 @@ 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; + 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 ??= [appendDelta(left.event)]; + left.deltaParts.push(appendDelta(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; } @@ -195,10 +211,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; - 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 { @@ -224,13 +265,20 @@ 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, -): 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/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..4a9ab24eda 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,24 @@ 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; + /** Zero-based UTF-16 code-unit offset where `inputTextDelta` begins. */ + inputTextOffset: number; + 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 +733,7 @@ export interface SessionCompletedStreamEvent { * consumers receive {@link MessageStreamEvent}. */ export type UnstampedMessageStreamEvent = + | ActionInputAppendedStreamEvent | ApprovalCandidateStreamEvent | ApprovalSettledStreamEvent | ContextClearedStreamEvent @@ -1081,6 +1100,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 inputTextOffset: number; + readonly sequence: number; + readonly stepIndex: number; + readonly toolName: string; + readonly turnId: string; +}): ActionInputAppendedStreamEvent { + return { + data: { + callId: input.callId, + inputTextDelta: input.inputTextDelta, + inputTextOffset: input.inputTextOffset, + 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">; 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,