From 2b2e8bbc19ce371a67202da11a4d2fc3a8ea256a Mon Sep 17 00:00:00 2001 From: Andrew Barba Date: Mon, 27 Jul 2026 12:07:41 -0400 Subject: [PATCH 01/12] feat(eve): stamp every stream event with a stable id MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Every session stream event now carries `meta.id`, a unique `evt_`-prefixed ULID minted once at emission. The id is stable across reconnects, rewinds, tail reads, and replays, so a consumer can key durable work on it — for example `INSERT ... ON CONFLICT (id) DO NOTHING` in a hook — instead of guessing identity from payload content. Events are stamped before the channel adapter runs, so the adapter, the persisted chunk, and hooks all observe the same id. Client consumers now use it. `EveAgentStore` no longer double-applies an `initialEvents` prefix that the live stream replays, which also protects the React, Vue, and Svelte bindings and any user-authored reducer. The dev TUI drops re-delivered chunks up front, which fixes a subagent section that rendered its child transcript twice when the child stream reopened, and retires four content-comparison replay heuristics that could silently drop a new model call whose text prefixed the previous one. Bumps the message stream version to 20 and the hook, dynamicTool, dynamicSkill, and dynamicInstructions capability epochs. All four changed only additively, so prior-epoch extensions stay compatible and are retained. Signed-off-by: Andrew Barba --- .changeset/dedupe-stream-events-by-id.md | 9 + .changeset/stable-stream-event-ids.md | 5 + docs/concepts/sessions-runs-and-streaming.md | 42 +++++ docs/guides/client/streaming.mdx | 10 +- docs/guides/frontend/overview.mdx | 2 + docs/guides/hooks.md | 33 +++- .../compatibility/dynamicInstructions/v1.ts | 17 ++ .../compatibility/dynamicSkill/v1.ts | 21 +++ .../compatibility/dynamicTool/v3.ts | 29 ++++ .../compatibility/hook/v2.ts | 32 ++++ .../reports/dynamicInstructions/v2.json | 7 + .../reports/dynamicSkill/v2.json | 7 + .../reports/dynamicTool/v4.json | 13 ++ .../extension-contracts/reports/hook/v3.json | 7 + packages/eve/src/channel/adapter.test.ts | 4 +- packages/eve/src/channel/adapter.ts | 13 +- packages/eve/src/channel/schedule.test.ts | 8 +- packages/eve/src/channel/send.test.ts | 20 ++- packages/eve/src/channel/session.ts | 4 +- packages/eve/src/channel/types.ts | 9 +- packages/eve/src/cli/dev/tui/runner.test.ts | 161 ++++++++++++++---- packages/eve/src/cli/dev/tui/runner.ts | 35 ++-- .../eve/src/cli/dev/tui/subagent-pump.test.ts | 90 ++++++++-- packages/eve/src/cli/dev/tui/subagent-pump.ts | 15 +- .../eve/src/client/eve-agent-store.test.ts | 126 ++++++++++++++ packages/eve/src/client/eve-agent-store.ts | 34 +++- packages/eve/src/client/index.ts | 2 + packages/eve/src/client/message-response.ts | 14 +- packages/eve/src/client/ndjson.ts | 10 +- packages/eve/src/client/open-stream.ts | 4 +- packages/eve/src/client/session.ts | 12 +- packages/eve/src/client/types.ts | 4 +- .../src/compiler/extension-compatibility.ts | 8 +- .../hook-lifecycle.integration.test.ts | 30 +++- packages/eve/src/context/hook-lifecycle.ts | 4 +- .../dispatch-runtime-actions-step.ts | 28 +-- .../execution/settle-cancelled-turn-step.ts | 12 +- .../src/execution/subagent-adapter.test.ts | 3 +- .../subagent-auth-proxy.integration.test.ts | 6 +- .../execution/subagent-event-proxy-step.ts | 12 +- .../subagent-hitl-proxy.integration.test.ts | 13 +- .../terminal-session-failure-step.ts | 8 +- .../workflow-entry.integration.test.ts | 67 +++++++- .../eve/src/execution/workflow-runtime.ts | 10 +- packages/eve/src/execution/workflow-steps.ts | 17 +- .../harness/ordered-stream-emitter.test.ts | 12 +- .../host/build-extension.scenario.test.ts | 2 +- packages/eve/src/internal/testing/events.ts | 39 ++++- .../eve/src/protocol/event-dedupe.test.ts | 82 +++++++++ packages/eve/src/protocol/event-dedupe.ts | 69 ++++++++ packages/eve/src/protocol/event-id.test.ts | 83 +++++++++ packages/eve/src/protocol/event-id.ts | 113 ++++++++++++ packages/eve/src/protocol/message.test.ts | 25 ++- packages/eve/src/protocol/message.ts | 53 ++++-- .../channels/chat-sdk/chatSdkChannel.test.ts | 12 +- .../channels/discord/discordChannel.test.ts | 16 +- packages/eve/src/public/channels/eve.test.ts | 15 +- .../channels/github/githubChannel.test.ts | 10 +- .../channels/linear/linearChannel.test.ts | 12 +- .../channels/slack/slackChannel.test.ts | 14 +- .../channels/telegram/telegramChannel.test.ts | 12 +- .../channels/twilio/twilioChannel.test.ts | 12 +- .../src/public/definitions/channel.test.ts | 13 +- packages/eve/src/public/definitions/hook.ts | 9 +- packages/eve/src/react/use-eve-agent.test.ts | 11 +- packages/eve/src/react/use-eve-agent.ts | 4 +- packages/eve/src/runtime/hooks/registry.ts | 4 +- packages/eve/src/runtime/types.ts | 4 +- packages/eve/src/svelte/use-eve-agent.test.ts | 14 +- packages/eve/src/svelte/use-eve-agent.ts | 8 +- packages/eve/src/vue/use-eve-agent.test.ts | 10 +- packages/eve/src/vue/use-eve-agent.ts | 6 +- .../eve/test/eve-run-stream-channel.test.ts | 39 +++-- .../schedule-trigger.scenario.test.ts | 6 +- .../tui-client/tui-connection-auth-states.ts | 24 ++- 75 files changed, 1464 insertions(+), 286 deletions(-) create mode 100644 .changeset/dedupe-stream-events-by-id.md create mode 100644 .changeset/stable-stream-event-ids.md create mode 100644 packages/eve/extension-contracts/compatibility/dynamicInstructions/v1.ts create mode 100644 packages/eve/extension-contracts/compatibility/dynamicSkill/v1.ts create mode 100644 packages/eve/extension-contracts/compatibility/dynamicTool/v3.ts create mode 100644 packages/eve/extension-contracts/compatibility/hook/v2.ts create mode 100644 packages/eve/extension-contracts/reports/dynamicInstructions/v2.json create mode 100644 packages/eve/extension-contracts/reports/dynamicSkill/v2.json create mode 100644 packages/eve/extension-contracts/reports/dynamicTool/v4.json create mode 100644 packages/eve/extension-contracts/reports/hook/v3.json create mode 100644 packages/eve/src/client/eve-agent-store.test.ts create mode 100644 packages/eve/src/protocol/event-dedupe.test.ts create mode 100644 packages/eve/src/protocol/event-dedupe.ts create mode 100644 packages/eve/src/protocol/event-id.test.ts create mode 100644 packages/eve/src/protocol/event-id.ts diff --git a/.changeset/dedupe-stream-events-by-id.md b/.changeset/dedupe-stream-events-by-id.md new file mode 100644 index 000000000..5843d896e --- /dev/null +++ b/.changeset/dedupe-stream-events-by-id.md @@ -0,0 +1,9 @@ +--- +"eve": patch +--- + +Stream consumers now drop re-delivered events by their stable `meta.id` +instead of guessing from payload content. `EveAgentStore` (and so the React, +Vue, and Svelte bindings) no longer double-applies an `initialEvents` prefix +that the live stream replays, and the dev TUI no longer renders a subagent's +transcript twice when its child stream reopens. diff --git a/.changeset/stable-stream-event-ids.md b/.changeset/stable-stream-event-ids.md new file mode 100644 index 000000000..166f3cfab --- /dev/null +++ b/.changeset/stable-stream-event-ids.md @@ -0,0 +1,5 @@ +--- +"eve": patch +--- + +Every session stream event now carries a stable `meta.id`. The id is a sortable, `evt_`-prefixed ULID minted once when the event is written to the durable stream, so reconnecting from a cursor, rewinding to `startIndex=0`, or replaying a finished session all return the same id for the same event — making it safe to use as a primary key when persisting events (`on conflict (id) do nothing`). Channel adapters, the durable stream, and authored hooks now all observe the same envelope, so a hook can use `event.meta.id` to make its own side effects idempotent. Events read from a stream are typed as the new `StampedHandleMessageStreamEvent`, which guarantees `meta` is present. diff --git a/docs/concepts/sessions-runs-and-streaming.md b/docs/concepts/sessions-runs-and-streaming.md index fef752a2d..400196cef 100644 --- a/docs/concepts/sessions-runs-and-streaming.md +++ b/docs/concepts/sessions-runs-and-streaming.md @@ -73,6 +73,46 @@ A delegated subagent publishes progress on its own child-session stream. The par `step.failed` and `turn.failed` carry `{ code, message, details? }` for the failed fragment or turn, and `session.failed` is the terminal session-level variant. `turn.cancelled` is not a failure: the cancelled turn ends without any failure event, `session.waiting` follows, and the session accepts the next message normally — whatever the turn streamed before cancellation stays on the stream, while durable history keeps only what had already settled. When a turn requested an output schema, the finalized payload lands on `result.completed` as `data.result` before the turn boundary. `authorization.required` carries the sign-in challenge (`data.authorization` may include `url`, `userCode`, `expiresAt`, `instructions`), and `authorization.completed` carries `data.outcome` (`"authorized" | "declined" | "failed" | "timed-out"`). +## The event envelope + +Alongside `type` and `data`, every event carries a `meta` envelope: + +```json +{ + "type": "message.completed", + "data": { + "message": "Sunny and 72°F.", + "finishReason": "stop", + "sequence": 0, + "stepIndex": 0, + "turnId": "turn_0" + }, + "meta": { "id": "evt_01K18VW2Q7B4M9XN3RTC5FDGHJ", "at": "2026-07-27T18:04:11.912Z" } +} +``` + +- **`meta.id`** uniquely identifies the event. It is a `evt_`-prefixed [ULID](https://github.com/ulid/spec), so sorting ids as strings reproduces emission order. +- **`meta.at`** is the ISO-8601 time the event was emitted. + +`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. + +That makes it the key for ingesting events into a database exactly once: + +```sql +insert into agent_events (id, session_id, type, data, emitted_at) +values ($1, $2, $3, $4, $5) +on conflict (id) do nothing; +``` + +Because ids sort chronologically, a `primary key (id)` also keeps inserts clustered and lets you page with `where id > $cursor order by id`. + +Two caveats worth knowing: + +- **Ids identify events, not intent.** Two events with identical payloads — the `step.failed` → `turn.failed` → `session.failed` cascade, or two identical text deltas in one step — are distinct events with distinct ids. Deduplicate on `meta.id` only; matching on content would drop real data. +- **A subagent's event is re-emitted, not shared.** When a parent forwards a child's event onto its own stream, the parent's copy is a separate event with its own id. Correlate the two streams through `subagent.called.data.childSessionId`. + +Authored [hooks](../guides/hooks) receive the same envelope, so a hook can use `event.meta.id` to make its own side effects idempotent across a retry. + ## Send a follow-up message Once the session is waiting (you'll see `session.waiting`), POST your follow-up to the session endpoint with `event.data.continuationToken`: @@ -112,6 +152,8 @@ Custom channel routes request the same cancellation without knowing the session The stream is durable. Every event is recorded before a step completes, so consumers can reconnect from their cursor when an HTTP connection ends. A nonnegative `startIndex` is an absolute event count: use it to pick up where you dropped off or pass `0` to rewind to the start. +If a reconnect overlaps events you already handled, [`meta.id`](#the-event-envelope) identifies the duplicates: it is unchanged across reconnects and rewinds, so a consumer keyed on it can replay safely. + ```bash curl "http://127.0.0.1:2000/eve/v1/session//stream?startIndex=" ``` diff --git a/docs/guides/client/streaming.mdx b/docs/guides/client/streaming.mdx index 8ba6fb3fd..a2e539f9f 100644 --- a/docs/guides/client/streaming.mdx +++ b/docs/guides/client/streaming.mdx @@ -60,13 +60,15 @@ for await (const event of response) { ## Handle event types -Import event types from `eve/client` when you want exhaustiveness or helpers: +Import event types from `eve/client` when you want exhaustiveness or helpers. Events read from a stream are `StampedHandleMessageStreamEvent`: the same union, with the `meta` envelope guaranteed present. ```ts -import type { HandleMessageStreamEvent } from "eve/client"; +import type { StampedHandleMessageStreamEvent } from "eve/client"; import { isCurrentTurnBoundaryEvent } from "eve/client"; -function handleEvent(event: HandleMessageStreamEvent) { +function handleEvent(event: StampedHandleMessageStreamEvent) { + console.log(event.meta.id, event.meta.at); + if (isCurrentTurnBoundaryEvent(event)) { console.log("turn settled:", event.type); } @@ -105,6 +107,8 @@ If you support refresh while an authorization prompt is pending, keep the sessio HTTP connections can end before a run does. The client reconnects from the number of events already consumed, so long turns continue without replaying events. It stops at a turn boundary, when aborted, or when the stream can no longer make progress. +If your consumer persists events, key on `event.meta.id`. It is stable across reconnects and rewinds, so an overlapping replay is safe to ingest twice. See [the event envelope](../../concepts/sessions-runs-and-streaming#the-event-envelope). + Set `streamReconnectPolicy: { reconnect: false }` when a relay or proxy owns the cursor and reconnection policy. This makes a single stream GET attempt and returns when that connection ends; it does not stop the server-side turn: ```ts diff --git a/docs/guides/frontend/overview.mdx b/docs/guides/frontend/overview.mdx index 06ee672bb..8353dd4d2 100644 --- a/docs/guides/frontend/overview.mdx +++ b/docs/guides/frontend/overview.mdx @@ -255,6 +255,8 @@ const agent = useEveAgent({ Store the full `session` object (`sessionId`, `continuationToken`, `streamIndex`), not a single field. The session cursor lets eve continue the durable conversation; the event log lets your UI render historical messages without replaying the whole stream. A database-backed chat app should usually persist stream events as they arrive with `onEvent` and then save a final snapshot in `onFinish`. +You do not have to make `initialEvents` line up exactly with where the stream resumes. Every event carries a stable [`meta.id`](/docs/concepts/sessions-runs-and-streaming#the-event-envelope), and the store drops any event whose id it has already applied — so a saved log that overlaps the replayed prefix renders once, and `onEvent` only fires for events your UI has not seen. + For multiple chat threads, keep one saved event log and session cursor per thread. `agent`, `host`, `reducer`, `session`, `initialEvents`, `initialSession`, `auth`, `headers`, and `optimistic` are read when the hook creates its store, so remount the chat component when switching threads, for example with `key={chat.id}`. If the user can refresh or navigate immediately after pressing send, create your app-level chat row and store the pending user message before calling `send()`. After the request starts, persist the session state as soon as it contains a `sessionId`, then reconnect an interrupted in-flight turn with `session.stream({ startIndex: savedEvents.length })` from the lower-level client. diff --git a/docs/guides/hooks.md b/docs/guides/hooks.md index 6a8132b91..9f1c00ed1 100644 --- a/docs/guides/hooks.md +++ b/docs/guides/hooks.md @@ -80,15 +80,44 @@ import { search } from "@acme/crm/tools"; const crmSearch = toolResultFrom(event.data.result, search); // typed; matches crm__search ``` +### Making a hook idempotent + +Every event carries a `meta` envelope with a stable `meta.id`. A hook may run more than once for the same event — a durable step can be retried, and a turn can be re-driven after a park — so use the id as the key for any side effect that must happen once: + +```ts title="agent/hooks/persist.ts" +import { defineHook } from "eve/hooks"; + +export default defineHook({ + events: { + async "*"(event, ctx) { + await db.query( + `insert into agent_events (id, session_id, type, data, emitted_at) + values ($1, $2, $3, $4, $5) + on conflict (id) do nothing`, + [ + event.meta.id, + ctx.session.id, + event.type, + "data" in event ? event.data : null, + event.meta.at, + ], + ); + }, + }, +}); +``` + +See [the event envelope](../concepts/sessions-runs-and-streaming#the-event-envelope) for the full contract, including what `meta.id` does and does not deduplicate. + ## Execution order When a stream event fires, three things happen in order: -1. Emit. The channel adapter handler runs, then the event is written to the durable stream. +1. Emit. The event is stamped with its `meta` envelope, the channel adapter handler runs, then the event is written to the durable stream. 2. Hooks. Stream-event hooks fire (typed handlers first, then the `*` wildcard). Return values are ignored. 3. Dynamic tool resolvers. Resolvers subscribed to the event type run and update the tool set. -Hooks always run after the event is durably recorded, so if a hook throws, the stream stays consistent. +Hooks always run after the event is durably recorded, so if a hook throws, the stream stays consistent. The adapter, the persisted event, and every hook observe the same `meta.id`. ## What happens when a hook throws diff --git a/packages/eve/extension-contracts/compatibility/dynamicInstructions/v1.ts b/packages/eve/extension-contracts/compatibility/dynamicInstructions/v1.ts new file mode 100644 index 000000000..e9cbfca49 --- /dev/null +++ b/packages/eve/extension-contracts/compatibility/dynamicInstructions/v1.ts @@ -0,0 +1,17 @@ +import { defineDynamic, defineInstructions } from "#public/instructions/index.js"; + +/** + * Epoch 1 resolves the per-session system prompt from `session.started`. The + * resolver reads only the resolve context, so the added `meta.id` on the event + * envelope must stay invisible to it. + */ +export default defineDynamic({ + events: { + "session.started": (_event, ctx) => { + const plan = ctx.session.auth.current?.attributes.plan ?? "free"; + return defineInstructions({ + markdown: `The caller is on the ${String(plan)} plan. Match the depth of your answers to it.`, + }); + }, + }, +}); diff --git a/packages/eve/extension-contracts/compatibility/dynamicSkill/v1.ts b/packages/eve/extension-contracts/compatibility/dynamicSkill/v1.ts new file mode 100644 index 000000000..6f6d0d5a4 --- /dev/null +++ b/packages/eve/extension-contracts/compatibility/dynamicSkill/v1.ts @@ -0,0 +1,21 @@ +import { defineDynamic, defineSkill } from "#public/skills/index.js"; + +/** + * Epoch 1 resolves a per-principal skill from `session.started` and + * `turn.started`. The resolver ignores the event body entirely, so the added + * `meta.id` on the envelope must stay invisible to it. + */ +export default defineDynamic({ + events: { + "session.started": (_event, ctx) => { + const team = ctx.session.auth.current?.attributes.team; + return typeof team === "string" + ? defineSkill({ + description: `Escalation playbook for the ${team} team.`, + markdown: `# ${team} playbook\n\nFollow the ${team} escalation path.`, + }) + : null; + }, + "turn.started": () => null, + }, +}); diff --git a/packages/eve/extension-contracts/compatibility/dynamicTool/v3.ts b/packages/eve/extension-contracts/compatibility/dynamicTool/v3.ts new file mode 100644 index 000000000..e49779829 --- /dev/null +++ b/packages/eve/extension-contracts/compatibility/dynamicTool/v3.ts @@ -0,0 +1,29 @@ +import { + defineDynamic, + defineTool, + type DynamicToolEvents, + type DynamicToolResult, +} from "#public/tools/index.js"; + +/** + * Epoch 3 resolvers take the stream event as `unknown` and read session + * identity from the resolve context. Adding `meta.id` to the event envelope + * must not disturb either. + */ +const events = { + "session.started": (_event, ctx): DynamicToolResult => ({ + inspect_session: defineTool({ + description: "Inspect the resolved session", + inputSchema: { type: "object", properties: {} }, + async execute(_input, toolContext) { + return { + resolverSessionId: ctx.session.id, + toolSessionId: toolContext.session.id, + }; + }, + }), + }), + "step.started": (): DynamicToolResult => null, +} satisfies DynamicToolEvents; + +export default defineDynamic({ events }); diff --git a/packages/eve/extension-contracts/compatibility/hook/v2.ts b/packages/eve/extension-contracts/compatibility/hook/v2.ts new file mode 100644 index 000000000..ebc128656 --- /dev/null +++ b/packages/eve/extension-contracts/compatibility/hook/v2.ts @@ -0,0 +1,32 @@ +import { defineHook } from "#public/hooks/index.js"; + +/** + * Epoch 2 predates `meta.id` on the event envelope, so authors keyed their own + * bookkeeping off the turn coordinates carried in `data` and read `meta` + * defensively. Both patterns must keep compiling now that `meta` is guaranteed. + */ +export default defineHook({ + events: { + "turn.started"(event, ctx) { + console.info("turn started", { + agentName: ctx.agent.name, + sequence: event.data.sequence, + sessionId: ctx.session.id, + turnId: event.data.turnId, + }); + }, + "action.result"(event) { + console.info("tool call settled", { + status: event.data.status, + stepIndex: event.data.stepIndex, + }); + }, + "*"(event, ctx) { + console.info("stream event", { + at: event.meta?.at, + channelKind: ctx.channel.kind, + type: event.type, + }); + }, + }, +}); diff --git a/packages/eve/extension-contracts/reports/dynamicInstructions/v2.json b/packages/eve/extension-contracts/reports/dynamicInstructions/v2.json new file mode 100644 index 000000000..8dc78f843 --- /dev/null +++ b/packages/eve/extension-contracts/reports/dynamicInstructions/v2.json @@ -0,0 +1,7 @@ +{ + "kind": "eve-extension-capability-contract", + "capability": "dynamicInstructions", + "epoch": 2, + "sha256": "69c67b56ca66fce4f9c53e181fc585ae23ecb1d184b4aef189ffa88c2da003c9", + "exports": ["defineDynamic"] +} diff --git a/packages/eve/extension-contracts/reports/dynamicSkill/v2.json b/packages/eve/extension-contracts/reports/dynamicSkill/v2.json new file mode 100644 index 000000000..c307cc5f9 --- /dev/null +++ b/packages/eve/extension-contracts/reports/dynamicSkill/v2.json @@ -0,0 +1,7 @@ +{ + "kind": "eve-extension-capability-contract", + "capability": "dynamicSkill", + "epoch": 2, + "sha256": "69c67b56ca66fce4f9c53e181fc585ae23ecb1d184b4aef189ffa88c2da003c9", + "exports": ["defineDynamic"] +} diff --git a/packages/eve/extension-contracts/reports/dynamicTool/v4.json b/packages/eve/extension-contracts/reports/dynamicTool/v4.json new file mode 100644 index 000000000..14d677f2d --- /dev/null +++ b/packages/eve/extension-contracts/reports/dynamicTool/v4.json @@ -0,0 +1,13 @@ +{ + "kind": "eve-extension-capability-contract", + "capability": "dynamicTool", + "epoch": 4, + "sha256": "0950b720c369c71cedf51a2983d3f764f7bd04d0d669c2592c00e34035a68bae", + "exports": [ + "DynamicToolEntry", + "DynamicToolEvents", + "DynamicToolResult", + "DynamicToolSet", + "defineDynamic" + ] +} diff --git a/packages/eve/extension-contracts/reports/hook/v3.json b/packages/eve/extension-contracts/reports/hook/v3.json new file mode 100644 index 000000000..a36c0331d --- /dev/null +++ b/packages/eve/extension-contracts/reports/hook/v3.json @@ -0,0 +1,7 @@ +{ + "kind": "eve-extension-capability-contract", + "capability": "hook", + "epoch": 3, + "sha256": "9f69dd72540de88770bb7c22bd8ad8197e5ad5c35d923f1c7da971abd19087a4", + "exports": ["defineHook"] +} diff --git a/packages/eve/src/channel/adapter.test.ts b/packages/eve/src/channel/adapter.test.ts index 32f0ff165..a0117dd2b 100644 --- a/packages/eve/src/channel/adapter.test.ts +++ b/packages/eve/src/channel/adapter.test.ts @@ -3,6 +3,7 @@ import { describe, expect, it } from "vitest"; import type { ChannelAdapter, ChannelAdapterContext, FetchFileResult } from "#channel/adapter.js"; import { callAdapterEventHandler, defaultDeliverResult, getAdapterKind } from "#channel/adapter.js"; import { createSessionWaitingEvent } from "#protocol/message.js"; +import { stampTestEvent } from "#internal/testing/events.js"; describe("ChannelAdapter (fetchFile field)", () => { it("treats the fetchFile field as optional", () => { @@ -126,12 +127,13 @@ describe("ChannelAdapter helpers", () => { const event = await callAdapterEventHandler( adapter, - createSessionWaitingEvent("slack:temporary"), + stampTestEvent(createSessionWaitingEvent("slack:temporary")), context, ); expect(event).toEqual({ data: { continuationToken: "C1:T1", wait: "next-user-message" }, + meta: { at: "2026-01-01T00:00:00.000Z", id: "evt_test_0000" }, type: "session.waiting", }); }); diff --git a/packages/eve/src/channel/adapter.ts b/packages/eve/src/channel/adapter.ts index e74219cfa..513452ae5 100644 --- a/packages/eve/src/channel/adapter.ts +++ b/packages/eve/src/channel/adapter.ts @@ -1,7 +1,10 @@ import type { ContextAccessor } from "#context/key.js"; import type { StepInput } from "#harness/types.js"; import { createLogger } from "#internal/logging.js"; -import type { HandleMessageStreamEvent } from "#protocol/message.js"; +import type { + HandleMessageStreamEvent, + StampedHandleMessageStreamEvent, +} from "#protocol/message.js"; import type { SessionHandle } from "#channel/session.js"; import type { DeliverPayload } from "#channel/types.js"; import type { FetchFileResult, FetchFileFunction } from "#shared/channel-definition.js"; @@ -227,14 +230,18 @@ export function getAdapterKind(adapter: ChannelAdapter): string { * runtime refreshes `session.waiting` with the live continuation token so a * handler that re-keyed the session publishes the new resume handle. * + * The event arrives already stamped, so the adapter, the persisted stream, and + * hooks all observe the same `meta.id`. The `session.waiting` refresh preserves + * that envelope. + * * Throwing handlers are logged and swallowed so a downstream delivery * failure does not corrupt the event stream write path. */ export async function callAdapterEventHandler( adapter: ChannelAdapter, - event: HandleMessageStreamEvent, + event: StampedHandleMessageStreamEvent, ctx: ChannelAdapterContext, -): Promise { +): Promise { const handler = adapter[event.type] as | ((data: unknown, ctx: ChannelAdapterContext) => void | Promise) | undefined; diff --git a/packages/eve/src/channel/schedule.test.ts b/packages/eve/src/channel/schedule.test.ts index 3cb6e58b5..736a4a724 100644 --- a/packages/eve/src/channel/schedule.test.ts +++ b/packages/eve/src/channel/schedule.test.ts @@ -2,7 +2,7 @@ import { describe, expect, it, vi } from "vitest"; import { CHANNEL_SENTINEL, type CompiledChannel } from "#channel/compiled-channel.js"; import { isCompiledChannel } from "#channel/compiled-channel.js"; -import type { HandleMessageStreamEvent } from "#protocol/message.js"; +import type { StampedHandleMessageStreamEvent } from "#protocol/message.js"; import { SCHEDULE_ADAPTER, SCHEDULE_ADAPTER_KIND, @@ -17,7 +17,7 @@ import type { ResolvedChannelDefinition } from "#runtime/types.js"; function createMockRunHandle(): RunHandle { return { continuationToken: "slack:C0123ABC:", - events: new ReadableStream(), + events: new ReadableStream(), sessionId: "mock-session-id", }; } @@ -28,7 +28,9 @@ function createMockRuntime(): Runtime { deliver: vi.fn().mockRejectedValue(new RuntimeNoActiveSessionError("schedule:token")), resolveSession: vi.fn(), run: vi.fn().mockResolvedValue(createMockRunHandle()), - getEventStream: vi.fn().mockResolvedValue(new ReadableStream()), + getEventStream: vi + .fn() + .mockResolvedValue(new ReadableStream()), terminateSession: vi.fn(), }; } diff --git a/packages/eve/src/channel/send.test.ts b/packages/eve/src/channel/send.test.ts index adf07c8fb..96d146299 100644 --- a/packages/eve/src/channel/send.test.ts +++ b/packages/eve/src/channel/send.test.ts @@ -4,12 +4,12 @@ import type { ChannelAdapter } from "#channel/adapter.js"; import { createSendFn } from "#channel/send.js"; import type { RunHandle, Runtime } from "#channel/types.js"; import { RuntimeNoActiveSessionError } from "#execution/runtime-errors.js"; -import type { HandleMessageStreamEvent } from "#protocol/message.js"; +import type { StampedHandleMessageStreamEvent } from "#protocol/message.js"; function createMockRunHandle(): RunHandle { return { continuationToken: "test:token", - events: new ReadableStream(), + events: new ReadableStream(), sessionId: "mock-session-id", }; } @@ -20,7 +20,9 @@ function createRuntime(deliverError: unknown): Runtime { deliver: vi.fn().mockRejectedValue(deliverError), resolveSession: vi.fn(), run: vi.fn().mockResolvedValue(createMockRunHandle()), - getEventStream: vi.fn().mockResolvedValue(new ReadableStream()), + getEventStream: vi + .fn() + .mockResolvedValue(new ReadableStream()), terminateSession: vi.fn(), }; } @@ -83,7 +85,9 @@ describe("createSendFn", () => { deliver: vi.fn().mockResolvedValue({ sessionId: "existing-session-id" }), resolveSession: vi.fn(), run: vi.fn().mockResolvedValue(createMockRunHandle()), - getEventStream: vi.fn().mockResolvedValue(new ReadableStream()), + getEventStream: vi + .fn() + .mockResolvedValue(new ReadableStream()), terminateSession: vi.fn(), }; @@ -117,7 +121,9 @@ describe("createSendFn", () => { deliver: vi.fn().mockResolvedValue({ sessionId: "existing-session-id" }), resolveSession: vi.fn(), run: vi.fn().mockResolvedValue(createMockRunHandle()), - getEventStream: vi.fn().mockResolvedValue(new ReadableStream()), + getEventStream: vi + .fn() + .mockResolvedValue(new ReadableStream()), terminateSession: vi.fn(), }; @@ -146,7 +152,9 @@ describe("createSendFn", () => { deliver: vi.fn().mockResolvedValue({ sessionId: "existing-session-id" }), resolveSession: vi.fn(), run: vi.fn().mockResolvedValue(createMockRunHandle()), - getEventStream: vi.fn().mockResolvedValue(new ReadableStream()), + getEventStream: vi + .fn() + .mockResolvedValue(new ReadableStream()), terminateSession: vi.fn(), }; diff --git a/packages/eve/src/channel/session.ts b/packages/eve/src/channel/session.ts index 89f946e6f..f6c3806ac 100644 --- a/packages/eve/src/channel/session.ts +++ b/packages/eve/src/channel/session.ts @@ -1,5 +1,5 @@ import type { ContextAccessor } from "#context/key.js"; -import type { HandleMessageStreamEvent } from "#protocol/message.js"; +import type { StampedHandleMessageStreamEvent } from "#protocol/message.js"; import type { CancelTurnResult, Runtime } from "#channel/types.js"; import type { SessionAuth } from "#context/keys.js"; import { AuthKey, ContinuationTokenKey, InitiatorAuthKey, SessionIdKey } from "#context/keys.js"; @@ -29,7 +29,7 @@ export interface Session { */ getEventStream(options?: { startIndex?: number; - }): Promise>; + }): Promise>; } /** diff --git a/packages/eve/src/channel/types.ts b/packages/eve/src/channel/types.ts index d57472b7c..bb193fa0f 100644 --- a/packages/eve/src/channel/types.ts +++ b/packages/eve/src/channel/types.ts @@ -1,6 +1,9 @@ import type { UserContent } from "ai"; -import type { HandleMessageStreamEvent } from "#protocol/message.js"; +import type { + HandleMessageStreamEvent, + StampedHandleMessageStreamEvent, +} from "#protocol/message.js"; import type { CancelTurnStatus } from "#protocol/cancel-turn.js"; import type { RunMode } from "#shared/run-mode.js"; import type { RuntimeActionResult } from "#runtime/actions/types.js"; @@ -371,7 +374,7 @@ export type RunResult = */ export interface RunHandle { readonly continuationToken: string; - readonly events: ReadableStream; + readonly events: ReadableStream; /** * Runtime-owned identifier for this session. Stream and inspection APIs * key on it: workflow-backed runs expose the workflow run id. @@ -425,7 +428,7 @@ export interface Runtime { getEventStream( sessionId: string, options?: GetEventStreamOptions, - ): Promise>; + ): Promise>; } /** diff --git a/packages/eve/src/cli/dev/tui/runner.test.ts b/packages/eve/src/cli/dev/tui/runner.test.ts index 2a82253fc..cebf3f832 100644 --- a/packages/eve/src/cli/dev/tui/runner.test.ts +++ b/packages/eve/src/cli/dev/tui/runner.test.ts @@ -6,9 +6,11 @@ import { MessageResponse, type AgentInfoResult, type ClientSession, - type HandleMessageStreamEvent, + type StampedHandleMessageStreamEvent, } from "#client/index.js"; +import { stampTestEvent } from "#internal/testing/events.js"; import { resolveTestVercelTarget } from "#internal/testing/verified-vercel-target.js"; +import type { HandleMessageStreamEvent } from "#protocol/message.js"; import { createDevelopmentCredentialGate } from "#services/dev-client/credential-gate.js"; import type { VercelDeploymentResolution } from "#setup/vercel-deployment.js"; @@ -54,17 +56,36 @@ function stubSession(): ClientSession { return stubClient().session(); } -/** Wraps literal stream events in a real `MessageResponse`. */ +/** + * Wraps literal stream events in a real `MessageResponse`. + * + * Each literal is stamped as its own emission, so a fixture that repeats an + * identical payload models two distinct events rather than a re-delivery. + * Re-delivery is modelled by reusing one stamped event — see + * {@link replayed}. + */ function messageResponseOf(events: readonly unknown[]): MessageResponse { + const stamped = events.map((event, index) => + isStamped(event) ? event : stampTestEvent(event as HandleMessageStreamEvent, index), + ); return new MessageResponse({ continuationToken: "eve:test", createStream: async function* () { - for (const event of events) yield event as HandleMessageStreamEvent; + for (const event of stamped) yield event; }, sessionId: "session_test", }); } +function isStamped(event: unknown): event is StampedHandleMessageStreamEvent { + return typeof (event as StampedHandleMessageStreamEvent).meta?.id === "string"; +} + +/** Models a re-delivered durable chunk: the same event, id included. */ +function replayed(event: unknown, index: number): StampedHandleMessageStreamEvent { + return isStamped(event) ? event : stampTestEvent(event as HandleMessageStreamEvent, index); +} + const AGENT_INFO: AgentInfoResult = { agent: { agentRoot: "/tmp/weather-agent/agent", @@ -558,7 +579,15 @@ describe("EveTUIRunner development session continuity", () => { start(controller) { controller.enqueue( encoder.encode( - '{"type":"session.waiting","data":{"continuationToken":"token-session-1","wait":"next-user-message"}}\n', + `${JSON.stringify( + stampTestEvent({ + type: "session.waiting", + data: { + continuationToken: "token-session-1", + wait: "next-user-message", + }, + } as HandleMessageStreamEvent), + )}\n`, ), ); controller.close(); @@ -965,14 +994,14 @@ describe("EveTUIRunner connection authorization", () => { { type: "session.waiting", data: { wait: "connection-authorization" } }, ]); vi.spyOn(session, "stream").mockImplementation(async function* () { - yield { + yield stampTestEvent({ type: "authorization.completed", data: { name: "linear", outcome: "authorized" }, - } as HandleMessageStreamEvent; - yield { + } as HandleMessageStreamEvent); + yield stampTestEvent({ type: "session.waiting", data: { wait: "next-user-message" }, - } as HandleMessageStreamEvent; + } as HandleMessageStreamEvent); }); const renderer: AgentTUIRenderer = { readPrompt: vi.fn(async () => prompts.shift()), @@ -1159,6 +1188,59 @@ describe("EveTUIRunner reused step indexes", () => { }); describe("EveTUIRunner replay guards", () => { + it("renders a re-delivered chunk once even after the next step reuses its key", async () => { + // The hard case: a reconnect re-sends a chunk from the completed step + // after `step.started` opened the next one. Coordinates alone cannot tell + // that apart from a fresh model call reusing `stepIndex`, so without the + // id it renders as a second assistant block. + const appended = replayed( + { + type: "message.appended", + data: { + messageDelta: "Sunny.", + messageSoFar: "Sunny.", + sequence: 0, + stepIndex: 0, + turnId: "turn_0", + }, + }, + 1, + ); + const prompts: Array = ["weather", undefined]; + const emitted: AgentTUIStreamEvent[] = []; + const session = sessionYielding([ + { type: "step.started", data: { sequence: 0, stepIndex: 0, turnId: "turn_0" } }, + appended, + { + type: "message.completed", + data: { + finishReason: "stop", + message: "Sunny.", + sequence: 0, + stepIndex: 0, + turnId: "turn_0", + }, + }, + { type: "step.started", data: { sequence: 1, stepIndex: 0, turnId: "turn_0" } }, + appended, + { type: "session.waiting", data: { wait: "next-user-message" } }, + ]); + const renderer: AgentTUIRenderer = { + readPrompt: vi.fn(async () => prompts.shift()), + renderStream: vi.fn(async (result) => { + for await (const event of result.events as AsyncIterable) { + emitted.push(event); + } + }), + }; + + const runner = new EveTUIRunner({ session, renderer, name: "Weather Agent" }); + await runner.run(); + + const deltas = emitted.filter((event) => event.type === "assistant-delta"); + expect(deltas).toEqual([{ type: "assistant-delta", id: "text:turn_0:0", delta: "Sunny." }]); + }); + it("deduplicates repeated call IDs and divergent text attempts in one turn", async () => { const prompts: Array = ["weather", undefined]; const emitted: AgentTUIStreamEvent[] = []; @@ -1907,20 +1989,26 @@ describe("EveTUIRunner renderer teardown", () => { vi.spyOn(client, "session").mockReturnValue(childSession); vi.spyOn(childSession, "stream").mockImplementation(() => ({ async *[Symbol.asyncIterator]() { - yield { - type: "message.completed", - data: { - finishReason: "stop", - message: "final answer", - sequence: 0, - stepIndex: 0, - turnId: "turn-child", - }, - } as HandleMessageStreamEvent; - yield { - type: "session.waiting", - data: { wait: "next-user-message" }, - } as HandleMessageStreamEvent; + yield stampTestEvent( + { + type: "message.completed", + data: { + finishReason: "stop", + message: "final answer", + sequence: 0, + stepIndex: 0, + turnId: "turn-child", + }, + } as HandleMessageStreamEvent, + 0, + ); + yield stampTestEvent( + { + type: "session.waiting", + data: { wait: "next-user-message" }, + } as HandleMessageStreamEvent, + 1, + ); }, })); @@ -1997,10 +2085,13 @@ describe("EveTUIRunner renderer teardown", () => { await new Promise((resolve) => { signal.addEventListener("abort", () => resolve(), { once: true }); }); - yield { - type: "session.waiting", - data: { wait: "next-user-message" }, - } as HandleMessageStreamEvent; + yield stampTestEvent( + { + type: "session.waiting", + data: { wait: "next-user-message" }, + } as HandleMessageStreamEvent, + 1, + ); }, }; }); @@ -2829,12 +2920,24 @@ describe("EveTUIRunner mid-turn message queue", () => { new MessageResponse({ continuationToken: "eve:test", createStream: async function* () { - yield { type: "turn.started", data: { turnId: "turn-1", sequence: 1 } } as never; + yield stampTestEvent( + { + type: "turn.started", + data: { turnId: "turn-1", sequence: 1 }, + } as HandleMessageStreamEvent, + 0, + ); // Hold the stream mid-turn until the retry loop has proven both // attempts, then settle the turn as cancelled. await gate.promise; - yield { type: "turn.cancelled", data: { turnId: "turn-1", sequence: 2 } } as never; - yield { type: "session.waiting" } as never; + yield stampTestEvent( + { + type: "turn.cancelled", + data: { turnId: "turn-1", sequence: 2 }, + } as HandleMessageStreamEvent, + 1, + ); + yield stampTestEvent({ type: "session.waiting" } as HandleMessageStreamEvent, 2); }, sessionId: "session_test", }), diff --git a/packages/eve/src/cli/dev/tui/runner.ts b/packages/eve/src/cli/dev/tui/runner.ts index 851402668..91807c9b7 100644 --- a/packages/eve/src/cli/dev/tui/runner.ts +++ b/packages/eve/src/cli/dev/tui/runner.ts @@ -5,7 +5,6 @@ import { type AuthorizationCompletedStreamEvent, type ConnectionAuthorizationOutcome, type AuthorizationRequiredStreamEvent, - type HandleMessageStreamEvent, type InputOption, type InputRequest, type InputRequestedStreamEvent, @@ -14,6 +13,7 @@ import { type ReasoningAppendedStreamEvent, type SessionFailedStreamEvent, type StepCompletedStreamEvent, + type StampedHandleMessageStreamEvent, type SubagentCalledStreamEvent, type SubagentCompletedStreamEvent, Client, @@ -21,6 +21,7 @@ import { } from "#client/index.js"; import { loadDevelopmentEnvironmentFiles } from "#cli/dev/environment.js"; import { subscribeDevelopmentSandboxPrewarmLogs } from "#execution/sandbox/development-prewarm.js"; +import { createEventDeduper } from "#protocol/event-dedupe.js"; import { createDevelopmentRuntimeArtifactRefresher, type DevelopmentRuntimeArtifactRefresher, @@ -1184,7 +1185,7 @@ export class EveTUIRunner { } #createTUIStreamResult( - events: AsyncIterable, + events: AsyncIterable, abort: () => void, ): AgentTUIStreamResult { const turnState = createTurnState(); @@ -1632,7 +1633,7 @@ function formatAgentUpdateNotice( } type EveStreamTranslatorInput = { - events: AsyncIterable; + events: AsyncIterable; pendingInputRequests: Map; turnState: AgentTUITurnState; onSubagentCalled?: (event: SubagentCalledStreamEvent) => void; @@ -1671,11 +1672,15 @@ async function* eveEventsToTUIStream( } = input; const textParts = new Map(); const reasoningParts = new Map(); + // Re-delivery of a durable chunk — a reconnect, or a rewind to replay the + // turn so far — carries the id it was emitted with. Dropping those here + // means every case below is a genuinely new emission. + const seenEvents = createEventDeduper(); // Counts `step.started` events. The harness reuses `stepIndex` across the // model calls of one turn (e.g. the post-subagent call restarts at the same - // index), so a part key alone cannot distinguish "new message under a - // reused key" from "replayed events of the finished message". A fresh - // `step.started` since the part completed is the discriminator. + // index), so a part key alone cannot distinguish a new message under a + // reused key from a re-emission of the finished one. A fresh `step.started` + // since the part completed is the discriminator. let stepEpoch = 0; const knownToolCalls = new Set(); const seenInputRequestIds = new Set(); @@ -1688,6 +1693,10 @@ async function* eveEventsToTUIStream( let latestStepUsage: StepCompletedStreamEvent["data"]["usage"] | undefined; for await (const event of events) { + if (seenEvents.isDuplicate(event)) { + continue; + } + if (visibleTurnCompleted && isPostTurnVisibleEvent(event)) { continue; } @@ -1726,10 +1735,9 @@ async function* eveEventsToTUIStream( const next = appended.data.messageSoFar; if (state.completed) { - // Replays of the finished message re-stream prefixes of it — drop. - if (state.text.startsWith(next)) break; - // Divergent text without an intervening `step.started` is a retry - // of the same model call — drop it rather than mixing attempts. + // Text under a completed key without an intervening `step.started` + // is a retry of the same model call — drop it rather than mixing + // attempts. if (stepEpoch <= state.completedEpoch) break; // A fresh model call reusing this part key (the harness restarts // `stepIndex` after a park/resume, e.g. post-subagent): open a new @@ -1755,7 +1763,7 @@ async function* eveEventsToTUIStream( const message = event.data.message; if (state.completed) { - if (message === null || message === state.text) break; + if (message === null) break; if (stepEpoch <= state.completedEpoch) break; // Channels that skip per-delta events: a new full message under a // reused key after a fresh model call. @@ -1802,7 +1810,6 @@ async function* eveEventsToTUIStream( const next = appended.data.reasoningSoFar; if (state.completed) { - if (state.text.startsWith(next)) break; if (stepEpoch <= state.completedEpoch) break; state.generation += 1; state.text = ""; @@ -1825,7 +1832,7 @@ async function* eveEventsToTUIStream( const next = event.data.reasoning; if (state.completed) { - if (next.length === 0 || next === state.text || state.text.startsWith(next)) break; + if (next.length === 0) break; if (stepEpoch <= state.completedEpoch) break; state.generation += 1; state.text = next; @@ -2144,7 +2151,7 @@ function* closeOpenParts( } } -function isPostTurnVisibleEvent(event: HandleMessageStreamEvent): boolean { +function isPostTurnVisibleEvent(event: StampedHandleMessageStreamEvent): boolean { switch (event.type) { case "actions.requested": case "authorization.completed": diff --git a/packages/eve/src/cli/dev/tui/subagent-pump.test.ts b/packages/eve/src/cli/dev/tui/subagent-pump.test.ts index d425c5089..40038759c 100644 --- a/packages/eve/src/cli/dev/tui/subagent-pump.test.ts +++ b/packages/eve/src/cli/dev/tui/subagent-pump.test.ts @@ -1,7 +1,8 @@ import { describe, expect, it, vi } from "vitest"; -import { Client, type HandleMessageStreamEvent } from "#client/index.js"; -import type { SubagentCalledStreamEvent } from "#protocol/message.js"; +import { Client, type StampedHandleMessageStreamEvent } from "#client/index.js"; +import { stampTestEvent } from "#internal/testing/events.js"; +import type { HandleMessageStreamEvent, SubagentCalledStreamEvent } from "#protocol/message.js"; import { SubagentPump, type SubagentView } from "./subagent-pump.js"; @@ -21,17 +22,17 @@ function fakeView(): SubagentView { * the pump must never reach the view. */ function pushableChildStream() { - const queue: HandleMessageStreamEvent[] = []; + const queue: StampedHandleMessageStreamEvent[] = []; let wake: (() => void) | undefined; let aborted = false; return { - push(event: HandleMessageStreamEvent) { + push(event: StampedHandleMessageStreamEvent) { queue.push(event); wake?.(); wake = undefined; }, - stream(options?: { signal?: AbortSignal }): AsyncIterable { + stream(options?: { signal?: AbortSignal }): AsyncIterable { const signal = options?.signal; return { async *[Symbol.asyncIterator]() { @@ -72,17 +73,29 @@ function subagentCalled(callId: string): SubagentCalledStreamEvent { } as SubagentCalledStreamEvent; } -function reasoningEvent(delta: string): HandleMessageStreamEvent { - return { - type: "reasoning.appended", - data: { - reasoningDelta: delta, - reasoningSoFar: delta, - sequence: 2, - stepIndex: 0, - turnId: "child-turn", - }, - } as HandleMessageStreamEvent; +let nextChildEventIndex = 0; + +function reasoningEvent(delta: string): StampedHandleMessageStreamEvent { + return stampTestEvent( + { + type: "reasoning.appended", + data: { + reasoningDelta: delta, + reasoningSoFar: delta, + sequence: 2, + stepIndex: 0, + turnId: "child-turn", + }, + } as HandleMessageStreamEvent, + nextChildEventIndex++, + ); +} + +function boundaryEvent(): StampedHandleMessageStreamEvent { + return stampTestEvent( + { type: "session.waiting", data: { wait: "next-user-message" } } as HandleMessageStreamEvent, + nextChildEventIndex++, + ); } async function settleAsyncWork(): Promise { @@ -125,3 +138,48 @@ describe("SubagentPump.settleAll", () => { expect(vi.mocked(view.complete).mock.calls.length).toBe(completions); }); }); + +describe("SubagentPump child stream replay", () => { + it("folds a replayed child transcript in once when the pump restarts", async () => { + // A pump is dropped from the registry when its stream ends, so a second + // `subagent.called` for the same call — an SSE-resume re-entry — reopens + // the child at `streamIndex: 0` while the run's sections survive. + const transcript = [reasoningEvent("looked up the forecast"), boundaryEvent()]; + let opened = 0; + const client = new Client({ host: "http://localhost:3000" }); + vi.spyOn(client, "session").mockReturnValue({ + stream: () => { + opened += 1; + return { + async *[Symbol.asyncIterator]() { + for (const event of transcript) yield event; + }, + }; + }, + } as never); + const view = fakeView(); + const pump = new SubagentPump({ client, view, formatActionResultError: () => "failed" }); + + const called = subagentCalled("call-1"); + pump.begin(called); + await settleAsyncWork(); + pump.begin(called); + await settleAsyncWork(); + + expect(opened).toBe(2); + const reasoning = vi + .mocked(view.upsertStep) + .mock.calls.map(([update]) => update.reasoning) + .filter((text): text is string => text !== undefined && text.length > 0); + expect(reasoning.length).toBeGreaterThan(0); + // Every paint carries the transcript once, and the replay opened no + // second section to paint it into again. + for (const text of reasoning) { + expect(text).toBe("looked up the forecast"); + } + const sectionKeys = new Set( + vi.mocked(view.upsertStep).mock.calls.map(([update]) => update.sectionKey), + ); + expect(sectionKeys.size).toBe(1); + }); +}); diff --git a/packages/eve/src/cli/dev/tui/subagent-pump.ts b/packages/eve/src/cli/dev/tui/subagent-pump.ts index 9fa8d4205..d863c1299 100644 --- a/packages/eve/src/cli/dev/tui/subagent-pump.ts +++ b/packages/eve/src/cli/dev/tui/subagent-pump.ts @@ -7,10 +7,11 @@ */ import type { Client } from "#client/index.js"; +import { createEventDeduper, type EventDeduper } from "#protocol/event-dedupe.js"; import { isCurrentTurnBoundaryEvent, type ActionResultStreamEvent, - type HandleMessageStreamEvent, + type StampedHandleMessageStreamEvent, type SubagentCalledStreamEvent, } from "#protocol/message.js"; import { toErrorMessage } from "#shared/errors.js"; @@ -87,6 +88,14 @@ export type SubagentRun = { /** Monotonic counter for new section keys. */ nextSectionKey: number; tools: Map; + /** + * Child events already folded into this run. A pump is removed from + * `#pumps` when its stream ends, so a later `subagent.called` for the same + * call — an SSE-resume re-entry, or a recovery after a dropped child + * stream — restarts the pump at `streamIndex: 0` while `steps` survives. + * Keyed on the durable event id, the replayed prefix folds in exactly once. + */ + seenChildEvents: EventDeduper; }; export type SubagentStepUpdate = { @@ -148,6 +157,7 @@ export class SubagentPump { currentSectionKey: null, nextSectionKey: 0, tools: new Map(), + seenChildEvents: createEventDeduper(), }); } else { existing.name = called.data.name; @@ -358,9 +368,10 @@ export class SubagentPump { } } - #applyChildEvent(callId: string, event: HandleMessageStreamEvent) { + #applyChildEvent(callId: string, event: StampedHandleMessageStreamEvent) { const run = this.#runs.get(callId); if (!run) return; + if (run.seenChildEvents.isDuplicate(event)) return; // A child event after settle is a HITL-parked turn resuming: reopen the // run explicitly (begin clears the header's Done mark) instead of // mutating a completed section by accident. diff --git a/packages/eve/src/client/eve-agent-store.test.ts b/packages/eve/src/client/eve-agent-store.test.ts new file mode 100644 index 000000000..9ff27f4f8 --- /dev/null +++ b/packages/eve/src/client/eve-agent-store.test.ts @@ -0,0 +1,126 @@ +import { afterEach, describe, expect, it, vi } from "vitest"; + +import { EveAgentStore } from "#client/eve-agent-store.js"; +import { defaultMessageReducer } from "#client/message-reducer.js"; +import { stampTestEvents } from "#internal/testing/events.js"; +import { + createMessageCompletedEvent, + createMessageReceivedEvent, + createSessionWaitingEvent, + EVE_SESSION_ID_HEADER, + type HandleMessageStreamEvent, + type StampedHandleMessageStreamEvent, +} from "#protocol/message.js"; + +function turnEvents(): StampedHandleMessageStreamEvent[] { + return stampTestEvents([ + createMessageReceivedEvent({ message: "Hello", sequence: 0, turnId: "turn_1" }), + createMessageCompletedEvent({ + finishReason: "stop", + message: "Hi there.", + sequence: 1, + stepIndex: 0, + turnId: "turn_1", + }), + createSessionWaitingEvent("http:session_1"), + ] as HandleMessageStreamEvent[]); +} + +function startedResponse(): Response { + return new Response( + JSON.stringify({ continuationToken: "http:session_1", ok: true, sessionId: "session_1" }), + { + headers: { "content-type": "application/json", [EVE_SESSION_ID_HEADER]: "session_1" }, + status: 202, + }, + ); +} + +function streamResponse(events: readonly StampedHandleMessageStreamEvent[]): Response { + const encoder = new TextEncoder(); + return new Response( + new ReadableStream({ + start(controller) { + for (const event of events) { + controller.enqueue(encoder.encode(`${JSON.stringify(event)}\n`)); + } + controller.close(); + }, + }), + ); +} + +afterEach(() => { + vi.restoreAllMocks(); + vi.unstubAllGlobals(); +}); + +describe("EveAgentStore stream overlap", () => { + it("folds an initialEvents prefix that the live stream re-delivers in once", async () => { + const events = turnEvents(); + vi.spyOn(globalThis, "fetch") + .mockResolvedValueOnce(startedResponse()) + // The server-rendered prefix is replayed ahead of the live tail. + .mockResolvedValueOnce(streamResponse(events)); + + const store = new EveAgentStore({ + initialEvents: events.slice(0, 2), + reducer: defaultMessageReducer(), + }); + + const seen: StampedHandleMessageStreamEvent[] = []; + store.setCallbacks({ onEvent: (event) => seen.push(event) }); + + await store.send({ message: "Hello" }); + + // Only the events the prefix did not already carry reach subscribers. + expect(seen.map((event) => event.meta.id)).toEqual([events[2]?.meta.id]); + expect(store.snapshot.events.map((event) => event.meta.id)).toEqual( + events.map((event) => event.meta.id), + ); + + const assistant = store.snapshot.data.messages.filter( + (message) => message.role === "assistant", + ); + expect(assistant).toHaveLength(1); + expect(assistant[0]?.parts).toEqual([ + { type: "step-start" }, + { state: "done", stepIndex: 0, text: "Hi there.", type: "text" }, + ]); + }); + + it("drops a duplicated initialEvents entry", () => { + const events = turnEvents(); + const store = new EveAgentStore({ + initialEvents: [...events, ...events], + reducer: defaultMessageReducer(), + }); + + expect(store.snapshot.events.map((event) => event.meta.id)).toEqual( + events.map((event) => event.meta.id), + ); + }); + + it("re-admits events after reset clears the window", async () => { + const events = turnEvents(); + const store = new EveAgentStore({ + initialEvents: events, + reducer: defaultMessageReducer(), + }); + expect(store.snapshot.events).toHaveLength(3); + + store.reset(); + expect(store.snapshot.events).toHaveLength(0); + + vi.spyOn(globalThis, "fetch") + .mockResolvedValueOnce(startedResponse()) + .mockResolvedValueOnce(streamResponse(events)); + + await store.send({ message: "Hello" }); + + // A fresh session must not have the retired ids held against it. + expect(store.snapshot.events.map((event) => event.meta.id)).toEqual( + events.map((event) => event.meta.id), + ); + }); +}); diff --git a/packages/eve/src/client/eve-agent-store.ts b/packages/eve/src/client/eve-agent-store.ts index 72c3e0edf..e9b87c216 100644 --- a/packages/eve/src/client/eve-agent-store.ts +++ b/packages/eve/src/client/eve-agent-store.ts @@ -1,7 +1,8 @@ import { Client } from "#client/client.js"; import type { EveAgentReducer, EveAgentReducerEvent } from "#client/reducer.js"; import type { ClientSession } from "#client/session.js"; -import type { HandleMessageStreamEvent } from "#protocol/message.js"; +import { createEventDeduper } from "#protocol/event-dedupe.js"; +import type { StampedHandleMessageStreamEvent } from "#protocol/message.js"; import { toError } from "#shared/errors.js"; import type { ClientAuth, HeadersValue, SendTurnPayload, SessionState } from "#client/types.js"; import type { UserContent } from "ai"; @@ -30,7 +31,7 @@ export type PrepareSend = (input: SendTurnPayload) => SendTurnPayload | Promise< export interface EveAgentStoreSnapshot { readonly data: TData; readonly error: Error | undefined; - readonly events: readonly HandleMessageStreamEvent[]; + readonly events: readonly StampedHandleMessageStreamEvent[]; readonly session: SessionState; readonly status: EveAgentStoreStatus; } @@ -44,7 +45,7 @@ export interface EveAgentStoreSnapshot { */ export interface EveAgentStoreCallbacks { readonly onError?: (error: Error) => void; - readonly onEvent?: (event: HandleMessageStreamEvent) => void; + readonly onEvent?: (event: StampedHandleMessageStreamEvent) => void; readonly onFinish?: (snapshot: EveAgentStoreSnapshot) => void; readonly onSessionChange?: (session: SessionState) => void; readonly prepareSend?: PrepareSend; @@ -66,7 +67,7 @@ export interface EveAgentStoreInit { readonly auth?: ClientAuth; readonly headers?: HeadersValue; readonly host?: string; - readonly initialEvents?: readonly HandleMessageStreamEvent[]; + readonly initialEvents?: readonly StampedHandleMessageStreamEvent[]; readonly initialSession?: SessionState; readonly optimistic?: boolean; readonly reducer: EveAgentReducer; @@ -97,11 +98,19 @@ export class EveAgentStore { readonly #reducer: EveAgentReducer; readonly #subscribers = new Set<() => void>(); + /** + * Events already folded into the projection. `initialEvents` is typically a + * server-rendered prefix that the live stream then overlaps, and a + * reconnect can re-deliver the chunk it resumed from — both re-deliver the + * ids they were emitted with. + */ + #seenEvents = createEventDeduper(); + #abortController: AbortController | undefined; #callbacks: EveAgentStoreCallbacks = {}; #data: TData; #error: Error | undefined; - #events: readonly HandleMessageStreamEvent[]; + #events: readonly StampedHandleMessageStreamEvent[]; #operationId = 0; #pendingMessageSubmission: PendingMessageSubmission | undefined; #projectionEvents: readonly EveAgentReducerEvent[]; @@ -118,7 +127,9 @@ export class EveAgentStore { headers: init.headers, host: init.host ?? "", }).session(init.initialSession); - this.#events = [...(init.initialEvents ?? [])]; + this.#events = (init.initialEvents ?? []).filter( + (event) => !this.#seenEvents.isDuplicate(event), + ); this.#projectionEvents = [...this.#events]; this.#optimistic = init.optimistic ?? true; this.#reducer = init.reducer; @@ -182,6 +193,10 @@ export class EveAgentStore { this.#status = "streaming"; } + if (this.#seenEvents.isDuplicate(event)) { + continue; + } + this.#events = [...this.#events, event]; this.#applyServerEvent(event); this.#callbacks.onEvent?.(event); @@ -228,6 +243,7 @@ export class EveAgentStore { this.#abortController = undefined; this.#session = this.#createSession?.() ?? this.#session; this.#events = []; + this.#seenEvents = createEventDeduper(); this.#pendingMessageSubmission = undefined; this.#projectionEvents = []; this.#data = this.#reducer.initial(); @@ -293,7 +309,7 @@ export class EveAgentStore { }); } - #applyServerEvent(event: HandleMessageStreamEvent): void { + #applyServerEvent(event: StampedHandleMessageStreamEvent): void { if (event.type === "message.received" && this.#pendingMessageSubmission !== undefined) { const submissionId = this.#pendingMessageSubmission.id; this.#pendingMessageSubmission = undefined; @@ -309,7 +325,7 @@ export class EveAgentStore { this.#appendProjectionEvent(event); } - #applyTerminalStreamFailure(event: HandleMessageStreamEvent): void { + #applyTerminalStreamFailure(event: StampedHandleMessageStreamEvent): void { const error = toTerminalStreamFailureError(event); if (error === undefined) { return; @@ -439,7 +455,7 @@ function isAbortError(error: unknown): boolean { return error instanceof Error && error.name === "AbortError"; } -function toTerminalStreamFailureError(event: HandleMessageStreamEvent): Error | undefined { +function toTerminalStreamFailureError(event: StampedHandleMessageStreamEvent): Error | undefined { if (event.type !== "session.failed") { return undefined; } diff --git a/packages/eve/src/client/index.ts b/packages/eve/src/client/index.ts index fa0c477b7..db97d4aa4 100644 --- a/packages/eve/src/client/index.ts +++ b/packages/eve/src/client/index.ts @@ -97,6 +97,7 @@ export type { ConnectionAuthorizationOutcome, AuthorizationRequiredStreamEvent, HandleMessageStreamEvent, + HandleMessageStreamEventMeta, InputRequestedStreamEvent, MessageAppendedStreamEvent, MessageCompletedStreamEvent, @@ -109,6 +110,7 @@ export type { SessionFailedStreamEvent, SessionStartedStreamEvent, SessionWaitingStreamEvent, + StampedHandleMessageStreamEvent, StepCompletedStreamEvent, StepFailedStreamEvent, StepStartedStreamEvent, diff --git a/packages/eve/src/client/message-response.ts b/packages/eve/src/client/message-response.ts index f3fe1573e..dd9c55641 100644 --- a/packages/eve/src/client/message-response.ts +++ b/packages/eve/src/client/message-response.ts @@ -1,4 +1,4 @@ -import type { HandleMessageStreamEvent } from "#protocol/message.js"; +import type { StampedHandleMessageStreamEvent } from "#protocol/message.js"; import { extractCompletedResult } from "#client/output-schema.js"; import { deriveResultStatus, @@ -12,7 +12,7 @@ import type { MessageResult } from "#client/types.js"; */ interface MessageResponseInput { readonly continuationToken?: string; - readonly createStream: () => AsyncGenerator; + readonly createStream: () => AsyncGenerator; readonly sessionId: string; } @@ -23,7 +23,9 @@ interface MessageResponseInput { * token) as soon as the POST completes. Collect the event stream via * {@link result} or iterate it with `for await...of`. */ -export class MessageResponse implements AsyncIterable { +export class MessageResponse< + TOutput = unknown, +> implements AsyncIterable { /** * Continuation token returned by the server for follow-up messages. */ @@ -35,7 +37,7 @@ export class MessageResponse implements AsyncIterable AsyncGenerator; + readonly #createStream: () => AsyncGenerator; /** @internal */ constructor(input: MessageResponseInput) { @@ -49,7 +51,7 @@ export class MessageResponse implements AsyncIterable> { - const events: HandleMessageStreamEvent[] = []; + const events: StampedHandleMessageStreamEvent[] = []; for await (const event of this) { events.push(event); @@ -70,7 +72,7 @@ export class MessageResponse implements AsyncIterable { + [Symbol.asyncIterator](): AsyncIterator { if (this.#consumed) { throw new Error("MessageResponse has already been consumed."); } diff --git a/packages/eve/src/client/ndjson.ts b/packages/eve/src/client/ndjson.ts index 3187a369a..9b36adbe9 100644 --- a/packages/eve/src/client/ndjson.ts +++ b/packages/eve/src/client/ndjson.ts @@ -1,4 +1,4 @@ -import type { HandleMessageStreamEvent } from "#protocol/message.js"; +import type { StampedHandleMessageStreamEvent } from "#protocol/message.js"; /** * Returns true when an error looks like a stream socket disconnection that @@ -27,7 +27,7 @@ export function isStreamDisconnectError(error: unknown): boolean { /** * Reads newline-delimited JSON events from a `ReadableStream`. * - * Yields one parsed {@link HandleMessageStreamEvent} per complete NDJSON line. + * Yields one parsed {@link StampedHandleMessageStreamEvent} per complete NDJSON line. * Handles partial lines across chunks via an internal buffer. * * All read errors — including socket disconnections — propagate to the caller. @@ -35,7 +35,7 @@ export function isStreamDisconnectError(error: unknown): boolean { */ export async function* readNdjsonStream( body: ReadableStream, -): AsyncGenerator { +): AsyncGenerator { const reader = body.getReader(); const decoder = new TextDecoder(); let buffer = ""; @@ -63,7 +63,7 @@ export async function* readNdjsonStream( buffer = buffer.slice(newlineIndex + 1); if (line.length > 0) { - yield JSON.parse(line) as HandleMessageStreamEvent; + yield JSON.parse(line) as StampedHandleMessageStreamEvent; } newlineIndex = buffer.indexOf("\n"); @@ -73,7 +73,7 @@ export async function* readNdjsonStream( // Yield any trailing content without a final newline. const trailing = buffer.trim(); if (trailing.length > 0) { - yield JSON.parse(trailing) as HandleMessageStreamEvent; + yield JSON.parse(trailing) as StampedHandleMessageStreamEvent; } } finally { if (!reachedEof) { diff --git a/packages/eve/src/client/open-stream.ts b/packages/eve/src/client/open-stream.ts index e5d01265c..819ce76be 100644 --- a/packages/eve/src/client/open-stream.ts +++ b/packages/eve/src/client/open-stream.ts @@ -1,4 +1,4 @@ -import type { HandleMessageStreamEvent } from "#protocol/message.js"; +import type { StampedHandleMessageStreamEvent } from "#protocol/message.js"; import { createEveMessageStreamRoutePath } from "#protocol/routes.js"; import { ClientError } from "#client/client-error.js"; import { isStreamDisconnectError, readNdjsonStream } from "#client/ndjson.js"; @@ -94,7 +94,7 @@ interface FollowStreamInput { */ export async function* followStreamIterable( input: FollowStreamInput, -): AsyncGenerator { +): AsyncGenerator { const retryPolicy = resolveStreamReconnectPolicy(input.streamReconnectPolicy); const idleRetryPolicy = retryPolicy.streamIdleReconnectPolicy; let startIndex = input.startIndex; diff --git a/packages/eve/src/client/session.ts b/packages/eve/src/client/session.ts index e07096433..c5fec60a0 100644 --- a/packages/eve/src/client/session.ts +++ b/packages/eve/src/client/session.ts @@ -1,4 +1,4 @@ -import type { HandleMessageStreamEvent } from "#protocol/message.js"; +import type { StampedHandleMessageStreamEvent } from "#protocol/message.js"; import { EVE_SESSION_ID_HEADER, isCurrentTurnBoundaryEvent } from "#protocol/message.js"; import { CancelTurnResponseSchema } from "#protocol/cancel-turn.js"; import { ResetResponseSchema } from "#protocol/reset-session.js"; @@ -232,7 +232,7 @@ export class ClientSession { * @throws {Error} If the session has no session ID (no message has been sent * yet). */ - stream(options?: StreamOptions): AsyncIterable { + stream(options?: StreamOptions): AsyncIterable { const sessionId = this.#state.sessionId; if (!sessionId) { @@ -302,8 +302,8 @@ export class ClientSession { continuationToken: string | undefined, initialState: SessionState, input: SendTurnPayload, - ): AsyncGenerator { - const events: HandleMessageStreamEvent[] = []; + ): AsyncGenerator { + const events: StampedHandleMessageStreamEvent[] = []; try { for await (const event of followStreamIterable({ @@ -336,10 +336,10 @@ export class ClientSession { async *#streamAndAdvance( sessionId: string, options?: StreamOptions, - ): AsyncGenerator { + ): AsyncGenerator { const initialState = this.#state; const streamIndex = options?.startIndex ?? initialState.streamIndex; - const events: HandleMessageStreamEvent[] = []; + const events: StampedHandleMessageStreamEvent[] = []; try { for await (const event of followStreamIterable({ diff --git a/packages/eve/src/client/types.ts b/packages/eve/src/client/types.ts index 9ed7d48ec..5092bf759 100644 --- a/packages/eve/src/client/types.ts +++ b/packages/eve/src/client/types.ts @@ -1,7 +1,7 @@ import type { UserContent } from "ai"; import type { StandardJSONSchemaV1 } from "#compiled/@standard-schema/spec/index.js"; -import type { HandleMessageStreamEvent } from "#protocol/message.js"; +import type { StampedHandleMessageStreamEvent } from "#protocol/message.js"; import type { CancelTurnStatus } from "#protocol/cancel-turn.js"; import type { ResetStatus } from "#protocol/reset-session.js"; import type { InputRequest, InputResponse } from "#runtime/input/types.js"; @@ -261,7 +261,7 @@ export interface MessageResult { /** * All events received during this turn. */ - readonly events: HandleMessageStreamEvent[]; + readonly events: StampedHandleMessageStreamEvent[]; /** * HITL input requests emitted during this turn. diff --git a/packages/eve/src/compiler/extension-compatibility.ts b/packages/eve/src/compiler/extension-compatibility.ts index 13a8d2d47..adbbc5c92 100644 --- a/packages/eve/src/compiler/extension-compatibility.ts +++ b/packages/eve/src/compiler/extension-compatibility.ts @@ -22,13 +22,13 @@ interface ExtensionCapabilityContract { const EXTENSION_CAPABILITY_CONTRACTS = { extension: { current: 1, supported: [1], dropped: {} }, tool: { current: 2, supported: [1, 2], dropped: {} }, - dynamicTool: { current: 3, supported: [1, 2, 3], dropped: {} }, + dynamicTool: { current: 4, supported: [1, 2, 3, 4], dropped: {} }, connection: { current: 2, supported: [1, 2], dropped: {} }, - hook: { current: 2, supported: [1, 2], dropped: {} }, + hook: { current: 3, supported: [1, 2, 3], dropped: {} }, skill: { current: 1, supported: [1], dropped: {} }, - dynamicSkill: { current: 1, supported: [1], dropped: {} }, + dynamicSkill: { current: 2, supported: [1, 2], dropped: {} }, instructions: { current: 1, supported: [1], dropped: {} }, - dynamicInstructions: { current: 1, supported: [1], dropped: {} }, + dynamicInstructions: { current: 2, supported: [1, 2], dropped: {} }, config: { current: 1, supported: [1], dropped: {} }, state: { current: 2, supported: [1, 2], dropped: {} }, } as const satisfies Record; diff --git a/packages/eve/src/context/hook-lifecycle.integration.test.ts b/packages/eve/src/context/hook-lifecycle.integration.test.ts index d52f33aca..0d202a153 100644 --- a/packages/eve/src/context/hook-lifecycle.integration.test.ts +++ b/packages/eve/src/context/hook-lifecycle.integration.test.ts @@ -3,6 +3,7 @@ import { describe, expect, it } from "vitest"; import { createRuntimeHookRegistry } from "#runtime/hooks/registry.js"; import type { ResolvedHookDefinition } from "#runtime/types.js"; import type { HandleMessageStreamEvent } from "#protocol/message.js"; +import { stampTestEvent } from "#internal/testing/events.js"; import { ContextContainer, contextStorage } from "./container.js"; import { dispatchStreamEventHooks } from "./hook-lifecycle.js"; import { @@ -77,7 +78,7 @@ describe("dispatchStreamEventHooks", () => { dispatchStreamEventHooks({ ctx, registry, - event: { type: "session.completed" } satisfies HandleMessageStreamEvent, + event: stampTestEvent({ type: "session.completed" }), }), ); expect(calls).toEqual(["typed", "wildcard:session.completed"]); @@ -96,9 +97,34 @@ describe("dispatchStreamEventHooks", () => { dispatchStreamEventHooks({ ctx, registry: brokenRegistry, - event: { type: "session.completed" } satisfies HandleMessageStreamEvent, + event: stampTestEvent({ type: "session.completed" }), }), ), ).rejects.toThrow(/event hook boom/); }); + + it("hands subscribers the stamped envelope so a hook can dedupe its own work", async () => { + const seen: string[] = []; + const registry = createRuntimeHookRegistry([ + hook("audit", { + events: { + // Typed and wildcard subscribers see one event, so both agree on the + // key a hook would write to a database. + "session.completed": async (event) => { + seen.push(event.meta.id); + }, + "*": async (event) => { + seen.push(event.meta.id); + }, + }, + }), + ]); + const ctx = buildCtx(); + const event = stampTestEvent({ type: "session.completed" }); + + await contextStorage.run(ctx, () => dispatchStreamEventHooks({ ctx, registry, event })); + + expect(seen).toEqual([event.meta.id, event.meta.id]); + expect(event.meta.id).toBe("evt_test_0000"); + }); }); diff --git a/packages/eve/src/context/hook-lifecycle.ts b/packages/eve/src/context/hook-lifecycle.ts index ea6ead904..ccad09295 100644 --- a/packages/eve/src/context/hook-lifecycle.ts +++ b/packages/eve/src/context/hook-lifecycle.ts @@ -1,5 +1,5 @@ import { getAdapterKind } from "#channel/adapter.js"; -import type { HandleMessageStreamEvent } from "#protocol/message.js"; +import type { StampedHandleMessageStreamEvent } from "#protocol/message.js"; import type { HookContext } from "#public/definitions/hook.js"; import type { RuntimeHookRegistry } from "#runtime/hooks/registry.js"; import { buildCallbackContext } from "#context/build-callback-context.js"; @@ -16,7 +16,7 @@ import { ContinuationTokenKey } from "./keys.js"; export async function dispatchStreamEventHooks(input: { readonly ctx: ContextContainer; readonly registry: RuntimeHookRegistry; - readonly event: HandleMessageStreamEvent; + readonly event: StampedHandleMessageStreamEvent; }): Promise { const typed = input.registry.streamEventsByType.get(input.event.type) ?? []; const wildcard = input.registry.streamEventsWildcard; diff --git a/packages/eve/src/execution/dispatch-runtime-actions-step.ts b/packages/eve/src/execution/dispatch-runtime-actions-step.ts index 370b11a75..435425f52 100644 --- a/packages/eve/src/execution/dispatch-runtime-actions-step.ts +++ b/packages/eve/src/execution/dispatch-runtime-actions-step.ts @@ -24,7 +24,7 @@ import { import { createSubagentCalledEvent, encodeMessageStreamEvent, - timestampHandleMessageStreamEvent, + stampMessageStreamEvent, } from "#protocol/message.js"; import type { RuntimeActionRequest, @@ -218,20 +218,22 @@ export async function dispatchRuntimeActionsStep(input: { const parentEvent = await callAdapterEventHandler( adapter, - createSubagentCalledEvent({ - callId: action.callId, - childSessionId, - name, - remote, - sequence: batch.event.sequence, - sessionId: session.sessionId, - toolName, - turnId: batch.event.turnId, - workflowId: workflowEntryReference.workflowId, - }), + stampMessageStreamEvent( + createSubagentCalledEvent({ + callId: action.callId, + childSessionId, + name, + remote, + sequence: batch.event.sequence, + sessionId: session.sessionId, + toolName, + turnId: batch.event.turnId, + workflowId: workflowEntryReference.workflowId, + }), + ), adapterCtx, ); - await writer.write(encodeMessageStreamEvent(timestampHandleMessageStreamEvent(parentEvent))); + await writer.write(encodeMessageStreamEvent(parentEvent)); } } finally { writer.releaseLock(); diff --git a/packages/eve/src/execution/settle-cancelled-turn-step.ts b/packages/eve/src/execution/settle-cancelled-turn-step.ts index eb6544458..158a08730 100644 --- a/packages/eve/src/execution/settle-cancelled-turn-step.ts +++ b/packages/eve/src/execution/settle-cancelled-turn-step.ts @@ -26,7 +26,7 @@ import { clearPendingWorkflowInterrupt } from "#harness/workflow-interrupt-state import { encodeMessageStreamEvent, type HandleMessageStreamEvent, - timestampHandleMessageStreamEvent, + stampMessageStreamEvent, } from "#protocol/message.js"; import { BundleKey, ChannelKey } from "#runtime/sessions/runtime-context-keys.js"; @@ -74,11 +74,13 @@ export async function settleCancelledTurnStep(input: { try { const scoped = await withContextScope(ctx, session, async (enrichedSession) => { const emit = async (event: HandleMessageStreamEvent): Promise => { - const transformed = await callAdapterEventHandler(adapter, event, adapterCtx); - setChannelContext(ctx, { ...adapter, state: { ...adapterCtx.state } }); - await writer.write( - encodeMessageStreamEvent(timestampHandleMessageStreamEvent(transformed)), + const transformed = await callAdapterEventHandler( + adapter, + stampMessageStreamEvent(event), + adapterCtx, ); + setChannelContext(ctx, { ...adapter, state: { ...adapterCtx.state } }); + await writer.write(encodeMessageStreamEvent(transformed)); await dispatchStreamEventHooks({ ctx, event: transformed, diff --git a/packages/eve/src/execution/subagent-adapter.test.ts b/packages/eve/src/execution/subagent-adapter.test.ts index 9daa0f835..30adb5560 100644 --- a/packages/eve/src/execution/subagent-adapter.test.ts +++ b/packages/eve/src/execution/subagent-adapter.test.ts @@ -7,6 +7,7 @@ import { ContextContainer } from "#context/container.js"; import { ContinuationTokenKey, SessionIdKey } from "#context/keys.js"; import type { InputRequest } from "#runtime/input/types.js"; import { SUBAGENT_ADAPTER } from "#execution/subagent-adapter.js"; +import { stampTestEvent } from "#internal/testing/events.js"; const SUBAGENT_INPUT_REQUESTED = SUBAGENT_ADAPTER["input.requested"]; const SUBAGENT_AUTHORIZATION_REQUIRED = SUBAGENT_ADAPTER["authorization.required"]; @@ -83,7 +84,7 @@ describe("SUBAGENT_ADAPTER authorization handlers", () => { await callAdapterEventHandler( SUBAGENT_ADAPTER, - { data, type: "authorization.required" }, + stampTestEvent({ data, type: "authorization.required" }), makeContext(), ); diff --git a/packages/eve/src/execution/subagent-auth-proxy.integration.test.ts b/packages/eve/src/execution/subagent-auth-proxy.integration.test.ts index de3805ba3..df899575f 100644 --- a/packages/eve/src/execution/subagent-auth-proxy.integration.test.ts +++ b/packages/eve/src/execution/subagent-auth-proxy.integration.test.ts @@ -10,7 +10,7 @@ import { AuthKey, ContinuationTokenKey, SessionIdKey } from "#context/keys.js"; import { emitProxiedSubagentEvent } from "#execution/subagent-event-proxy-step.js"; import { projectToDurableSession } from "#execution/session.js"; import type { HarnessSession } from "#harness/types.js"; -import type { TimedHandleMessageStreamEvent } from "#protocol/message.js"; +import type { StampedHandleMessageStreamEvent } from "#protocol/message.js"; import { deserializeRuntimeAdapter } from "#runtime/channels/registry.js"; import { createEmptyHookRegistry } from "#runtime/hooks/registry.js"; import { @@ -131,8 +131,8 @@ function createCapturingWritable(chunks: Uint8Array[]): WritableStream { diff --git a/packages/eve/src/execution/subagent-event-proxy-step.ts b/packages/eve/src/execution/subagent-event-proxy-step.ts index 8a6604e6d..e959c5f03 100644 --- a/packages/eve/src/execution/subagent-event-proxy-step.ts +++ b/packages/eve/src/execution/subagent-event-proxy-step.ts @@ -21,7 +21,7 @@ import { emitProxiedInputRequest } from "#execution/subagent-hitl-proxy.js"; import { upsertProxyInputRequests } from "#harness/proxy-input-requests.js"; import type { HarnessSession } from "#harness/types.js"; import type { HandleMessageStreamEvent } from "#protocol/message.js"; -import { encodeMessageStreamEvent, timestampHandleMessageStreamEvent } from "#protocol/message.js"; +import { encodeMessageStreamEvent, stampMessageStreamEvent } from "#protocol/message.js"; import { BundleKey, ChannelKey } from "#runtime/sessions/runtime-context-keys.js"; type SubagentEventHookPayload = @@ -81,9 +81,15 @@ export async function emitProxiedSubagentEvent(input: { let proxyEntries: ProxyInputRequestEntries | undefined; let scopedSession: HarnessSession; try { + // A child event re-emitted here becomes a distinct event on the parent + // stream, so it is stamped with its own id rather than reusing the child's. const emit = async (event: HandleMessageStreamEvent): Promise => { - const transformed = await callAdapterEventHandler(adapter, event, adapterCtx); - await writer.write(encodeMessageStreamEvent(timestampHandleMessageStreamEvent(transformed))); + const transformed = await callAdapterEventHandler( + adapter, + stampMessageStreamEvent(event), + adapterCtx, + ); + await writer.write(encodeMessageStreamEvent(transformed)); }; const scopeResult = await withContextScope(ctx, session, async (enrichedSession) => { diff --git a/packages/eve/src/execution/subagent-hitl-proxy.integration.test.ts b/packages/eve/src/execution/subagent-hitl-proxy.integration.test.ts index 95b759c6e..a64855cbc 100644 --- a/packages/eve/src/execution/subagent-hitl-proxy.integration.test.ts +++ b/packages/eve/src/execution/subagent-hitl-proxy.integration.test.ts @@ -11,7 +11,8 @@ import { BundleKey, ChannelKey } from "#runtime/sessions/runtime-context-keys.js import { serializeContext } from "#context/serialize.js"; import { hasProxyInputRequests, upsertProxyInputRequests } from "#harness/proxy-input-requests.js"; import type { HarnessEmitFn, HarnessSession } from "#harness/types.js"; -import type { HandleMessageStreamEvent } from "#protocol/message.js"; +import type { StampedHandleMessageStreamEvent } from "#protocol/message.js"; +import { stampMessageStreamEvent } from "#protocol/message.js"; import type { InputRequest } from "#runtime/input/types.js"; import { createRuntimeAdapterRegistry } from "#runtime/channels/registry.js"; import type { RuntimeCompiledArtifactsSource } from "#runtime/compiled-artifacts-source.js"; @@ -175,14 +176,18 @@ function buildEmptySession(continuationToken: string, sessionId: string): Harnes */ function buildCapturingEmit(ctx: ContextContainer): { readonly emit: HarnessEmitFn; - readonly events: HandleMessageStreamEvent[]; + readonly events: StampedHandleMessageStreamEvent[]; readonly persistAdapterState: () => void; } { - const events: HandleMessageStreamEvent[] = []; + const events: StampedHandleMessageStreamEvent[] = []; const adapter = ctx.require(ChannelKey); const adapterCtx = buildAdapterContext(adapter, ctx); const emit: HarnessEmitFn = async (event) => { - const transformed = await callAdapterEventHandler(adapter, event, adapterCtx); + const transformed = await callAdapterEventHandler( + adapter, + stampMessageStreamEvent(event), + adapterCtx, + ); events.push(transformed); }; const persistAdapterState = () => { diff --git a/packages/eve/src/execution/terminal-session-failure-step.ts b/packages/eve/src/execution/terminal-session-failure-step.ts index 080fcea36..8b886126b 100644 --- a/packages/eve/src/execution/terminal-session-failure-step.ts +++ b/packages/eve/src/execution/terminal-session-failure-step.ts @@ -6,7 +6,7 @@ import { createLogger, formatError } from "#internal/logging.js"; import { createSessionFailedEvent, encodeMessageStreamEvent, - timestampHandleMessageStreamEvent, + stampMessageStreamEvent, } from "#protocol/message.js"; import { ChannelKey } from "#runtime/sessions/runtime-context-keys.js"; @@ -45,7 +45,9 @@ export async function emitTerminalSessionFailureStep(input: { code, }); - const event = createSessionFailedEvent({ code, details, message, sessionId }); + const event = stampMessageStreamEvent( + createSessionFailedEvent({ code, details, message, sessionId }), + ); // Best-effort: invoke the adapter handler so channels surface the // failure. Errors are logged, never rethrown — the outer workflow @@ -71,7 +73,7 @@ export async function emitTerminalSessionFailureStep(input: { try { const writer = input.parentWritable.getWriter(); try { - await writer.write(encodeMessageStreamEvent(timestampHandleMessageStreamEvent(event))); + await writer.write(encodeMessageStreamEvent(event)); } finally { writer.releaseLock(); } diff --git a/packages/eve/src/execution/workflow-entry.integration.test.ts b/packages/eve/src/execution/workflow-entry.integration.test.ts index 7fa48d882..b86429f37 100644 --- a/packages/eve/src/execution/workflow-entry.integration.test.ts +++ b/packages/eve/src/execution/workflow-entry.integration.test.ts @@ -15,7 +15,11 @@ import { createWorkflowRuntime } from "#execution/workflow-runtime.js"; import { normalizeEveAttributes } from "#runtime/attributes/normalize.js"; import { ROOT_COMPILED_AGENT_NODE_ID } from "#compiler/manifest.js"; import { ConnectionAuthorizationRequiredError } from "#public/connections/errors.js"; -import type { HandleMessageStreamEvent } from "#protocol/message.js"; +import type { + HandleMessageStreamEvent, + StampedHandleMessageStreamEvent, +} from "#protocol/message.js"; +import { isEventId } from "#protocol/event-id.js"; import type { ToolContext } from "#public/definitions/tool.js"; import type { AuthorizationDefinition, TokenResult } from "#runtime/connections/types.js"; import type { ResolvedToolDefinition } from "#runtime/types.js"; @@ -295,6 +299,67 @@ describe("workflowEntry integration", () => { }); }); + it("stamps every stream event with an id that survives a rewind", async () => { + const runtime = createTestRuntime({ agent: { name: "workflow-entry-event-ids" } }); + const continuationToken = "http:workflow-entry-event-ids"; + + await runtime.run(async () => { + const run = await start(workflowEntry, [ + { + input: { message: "identify these events" }, + serializedContext: buildSerializedContext({ + channelKind: "http", + continuationToken, + mode: "conversation", + }), + }, + ]); + + const stream = captureTurnEvents(run); + let firstTurn: readonly StampedHandleMessageStreamEvent[]; + try { + firstTurn = await stream.nextTurn(); + } finally { + stream.dispose(); + } + + try { + expect(firstTurn.length).toBeGreaterThan(1); + // Every event carries a well-formed id, and no two events share one — + // including the append events that share `(turnId, sequence, stepIndex)`. + expect(firstTurn.every((event) => isEventId(event.meta.id))).toBe(true); + expect(new Set(firstTurn.map((event) => event.meta.id)).size).toBe(firstTurn.length); + + // Ids are minted in emission order, so a consumer can sort by them. + const ids = firstTurn.map((event) => event.meta.id); + expect(ids).toEqual([...ids].sort()); + + // The contract that makes DB ingestion idempotent: re-reading the same + // durable stream returns the same ids, so a reconnect or a full rewind + // re-delivers events a consumer has already stored under those keys. + const workflowRuntime = createWorkflowRuntime({ + compiledArtifactsSource: createBundledRuntimeCompiledArtifactsSource(), + }); + const replayed = await workflowRuntime.getEventStream(run.runId, { startIndex: 0 }); + const replayedIds: string[] = []; + const reader = replayed.getReader(); + try { + while (replayedIds.length < firstTurn.length) { + const { done, value } = await reader.read(); + if (done) break; + replayedIds.push(value.meta.id); + } + } finally { + await reader.cancel(); + } + + expect(replayedIds).toEqual(ids); + } finally { + await run.cancel(); + } + }); + }); + it("fails a competing continuation owner before its first turn", async () => { const runtime = createTestRuntime({ agent: { name: "workflow-entry-hook-owner" } }); const continuationToken = "http:workflow-entry-hook-owner"; diff --git a/packages/eve/src/execution/workflow-runtime.ts b/packages/eve/src/execution/workflow-runtime.ts index f4163c1bf..501fab32b 100644 --- a/packages/eve/src/execution/workflow-runtime.ts +++ b/packages/eve/src/execution/workflow-runtime.ts @@ -38,7 +38,7 @@ import { type WorkflowFunction, type WorkflowMetadata, } from "#internal/workflow/runtime.js"; -import type { HandleMessageStreamEvent } from "#protocol/message.js"; +import type { StampedHandleMessageStreamEvent } from "#protocol/message.js"; import type { RuntimeCompiledArtifactsSource } from "#runtime/compiled-artifacts-source.js"; import { ROOT_RUNTIME_AGENT_NODE_ID } from "#runtime/graph.js"; import { normalizeEveAttributes } from "#runtime/attributes/normalize.js"; @@ -158,9 +158,9 @@ export function createWorkflowRuntime(config: { throw error; } - let events: ReadableStream | undefined; + let events: ReadableStream | undefined; const getEvents = () => { - events ??= parseNdjsonStream(() => + events ??= parseNdjsonStream(() => getRun(run.runId).getReadable(), ); return events; @@ -219,8 +219,8 @@ export function createWorkflowRuntime(config: { async getEventStream( sessionId: string, options?: GetEventStreamOptions, - ): Promise> { - return parseNdjsonStream(() => + ): Promise> { + return parseNdjsonStream(() => getRun(sessionId).getReadable({ startIndex: options?.startIndex }), ); }, diff --git a/packages/eve/src/execution/workflow-steps.ts b/packages/eve/src/execution/workflow-steps.ts index c4b98ff26..449c3607f 100644 --- a/packages/eve/src/execution/workflow-steps.ts +++ b/packages/eve/src/execution/workflow-steps.ts @@ -44,7 +44,8 @@ import { createSessionStartedEvent, encodeMessageStreamEvent, type HandleMessageStreamEvent, - timestampHandleMessageStreamEvent, + stampMessageStreamEvent, + type StampedHandleMessageStreamEvent, } from "#protocol/message.js"; import { CallbackBaseUrlKey, @@ -290,10 +291,18 @@ export async function turnStep(rawInput: TurnStepInput): Promise => { - const toEmit = await callAdapterEventHandler(adapter, event, adapterCtx); + // Stamp before the adapter runs so the adapter, the persisted chunk, and the + // hooks below all observe one `meta.id` for this event. + const emit = async ( + event: HandleMessageStreamEvent, + ): Promise => { + const toEmit = await callAdapterEventHandler( + adapter, + stampMessageStreamEvent(event), + adapterCtx, + ); setChannelContext(ctx, { ...adapter, state: { ...adapterCtx.state } }); - await writer.write(encodeMessageStreamEvent(timestampHandleMessageStreamEvent(toEmit))); + await writer.write(encodeMessageStreamEvent(toEmit)); return toEmit; }; diff --git a/packages/eve/src/harness/ordered-stream-emitter.test.ts b/packages/eve/src/harness/ordered-stream-emitter.test.ts index d425f1d95..bc5d558a7 100644 --- a/packages/eve/src/harness/ordered-stream-emitter.test.ts +++ b/packages/eve/src/harness/ordered-stream-emitter.test.ts @@ -47,8 +47,14 @@ describe("createOrderedStreamEmitter", () => { const emitter = createOrderedStreamEmitter(emitFn); await emitter.emit(message("A", "A")); - await emitter.emit({ ...message("B", "AB"), meta: { at: "2026-07-10T18:00:00.000Z" } }); - await emitter.emit({ ...message("C", "ABC"), meta: { at: "2026-07-10T18:00:01.000Z" } }); + await emitter.emit({ + ...message("B", "AB"), + meta: { at: "2026-07-10T18:00:00.000Z", id: "evt_test_0000" }, + }); + await emitter.emit({ + ...message("C", "ABC"), + meta: { at: "2026-07-10T18:00:01.000Z", id: "evt_test_0001" }, + }); expect(emitFn).toHaveBeenCalledTimes(1); firstWrite.resolve(); @@ -56,7 +62,7 @@ describe("createOrderedStreamEmitter", () => { expect(events).toEqual([ message("A", "A"), - { ...message("BC", "ABC"), meta: { at: "2026-07-10T18:00:01.000Z" } }, + { ...message("BC", "ABC"), meta: { at: "2026-07-10T18:00:01.000Z", id: "evt_test_0001" } }, ]); }); diff --git a/packages/eve/src/internal/nitro/host/build-extension.scenario.test.ts b/packages/eve/src/internal/nitro/host/build-extension.scenario.test.ts index 1e150589e..1a42dea78 100644 --- a/packages/eve/src/internal/nitro/host/build-extension.scenario.test.ts +++ b/packages/eve/src/internal/nitro/host/build-extension.scenario.test.ts @@ -226,7 +226,7 @@ describe("extension build output", () => { expect(manifest.requires).toEqual({ extension: 1, tool: 2, - dynamicTool: 3, + dynamicTool: 4, skill: 1, config: 1, state: 2, diff --git a/packages/eve/src/internal/testing/events.ts b/packages/eve/src/internal/testing/events.ts index 8f699c96d..3aaf0d866 100644 --- a/packages/eve/src/internal/testing/events.ts +++ b/packages/eve/src/internal/testing/events.ts @@ -1,5 +1,8 @@ -import type { HandleMessageStreamEvent } from "#protocol/message.js"; -import { isCurrentTurnBoundaryEvent } from "#protocol/message.js"; +import type { + HandleMessageStreamEvent, + StampedHandleMessageStreamEvent, +} from "#protocol/message.js"; +import { isCurrentTurnBoundaryEvent, stampMessageStreamEvent } from "#protocol/message.js"; /** * Minimal, duck-typed handle to one workflow `Run`'s readable stream. @@ -27,7 +30,7 @@ export interface CapturedTurnStream { * `session.completed`, or `session.failed`) and returns every event * observed in that turn. */ - nextTurn(): Promise; + nextTurn(): Promise; /** Releases the reader lock on the underlying `ReadableStream`. */ dispose(): void; } @@ -133,8 +136,8 @@ async function readUntilBoundary( reader: ReadableStreamDefaultReader, state: StreamState, decoder: InstanceType, -): Promise { - const events: HandleMessageStreamEvent[] = []; +): Promise { + const events: StampedHandleMessageStreamEvent[] = []; while (true) { const { done, value } = await reader.read(); @@ -157,7 +160,7 @@ async function readUntilBoundary( continue; } - const event = JSON.parse(line) as HandleMessageStreamEvent; + const event = JSON.parse(line) as StampedHandleMessageStreamEvent; events.push(event); if (isCurrentTurnBoundaryEvent(event)) { @@ -166,3 +169,27 @@ async function readUntilBoundary( } } } + +/** + * Stamps a constructed event so a test fixture satisfies the stamped stream + * contract without going through a real emit seam. + * + * Ids are sequential and deterministic so failure output stays readable. Use + * the real {@link stampMessageStreamEvent} when a test asserts on id format. + */ +export function stampTestEvent( + event: HandleMessageStreamEvent, + index = 0, +): StampedHandleMessageStreamEvent { + return stampMessageStreamEvent(event, { + at: new Date(Date.UTC(2026, 0, 1) + index).toISOString(), + id: `evt_test_${String(index).padStart(4, "0")}`, + }); +} + +/** Stamps every event in a fixture list. See {@link stampTestEvent}. */ +export function stampTestEvents( + events: readonly HandleMessageStreamEvent[], +): StampedHandleMessageStreamEvent[] { + return events.map((event, index) => stampTestEvent(event, index)); +} diff --git a/packages/eve/src/protocol/event-dedupe.test.ts b/packages/eve/src/protocol/event-dedupe.test.ts new file mode 100644 index 000000000..b23be2e9d --- /dev/null +++ b/packages/eve/src/protocol/event-dedupe.test.ts @@ -0,0 +1,82 @@ +import { describe, expect, it } from "vitest"; + +import { stampTestEvent } from "#internal/testing/events.js"; +import { createEventDeduper, DEFAULT_EVENT_DEDUPE_CAPACITY } from "#protocol/event-dedupe.js"; +import { stampMessageStreamEvent } from "#protocol/message.js"; + +function sessionStarted(index: number) { + return stampTestEvent({ type: "session.started", data: {} }, index); +} + +describe("createEventDeduper", () => { + it("admits an event once and rejects every re-delivery of it", () => { + const deduper = createEventDeduper(); + const event = sessionStarted(0); + + expect(deduper.isDuplicate(event)).toBe(false); + expect(deduper.isDuplicate(event)).toBe(true); + expect(deduper.isDuplicate(event)).toBe(true); + expect(deduper.size).toBe(1); + }); + + it("keys on the event id rather than the payload", () => { + const deduper = createEventDeduper(); + const data = {} as const; + + // Byte-identical payloads stamped as two distinct emissions are two + // events, not a replay: both must pass. + expect(deduper.isDuplicate(stampMessageStreamEvent({ type: "session.started", data }))).toBe( + false, + ); + expect(deduper.isDuplicate(stampMessageStreamEvent({ type: "session.started", data }))).toBe( + false, + ); + expect(deduper.size).toBe(2); + }); + + it("survives an overlapping replay of a stream prefix", () => { + const deduper = createEventDeduper(); + const stream = [0, 1, 2, 3].map(sessionStarted); + + const first = stream.filter((event) => !deduper.isDuplicate(event)); + // A reconnect that rewinds past the cursor re-delivers the whole prefix. + const second = stream.filter((event) => !deduper.isDuplicate(event)); + + expect(first).toHaveLength(4); + expect(second).toHaveLength(0); + }); + + it("evicts the oldest id once the window is full", () => { + const deduper = createEventDeduper(2); + const [a, b, c] = [sessionStarted(0), sessionStarted(1), sessionStarted(2)]; + + expect(deduper.isDuplicate(a)).toBe(false); + expect(deduper.isDuplicate(b)).toBe(false); + expect(deduper.isDuplicate(c)).toBe(false); + expect(deduper.size).toBe(2); + + // `a` fell out of the window; `b` and `c` are still remembered. + expect(deduper.isDuplicate(b)).toBe(true); + expect(deduper.isDuplicate(c)).toBe(true); + expect(deduper.isDuplicate(a)).toBe(false); + }); + + it("never exceeds its capacity", () => { + const deduper = createEventDeduper(8); + for (let index = 0; index < 100; index += 1) { + deduper.isDuplicate(sessionStarted(index)); + } + expect(deduper.size).toBe(8); + }); + + it("rejects a capacity that cannot hold an event", () => { + expect(() => createEventDeduper(0)).toThrow(TypeError); + expect(() => createEventDeduper(-1)).toThrow(TypeError); + expect(() => createEventDeduper(1.5)).toThrow(TypeError); + }); + + it("defaults to a bounded window", () => { + expect(DEFAULT_EVENT_DEDUPE_CAPACITY).toBeGreaterThan(0); + expect(createEventDeduper().size).toBe(0); + }); +}); diff --git a/packages/eve/src/protocol/event-dedupe.ts b/packages/eve/src/protocol/event-dedupe.ts new file mode 100644 index 000000000..4d98fc902 --- /dev/null +++ b/packages/eve/src/protocol/event-dedupe.ts @@ -0,0 +1,69 @@ +import type { StampedHandleMessageStreamEvent } from "#protocol/message.js"; + +/** + * Default number of event ids a deduper remembers before evicting the oldest. + * + * Duplicates only arise from stream overlap — a reconnect, a rewind, or an + * initial replay stitched to a live tail — so the window that matters is the + * recent past. The cap bounds memory on long-lived sessions. + */ +export const DEFAULT_EVENT_DEDUPE_CAPACITY = 10_000; + +/** + * Remembers which session-stream events have already been consumed. + * + * Every event carries a stable `meta.id` that survives reconnects, rewinds, + * and replays, so re-delivery of the same durable chunk is exactly an id we + * have already recorded. + */ +export type EventDeduper = { + /** + * Records `event` and reports whether it had already been recorded. + * + * Mutates: the first call for an id returns `false` and remembers it, and + * every later call for that id returns `true`. + */ + isDuplicate(event: StampedHandleMessageStreamEvent): boolean; + /** Number of ids currently remembered. */ + readonly size: number; +}; + +/** + * Creates an {@link EventDeduper} that remembers up to `capacity` event ids. + * + * Eviction is insertion-ordered: once the window is full, remembering a new id + * forgets the oldest. Events arrive in stream order, so the forgotten ids are + * the ones a reconnect can no longer re-deliver. + * + * @example + * ```ts + * const deduper = createEventDeduper(); + * for await (const event of session.stream()) { + * if (deduper.isDuplicate(event)) continue; + * render(event); + * } + * ``` + */ +export function createEventDeduper(capacity: number = DEFAULT_EVENT_DEDUPE_CAPACITY): EventDeduper { + if (!Number.isInteger(capacity) || capacity < 1) { + throw new TypeError(`event deduper capacity must be a positive integer, received ${capacity}`); + } + + const seen = new Set(); + + return { + isDuplicate(event) { + const id = event.meta.id; + if (seen.has(id)) return true; + seen.add(id); + if (seen.size > capacity) { + const oldest = seen.values().next(); + if (!oldest.done) seen.delete(oldest.value); + } + return false; + }, + get size() { + return seen.size; + }, + }; +} diff --git a/packages/eve/src/protocol/event-id.test.ts b/packages/eve/src/protocol/event-id.test.ts new file mode 100644 index 000000000..63ea255f3 --- /dev/null +++ b/packages/eve/src/protocol/event-id.test.ts @@ -0,0 +1,83 @@ +import { afterEach, describe, expect, it, vi } from "vitest"; + +import { + createEventId, + EVENT_ID_BODY_LENGTH, + EVENT_ID_PREFIX, + isEventId, +} from "#protocol/event-id.js"; + +const CROCKFORD_BODY = /^[0-9A-HJKMNP-TV-Z]{26}$/; + +afterEach(() => { + vi.useRealTimers(); +}); + +describe("createEventId", () => { + it("produces a prefixed 26-character Crockford base32 body", () => { + const id = createEventId(); + + expect(id.startsWith(EVENT_ID_PREFIX)).toBe(true); + const body = id.slice(EVENT_ID_PREFIX.length); + expect(body).toHaveLength(EVENT_ID_BODY_LENGTH); + expect(body).toMatch(CROCKFORD_BODY); + }); + + it("never repeats an id across a large batch", () => { + const ids = new Set(); + for (let index = 0; index < 100_000; index += 1) { + ids.add(createEventId()); + } + + expect(ids.size).toBe(100_000); + }); + + it("sorts in emission order even when many ids share one millisecond", () => { + // A single turn emits far more than one event per millisecond, so a + // non-monotonic generator would scramble `ORDER BY id` for exactly the + // events consumers most want ordered. + vi.useFakeTimers(); + vi.setSystemTime(new Date("2026-07-27T00:00:00.000Z")); + + const ids = Array.from({ length: 500 }, () => createEventId()); + + expect(ids).toEqual([...ids].sort()); + expect(new Set(ids).size).toBe(ids.length); + }); + + it("sorts later timestamps after earlier ones", () => { + vi.useFakeTimers(); + + vi.setSystemTime(new Date("2026-07-27T00:00:00.000Z")); + const earlier = createEventId(); + vi.setSystemTime(new Date("2026-07-27T00:00:01.000Z")); + const later = createEventId(); + + expect(earlier < later).toBe(true); + }); + + it("stays monotonic when the clock moves backwards", () => { + vi.useFakeTimers(); + + vi.setSystemTime(new Date("2026-07-27T00:00:05.000Z")); + const first = createEventId(); + vi.setSystemTime(new Date("2026-07-27T00:00:04.000Z")); + const second = createEventId(); + + expect(first < second).toBe(true); + }); +}); + +describe("isEventId", () => { + it("accepts ids this module mints", () => { + expect(isEventId(createEventId())).toBe(true); + }); + + it("rejects a missing prefix, a wrong length, or a non-Crockford character", () => { + const body = createEventId().slice(EVENT_ID_PREFIX.length); + + expect(isEventId(body)).toBe(false); + expect(isEventId(`${EVENT_ID_PREFIX}${body.slice(1)}`)).toBe(false); + expect(isEventId(`${EVENT_ID_PREFIX}U${body.slice(1)}`)).toBe(false); + }); +}); diff --git a/packages/eve/src/protocol/event-id.ts b/packages/eve/src/protocol/event-id.ts new file mode 100644 index 000000000..dbc21bff3 --- /dev/null +++ b/packages/eve/src/protocol/event-id.ts @@ -0,0 +1,113 @@ +/** Prefix stamped on every eve session stream event id. */ +export const EVENT_ID_PREFIX = "evt_"; + +/** Character count of the ULID body that follows {@link EVENT_ID_PREFIX}. */ +export const EVENT_ID_BODY_LENGTH = 26; + +// Crockford base32: no I, L, O, or U, so a transcribed id cannot be confused +// with 1/0 and the alphabet stays sort-order-compatible with the raw bits. +const ENCODING = "0123456789ABCDEFGHJKMNPQRSTVWXYZ"; +const TIME_CHARS = 10; +const RANDOM_CHARS = 16; +const RANDOM_BYTES = 10; + +let lastTimeMs = -1; +const lastRandom = new Uint8Array(RANDOM_BYTES); + +/** + * Mints a unique, lexicographically sortable event id. + * + * The value is `evt_` followed by a 26-character ULID: a 48-bit millisecond + * timestamp then 80 bits of randomness, both in Crockford base32. Sorting ids + * as strings therefore reproduces emission order, which keeps a primary key + * built on this id chronologically clustered in a database. + * + * Ids minted within the same millisecond increment the random component rather + * than re-randomizing it, so they still sort in emission order. A clock that + * moves backwards is pinned to the last observed millisecond for the same + * reason: monotonicity outranks timestamp precision, because consumers order by + * the id and read the exact emission time from `meta.at`. + */ +export function createEventId(): string { + const now = Date.now(); + + if (now > lastTimeMs) { + lastTimeMs = now; + randomFill(lastRandom); + } else if (!incrementRandom(lastRandom)) { + // 2^80 ids inside one millisecond. Unreachable in practice; advancing the + // logical clock keeps the generator monotonic instead of blocking on it. + lastTimeMs += 1; + randomFill(lastRandom); + } + + return `${EVENT_ID_PREFIX}${encodeTime(lastTimeMs)}${encodeRandom(lastRandom)}`; +} + +/** + * Returns true when `value` has the shape {@link createEventId} produces. + * + * Shape-only: this does not prove eve minted the id. + */ +export function isEventId(value: string): boolean { + if (!value.startsWith(EVENT_ID_PREFIX)) return false; + const body = value.slice(EVENT_ID_PREFIX.length); + if (body.length !== EVENT_ID_BODY_LENGTH) return false; + for (const character of body) { + if (!ENCODING.includes(character)) return false; + } + return true; +} + +function randomFill(target: Uint8Array): void { + // Web Crypto rather than `node:crypto`, because this module is bundled into + // browser clients and into the workflow step sandbox. `getRandomValues` is + // also available in insecure browser contexts, unlike `randomUUID`. + const webCrypto = globalThis.crypto; + if (typeof webCrypto?.getRandomValues !== "function") { + throw new Error("Cannot mint an event id: globalThis.crypto.getRandomValues is unavailable."); + } + webCrypto.getRandomValues(target); +} + +function encodeTime(timeMs: number): string { + let remaining = timeMs; + let encoded = ""; + for (let index = 0; index < TIME_CHARS; index += 1) { + encoded = ENCODING[remaining % 32] + encoded; + remaining = Math.floor(remaining / 32); + } + return encoded; +} + +function encodeRandom(bytes: Uint8Array): string { + // 80 bits divide evenly into 16 five-bit groups, so no padding is needed. + let buffer = 0; + let bufferedBits = 0; + let encoded = ""; + + for (const byte of bytes) { + buffer = (buffer << 8) | byte; + bufferedBits += 8; + while (bufferedBits >= 5) { + bufferedBits -= 5; + encoded += ENCODING[(buffer >>> bufferedBits) & 31]; + } + buffer &= (1 << bufferedBits) - 1; + } + + return encoded.padStart(RANDOM_CHARS, ENCODING[0]); +} + +/** Adds one to a big-endian counter in place. Returns false on overflow. */ +function incrementRandom(bytes: Uint8Array): boolean { + for (let index = bytes.length - 1; index >= 0; index -= 1) { + const byte = bytes[index] ?? 0; + if (byte < 0xff) { + bytes[index] = byte + 1; + return true; + } + bytes[index] = 0; + } + return false; +} diff --git a/packages/eve/src/protocol/message.test.ts b/packages/eve/src/protocol/message.test.ts index 7b37347e6..3d4b9d82c 100644 --- a/packages/eve/src/protocol/message.test.ts +++ b/packages/eve/src/protocol/message.test.ts @@ -14,13 +14,14 @@ import { createStepStartedEvent, createTurnCancelledEvent, encodeMessageStreamEvent, - timestampHandleMessageStreamEvent, + stampMessageStreamEvent, } from "#protocol/message.js"; +import { isEventId } from "#protocol/event-id.js"; 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("19"); + expect(EVE_MESSAGE_STREAM_VERSION).toBe("20"); }); it("publishes the channel-local continuation token on session.waiting", () => { @@ -59,18 +60,19 @@ describe("message stream protocol", () => { }); }); - it("stamps durable timing metadata and preserves it through encoding", () => { - const timed = timestampHandleMessageStreamEvent( + it("stamps durable envelope metadata and preserves it through encoding", () => { + const timed = stampMessageStreamEvent( createStepStartedEvent({ sequence: 0, stepIndex: 1, turnId: "turn_0", }), - "2026-04-17T10:14:22.123Z", + { at: "2026-04-17T10:14:22.123Z", id: "evt_fixed" }, ); expect(timed.meta).toEqual({ at: "2026-04-17T10:14:22.123Z", + id: "evt_fixed", }); const encoded = encodeMessageStreamEvent(timed); @@ -79,6 +81,19 @@ describe("message stream protocol", () => { expect(decoded).toEqual(timed); }); + it("mints a unique event id per stamp", () => { + const event = createStepStartedEvent({ sequence: 0, stepIndex: 0, turnId: "turn_0" }); + + const first = stampMessageStreamEvent(event); + const second = stampMessageStreamEvent(event); + + expect(isEventId(first.meta.id)).toBe(true); + expect(isEventId(second.meta.id)).toBe(true); + // Two emissions of an identical payload are two distinct stream events, so + // a consumer keyed on `meta.id` must store both rather than collapse them. + expect(first.meta.id).not.toBe(second.meta.id); + }); + it("builds authorization.required with optional challenge and webhookUrl", () => { const bare = createAuthorizationRequiredEvent({ name: "linear", diff --git a/packages/eve/src/protocol/message.ts b/packages/eve/src/protocol/message.ts index 2ba630452..64acb2dd9 100644 --- a/packages/eve/src/protocol/message.ts +++ b/packages/eve/src/protocol/message.ts @@ -6,6 +6,7 @@ import { isSerializedUrlFilePart, } from "#internal/attachments/url-refs.js"; import { decodeSandboxRef, isSandboxRefUrl } from "#internal/attachments/sandbox-refs.js"; +import { createEventId } from "#protocol/event-id.js"; import type { ConnectionAuthorizationChallenge } from "#public/connections/errors.js"; import type { RuntimeActionRequest, RuntimeActionResult } from "#runtime/actions/types.js"; import type { InputRequest, InputResponse } from "#runtime/input/types.js"; @@ -17,7 +18,7 @@ export const EVE_STREAM_FORMAT_HEADER = "x-eve-stream-format"; 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 = "19"; +export const EVE_MESSAGE_STREAM_VERSION = "20"; /** * eve-owned finish reason for one completed assistant step. @@ -46,11 +47,22 @@ export interface StepCompletedProviderMetadata { /** * Durable metadata attached to one persisted session stream event. * - * Runtime code stamps this immediately before writing the event to the - * workflow-owned stream so replay preserves the original timing. + * Runtime code stamps this once, before the channel adapter and hooks observe + * the event and before it is written to the workflow-owned stream, so every + * observer of one event agrees on its identity and timing. + * + * `id` identifies the persisted event. It is minted once and stored with the + * event, so re-reading the stream — reconnecting from a cursor, rewinding to + * `startIndex: 0`, or replaying a finished session — always yields the same id + * for the same event. Consumers persisting events can use it as a primary key + * to make ingestion idempotent. + * + * `at` is the ISO-8601 emission time. Replays preserve the original value + * instead of recomputing it. */ export interface HandleMessageStreamEventMeta { readonly at: string; + readonly id: string; } /** @@ -626,10 +638,13 @@ export type TurnFailureStreamEvent = /** * One public session stream event after runtime metadata has been stamped. * - * Runtime/execution code owns this stamping boundary. Replays must preserve the - * original `meta.at` value instead of recomputing it. + * Every event read from a session stream is stamped, so this — not + * {@link HandleMessageStreamEvent} — is the type consumers receive. + * + * Runtime/execution code owns the stamping boundary. Replays must preserve the + * original `meta` values instead of recomputing them. */ -export type TimedHandleMessageStreamEvent = HandleMessageStreamEvent & { +export type StampedHandleMessageStreamEvent = HandleMessageStreamEvent & { readonly meta: HandleMessageStreamEventMeta; }; @@ -1404,21 +1419,27 @@ export function createSessionCompletedEvent(): SessionCompletedStreamEvent { } /** - * Stamps one session event with durable timing metadata immediately before it - * is written to the workflow-owned stream. + * Stamps one session event with its durable identity and emission time. + * + * Only runtime/execution code should call this, once per event, at the top of + * an emit seam. Keeping one stamping seam ensures the channel adapter, the + * persisted stream, and authored hooks all observe the same `meta.id`, and that + * replay never invents a new id or timestamp. * - * Only runtime/execution code should call this. Keeping one stamping seam - * ensures every persisted event shares the same clock contract and replay never - * invents new timestamps. + * `overrides` exists for tests that need a fixed envelope. */ -export function timestampHandleMessageStreamEvent( +export function stampMessageStreamEvent( event: HandleMessageStreamEvent, - at = new Date().toISOString(), -): TimedHandleMessageStreamEvent { + overrides?: { + readonly at?: string; + readonly id?: string; + }, +): StampedHandleMessageStreamEvent { return { ...event, meta: { - at, + at: overrides?.at ?? new Date().toISOString(), + id: overrides?.id ?? createEventId(), }, }; } @@ -1426,7 +1447,7 @@ export function timestampHandleMessageStreamEvent( /** * Encodes one message stream event as newline-delimited JSON. */ -export function encodeMessageStreamEvent(event: TimedHandleMessageStreamEvent): Uint8Array { +export function encodeMessageStreamEvent(event: StampedHandleMessageStreamEvent): Uint8Array { return textEncoder.encode(`${JSON.stringify(event)}\n`); } diff --git a/packages/eve/src/public/channels/chat-sdk/chatSdkChannel.test.ts b/packages/eve/src/public/channels/chat-sdk/chatSdkChannel.test.ts index b47f30337..063be447b 100644 --- a/packages/eve/src/public/channels/chat-sdk/chatSdkChannel.test.ts +++ b/packages/eve/src/public/channels/chat-sdk/chatSdkChannel.test.ts @@ -6,7 +6,11 @@ import { isCompiledChannel, type CompiledChannel } from "#channel/compiled-chann import { isHttpRouteDefinition } from "#channel/routes.js"; import { ContextContainer, contextStorage } from "#context/container.js"; import { SessionKey } from "#context/keys.js"; -import type { HandleMessageStreamEvent } from "#protocol/message.js"; +import type { + HandleMessageStreamEvent, + StampedHandleMessageStreamEvent, +} from "#protocol/message.js"; +import { stampTestEvent } from "#internal/testing/events.js"; import { chatSdkChannel, isNotImplemented, @@ -72,8 +76,10 @@ function callEvent( adapter: ChannelAdapter, event: HandleMessageStreamEvent, ctx: any, -): Promise { - return contextStorage.run(stubAlsContext, () => callAdapterEventHandler(adapter, event, ctx)); +): Promise { + return contextStorage.run(stubAlsContext, () => + callAdapterEventHandler(adapter, stampTestEvent(event), ctx), + ); } function makeEvent( diff --git a/packages/eve/src/public/channels/discord/discordChannel.test.ts b/packages/eve/src/public/channels/discord/discordChannel.test.ts index 2ac1c4cef..ebd6517d5 100644 --- a/packages/eve/src/public/channels/discord/discordChannel.test.ts +++ b/packages/eve/src/public/channels/discord/discordChannel.test.ts @@ -8,7 +8,11 @@ import { isCompiledChannel, type CompiledChannel } from "#channel/compiled-chann import { isHttpRouteDefinition } from "#channel/routes.js"; import { ContextContainer, contextStorage } from "#context/container.js"; import { SessionKey } from "#context/keys.js"; -import type { HandleMessageStreamEvent } from "#protocol/message.js"; +import type { + HandleMessageStreamEvent, + StampedHandleMessageStreamEvent, +} from "#protocol/message.js"; +import { stampTestEvent } from "#internal/testing/events.js"; import { DISCORD_HITL_FREEFORM_TEXT_INPUT_ID, renderInputRequestComponents, @@ -47,8 +51,10 @@ function callEvent( adapter: ChannelAdapter, event: HandleMessageStreamEvent, ctx: any, -): Promise { - return contextStorage.run(stubAlsContext, () => callAdapterEventHandler(adapter, event, ctx)); +): Promise { + return contextStorage.run(stubAlsContext, () => + callAdapterEventHandler(adapter, stampTestEvent(event), ctx), + ); } function captureAccessor(initialContinuationToken: string): { @@ -422,8 +428,8 @@ describe("discordChannel() default event handlers", () => { const ctx = buildAdapterContext(adapter, accessor); await expect( - callAdapterEventHandler(adapter, makeEvent("turn.started", {}), ctx), - ).resolves.toEqual(makeEvent("turn.started", {})); + callAdapterEventHandler(adapter, stampTestEvent(makeEvent("turn.started", {})), ctx), + ).resolves.toEqual(stampTestEvent(makeEvent("turn.started", {}))); }); it("uses the environment bot token for proactive messages", async () => { diff --git a/packages/eve/src/public/channels/eve.test.ts b/packages/eve/src/public/channels/eve.test.ts index ef77801ef..3b2463093 100644 --- a/packages/eve/src/public/channels/eve.test.ts +++ b/packages/eve/src/public/channels/eve.test.ts @@ -3,6 +3,7 @@ import { describe, expect, it, vi } from "vitest"; import { buildAdapterContext } from "#channel/adapter-context.js"; import { callAdapterEventHandler, type ChannelAdapter } from "#channel/adapter.js"; +import { stampTestEvent } from "#internal/testing/events.js"; import { isCompiledChannel } from "#channel/compiled-channel.js"; import { createJsonMessageRequest, createMockAgent } from "#internal/testing/route-harness.js"; import { attachRouteAgent } from "#internal/nitro/routes/channel-route-context.js"; @@ -402,12 +403,14 @@ describe("eveChannel — events", () => { await contextStorage.run(ctx, async () => { await callAdapterEventHandler( adapter, - createMessageCompletedEvent({ - message: "done", - sequence: 1, - stepIndex: 0, - turnId: "turn-1", - }), + stampTestEvent( + createMessageCompletedEvent({ + message: "done", + sequence: 1, + stepIndex: 0, + turnId: "turn-1", + }), + ), adapterCtx, ); }); diff --git a/packages/eve/src/public/channels/github/githubChannel.test.ts b/packages/eve/src/public/channels/github/githubChannel.test.ts index bd72030ce..2d789bbc0 100644 --- a/packages/eve/src/public/channels/github/githubChannel.test.ts +++ b/packages/eve/src/public/channels/github/githubChannel.test.ts @@ -7,7 +7,11 @@ import { isHttpRouteDefinition } from "#channel/routes.js"; import { ContextContainer, contextStorage } from "#context/container.js"; import { SandboxKey, SessionKey } from "#context/keys.js"; import { mockSandbox, type MockSandbox } from "#internal/testing/mocks/mock-sandbox.js"; -import type { HandleMessageStreamEvent } from "#protocol/message.js"; +import type { + HandleMessageStreamEvent, + StampedHandleMessageStreamEvent, +} from "#protocol/message.js"; +import { stampTestEvent } from "#internal/testing/events.js"; import { clearGitHubInstallationTokenCache, seedGitHubInstallationTokenForTests, @@ -84,10 +88,10 @@ function callEvent( event: HandleMessageStreamEvent, ctx: any, sandbox?: MockSandbox, -): Promise { +): Promise { return contextStorage.run( sandbox === undefined ? stubAlsContext : createAlsContext(sandbox), - () => callAdapterEventHandler(adapter, event, ctx), + () => callAdapterEventHandler(adapter, stampTestEvent(event), ctx), ); } diff --git a/packages/eve/src/public/channels/linear/linearChannel.test.ts b/packages/eve/src/public/channels/linear/linearChannel.test.ts index eff91db1f..43d924426 100644 --- a/packages/eve/src/public/channels/linear/linearChannel.test.ts +++ b/packages/eve/src/public/channels/linear/linearChannel.test.ts @@ -6,7 +6,11 @@ import { isCompiledChannel, type CompiledChannel } from "#channel/compiled-chann import { isHttpRouteDefinition } from "#channel/routes.js"; import { ContextContainer, contextStorage } from "#context/container.js"; import { SessionKey } from "#context/keys.js"; -import type { HandleMessageStreamEvent } from "#protocol/message.js"; +import type { + HandleMessageStreamEvent, + StampedHandleMessageStreamEvent, +} from "#protocol/message.js"; +import { stampTestEvent } from "#internal/testing/events.js"; import { linearChannel, type LinearChannelState } from "#public/channels/linear/linearChannel.js"; import { signLinearWebhookBody } from "#public/channels/linear/verify.js"; import type { InputRequest } from "#runtime/input/types.js"; @@ -47,8 +51,10 @@ function callEvent( adapter: ChannelAdapter, event: HandleMessageStreamEvent, ctx: any, -): Promise { - return contextStorage.run(stubAlsContext, () => callAdapterEventHandler(adapter, event, ctx)); +): Promise { + return contextStorage.run(stubAlsContext, () => + callAdapterEventHandler(adapter, stampTestEvent(event), ctx), + ); } function makeEvent( diff --git a/packages/eve/src/public/channels/slack/slackChannel.test.ts b/packages/eve/src/public/channels/slack/slackChannel.test.ts index 630e27a9f..5b5d96e36 100644 --- a/packages/eve/src/public/channels/slack/slackChannel.test.ts +++ b/packages/eve/src/public/channels/slack/slackChannel.test.ts @@ -8,7 +8,11 @@ import { isCompiledChannel, type CompiledChannel } from "#channel/compiled-chann import { isHttpRouteDefinition } from "#channel/routes.js"; import { ContextContainer, contextStorage } from "#context/container.js"; import { SessionKey } from "#context/keys.js"; -import type { HandleMessageStreamEvent } from "#protocol/message.js"; +import type { + HandleMessageStreamEvent, + StampedHandleMessageStreamEvent, +} from "#protocol/message.js"; +import { stampTestEvent } from "#internal/testing/events.js"; import { decodeSlackApiBody } from "#public/channels/slack/api-encoding.js"; import { HITL_ACTION_PREFIX, @@ -80,8 +84,10 @@ function callEvent( adapter: ChannelAdapter, event: HandleMessageStreamEvent, ctx: any, -): Promise { - return contextStorage.run(stubAlsContext, () => callAdapterEventHandler(adapter, event, ctx)); +): Promise { + return contextStorage.run(stubAlsContext, () => + callAdapterEventHandler(adapter, stampTestEvent(event), ctx), + ); } /** @@ -684,7 +690,7 @@ describe("slackChannel() default event handlers", () => { await contextStorage.run(callerContext, () => callAdapterEventHandler( adapter, - makeEvent("turn.started", { sequence: 1, stepIndex: 0, turnId: "t2" }), + stampTestEvent(makeEvent("turn.started", { sequence: 1, stepIndex: 0, turnId: "t2" })), ctx, ), ); diff --git a/packages/eve/src/public/channels/telegram/telegramChannel.test.ts b/packages/eve/src/public/channels/telegram/telegramChannel.test.ts index bb9b9ddb7..a342e4ddf 100644 --- a/packages/eve/src/public/channels/telegram/telegramChannel.test.ts +++ b/packages/eve/src/public/channels/telegram/telegramChannel.test.ts @@ -6,7 +6,11 @@ import { isCompiledChannel, type CompiledChannel } from "#channel/compiled-chann import { isHttpRouteDefinition } from "#channel/routes.js"; import { ContextContainer, contextStorage } from "#context/container.js"; import { SessionKey } from "#context/keys.js"; -import type { HandleMessageStreamEvent } from "#protocol/message.js"; +import type { + HandleMessageStreamEvent, + StampedHandleMessageStreamEvent, +} from "#protocol/message.js"; +import { stampTestEvent } from "#internal/testing/events.js"; import { telegramChannel, type TelegramChannelState } from "#public/channels/telegram/index.js"; const SECRET = "telegram-secret"; @@ -43,8 +47,10 @@ function callEvent( adapter: ChannelAdapter, event: HandleMessageStreamEvent, ctx: any, -): Promise { - return contextStorage.run(stubAlsContext, () => callAdapterEventHandler(adapter, event, ctx)); +): Promise { + return contextStorage.run(stubAlsContext, () => + callAdapterEventHandler(adapter, stampTestEvent(event), ctx), + ); } function makeEvent( diff --git a/packages/eve/src/public/channels/twilio/twilioChannel.test.ts b/packages/eve/src/public/channels/twilio/twilioChannel.test.ts index c771cdd68..50444a5df 100644 --- a/packages/eve/src/public/channels/twilio/twilioChannel.test.ts +++ b/packages/eve/src/public/channels/twilio/twilioChannel.test.ts @@ -6,7 +6,11 @@ import { isCompiledChannel, type CompiledChannel } from "#channel/compiled-chann import { isHttpRouteDefinition } from "#channel/routes.js"; import { ContextContainer, contextStorage } from "#context/container.js"; import { SessionKey } from "#context/keys.js"; -import type { HandleMessageStreamEvent } from "#protocol/message.js"; +import type { + HandleMessageStreamEvent, + StampedHandleMessageStreamEvent, +} from "#protocol/message.js"; +import { stampTestEvent } from "#internal/testing/events.js"; import type { TwilioTextMessage } from "#public/channels/twilio/inbound.js"; import { twilioChannel, type TwilioContext } from "#public/channels/twilio/twilioChannel.js"; import { signTwilioRequest } from "#public/channels/twilio/verify.js"; @@ -49,8 +53,10 @@ function callEvent( adapter: ChannelAdapter, event: HandleMessageStreamEvent, ctx: any, -): Promise { - return contextStorage.run(stubAlsContext, () => callAdapterEventHandler(adapter, event, ctx)); +): Promise { + return contextStorage.run(stubAlsContext, () => + callAdapterEventHandler(adapter, stampTestEvent(event), ctx), + ); } function makeEvent( diff --git a/packages/eve/src/public/definitions/channel.test.ts b/packages/eve/src/public/definitions/channel.test.ts index a9db21800..87cd5d2dd 100644 --- a/packages/eve/src/public/definitions/channel.test.ts +++ b/packages/eve/src/public/definitions/channel.test.ts @@ -2,6 +2,7 @@ import { describe, expect, it } from "vitest"; import { buildAdapterContext } from "#channel/adapter-context.js"; import { callAdapterEventHandler, type ChannelAdapter } from "#channel/adapter.js"; +import { stampTestEvent } from "#internal/testing/events.js"; import { isCompiledChannel } from "#channel/compiled-channel.js"; import type { InferReceiveTarget } from "#channel/receive-target.js"; import { ContextContainer, contextStorage } from "#context/container.js"; @@ -459,7 +460,7 @@ describe("defineChannel", () => { await contextStorage.run(ctx, async () => { await callAdapterEventHandler( adapter, - { + stampTestEvent({ type: "reasoning.appended", data: { reasoningDelta: "Need", @@ -468,12 +469,12 @@ describe("defineChannel", () => { stepIndex: 0, turnId: "turn-1", }, - }, + }), adapterCtx, ); await callAdapterEventHandler( adapter, - { + stampTestEvent({ type: "reasoning.completed", data: { reasoning: "Need to inspect the repo.", @@ -481,7 +482,7 @@ describe("defineChannel", () => { stepIndex: 0, turnId: "turn-1", }, - }, + }), adapterCtx, ); }); @@ -528,10 +529,10 @@ describe("defineChannel", () => { await callAdapterEventHandler( adapter, - { + stampTestEvent({ type: "session.failed", data: { code: "INTERNAL", message: "boom", sessionId: "sess-1" }, - }, + }), adapterCtx, ); diff --git a/packages/eve/src/public/definitions/hook.ts b/packages/eve/src/public/definitions/hook.ts index 2a9abb4dd..1975b5e12 100644 --- a/packages/eve/src/public/definitions/hook.ts +++ b/packages/eve/src/public/definitions/hook.ts @@ -1,9 +1,14 @@ -import type { HandleMessageStreamEvent } from "../../protocol/message.js"; +import type { + HandleMessageStreamEvent, + StampedHandleMessageStreamEvent, +} from "../../protocol/message.js"; import type { SessionContext } from "./callback-context.js"; import type { ExactDefinition } from "./exact.js"; +// Stamped, so a hook can read `event.meta.id` unguarded — the stable key for +// making a hook's own side effects idempotent. type ProtocolEvent = Extract< - HandleMessageStreamEvent, + StampedHandleMessageStreamEvent, { type: TType } >; diff --git a/packages/eve/src/react/use-eve-agent.test.ts b/packages/eve/src/react/use-eve-agent.test.ts index 37fdec7ea..4563ff711 100644 --- a/packages/eve/src/react/use-eve-agent.test.ts +++ b/packages/eve/src/react/use-eve-agent.test.ts @@ -14,6 +14,7 @@ import { createTurnFailedEvent, type HandleMessageStreamEvent, } from "#protocol/message.js"; +import { stampTestEvents } from "#internal/testing/events.js"; import type { SessionState } from "#client/types.js"; function createStartedMessageResponse(sessionId: string, continuationToken: string): Response { @@ -26,13 +27,13 @@ function createStartedMessageResponse(sessionId: string, continuationToken: stri }); } +/** Serializes fixture events the way the server does: stamped, one per line. */ function createEagerStreamResponse(events: readonly HandleMessageStreamEvent[]): Response { const encoder = new TextEncoder(); - return new Response( new ReadableStream({ start(controller) { - for (const event of events) { + for (const event of stampTestEvents(events)) { controller.enqueue(encoder.encode(`${JSON.stringify(event)}\n`)); } controller.close(); @@ -199,7 +200,7 @@ describe("useEveAgent", () => { }); expect(fetchMock).toHaveBeenCalledTimes(2); - expect(seenEvents).toEqual(events); + expect(seenEvents).toEqual(stampTestEvents(events)); expect(seenSessions).toEqual([ { continuationToken: "http:session_1", @@ -481,7 +482,7 @@ describe("useEveAgent", () => { expect(seenErrors.map((error) => error.name)).toEqual(["MODEL_CALL_FAILED"]); expect(helpers?.status).toBe("error"); expect(helpers?.error?.message).toBe("Bad Request"); - expect(helpers?.events).toEqual(events); + expect(helpers?.events).toEqual(stampTestEvents(events)); expect(helpers?.data).toEqual( completedTurnData({ turnId: "turn_1", @@ -548,7 +549,7 @@ describe("useEveAgent", () => { expect(seenErrors).toEqual([]); expect(helpers?.status).toBe("ready"); expect(helpers?.error).toBeUndefined(); - expect(helpers?.events).toEqual(events); + expect(helpers?.events).toEqual(stampTestEvents(events)); }); it("projects input responses before the resumed stream returns", async () => { diff --git a/packages/eve/src/react/use-eve-agent.ts b/packages/eve/src/react/use-eve-agent.ts index 5b8e9ffaa..8c725b9c7 100644 --- a/packages/eve/src/react/use-eve-agent.ts +++ b/packages/eve/src/react/use-eve-agent.ts @@ -11,7 +11,7 @@ import { resolveEveAgentHost } from "#client/agent-host.js"; import type { EveAgentReducer } from "#client/reducer.js"; import type { ClientSession } from "#client/session.js"; import { defaultMessageReducer, type EveMessageData } from "#client/message-reducer.js"; -import type { HandleMessageStreamEvent } from "#protocol/message.js"; +import type { StampedHandleMessageStreamEvent } from "#protocol/message.js"; import type { ClientAuth, HeadersValue, SendTurnPayload, SessionState } from "#client/types.js"; export type { PrepareSend }; @@ -75,7 +75,7 @@ export interface UseEveAgentOptions extends EveAgentStoreCallbacks * @default "" */ readonly host?: string; - readonly initialEvents?: readonly HandleMessageStreamEvent[]; + readonly initialEvents?: readonly StampedHandleMessageStreamEvent[]; readonly initialSession?: SessionState; /** * Project submitted user messages before eve confirms them with a diff --git a/packages/eve/src/runtime/hooks/registry.ts b/packages/eve/src/runtime/hooks/registry.ts index b223a7c68..03d3a4370 100644 --- a/packages/eve/src/runtime/hooks/registry.ts +++ b/packages/eve/src/runtime/hooks/registry.ts @@ -1,4 +1,4 @@ -import type { HandleMessageStreamEvent } from "#protocol/message.js"; +import type { StampedHandleMessageStreamEvent } from "#protocol/message.js"; import type { StreamEventHook } from "../../public/definitions/hook.js"; import type { ResolvedHookDefinition } from "../types.js"; @@ -10,7 +10,7 @@ import type { ResolvedHookDefinition } from "../types.js"; */ interface RuntimeStreamEventHookEntry { readonly slug: string; - readonly handler: StreamEventHook; + readonly handler: StreamEventHook; readonly eventType: string; } diff --git a/packages/eve/src/runtime/types.ts b/packages/eve/src/runtime/types.ts index 8e10abea7..2d8287204 100644 --- a/packages/eve/src/runtime/types.ts +++ b/packages/eve/src/runtime/types.ts @@ -3,7 +3,7 @@ import type { CompiledChannel } from "#channel/compiled-channel.js"; import type { NormalizedChannelCorsOptions } from "#channel/cors.js"; import type { HeadersValue } from "#client/types.js"; import type { DiscoverDiagnosticsSummary } from "#discover/diagnostics.js"; -import type { HandleMessageStreamEvent } from "#protocol/message.js"; +import type { StampedHandleMessageStreamEvent } from "#protocol/message.js"; import type { ChannelRouteMethod, RouteContext } from "#public/definitions/channel.js"; import type { RouteHandler, WebSocketRouteHandler } from "#channel/routes.js"; import type { OutboundAuthFn } from "#public/agents/auth.js"; @@ -210,7 +210,7 @@ export interface ResolvedHookDefinition extends ResolvedModuleSourceRef { * wildcard if declared. Unknown keys are accepted at resolve time * and ignored at dispatch time. */ - readonly events: Readonly>>; + readonly events: Readonly>>; } /** diff --git a/packages/eve/src/svelte/use-eve-agent.test.ts b/packages/eve/src/svelte/use-eve-agent.test.ts index fee4f8a00..23af9a301 100644 --- a/packages/eve/src/svelte/use-eve-agent.test.ts +++ b/packages/eve/src/svelte/use-eve-agent.test.ts @@ -8,6 +8,7 @@ import { createSessionWaitingEvent, type HandleMessageStreamEvent, } from "#protocol/message.js"; +import { stampTestEvents } from "#internal/testing/events.js"; function createStartedMessageResponse(sessionId: string, continuationToken: string): Response { return new Response(JSON.stringify({ continuationToken, ok: true, sessionId }), { @@ -19,12 +20,13 @@ function createStartedMessageResponse(sessionId: string, continuationToken: stri }); } +/** Serializes fixture events the way the server does: stamped, one per line. */ function createEagerStreamResponse(events: readonly HandleMessageStreamEvent[]): Response { const encoder = new TextEncoder(); return new Response( new ReadableStream({ start(controller) { - for (const event of events) { + for (const event of stampTestEvents(events)) { controller.enqueue(encoder.encode(`${JSON.stringify(event)}\n`)); } controller.close(); @@ -41,7 +43,7 @@ afterEach(() => { describe("useEveAgent (Svelte rune binding)", () => { it("renders the initial projection through plain reactive properties", () => { const agent = useEveAgent({ - initialEvents: [ + initialEvents: stampTestEvents([ createMessageReceivedEvent({ message: "Hello", sequence: 0, turnId: "turn_1" }), createMessageCompletedEvent({ message: "Hi there.", @@ -49,7 +51,7 @@ describe("useEveAgent (Svelte rune binding)", () => { stepIndex: 0, turnId: "turn_1", }), - ], + ]), initialSession: { continuationToken: "http:session_1", sessionId: "session_1", @@ -68,9 +70,9 @@ describe("useEveAgent (Svelte rune binding)", () => { it("does not update the visible snapshot without browser reactivity", async () => { const agent = useEveAgent({ - initialEvents: [ + initialEvents: stampTestEvents([ createMessageReceivedEvent({ message: "Hello", sequence: 0, turnId: "turn_1" }), - ], + ]), }); const dataBeforeSend = agent.data; @@ -101,6 +103,6 @@ describe("useEveAgent (Svelte rune binding)", () => { await agent.send({ message: "Hello" }); - expect(seenEvents).toEqual(events); + expect(seenEvents).toEqual(stampTestEvents(events)); }); }); diff --git a/packages/eve/src/svelte/use-eve-agent.ts b/packages/eve/src/svelte/use-eve-agent.ts index 78face8e3..53f0c2741 100644 --- a/packages/eve/src/svelte/use-eve-agent.ts +++ b/packages/eve/src/svelte/use-eve-agent.ts @@ -12,7 +12,7 @@ import { defaultMessageReducer, type EveMessageData } from "#client/message-redu import type { EveAgentReducer } from "#client/reducer.js"; import type { ClientSession } from "#client/session.js"; import type { ClientAuth, HeadersValue, SendTurnPayload, SessionState } from "#client/types.js"; -import type { HandleMessageStreamEvent } from "#protocol/message.js"; +import type { StampedHandleMessageStreamEvent } from "#protocol/message.js"; export type { PrepareSend }; @@ -43,7 +43,7 @@ export interface UseEveAgentReturn { /** Last transport-level error, or `undefined` when healthy. */ readonly error: Error | undefined; /** Raw server events received during this session (authoritative stream). */ - readonly events: readonly HandleMessageStreamEvent[]; + readonly events: readonly StampedHandleMessageStreamEvent[]; /** Clear all state and start a new session. */ readonly reset: () => void; /** Send a turn with full structured input (message, attachments, input responses). */ @@ -91,7 +91,7 @@ export interface UseEveAgentOptions extends EveAgentStoreCallbacks */ readonly host?: string; /** Seed events for resuming a prior conversation. */ - readonly initialEvents?: readonly HandleMessageStreamEvent[]; + readonly initialEvents?: readonly StampedHandleMessageStreamEvent[]; /** Seed session identity and stream cursor for resuming a prior conversation. */ readonly initialSession?: SessionState; /** @@ -148,7 +148,7 @@ class SvelteEveAgent implements UseEveAgentReturn { return this.#snapshot.error; } - get events(): readonly HandleMessageStreamEvent[] { + get events(): readonly StampedHandleMessageStreamEvent[] { this.#subscribe(); return this.#snapshot.events; } diff --git a/packages/eve/src/vue/use-eve-agent.test.ts b/packages/eve/src/vue/use-eve-agent.test.ts index a3e4baffa..c4ad102c9 100644 --- a/packages/eve/src/vue/use-eve-agent.test.ts +++ b/packages/eve/src/vue/use-eve-agent.test.ts @@ -12,6 +12,7 @@ import { createSessionWaitingEvent, type HandleMessageStreamEvent, } from "#protocol/message.js"; +import { stampTestEvents } from "#internal/testing/events.js"; import { defaultMessageReducer } from "#client/message-reducer.js"; import type { SessionState } from "#client/types.js"; @@ -25,12 +26,13 @@ function createStartedMessageResponse(sessionId: string, continuationToken: stri }); } +/** Serializes fixture events the way the server does: stamped, one per line. */ function createEagerStreamResponse(events: readonly HandleMessageStreamEvent[]): Response { const encoder = new TextEncoder(); return new Response( new ReadableStream({ start(controller) { - for (const event of events) { + for (const event of stampTestEvents(events)) { controller.enqueue(encoder.encode(`${JSON.stringify(event)}\n`)); } controller.close(); @@ -179,7 +181,7 @@ describe("EveAgentStore (Vue composable backing store)", () => { startResponse.resolve(createStartedMessageResponse("session_1", "http:session_1")); await sendPromise; - expect(seenEvents).toEqual(events); + expect(seenEvents).toEqual(stampTestEvents(events)); expect(store.snapshot.status).toBe("ready"); expect(store.snapshot.data).toEqual( completedTurnData({ @@ -403,7 +405,7 @@ describe("useEveAgent (Vue composable wiring)", () => { it("renders initial projection without subscribing during SSR", async () => { const agent = useEveAgent({ - initialEvents: [ + initialEvents: stampTestEvents([ createMessageReceivedEvent({ message: "Hello", sequence: 0, turnId: "turn_1" }), createMessageCompletedEvent({ message: "Hi there.", @@ -411,7 +413,7 @@ describe("useEveAgent (Vue composable wiring)", () => { stepIndex: 0, turnId: "turn_1", }), - ], + ]), initialSession: { continuationToken: "http:session_1", sessionId: "session_1", diff --git a/packages/eve/src/vue/use-eve-agent.ts b/packages/eve/src/vue/use-eve-agent.ts index 521e7fc84..ca57797b0 100644 --- a/packages/eve/src/vue/use-eve-agent.ts +++ b/packages/eve/src/vue/use-eve-agent.ts @@ -11,7 +11,7 @@ import { resolveEveAgentHost } from "#client/agent-host.js"; import type { EveAgentReducer } from "#client/reducer.js"; import type { ClientSession } from "#client/session.js"; import { defaultMessageReducer, type EveMessageData } from "#client/message-reducer.js"; -import type { HandleMessageStreamEvent } from "#protocol/message.js"; +import type { StampedHandleMessageStreamEvent } from "#protocol/message.js"; import type { ClientAuth, HeadersValue, SendTurnPayload, SessionState } from "#client/types.js"; export type { PrepareSend }; @@ -41,7 +41,7 @@ export interface UseEveAgentReturn { /** Last transport-level error, or `undefined` when healthy. */ readonly error: ComputedRef; /** Raw server events from this session (authoritative stream). */ - readonly events: ComputedRef; + readonly events: ComputedRef; /** Clear all state and start a new session. */ readonly reset: () => void; /** Send a turn with full structured input (message, attachments, input responses). */ @@ -89,7 +89,7 @@ export interface UseEveAgentOptions extends EveAgentStoreCallbacks */ readonly host?: string; /** Prior stream events to rehydrate the projected state from on mount. */ - readonly initialEvents?: readonly HandleMessageStreamEvent[]; + readonly initialEvents?: readonly StampedHandleMessageStreamEvent[]; /** Prior session cursor to resume from on mount. */ readonly initialSession?: SessionState; /** diff --git a/packages/eve/test/eve-run-stream-channel.test.ts b/packages/eve/test/eve-run-stream-channel.test.ts index 2c5f63952..e411af9ca 100644 --- a/packages/eve/test/eve-run-stream-channel.test.ts +++ b/packages/eve/test/eve-run-stream-channel.test.ts @@ -1,6 +1,7 @@ import { describe, expect, it, vi } from "vitest"; -import type { HandleMessageStreamEvent } from "../src/protocol/message.js"; +import type { StampedHandleMessageStreamEvent } from "../src/protocol/message.js"; +import { stampTestEvents } from "../src/internal/testing/events.js"; import type { RouteHandlerArgs, GetSessionFn } from "../src/channel/routes.js"; import type { Session } from "../src/channel/session.js"; import { EVE_MESSAGE_STREAM_ROUTE_PATTERN } from "../src/protocol/routes.js"; @@ -30,18 +31,20 @@ function createGetHandler() { describe("eveChannel GET stream", () => { it("forwards the startIndex query parameter into getSession/getEventStream", async () => { const getRoute = createGetHandler(); - const events = createEvents([ - { - type: "message.completed", - data: { - finishReason: "stop", - message: "second turn reply", - sequence: 0, - stepIndex: 0, - turnId: "turn-1", + const events = createEvents( + stampTestEvents([ + { + type: "message.completed", + data: { + finishReason: "stop", + message: "second turn reply", + sequence: 0, + stepIndex: 0, + turnId: "turn-1", + }, }, - }, - ]); + ]), + ); const getSession = createMockGetSession(events); const response = await (getRoute as any).handler( @@ -116,7 +119,7 @@ describe("eveChannel GET stream", () => { it("re-serializes the parsed event stream as NDJSON bytes", async () => { const getRoute = createGetHandler(); - const events: HandleMessageStreamEvent[] = [ + const events = stampTestEvents([ { type: "message.completed", data: { @@ -131,7 +134,7 @@ describe("eveChannel GET stream", () => { type: "session.waiting", data: { continuationToken: "eve:test", wait: "next-user-message" }, }, - ]; + ]); const getSession = createMockGetSession(createEvents(events)); const response = await (getRoute as any).handler( @@ -154,9 +157,9 @@ describe("eveChannel GET stream", () => { }); function createEvents( - events: readonly HandleMessageStreamEvent[], -): ReadableStream { - return new ReadableStream({ + events: readonly StampedHandleMessageStreamEvent[], +): ReadableStream { + return new ReadableStream({ start(controller) { for (const event of events) { controller.enqueue(event); @@ -166,7 +169,7 @@ function createEvents( }); } -function createMockGetSession(events: ReadableStream) { +function createMockGetSession(events: ReadableStream) { return vi.fn().mockReturnValue({ id: "session_xyz", continuationToken: "", diff --git a/packages/eve/test/scenarios/schedule-trigger.scenario.test.ts b/packages/eve/test/scenarios/schedule-trigger.scenario.test.ts index 3b9176e51..795137798 100644 --- a/packages/eve/test/scenarios/schedule-trigger.scenario.test.ts +++ b/packages/eve/test/scenarios/schedule-trigger.scenario.test.ts @@ -5,7 +5,7 @@ import { describe, expect, it } from "vitest"; import type { ChannelAdapter } from "#channel/adapter.js"; import { expectScheduleRun, SCHEDULE_ADAPTER_KIND, ScheduleDispatcher } from "#channel/schedule.js"; -import type { HandleMessageStreamEvent } from "#protocol/message.js"; +import type { StampedHandleMessageStreamEvent } from "#protocol/message.js"; import type { RunHandle, Runtime } from "#channel/types.js"; import { compileAgent } from "#compiler/compile-agent.js"; import { ContextContainer } from "#context/container.js"; @@ -71,7 +71,7 @@ function createCapturingRuntime(captured: CapturedRun[]): Runtime { const handle: RunHandle = { continuationToken: "scenario-token", - events: new ReadableStream(), + events: new ReadableStream(), sessionId: "scenario-session", }; return handle; @@ -80,7 +80,7 @@ function createCapturingRuntime(captured: CapturedRun[]): Runtime { throw new Error("deliver should not be called in this scenario"); }, async getEventStream() { - return new ReadableStream(); + return new ReadableStream(); }, async terminateSession() { throw new Error("terminateSession should not be called in this scenario"); diff --git a/packages/eve/test/tui-client/tui-connection-auth-states.ts b/packages/eve/test/tui-client/tui-connection-auth-states.ts index 635da139e..72cd9450d 100644 --- a/packages/eve/test/tui-client/tui-connection-auth-states.ts +++ b/packages/eve/test/tui-client/tui-connection-auth-states.ts @@ -1,6 +1,11 @@ import { setTimeout as sleep } from "node:timers/promises"; -import { ClientSession, MessageResponse, type HandleMessageStreamEvent } from "eve/client"; +import { + ClientSession, + MessageResponse, + type HandleMessageStreamEvent, + type StampedHandleMessageStreamEvent, +} from "eve/client"; import { EveTUIRunner, MockScreen, MockUserInput } from "./lib/tui.ts"; import { theme } from "./lib/theme.ts"; @@ -61,18 +66,29 @@ class FakeSession extends ClientSession { }); } - override stream(): AsyncIterable { + override stream(): AsyncIterable { const events = this.#continuations[this.#continuationIndex] ?? []; this.#continuationIndex += 1; return pacedEvents(events); } } +let nextEventIndex = 0; + +/** Stamps a fixture event the way an emit seam does before it hits the wire. */ +function stamp(event: HandleMessageStreamEvent): StampedHandleMessageStreamEvent { + nextEventIndex += 1; + return { + ...event, + meta: { at: new Date().toISOString(), id: `evt_smoke_${nextEventIndex}` }, + } as StampedHandleMessageStreamEvent; +} + async function* pacedEvents( events: readonly HandleMessageStreamEvent[], -): AsyncGenerator { +): AsyncGenerator { for (const event of events) { - yield event; + yield stamp(event); // Pacing gives the renderer and smoke assertions a chance to observe // each lifecycle state between events, as the HTTP transport does live. await sleep(200); From 0bba334753e952f1d8a7806ae9e99fd765a90d2a Mon Sep 17 00:00:00 2001 From: Andrew Barba Date: Mon, 27 Jul 2026 14:58:58 -0400 Subject: [PATCH 02/12] refactor(eve): tighten the stream event id surface Review pass over the stamped-event-id work. Correctness: - Remove the event deduper's bounded window. Once a stream exceeded the capacity, rewinding to the start re-admitted the evicted id, which evicted the next, cascading until the whole replay was applied again. Every caller already retains an object per event, so the cap bought nothing. - Tolerate a missing envelope. Events written before stream version 20 have no `meta`, so rewinding into an older part of a live session threw instead of admitting the event. - Stop promising that `meta.id` makes hook side effects idempotent. A durable step that is interrupted re-runs and re-emits with fresh ids, so `on conflict (id) do nothing` does not dedupe a retry. The guarantee is scoped to re-reading an already persisted stream. - Stop promising a total order. Ids are minted per process, so two steps on different machines can sort either way; drop the `where id > $cursor` pagination guidance in favour of stream position. - Narrow `RouteContext.agent.getEventStream()`, which still returned unstamped events. - Mark the changeset `minor`: `initialEvents` now requires stamped events. Structure: - Extract the ULID generator to `#shared/ulid.ts` with a `createUlidFactory()` for callers needing isolated monotonic state, leaving `event-id.ts` as the `evt_` prefix layer. Records why the generator is in-repo rather than an npm package, and pins the timestamp encoding with golden vectors cross-checked against the reference implementation. - Revert `callAdapterEventHandler` to its original signature and stamp after the adapter runs. Adapter handlers only ever receive `event.data`, so they never observed the envelope the comments claimed; the parameter change was churn across thirteen files. Also drops development notes, redundant assertions, and an unreachable `padStart`, and removes the test-only override parameter from `stampMessageStreamEvent`. Signed-off-by: Andrew Barba --- .changeset/dedupe-stream-events-by-id.md | 9 -- .changeset/stable-stream-event-ids.md | 8 +- docs/concepts/sessions-runs-and-streaming.md | 13 +- docs/guides/frontend/overview.mdx | 4 +- docs/guides/frontend/use-eve-agent-svelte.mdx | 20 +-- docs/guides/frontend/use-eve-agent-vue.mdx | 20 +-- docs/guides/hooks.md | 12 +- packages/eve/src/channel/adapter.test.ts | 4 +- packages/eve/src/channel/adapter.ts | 13 +- packages/eve/src/cli/dev/tui/runner.test.ts | 20 +-- packages/eve/src/cli/dev/tui/runner.ts | 5 +- .../eve/src/cli/dev/tui/subagent-pump.test.ts | 21 +-- packages/eve/src/cli/dev/tui/subagent-pump.ts | 4 +- .../eve/src/client/eve-agent-store.test.ts | 12 -- packages/eve/src/client/eve-agent-store.ts | 5 +- .../hook-lifecycle.integration.test.ts | 25 --- .../dispatch-runtime-actions-step.ts | 26 ++-- .../execution/settle-cancelled-turn-step.ts | 12 +- .../src/execution/subagent-adapter.test.ts | 3 +- .../execution/subagent-event-proxy-step.ts | 12 +- .../subagent-hitl-proxy.integration.test.ts | 13 +- .../terminal-session-failure-step.ts | 6 +- .../workflow-entry.integration.test.ts | 9 +- packages/eve/src/execution/workflow-steps.ts | 15 +- packages/eve/src/internal/testing/events.ts | 15 +- .../eve/src/protocol/event-dedupe.test.ts | 97 ++++-------- packages/eve/src/protocol/event-dedupe.ts | 59 ++----- packages/eve/src/protocol/event-id.test.ts | 68 +------- packages/eve/src/protocol/event-id.ts | 106 ++----------- packages/eve/src/protocol/message.test.ts | 32 ++-- packages/eve/src/protocol/message.ts | 36 ++--- .../channels/chat-sdk/chatSdkChannel.test.ts | 12 +- .../channels/discord/discordChannel.test.ts | 16 +- packages/eve/src/public/channels/eve.test.ts | 15 +- .../channels/github/githubChannel.test.ts | 10 +- .../channels/linear/linearChannel.test.ts | 12 +- .../channels/slack/slackChannel.test.ts | 14 +- .../channels/telegram/telegramChannel.test.ts | 12 +- .../channels/twilio/twilioChannel.test.ts | 12 +- .../src/public/definitions/channel.test.ts | 13 +- .../eve/src/public/definitions/channel.ts | 7 +- packages/eve/src/public/definitions/hook.ts | 3 +- packages/eve/src/react/use-eve-agent.test.ts | 2 +- packages/eve/src/shared/ulid.test.ts | 107 +++++++++++++ packages/eve/src/shared/ulid.ts | 146 ++++++++++++++++++ packages/eve/src/svelte/use-eve-agent.test.ts | 1 - packages/eve/src/vue/use-eve-agent.test.ts | 1 - 47 files changed, 491 insertions(+), 596 deletions(-) delete mode 100644 .changeset/dedupe-stream-events-by-id.md create mode 100644 packages/eve/src/shared/ulid.test.ts create mode 100644 packages/eve/src/shared/ulid.ts diff --git a/.changeset/dedupe-stream-events-by-id.md b/.changeset/dedupe-stream-events-by-id.md deleted file mode 100644 index 5843d896e..000000000 --- a/.changeset/dedupe-stream-events-by-id.md +++ /dev/null @@ -1,9 +0,0 @@ ---- -"eve": patch ---- - -Stream consumers now drop re-delivered events by their stable `meta.id` -instead of guessing from payload content. `EveAgentStore` (and so the React, -Vue, and Svelte bindings) no longer double-applies an `initialEvents` prefix -that the live stream replays, and the dev TUI no longer renders a subagent's -transcript twice when its child stream reopens. diff --git a/.changeset/stable-stream-event-ids.md b/.changeset/stable-stream-event-ids.md index 166f3cfab..9195c315d 100644 --- a/.changeset/stable-stream-event-ids.md +++ b/.changeset/stable-stream-event-ids.md @@ -1,5 +1,9 @@ --- -"eve": patch +"eve": minor --- -Every session stream event now carries a stable `meta.id`. The id is a sortable, `evt_`-prefixed ULID minted once when the event is written to the durable stream, so reconnecting from a cursor, rewinding to `startIndex=0`, or replaying a finished session all return the same id for the same event — making it safe to use as a primary key when persisting events (`on conflict (id) do nothing`). Channel adapters, the durable stream, and authored hooks now all observe the same envelope, so a hook can use `event.meta.id` to make its own side effects idempotent. Events read from a stream are typed as the new `StampedHandleMessageStreamEvent`, which guarantees `meta` is present. +Every session stream event now carries a `meta.id`: a unique, `evt_`-prefixed ULID minted once when the event is written to the durable stream. Re-reading a stream — reconnecting from a cursor, rewinding to `startIndex=0`, or replaying a finished session — returns the same id for the same event, so it is safe to use as a primary key when persisting events (`on conflict (id) do nothing`). Authored hooks receive the same envelope. + +Stream consumers now drop re-delivered events by id instead of guessing from payload content. `EveAgentStore` (and so the React, Vue, and Svelte bindings) no longer double-applies an `initialEvents` prefix that the live stream replays, and the dev TUI no longer renders a subagent's transcript twice when its child stream reopens. + +**Breaking:** events read from a stream are now typed as `StampedHandleMessageStreamEvent`, which guarantees `meta` is present. `initialEvents` on `useEveAgent`/`EveAgentStore` requires that type, so a saved log typed as `HandleMessageStreamEvent[]` no longer typechecks — widen it to `StampedHandleMessageStreamEvent[]`. Events persisted by an earlier version have no `meta`; sessions started before this release will yield unstamped events when rewound, and consumers that read `event.meta` must tolerate that until those sessions age out. diff --git a/docs/concepts/sessions-runs-and-streaming.md b/docs/concepts/sessions-runs-and-streaming.md index 400196cef..c56a9b133 100644 --- a/docs/concepts/sessions-runs-and-streaming.md +++ b/docs/concepts/sessions-runs-and-streaming.md @@ -91,12 +91,14 @@ Alongside `type` and `data`, every event carries a `meta` envelope: } ``` -- **`meta.id`** uniquely identifies the event. It is a `evt_`-prefixed [ULID](https://github.com/ulid/spec), so sorting ids as strings reproduces emission order. +- **`meta.id`** uniquely identifies the event. It is an `evt_`-prefixed [ULID](https://github.com/ulid/spec): a millisecond timestamp followed by random bits, so ids are broadly time-ordered. - **`meta.at`** is the ISO-8601 time the event was emitted. `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. -That makes it the key for ingesting events into a database exactly once: +The envelope arrived in stream version 20. Events written by an earlier version are stored without it, so a session that started before you upgraded yields events with no `meta` when you rewind into that part of its stream. Guard the read (`event.meta?.id`) if your agent has live sessions that predate the upgrade. + +That makes it the key for ingesting a stream into a database without duplicating rows when you re-read it: ```sql insert into agent_events (id, session_id, type, data, emitted_at) @@ -104,14 +106,15 @@ values ($1, $2, $3, $4, $5) on conflict (id) do nothing; ``` -Because ids sort chronologically, a `primary key (id)` also keeps inserts clustered and lets you page with `where id > $cursor order by id`. +Because ids lead with a timestamp, a `primary key (id)` stays roughly append-ordered and keeps inserts clustered. -Two caveats worth knowing: +Three caveats worth knowing: +- **Ids are time-ordered, not a total order.** The turn steps of one session can run in different processes, each generating ids from its own clock and its own random bits. Two events emitted in the same millisecond by different steps may sort either way, and clock skew between machines can invert neighbours. Record your own ingestion sequence, or read the stream in order and store the index, when you need an exact ordering to page against — do not use `where id > $cursor` as a lossless cursor. The stream itself is authoritative: `startIndex` is an absolute event count. - **Ids identify events, not intent.** Two events with identical payloads — the `step.failed` → `turn.failed` → `session.failed` cascade, or two identical text deltas in one step — are distinct events with distinct ids. Deduplicate on `meta.id` only; matching on content would drop real data. - **A subagent's event is re-emitted, not shared.** When a parent forwards a child's event onto its own stream, the parent's copy is a separate event with its own id. Correlate the two streams through `subagent.called.data.childSessionId`. -Authored [hooks](../guides/hooks) receive the same envelope, so a hook can use `event.meta.id` to make its own side effects idempotent across a retry. +Authored [hooks](../guides/hooks) receive the same envelope. Note that a hook observes each event as it is emitted, not as it is read: if an interrupted step re-runs, the turn re-emits its events as new events with new ids, so `meta.id` is a key for a stored row rather than a retry guard. ## Send a follow-up message diff --git a/docs/guides/frontend/overview.mdx b/docs/guides/frontend/overview.mdx index 8353dd4d2..e67992843 100644 --- a/docs/guides/frontend/overview.mdx +++ b/docs/guides/frontend/overview.mdx @@ -226,10 +226,10 @@ const agent = useEveAgent({ reducer: toolCounter }); The browser conversation lives durably on the server. Persist both the rendered event log and the `session` cursor to pick it back up after a reload: ```tsx -import type { HandleMessageStreamEvent, SessionState } from "eve/client"; +import type { SessionState, StampedHandleMessageStreamEvent } from "eve/client"; type SavedEveChat = { - events?: readonly HandleMessageStreamEvent[]; + events?: readonly StampedHandleMessageStreamEvent[]; session?: SessionState; }; diff --git a/docs/guides/frontend/use-eve-agent-svelte.mdx b/docs/guides/frontend/use-eve-agent-svelte.mdx index 99127adf4..775606a2c 100644 --- a/docs/guides/frontend/use-eve-agent-svelte.mdx +++ b/docs/guides/frontend/use-eve-agent-svelte.mdx @@ -23,16 +23,16 @@ Import from `eve/svelte` and read the reactive getters directly. No `$` prefix: ## What it returns -| Property | Type | Description | -| --------- | ------------------------------------------- | ------------------------------------------------------------------------- | -| `data` | `TData` | Projected state. With the default reducer, `EveMessageData` (`messages`). | -| `status` | `UseEveAgentStatus` | `"ready"`, `"submitted"`, `"streaming"`, or `"error"`. | -| `error` | `Error \| undefined` | Last transport-level error. | -| `events` | `readonly HandleMessageStreamEvent[]` | Raw server events for this session. | -| `session` | `SessionState` | Snapshot of session state. | -| `send` | `(input: SendTurnPayload) => Promise` | Send text or a full turn (multi-part, attachments, HITL responses). | -| `stop` | `() => void` | Abort the client's in-flight stream. | -| `reset` | `() => void` | Clear state and start a new session. | +| Property | Type | Description | +| --------- | -------------------------------------------- | ------------------------------------------------------------------------- | +| `data` | `TData` | Projected state. With the default reducer, `EveMessageData` (`messages`). | +| `status` | `UseEveAgentStatus` | `"ready"`, `"submitted"`, `"streaming"`, or `"error"`. | +| `error` | `Error \| undefined` | Last transport-level error. | +| `events` | `readonly StampedHandleMessageStreamEvent[]` | Raw server events for this session. | +| `session` | `SessionState` | Snapshot of session state. | +| `send` | `(input: SendTurnPayload) => Promise` | Send text or a full turn (multi-part, attachments, HITL responses). | +| `stop` | `() => void` | Abort the client's in-flight stream. | +| `reset` | `() => void` | Clear state and start a new session. | These state fields are reactive getters, so read them straight from templates, `$derived`, or `$effect`. They are not stores, so don't prefix them with `$`. diff --git a/docs/guides/frontend/use-eve-agent-vue.mdx b/docs/guides/frontend/use-eve-agent-vue.mdx index 2cb6adab6..5a6f17eb7 100644 --- a/docs/guides/frontend/use-eve-agent-vue.mdx +++ b/docs/guides/frontend/use-eve-agent-vue.mdx @@ -25,16 +25,16 @@ const { data } = useEveAgent(); ## What it returns -| Property | Type | Description | -| --------- | -------------------------------------------------- | ------------------------------------------------------------------------- | -| `data` | `ComputedRef` | Projected state. With the default reducer, `EveMessageData` (`messages`). | -| `status` | `ComputedRef` | `"ready"`, `"submitted"`, `"streaming"`, or `"error"`. | -| `error` | `ComputedRef` | Last transport-level error. | -| `events` | `ComputedRef` | Raw server events for this session. | -| `session` | `ComputedRef` | Snapshot of session state. | -| `send` | `(input: SendTurnPayload) => Promise` | Send text or a full turn (multi-part, attachments, HITL responses). | -| `stop` | `() => void` | Abort the client's in-flight stream. | -| `reset` | `() => void` | Clear state and start a new session. | +| Property | Type | Description | +| --------- | --------------------------------------------------------- | ------------------------------------------------------------------------- | +| `data` | `ComputedRef` | Projected state. With the default reducer, `EveMessageData` (`messages`). | +| `status` | `ComputedRef` | `"ready"`, `"submitted"`, `"streaming"`, or `"error"`. | +| `error` | `ComputedRef` | Last transport-level error. | +| `events` | `ComputedRef` | Raw server events for this session. | +| `session` | `ComputedRef` | Snapshot of session state. | +| `send` | `(input: SendTurnPayload) => Promise` | Send text or a full turn (multi-part, attachments, HITL responses). | +| `stop` | `() => void` | Abort the client's in-flight stream. | +| `reset` | `() => void` | Clear state and start a new session. | The first five are `ComputedRef`s; the rest are methods. Destructure whatever you need, since refs keep their reactivity through destructuring. Read them with `.value` in `