diff --git a/.changeset/background-task-web-chat.md b/.changeset/background-task-web-chat.md new file mode 100644 index 0000000000..8213c0cb4d --- /dev/null +++ b/.changeset/background-task-web-chat.md @@ -0,0 +1,5 @@ +--- +"eve": patch +--- + +Keep frontend agent sessions attached for background task results, and prevent framework-authored task wake-ups from appearing as user messages or corrupting optimistic message reconciliation. diff --git a/docs/guides/frontend/overview.mdx b/docs/guides/frontend/overview.mdx index 879eabd868..7778e6a4b4 100644 --- a/docs/guides/frontend/overview.mdx +++ b/docs/guides/frontend/overview.mdx @@ -87,7 +87,7 @@ Most chat UIs only need `data.messages` and `status`. Drop down to `events` when `data.messages` are eve-owned `EveMessage[]`. Common text, reasoning, file, and dynamic-tool parts follow the [AI SDK `UIMessage`](https://ai-sdk.dev/docs/reference/ai-sdk-core/ui-message) rendering convention, but the types are not interchangeable. eve also exposes authorization and HITL metadata, and a file part's URL can be absent. Adapt those parts before passing messages to an API typed as `UIMessage[]`. -When the root agent delegates, its stream emits `subagent.called` with the child's `childSessionId`, then `subagent.completed` after admission with a working task receipt. Later task notifications wake the parent with updates or the final result. Detailed child progress lives on the child session's stream instead of being flattened into the root `data.messages`. Use the lower-level [TypeScript client](../client/overview#sessions) to attach to that ID when your UI needs live subagent activity. See [What the parent sees](../../subagents#what-the-parent-sees) for the complete contract. +When the root agent delegates, its stream emits `subagent.called` with the child's `childSessionId`, then `subagent.completed` after admission with a working task receipt. Later task notifications wake the parent with updates or the final result. The frontend helpers keep following these parent turns and project their assistant responses without exposing framework-authored task state as user messages. Detailed child progress lives on the child session's stream instead of being flattened into the root `data.messages`. Use the lower-level [TypeScript client](../client/overview#sessions) to attach to that ID when your UI needs live subagent activity. See [What the parent sees](../../subagents#what-the-parent-sees) for the complete contract. ## Sending and streaming diff --git a/packages/eve/src/client/background-task-follower.ts b/packages/eve/src/client/background-task-follower.ts new file mode 100644 index 0000000000..95e060808e --- /dev/null +++ b/packages/eve/src/client/background-task-follower.ts @@ -0,0 +1,90 @@ +import type { ClientSession } from "#client/session.js"; +import { isAbortError } from "#client/eve-agent-store-helpers.js"; +import type { MessageStreamEvent } from "#protocol/message.js"; + +interface BackgroundTaskFollowerCallbacks { + readonly acceptEvent: (event: MessageStreamEvent) => void; + readonly getSession: () => ClientSession | undefined; + readonly onError: (error: unknown) => void; + readonly onWaiting: (session: ClientSession) => void; +} + +export class BackgroundTaskFollower { + readonly #callbacks: BackgroundTaskFollowerCallbacks; + #controller: AbortController | undefined; + #enabled = false; + #promise: Promise | undefined; + + constructor(callbacks: BackgroundTaskFollowerCallbacks) { + this.#callbacks = callbacks; + } + + observe(event: MessageStreamEvent): void { + if (isBackgroundTaskReceiptEvent(event)) this.#enabled = true; + } + + seed(events: readonly MessageStreamEvent[]): void { + this.#enabled = events.some(isBackgroundTaskReceiptEvent); + } + + stop(): Promise | undefined { + this.#controller?.abort(); + return this.#promise; + } + + reset(): void { + this.#enabled = false; + this.#controller?.abort(); + this.#controller = undefined; + this.#promise = undefined; + } + + start(): void { + const session = this.#callbacks.getSession(); + if (!this.#enabled || session === undefined || this.#controller !== undefined) return; + + const controller = new AbortController(); + this.#controller = controller; + let promise!: Promise; + promise = this.#follow(session, controller).finally(() => { + if (this.#controller === controller) this.#controller = undefined; + if (this.#promise === promise) this.#promise = undefined; + }); + this.#promise = promise; + } + + async #follow(session: ClientSession, controller: AbortController): Promise { + try { + while (this.#enabled && !controller.signal.aborted) { + for await (const event of session.stream({ signal: controller.signal })) { + if (this.#controller !== controller) return; + this.#callbacks.acceptEvent(event); + if (event.type === "session.waiting") { + this.#callbacks.onWaiting(session); + } else if (event.type === "session.completed" || event.type === "session.failed") { + this.#enabled = false; + } + } + } + } catch (error) { + if (!isAbortError(error)) this.#callbacks.onError(error); + } + } +} + +export function isBackgroundTaskReceiptEvent(event: MessageStreamEvent): boolean { + if (event.type === "subagent.completed") { + return event.data.backgroundTask?.status === "working"; + } + if (event.type !== "action.result" || event.data.result.kind !== "tool-result") return false; + + const output = event.data.result.output; + return ( + typeof output === "object" && + output !== null && + "status" in output && + output.status === "working" && + "taskId" in output && + typeof output.taskId === "string" + ); +} diff --git a/packages/eve/src/client/eve-agent-store-helpers.ts b/packages/eve/src/client/eve-agent-store-helpers.ts index 970cdbfbc3..6570ae102d 100644 --- a/packages/eve/src/client/eve-agent-store-helpers.ts +++ b/packages/eve/src/client/eve-agent-store-helpers.ts @@ -1,6 +1,24 @@ -import type { SendTurnPayload } from "#client/types.js"; +import type { MessageResponse } from "#client/message-response.js"; +import type { CancelSessionResult, SendTurnPayload } from "#client/types.js"; import { isCurrentTurnBoundaryEvent, type MessageStreamEvent } from "#protocol/message.js"; -import type { UserContent } from "ai"; + +export interface ActiveTurn { + readonly abortController: AbortController; + acceptedFollowUps: number; + readonly cancel: () => Promise; + readonly completion: Promise; + readonly followUpDispatches: Set>; + receivedFollowUps: number; + readonly resolveCompletion: () => void; + readonly response: Promise; + readonly resolveResponse: (response: MessageResponse | undefined) => void; +} + +export interface PendingMessageSubmission { + readonly createdAt: number; + readonly id: string; + readonly message: string; +} export function isSettledSessionTail(events: readonly MessageStreamEvent[]): boolean { const tail = events.at(-1); @@ -52,20 +70,6 @@ export function createAbortSignal( return first ? AbortSignal.any([first, second]) : second; } -export function summarizeUserContent(message: string | UserContent): string { - if (typeof message === "string") return message; - - const parts: string[] = []; - for (const part of message) { - if (part.type === "text") { - parts.push(part.text); - } else if (part.type === "file") { - parts.push(part.filename ? `[file: ${part.filename}]` : "[file]"); - } - } - return parts.join("\n"); -} - export function isAbortError(error: unknown): boolean { return error instanceof Error && error.name === "AbortError"; } diff --git a/packages/eve/src/client/eve-agent-store.test.ts b/packages/eve/src/client/eve-agent-store.test.ts index bfaaf051a4..7424ee6e33 100644 --- a/packages/eve/src/client/eve-agent-store.test.ts +++ b/packages/eve/src/client/eve-agent-store.test.ts @@ -4,6 +4,7 @@ import { detachEveAgentStore, EveAgentStore } from "#client/eve-agent-store.js"; import { defaultMessageReducer } from "#client/message-reducer.js"; import { stampTestEvents } from "#internal/testing/events.js"; import { + createActionResultEvent, createMessageAppendedEvent, createMessageCompletedEvent, createMessageReceivedEvent, @@ -781,6 +782,157 @@ describe("EveAgentStore steering", () => { }); }); +describe("EveAgentStore background tasks", () => { + it("keeps following after a background tool receipt", async () => { + const initialEvents = stampTestEvents([ + createMessageReceivedEvent({ message: "Hello", sequence: 0, turnId: "turn_0" }), + createActionResultEvent({ + result: { + callId: "call_1", + kind: "tool-result", + output: { status: "working", taskId: "task_1" }, + toolName: "write_later", + }, + sequence: 0, + stepIndex: 0, + turnId: "turn_0", + }), + createMessageCompletedEvent({ + finishReason: "stop", + message: "The background task started.", + sequence: 0, + stepIndex: 1, + turnId: "turn_0", + }), + createSessionWaitingEvent(), + ] as UnstampedMessageStreamEvent[]); + const callbackStream = controlledStreamResponse(); + const [callbackStarted, callbackCompleted, callbackWaiting] = stampTestEvents([ + createTurnStartedEvent({ sequence: 1, turnId: "turn_1" }), + createMessageCompletedEvent({ + finishReason: "stop", + message: "The background task finished.", + sequence: 1, + stepIndex: 0, + turnId: "turn_1", + }), + createSessionWaitingEvent(), + ] as UnstampedMessageStreamEvent[]).map((event, index) => ({ + ...event, + meta: { ...event.meta, id: `callback_${index}` }, + })); + const fetchMock = vi + .spyOn(globalThis, "fetch") + .mockResolvedValueOnce(startedResponse()) + .mockResolvedValueOnce(streamResponse(initialEvents)) + .mockResolvedValueOnce(callbackStream.response); + const store = new EveAgentStore({ reducer: defaultMessageReducer() }); + + await store.send({ message: "Hello" }); + await vi.waitFor(() => expect(fetchMock).toHaveBeenCalledTimes(3)); + callbackStream.emit(callbackStarted!); + callbackStream.emit(callbackCompleted!); + callbackStream.emit(callbackWaiting!); + + await vi.waitFor(() => + expect(store.snapshot.data.messages.at(-1)?.parts).toContainEqual({ + state: "done", + stepIndex: 0, + text: "The background task finished.", + type: "text", + }), + ); + expect(store.snapshot.status).toBe("ready"); + detachEveAgentStore(store); + }); + + it("recognizes background subagent receipts", async () => { + const initialEvents = stampTestEvents([ + createMessageReceivedEvent({ message: "Hello", sequence: 0, turnId: "turn_0" }), + { + data: { + backgroundTask: { status: "working", taskId: "task_1" }, + callId: "call_1", + output: "Started.", + subagentName: "researcher", + }, + type: "subagent.completed", + }, + createSessionWaitingEvent(), + ] as UnstampedMessageStreamEvent[]); + const callbackStream = controlledStreamResponse(); + const fetchMock = vi + .spyOn(globalThis, "fetch") + .mockResolvedValueOnce(startedResponse()) + .mockResolvedValueOnce(streamResponse(initialEvents)) + .mockResolvedValueOnce(callbackStream.response); + const store = new EveAgentStore({ reducer: defaultMessageReducer() }); + + await store.send({ message: "Hello" }); + await vi.waitFor(() => expect(fetchMock).toHaveBeenCalledTimes(3)); + + detachEveAgentStore(store); + }); + + it("reconciles an optimistic submission containing a file", async () => { + const message = [ + { text: "Review this document", type: "text" as const }, + { + data: "data:text/plain;base64,SGVsbG8=", + filename: "notes.txt", + mediaType: "text/plain", + type: "file" as const, + }, + ]; + const events = stampTestEvents([ + createMessageReceivedEvent({ message, sequence: 0, turnId: "turn_1" }), + createSessionWaitingEvent(), + ] as UnstampedMessageStreamEvent[]); + vi.spyOn(globalThis, "fetch") + .mockResolvedValueOnce(startedResponse()) + .mockResolvedValueOnce(streamResponse(events)); + const store = new EveAgentStore({ reducer: defaultMessageReducer() }); + + await store.send({ message }); + + const userMessages = store.snapshot.data.messages.filter( + (candidate) => candidate.role === "user", + ); + expect(userMessages).toHaveLength(1); + expect(userMessages[0]?.metadata?.optimistic).toBeUndefined(); + expect(userMessages[0]?.parts).toMatchObject([ + { state: "done", text: "Review this document", type: "text" }, + { filename: "notes.txt", mediaType: "text/plain", type: "file" }, + ]); + }); + + it("only reconciles an optimistic submission with its matching server message", async () => { + const events = stampTestEvents([ + createMessageReceivedEvent({ + message: "Framework-authored task state", + sequence: 0, + turnId: "turn_internal", + }), + createMessageReceivedEvent({ message: "Hello", sequence: 1, turnId: "turn_1" }), + createSessionWaitingEvent(), + ] as UnstampedMessageStreamEvent[]); + vi.spyOn(globalThis, "fetch") + .mockResolvedValueOnce(startedResponse()) + .mockResolvedValueOnce(streamResponse(events)); + const store = new EveAgentStore({ reducer: defaultMessageReducer() }); + + await store.send({ message: "Hello" }); + + const userMessages = store.snapshot.data.messages.filter((message) => message.role === "user"); + expect(userMessages).toHaveLength(2); + expect(userMessages.map((message) => message.parts[0])).toEqual([ + { state: "done", text: "Hello", type: "text" }, + { state: "done", text: "Framework-authored task state", type: "text" }, + ]); + expect(userMessages[0]?.metadata?.optimistic).toBeUndefined(); + }); +}); + describe("EveAgentStore terminal failure", () => { it("publishes a live terminal failure with error status", async () => { const failed = stampTestEvents([ diff --git a/packages/eve/src/client/eve-agent-store.ts b/packages/eve/src/client/eve-agent-store.ts index eb331421bc..f396dd2a9a 100644 --- a/packages/eve/src/client/eve-agent-store.ts +++ b/packages/eve/src/client/eve-agent-store.ts @@ -1,17 +1,23 @@ +import { BackgroundTaskFollower } from "#client/background-task-follower.js"; import { Client } from "#client/client.js"; import type { MessageResponse } from "#client/message-response.js"; import type { EveAgentReducer, EveAgentReducerEvent } from "#client/reducer.js"; import type { ClientSession } from "#client/session.js"; import { createEventDeduper } from "#protocol/event-dedupe.js"; -import { isCurrentTurnBoundaryEvent, type MessageStreamEvent } from "#protocol/message.js"; import { + isCurrentTurnBoundaryEvent, + summarizeUserContent, + type MessageStreamEvent, +} from "#protocol/message.js"; +import { + type ActiveTurn, assertExclusiveTurnInput, collectPendingAuthorizations, createAbortSignal, createSubmissionId, isAbortError, isSettledSessionTail, - summarizeUserContent, + type PendingMessageSubmission, toTerminalStreamFailureError, updatePendingAuthorizations, } from "#client/eve-agent-store-helpers.js"; @@ -95,26 +101,8 @@ export interface EveAgentStoreInit { readonly session?: ClientSession; } -interface PendingMessageSubmission { - readonly createdAt: number; - readonly id: string; - readonly message: string; -} - const detachStore = Symbol("detachEveAgentStore"); -interface ActiveTurn { - readonly abortController: AbortController; - acceptedFollowUps: number; - readonly cancel: () => Promise; - readonly completion: Promise; - readonly followUpDispatches: Set>; - receivedFollowUps: number; - readonly resolveCompletion: () => void; - readonly response: Promise; - readonly resolveResponse: (response: MessageResponse | undefined) => void; -} - /** * Framework-agnostic state machine for an eve agent session. * @@ -135,10 +123,9 @@ export class EveAgentStore { readonly #reducer: EveAgentReducer; readonly #subscribers = new Set<() => void>(); - /** Ids already folded into the projection: `initialEvents` and a reconnect can overlap. */ #seenEvents = createEventDeduper(); - #activeTurn: ActiveTurn | undefined; + readonly #backgroundTaskFollower: BackgroundTaskFollower; #callbacks: EveAgentStoreCallbacks = {}; #data: TData; #error: Error | undefined; @@ -159,8 +146,6 @@ export class EveAgentStore { headers: init.headers, host: init.host ?? "", }); - // Seed the deduper from the saved log so a live stream that replays the - // same prefix does not double-apply it. const initialEvents: MessageStreamEvent[] = []; for (const event of init.initialEvents ?? []) { if (this.#seenEvents.admit(event)) initialEvents.push(event); @@ -179,6 +164,23 @@ export class EveAgentStore { this.#data = this.#reduceProjectionEvents(this.#projectionEvents); this.#snapshot = this.#createSnapshot(); + this.#backgroundTaskFollower = new BackgroundTaskFollower({ + acceptEvent: (event) => this.#acceptServerEvent(event), + getSession: () => (this.#activeTurn === undefined ? this.#session : undefined), + onError: (error) => { + this.#error = toError(error); + this.#status = "error"; + this.#callbacks.onError?.(this.#error); + this.#publish(); + }, + onWaiting: (session) => { + this.#status = "ready"; + this.#callbacks.onSessionChange?.(session.state); + this.#publish(); + this.#callbacks.onFinish?.(this.#snapshot); + }, + }); + this.#backgroundTaskFollower.seed(initialEvents); } get snapshot(): EveAgentStoreSnapshot { @@ -197,6 +199,8 @@ export class EveAgentStore { } async send(input: SendTurnPayload): Promise { + const stoppedBackgroundFollower = this.#backgroundTaskFollower.stop(); + if (stoppedBackgroundFollower !== undefined) await stoppedBackgroundFollower; if (this.#activeTurn !== undefined) { if (this.#status === "resuming") { throw new Error("eve session is resuming."); @@ -284,6 +288,7 @@ export class EveAgentStore { this.#publish(); this.#callbacks.onFinish?.(this.#snapshot); turn.resolveCompletion(); + this.#backgroundTaskFollower.start(); } } } @@ -302,6 +307,8 @@ export class EveAgentStore { } async #resume(): Promise { + const stoppedBackgroundFollower = this.#backgroundTaskFollower.stop(); + if (stoppedBackgroundFollower !== undefined) await stoppedBackgroundFollower; if ( this.#status === "resuming" || this.#status === "streaming" || @@ -394,6 +401,7 @@ export class EveAgentStore { this.#publish(); this.#callbacks.onFinish?.(this.#snapshot); turn.resolveCompletion(); + this.#backgroundTaskFollower.start(); } } } @@ -412,11 +420,13 @@ export class EveAgentStore { [detachStore](): void { this.#activeTurn?.abortController.abort(); + void this.#backgroundTaskFollower.stop(); } reset(): void { const turn = this.#activeTurn; this.#activeTurn = undefined; + this.#backgroundTaskFollower.reset(); turn?.resolveResponse(undefined); turn?.resolveCompletion(); turn?.abortController.abort(); @@ -568,6 +578,7 @@ export class EveAgentStore { ): void { if (!this.#seenEvents.admit(event)) return; this.#events = [...this.#events, event]; + this.#backgroundTaskFollower.observe(event); this.#applyServerEvent(event); this.#callbacks.onEvent?.(event); if (options.transitionToStreaming ?? true) this.#status = "streaming"; @@ -577,7 +588,7 @@ export class EveAgentStore { #applyServerEvent(event: MessageStreamEvent): void { const pendingSubmission = this.#pendingMessageSubmissions[0]; - if (event.type === "message.received" && pendingSubmission !== undefined) { + if (event.type === "message.received" && event.data.message === pendingSubmission?.message) { const submissionId = pendingSubmission.id; this.#pendingMessageSubmissions = this.#pendingMessageSubmissions.slice(1); this.#replaceProjectionEvent( diff --git a/packages/eve/src/harness/emission.test.ts b/packages/eve/src/harness/emission.test.ts index 488e3e734d..bdde977a6d 100644 --- a/packages/eve/src/harness/emission.test.ts +++ b/packages/eve/src/harness/emission.test.ts @@ -13,6 +13,8 @@ import { getTurnClientContextState, setTurnClientContextState, } from "#harness/turn-client-context.js"; +import { ContextContainer, contextStorage } from "#context/container.js"; +import { TurnTaskDeliveryKey } from "#context/keys.js"; import type { HarnessEmitFn, HarnessSession } from "#harness/types.js"; import { EMPTY_DELIVERY_SENTINEL } from "#shared/empty-delivery.js"; @@ -150,6 +152,45 @@ describe("setHarnessEmissionState", () => { }); describe("emitTurnPreamble", () => { + it.each(["pending", "settled"] as const)( + "does not expose a %s task-delivery prompt as a received user message", + async (phase) => { + const events: Array[0]> = []; + const ctx = new ContextContainer(); + ctx.set(TurnTaskDeliveryKey, phase); + + await contextStorage.run(ctx, () => + emitTurnPreamble( + async (event) => { + events.push(event); + }, + { message: "Framework-authored task state" }, + { sequence: 1, sessionStarted: true, stepIndex: 0, turnId: "" }, + ), + ); + + expect(events).toEqual([{ data: { sequence: 1, turnId: "turn_1" }, type: "turn.started" }]); + }, + ); + + it("keeps the initiating human message visible", async () => { + const events: Array[0]> = []; + const ctx = new ContextContainer(); + ctx.set(TurnTaskDeliveryKey, "initiating"); + + await contextStorage.run(ctx, () => + emitTurnPreamble( + async (event) => { + events.push(event); + }, + { message: "Start background work" }, + { sequence: 0, sessionStarted: true, stepIndex: 0, turnId: "" }, + ), + ); + + expect(events.map((event) => event.type)).toEqual(["turn.started", "message.received"]); + }); + it("attaches one trace context to the session and turn start events", async () => { const events: Array[0]> = []; const trace = { diff --git a/packages/eve/src/harness/emission.ts b/packages/eve/src/harness/emission.ts index 2cdb30fffb..19b4a54069 100644 --- a/packages/eve/src/harness/emission.ts +++ b/packages/eve/src/harness/emission.ts @@ -15,6 +15,8 @@ import type { RuntimeIdentity, RuntimeTraceContext, } from "#protocol/message.js"; +import { contextStorage } from "#context/container.js"; +import { TurnTaskDeliveryKey } from "#context/keys.js"; import { createActionsRequestedEvent, createActionInputAppendedEvent, @@ -88,7 +90,12 @@ export async function emitTurnPreamble( await emitFn(createTurnStartedEvent({ sequence: state.sequence, trace: traceContext, turnId })); - if (input.message !== undefined) { + const taskDeliveryPhase = contextStorage.getStore()?.get(TurnTaskDeliveryKey); + if ( + input.message !== undefined && + taskDeliveryPhase !== "pending" && + taskDeliveryPhase !== "settled" + ) { await emitFn( createMessageReceivedEvent({ message: input.message, diff --git a/packages/eve/src/protocol/message.ts b/packages/eve/src/protocol/message.ts index fc45f348f4..3f371e01b1 100644 --- a/packages/eve/src/protocol/message.ts +++ b/packages/eve/src/protocol/message.ts @@ -892,7 +892,7 @@ export function createMessageReceivedEvent(input: { }; } -function summarizeUserContent(message: string | UserContent): string { +export function summarizeUserContent(message: string | UserContent): string { if (typeof message === "string") { return message; }