From 614404f65ba55de51dab11dad1ca8e833d599d2c Mon Sep 17 00:00:00 2001 From: "Builder.io" Date: Wed, 9 Sep 2026 17:08:53 +0000 Subject: [PATCH] feat(agentkit): carry runs over the AG-UI wire with a typed extension profile Implements Option C from the AgentKit protocol evaluation: AG-UI (@ag-ui/core 0.0.59) becomes the wire format, and the Builder events it has no counterpart for travel as a versioned typed profile over CUSTOM. The domain AgentEvent survives in memory and is rebuilt by a Builder-side parser, so the client reducer, its exhaustive switch, and the React layer are untouched. Encoding holds one domain event to one AG-UI frame, which keeps the sequence contract, gap rejection, and reconnect cursor intact. Sequencing, event ids, and profile version are defined as explicit extensions because AG-UI specifies none of them; they ride in one owned metadata key, since @ag-ui/core reserves "ag-ui" for itself. Approvals move to AG-UI's interrupt model: the transport exposes resumeRun with resume entries, replacing resolveApproval and its bespoke route. Conformance now asserts that every event from Core's adapter round-trips through AG-UI unchanged and that each frame parses under AG-UI's own schemas. --- .changeset/agentkit-agui-wire-profile.md | 22 + packages/agentkit/ARCHITECTURE.md | 15 + packages/agentkit/README.md | 18 + packages/agentkit/package.json | 1 + packages/agentkit/src/adapters/http.spec.ts | 143 +++++- packages/agentkit/src/adapters/http.ts | 114 +++-- packages/agentkit/src/client/client.spec.ts | 185 +++++++ packages/agentkit/src/client/client.ts | 74 ++- packages/agentkit/src/client/state.ts | 2 +- .../agentkit/src/conformance/index.spec.ts | 91 +++- packages/agentkit/src/conformance/index.ts | 166 +++++-- .../agentkit/src/protocol/agui-codec.spec.ts | 468 ++++++++++++++++++ packages/agentkit/src/protocol/agui-codec.ts | 338 +++++++++++++ packages/agentkit/src/protocol/agui.ts | 388 +++++++++++++++ .../agentkit/src/protocol/hardening.spec.ts | 10 + packages/agentkit/src/protocol/index.ts | 41 ++ .../agentkit/src/protocol/validation.spec.ts | 17 +- packages/agentkit/src/protocol/validation.ts | 48 ++ packages/agentkit/src/protocol/version.ts | 2 +- .../src/client/chat/agentkit-protocol.spec.ts | 260 ++++++++-- .../core/src/client/chat/agentkit-protocol.ts | 173 ++++++- pnpm-lock.yaml | 10 + 22 files changed, 2414 insertions(+), 172 deletions(-) create mode 100644 .changeset/agentkit-agui-wire-profile.md create mode 100644 packages/agentkit/src/protocol/agui-codec.spec.ts create mode 100644 packages/agentkit/src/protocol/agui-codec.ts create mode 100644 packages/agentkit/src/protocol/agui.ts diff --git a/.changeset/agentkit-agui-wire-profile.md b/.changeset/agentkit-agui-wire-profile.md new file mode 100644 index 00000000000..025a6651f15 --- /dev/null +++ b/.changeset/agentkit-agui-wire-profile.md @@ -0,0 +1,22 @@ +--- +"@agent-native/core": patch +"@agent-native/agentkit": minor +--- + +**Breaking (pre-1.0):** AgentKit protocol v2 intentionally rejects v1-only +peers because the AG-UI envelope is not wire-compatible with the original +Builder envelope. Upgrade the AgentKit client and server together, then rerun +transport conformance before deploying a custom adapter. The deprecated +`resolveApproval` API remains only as a source-compatibility bridge after both +peers are on v2. + +Carry AgentKit runs over the AG-UI wire format instead of a Builder-only +envelope. Overlapping events map onto native AG-UI event types, and the +Builder-specific events travel as a versioned typed extension profile over +`CUSTOM`, so a stock AG-UI client can read the stream while AgentKit consumers +still receive fully typed domain events. Sequencing, replay cursors, and profile +version negotiation are defined as explicit extensions because AG-UI specifies +none of them. Approvals now use AG-UI's interrupt model: the Core transport +exposes `resumeRun` with `resume` entries in place of `resolveApproval`. An +approval interrupt terminally closes its protocol run, and `resumeRun` returns +the distinct replacement run that carries the resolution and continued work. diff --git a/packages/agentkit/ARCHITECTURE.md b/packages/agentkit/ARCHITECTURE.md index f05feb569d5..5362635afa9 100644 --- a/packages/agentkit/ARCHITECTURE.md +++ b/packages/agentkit/ARCHITECTURE.md @@ -267,6 +267,21 @@ distinguishes available, degraded, unavailable, unsupported, and omitted unknown state. During the pre-1.0 period, consumers should review minor-release notes before upgrading a custom adapter. +Protocol v2 is one such explicitly breaking pre-1.0 minor. Its AG-UI envelope +is not wire-compatible with v1, so clients and servers must upgrade together +and custom adapters must rerun conformance before deployment. V2-only peers +reject v1 at negotiation instead of guessing a fallback; the deprecated +`resolveApproval` API is only a source-compatibility bridge after both peers +have upgraded. + +Approvals follow AG-UI's terminal interrupt lifecycle. The request closes the +interrupted protocol run, and `resumeRun()` returns a distinct replacement run +that contains the resolution and continued events. + +Conformance enforces `resumeRun()` for transports that advertise protocol v2. +Unversioned compatibility transports may temporarily use the deprecated +`resolveApproval()` bridge while custom adapters migrate. + Conformance follows negotiated capabilities. A transport that declares `resumableRuns` must prove cursor replay. A transport that declares `durableThreadSnapshots` must return a value accepted by diff --git a/packages/agentkit/README.md b/packages/agentkit/README.md index 6052551c686..983c596654c 100644 --- a/packages/agentkit/README.md +++ b/packages/agentkit/README.md @@ -774,6 +774,13 @@ where `true` means available, `false` means unsupported, and omission means unknown. Breaking required wire changes add a protocol version instead of guessing a fallback. +Protocol v2 is an explicitly breaking pre-1.0 minor: its AG-UI envelope is not +wire-compatible with v1. Upgrade AgentKit clients and servers together, then +rerun transport conformance before deploying a custom adapter. V2-only peers +reject v1 rather than silently decoding it, and the deprecated +`resolveApproval` API is only a source-compatibility bridge once both peers use +v2. + New transports implement `discoverCapabilities(input)` and return an `AgentCapabilitiesDiscovery` descriptor for every requested capability. `degraded` and `unavailable` descriptors carry a typed `capability_unavailable` @@ -787,6 +794,17 @@ apps pin it through Core's dependency rather than resolving a `latest` tag. Read release notes for minor updates, and run transport conformance after upgrading a custom adapter. +Approval requests are terminal interrupts on the wire. The interrupted run +closes after the request; `resumeRun()` returns a distinct replacement run id +whose stream begins with the approval resolution and carries the continued +work. Consumers must subscribe to that returned run instead of waiting for more +events on the interrupted run. + +Transport conformance requires `resumeRun()` when a transport advertises +protocol v2. An unversioned compatibility transport may temporarily advertise +approvals through the deprecated `resolveApproval()` bridge, allowing custom +adapters to migrate without weakening the v2 lifecycle contract. + ## Migrate an existing Core chat surface Keep the Core runtime and replace the presentation boundary in one pass: diff --git a/packages/agentkit/package.json b/packages/agentkit/package.json index 18c6ce206a8..81d4b20baeb 100644 --- a/packages/agentkit/package.json +++ b/packages/agentkit/package.json @@ -83,6 +83,7 @@ "prepublishOnly": "npm run build" }, "dependencies": { + "@ag-ui/core": "0.0.59", "@agent-native/toolkit": "workspace:*", "@tabler/icons-react": "catalog:", "react-markdown": "^10.1.0", diff --git a/packages/agentkit/src/adapters/http.spec.ts b/packages/agentkit/src/adapters/http.spec.ts index 3bb2e3f9a48..83efc33add8 100644 --- a/packages/agentkit/src/adapters/http.spec.ts +++ b/packages/agentkit/src/adapters/http.spec.ts @@ -3,8 +3,11 @@ import { describe, expect, it, vi } from "vitest"; import type { AgentEvent, AgentTransport } from "../protocol/index.js"; import { AgentKitProtocolError, + AGENTKIT_PROFILE_HEADER, + AGENTKIT_PROFILE_ID, createAgentKitProtocolVersionOffer, createAgentProtocolEnvelope, + encodeAgentEvent, } from "../protocol/index.js"; import { AgentKitHttpError, @@ -12,6 +15,17 @@ import { createAgentKitHttpTransport, } from "./http.js"; +const sseEvent = (sequence: number): AgentEvent => ({ + id: `event-${sequence}`, + threadId: "thread-1", + runId: "run-1", + sequence, + occurredAt: "2026-08-29T00:00:00.000Z", + type: "message.delta", + messageId: "message-1", + text: String(sequence), +}); + describe("AgentKit HTTP adapter", () => { it("round-trips commands and resumable SSE through Fetch primitives", async () => { const events: AgentEvent[] = [ @@ -128,6 +142,47 @@ describe("AgentKit HTTP adapter", () => { ); }); + it("keeps the deprecated approval route as a resume bridge", async () => { + const resumeRun = vi.fn(async () => ({ runId: "run-1" })); + const handler = createAgentKitHttpHandler({ + transport: { + capabilities: { approvals: true }, + async startRun() { + return { runId: "run-1" }; + }, + async *subscribeToRun() {}, + async cancelRun() {}, + resumeRun, + }, + }); + const transport = createAgentKitHttpTransport({ + baseUrl: "https://agentkit.test/agentkit", + fetch: (input, init) => handler(new Request(input, init)), + }); + + await transport.resolveApproval?.({ + threadId: "thread-1", + runId: "run-1", + approvalId: "approval-1", + response: { decision: "deny" }, + }); + + expect(resumeRun).toHaveBeenCalledWith( + { + threadId: "thread-1", + runId: "run-1", + resume: [ + { + interruptId: "approval-1", + status: "resolved", + payload: { decision: "deny" }, + }, + ], + }, + expect.objectContaining({ signal: expect.any(AbortSignal) }), + ); + }); + it("rejects malformed commands and route identity mismatches", async () => { const queueMessage = vi.fn(); const handler = createAgentKitHttpHandler({ @@ -356,7 +411,10 @@ describe("AgentKit HTTP adapter", () => { const fetcher = vi.fn( async () => new Response("not an event stream", { - headers: { "content-type": "application/json" }, + headers: { + "content-type": "application/json", + [AGENTKIT_PROFILE_HEADER]: AGENTKIT_PROFILE_ID, + }, }), ); const transport = createAgentKitHttpTransport({ @@ -388,23 +446,45 @@ describe("AgentKit HTTP adapter", () => { }); }); - it("rejects non-event envelopes inside an SSE response", async () => { + it("rejects a frame that is not a valid AG-UI event", async () => { const transport = createAgentKitHttpTransport({ baseUrl: "https://agentkit.test/agentkit", createCorrelationId: () => "request-42", + fetch: async () => + new Response(`data: ${JSON.stringify({ runId: "run-1" })}\n\n`, { + headers: { + "content-type": "text/event-stream; charset=utf-8", + "x-agentkit-correlation-id": "request-42", + [AGENTKIT_PROFILE_HEADER]: AGENTKIT_PROFILE_ID, + }, + }), + }); + + await expect( + transport + .subscribeToRun({ threadId: "thread-1", runId: "run-1" }) + [Symbol.asyncIterator]() + .next(), + ).rejects.toThrow("is not a valid AG-UI event"); + }); + + it("skips a frame emitted by a producer outside the profile", async () => { + const foreign = JSON.stringify({ + type: "CUSTOM", + name: "someone-else/thing", + value: {}, + }); + const transport = createAgentKitHttpTransport({ + baseUrl: "https://agentkit.test/agentkit", fetch: async () => new Response( - `data: ${JSON.stringify( - createAgentProtocolEnvelope( - "response", - { runId: "run-1" }, - "request-42", - ), + `id: opaque-foreign\ndata: ${foreign}\n\nid: 1\ndata: ${JSON.stringify( + encodeAgentEvent(sseEvent(1)), )}\n\n`, { headers: { - "content-type": "text/event-stream; charset=utf-8", - "x-agentkit-correlation-id": "request-42", + "content-type": "text/event-stream", + [AGENTKIT_PROFILE_HEADER]: AGENTKIT_PROFILE_ID, }, }, ), @@ -415,27 +495,22 @@ describe("AgentKit HTTP adapter", () => { .subscribeToRun({ threadId: "thread-1", runId: "run-1" }) [Symbol.asyncIterator]() .next(), - ).rejects.toThrow("expected an event envelope"); + ).resolves.toMatchObject({ value: { sequence: 1 } }); }); it("rejects an SSE sequence gap before advancing the replay cursor", async () => { - const event = (sequence: number) => - createAgentProtocolEnvelope("event", { - id: `event-${sequence}`, - threadId: "thread-1", - runId: "run-1", - sequence, - occurredAt: "2026-08-29T00:00:00.000Z", - type: "message.delta", - messageId: "message-1", - text: String(sequence), - } satisfies AgentEvent); + const event = (sequence: number) => encodeAgentEvent(sseEvent(sequence)); const transport = createAgentKitHttpTransport({ baseUrl: "https://agentkit.test/agentkit", fetch: async () => new Response( `id: 1\ndata: ${JSON.stringify(event(1))}\n\nid: 3\ndata: ${JSON.stringify(event(3))}\n\n`, - { headers: { "content-type": "text/event-stream" } }, + { + headers: { + "content-type": "text/event-stream", + [AGENTKIT_PROFILE_HEADER]: AGENTKIT_PROFILE_ID, + }, + }, ), }); @@ -450,6 +525,24 @@ describe("AgentKit HTTP adapter", () => { ); }); + it("rejects an event stream that omits the negotiated profile", async () => { + const transport = createAgentKitHttpTransport({ + baseUrl: "https://agentkit.test/agentkit", + fetch: async () => + new Response( + `id: 1\ndata: ${JSON.stringify(encodeAgentEvent(sseEvent(1)))}\n\n`, + { headers: { "content-type": "text/event-stream" } }, + ), + }); + + await expect( + transport + .subscribeToRun({ threadId: "thread-1", runId: "run-1" }) + [Symbol.asyncIterator]() + .next(), + ).rejects.toMatchObject({ code: "unsupported_profile" }); + }); + it("rejects a backend sequence gap before writing an invalid SSE cursor", async () => { const handler = createAgentKitHttpHandler({ transport: { @@ -1155,12 +1248,12 @@ describe("AgentKit HTTP adapter", () => { }); const discovery = await transport.discoverCapabilities?.({ - protocol: { protocol: "agentkit", versions: [1] }, + protocol: createAgentKitProtocolVersionOffer(), }); expect(discovery?.protocol).toMatchObject({ status: "compatible", - selectedVersion: 1, + selectedVersion: 2, }); expect(discovery?.capabilities).toEqual( expect.arrayContaining([ diff --git a/packages/agentkit/src/adapters/http.ts b/packages/agentkit/src/adapters/http.ts index 7bc00cb8fcb..28b79dcfb42 100644 --- a/packages/agentkit/src/adapters/http.ts +++ b/packages/agentkit/src/adapters/http.ts @@ -1,4 +1,5 @@ import type { + AgentKitStreamContext, AgentActionResult, AgentCapabilities, AgentCapabilitiesDiscovery, @@ -11,6 +12,10 @@ import type { FilePart, } from "../protocol/index.js"; import { + AGENTKIT_PROFILE_HEADER, + AGENTKIT_PROFILE_ID, + decodeAgUiEvent, + encodeAgentEvent, AgentProtocolValidationError, AgentKitProtocolError, createCapabilityUnsupportedError, @@ -40,6 +45,7 @@ import { parseQueueMessageInput, parseQueueMessageResult, parseResolveApprovalInput, + parseResumeRunInput, parseResolveConnectionRequestInput, parseStartRunInput, parseStartRunResult, @@ -54,6 +60,7 @@ import { parseUpdateThreadInput, parseDiscoverCapabilitiesInput, parseVoidResponse, + resumeEntryFromApproval, } from "../protocol/index.js"; export interface AgentKitHttpTransportOptions { @@ -236,6 +243,7 @@ async function* parseEventStream( response: Response, expectedCorrelationId: string, afterSequence: number, + context: AgentKitStreamContext, ): AsyncGenerator { if (!response.ok) { await parseResponse(response, parseVoidResponse, expectedCorrelationId); @@ -245,6 +253,19 @@ async function* parseEventStream( ?.split(";", 1)[0] ?.trim() .toLowerCase(); + const profile = response.headers.get(AGENTKIT_PROFILE_HEADER); + if (profile !== AGENTKIT_PROFILE_ID) { + throw new AgentKitHttpError( + response.status, + profile === null + ? `The event stream did not declare the required ${AGENTKIT_PROFILE_HEADER} profile.` + : `The event stream declared profile ${profile}, which this client cannot read.`, + undefined, + "unsupported_profile", + responseCorrelationId(response, undefined, expectedCorrelationId), + false, + ); + } if (contentType !== "text/event-stream") { throw new AgentKitHttpError( response.status, @@ -265,11 +286,9 @@ async function* parseEventStream( false, ); } - const correlationId = responseCorrelationId( - response, - undefined, - expectedCorrelationId, - ); + // AG-UI frames carry no envelope, so correlation is proven once for the whole + // stream here instead of per event. This throws on a mismatched response. + responseCorrelationId(response, undefined, expectedCorrelationId); const reader = response.body.pipeThrough(new TextDecoderStream()).getReader(); let buffer = ""; let lastSequence = afterSequence; @@ -291,6 +310,8 @@ async function* parseEventStream( .find((line) => line.startsWith("id:")) ?.slice(3) .trim(); + const event = decodeAgUiEvent(JSON.parse(data) as unknown, context); + if (event === undefined) continue; if ( id !== undefined && (!/^\d+$/.test(id) || !Number.isSafeInteger(Number(id))) @@ -300,25 +321,6 @@ async function* parseEventStream( "expected a non-negative safe integer cursor", ); } - const envelope = parseAgentProtocolEnvelope( - JSON.parse(data) as unknown, - ); - if (envelope.kind !== "event") { - throw new AgentProtocolValidationError( - "envelope.kind", - "expected an event envelope", - ); - } - if ( - envelope.correlationId && - envelope.correlationId !== correlationId - ) { - throw new AgentProtocolValidationError( - "envelope.correlationId", - "must match the event stream correlation id", - ); - } - const event = parseAgentEvent(envelope.payload, "envelope.payload"); if (id !== undefined && Number(id) !== event.sequence) { throw new AgentProtocolValidationError( "event.id", @@ -562,7 +564,10 @@ export function createAgentKitHttpTransport( `${baseUrl}/runs/${encodeURIComponent(input.runId)}/events?${query}`, { headers, signal: combinedSignal.signal }, ); - yield* parseEventStream(response, correlationId, afterSequence); + yield* parseEventStream(response, correlationId, afterSequence, { + threadId: input.threadId, + runId: input.runId, + }); } finally { combinedSignal.release(); } @@ -587,6 +592,17 @@ export function createAgentKitHttpTransport( context, ); }, + resumeRun(input, context) { + return request( + `/runs/${encodeURIComponent(input.runId)}/resume`, + parseStartRunResult, + { + method: "POST", + body: JSON.stringify(createAgentProtocolEnvelope("request", input)), + }, + context, + ); + }, async resolveApproval(input, context) { await request( `/runs/${encodeURIComponent(input.runId)}/approvals/${encodeURIComponent(input.approvalId)}`, @@ -889,7 +905,7 @@ function eventStream( lastSequence = event.sequence; controller.enqueue( encoder.encode( - `id: ${event.sequence}\ndata: ${JSON.stringify(createAgentProtocolEnvelope("event", event, correlationId))}\n\n`, + `id: ${event.sequence}\ndata: ${JSON.stringify(encodeAgentEvent(event))}\n\n`, ), ); } catch (error) { @@ -910,6 +926,7 @@ function eventStream( headers: { "cache-control": "no-cache, no-transform", "content-type": "text/event-stream", + [AGENTKIT_PROFILE_HEADER]: AGENTKIT_PROFILE_ID, ...(correlationId ? { "x-agentkit-correlation-id": correlationId } : {}), @@ -1293,15 +1310,25 @@ export function createAgentKitHttpHandler( await invoke((context) => transport.cancelRun(input, context)); return respond(undefined); } + const resumeMatch = path.match(/^\/runs\/([^/]+)\/resume$/); + if (request.method === "POST" && resumeMatch) { + const resumeRun = transport.resumeRun; + if (!resumeRun) { + return notSupported("resumeRun", "Approvals are not supported."); + } + const input = await readPayload(parseResumeRunInput); + assertRouteIdentifier( + decodeURIComponent(resumeMatch[1] ?? ""), + input.runId, + "runId", + ); + return respond( + await invoke((context) => resumeRun(input, context)), + 201, + ); + } const approvalMatch = path.match(/^\/runs\/([^/]+)\/approvals\/([^/]+)$/); if (request.method === "POST" && approvalMatch) { - const resolveApproval = transport.resolveApproval; - if (!resolveApproval) { - return notSupported( - "resolveApproval", - "Approvals are not supported.", - ); - } const input = await readPayload(parseResolveApprovalInput); assertRouteIdentifier( decodeURIComponent(approvalMatch[1] ?? ""), @@ -1313,7 +1340,26 @@ export function createAgentKitHttpHandler( input.approvalId, "approvalId", ); - await invoke((context) => resolveApproval(input, context)); + if (transport.resolveApproval) { + await invoke((context) => transport.resolveApproval!(input, context)); + return respond(undefined); + } + if (!transport.resumeRun) { + return notSupported( + "resolveApproval", + "Approvals are not supported.", + ); + } + await invoke((context) => + transport.resumeRun!( + { + threadId: input.threadId, + runId: input.runId, + resume: [resumeEntryFromApproval(input)], + }, + context, + ), + ); return respond(undefined); } const connectionMatch = path.match( diff --git a/packages/agentkit/src/client/client.spec.ts b/packages/agentkit/src/client/client.spec.ts index 533e6c6abc8..c86e7c9ff62 100644 --- a/packages/agentkit/src/client/client.spec.ts +++ b/packages/agentkit/src/client/client.spec.ts @@ -8,6 +8,7 @@ import type { } from "../protocol/index.js"; import type { AgentStreamIntegrityReport } from "../protocol/index.js"; import { + AgentProtocolValidationError, AgentKitProtocolError, createAgentKitProtocolVersionOffer, createCapabilityUnsupportedError, @@ -465,6 +466,41 @@ describe("AgentKitClient", () => { expect(client.getThread("thread-1").runs["run-1"]?.lastSequence).toBe(2); }); + it("preserves an approval interrupt across a dropped stream", async () => { + let subscriptions = 0; + const reports: AgentStreamIntegrityReport[] = []; + const transport = createTransport([]); + transport.subscribeToRun = async function* (input) { + subscriptions += 1; + if (subscriptions === 1) { + yield protocolEvent(1, { type: "run.started" }); + yield protocolEvent(2, { + type: "approval.requested", + request: { id: "approval-1", title: "Continue?" }, + }); + throw new Error("connection dropped"); + } + expect(input.afterSequence).toBe(2); + }; + const client = new AgentKitClient({ + transport, + reconnect: { attempts: 1, delayMs: () => 0 }, + onIntegrityReport: (report) => reports.push(report), + }); + + const run = await client.sendMessage({ threadId: "thread-1", text: "Go" }); + await run.completed; + + expect(subscriptions).toBe(2); + expect(client.getThread("thread-1").runs["run-1"]?.status).toBe( + "awaiting_approval", + ); + expect(reports).not.toContainEqual( + expect.objectContaining({ code: "run_missing_terminal" }), + ); + expect(client.getSnapshot().connection).toBe("connected"); + }); + it("rejects a sequence gap before advancing the durable resume cursor", async () => { const subscribeToRun = vi.fn(async function* () { yield protocolEvent(1, { type: "run.started" }); @@ -1374,6 +1410,155 @@ describe("AgentKitClient", () => { expect(client.getThread("thread-1").activeRunIds).toEqual([]); }); + it("retires an interrupted run when approval resumes on a new run", async () => { + const transport: AgentTransport = { + capabilities: { approvals: true }, + async startRun() { + return { runId: "run-interrupted" }; + }, + async *subscribeToRun({ runId }) { + if (runId === "run-interrupted") { + yield { ...protocolEvent(1, { type: "run.started" }), runId }; + yield { + ...protocolEvent(2, { + type: "approval.requested", + request: { id: "approval-1", title: "Continue?" }, + }), + runId, + }; + return; + } + yield { ...protocolEvent(1, { type: "run.started" }), runId }; + yield { + ...protocolEvent(2, { + type: "approval.resolved", + approvalId: "approval-1", + response: { decision: "approve" }, + }), + runId, + }; + yield { ...protocolEvent(3, { type: "run.completed" }), runId }; + }, + async cancelRun() {}, + async resumeRun() { + return { runId: "run-resumed" }; + }, + }; + const client = new AgentKitClient({ transport }); + const run = await client.sendMessage({ + threadId: "thread-1", + text: "Start", + }); + await run.completed; + + expect(client.getThread("thread-1").activeRunIds).toEqual([ + "run-interrupted", + ]); + expect(client.getThread("thread-1").runs["run-interrupted"]?.status).toBe( + "awaiting_approval", + ); + + await client.resolveApproval({ + threadId: "thread-1", + runId: "run-interrupted", + approvalId: "approval-1", + response: { decision: "approve" }, + }); + + await vi.waitFor(() => + expect(client.getThread("thread-1").activeRunIds).toEqual([]), + ); + expect(client.getThread("thread-1").runs["run-resumed"]?.status).toBe( + "completed", + ); + }); + + it("reports replacement-run failures and permits an explicit reattach", async () => { + const onError = vi.fn(); + let replacementSubscriptions = 0; + const transport: AgentTransport = { + capabilities: { approvals: true }, + async startRun() { + return { runId: "run-interrupted" }; + }, + async *subscribeToRun({ runId }) { + if (runId === "run-interrupted") { + yield { ...protocolEvent(1, { type: "run.started" }), runId }; + yield { + ...protocolEvent(2, { + type: "approval.requested", + request: { id: "approval-1", title: "Continue?" }, + }), + runId, + }; + return; + } + replacementSubscriptions += 1; + yield { ...protocolEvent(1, { type: "run.started" }), runId }; + throw new AgentProtocolValidationError("replacement", "stream failed"); + }, + async cancelRun() {}, + async resumeRun() { + return { runId: "run-resumed" }; + }, + }; + const client = new AgentKitClient({ + transport, + reconnect: { attempts: 0 }, + onError, + }); + const run = await client.sendMessage({ + threadId: "thread-1", + text: "Start", + }); + await run.completed; + + await client.resolveApproval({ + threadId: "thread-1", + runId: "run-interrupted", + approvalId: "approval-1", + response: { decision: "approve" }, + }); + + await vi.waitFor(() => + expect(onError).toHaveBeenCalledWith( + expect.objectContaining({ code: "run_stream_failed" }), + ), + ); + expect(client.getSnapshot().error).toMatchObject({ + code: "run_stream_failed", + }); + await expect( + client.resubscribeRun("thread-1", "run-resumed"), + ).rejects.toThrow("replacement: stream failed"); + expect(replacementSubscriptions).toBe(2); + }); + + it("keeps the deprecated approval transport bridge operational", async () => { + const resolveApproval = vi.fn(async () => undefined); + const client = new AgentKitClient({ + transport: { + capabilities: { approvals: true }, + async startRun() { + return { runId: "run-1" }; + }, + async *subscribeToRun() {}, + async cancelRun() {}, + resolveApproval, + }, + }); + const input = { + threadId: "thread-1", + runId: "run-1", + approvalId: "approval-1", + response: { decision: "approve" as const }, + }; + + await client.resolveApproval(input); + + expect(resolveApproval).toHaveBeenCalledWith(input, expect.any(Object)); + }); + it("publishes a promoted queued run before its stream produces an event", async () => { const releaseStream = Promise.withResolvers(); const transport: AgentTransport = { diff --git a/packages/agentkit/src/client/client.ts b/packages/agentkit/src/client/client.ts index 6fe0881255f..16e0493adb3 100644 --- a/packages/agentkit/src/client/client.ts +++ b/packages/agentkit/src/client/client.ts @@ -37,6 +37,7 @@ import { createAgentKitProtocolVersionOffer, parseAgentEvent, projectAgentCapabilities, + resumeEntryFromApproval, } from "../protocol/index.js"; import { classifyAgentEvent, @@ -810,7 +811,7 @@ export class AgentKitClient implements AgentKitController { } this.markRunStarted(input.threadId, result.runId); const completed = this.consume(input.threadId, result.runId); - this.consumers.set(this.runKey(input.threadId, result.runId), completed); + this.trackConsumer(input.threadId, result.runId, completed); return { runId: result.runId, completed, @@ -840,7 +841,7 @@ export class AgentKitClient implements AgentKitController { const existing = this.consumers.get(key); if (existing) return existing; const consumer = this.consume(threadId, runId); - this.consumers.set(key, consumer); + this.trackConsumer(threadId, runId, consumer); return consumer; } @@ -857,14 +858,43 @@ export class AgentKitClient implements AgentKitController { this.assertActive(); const requestContext = this.createRequestContext(context); await this.requireCapability("approvals", requestContext); + const resumeRun = this.transport.resumeRun; const resolveApproval = this.transport.resolveApproval; - if (!resolveApproval) { + if (!resumeRun && !resolveApproval) { throw new AgentKitCapabilityError("approvals"); } - await this.invokeRequest(requestContext, (context) => - resolveApproval(input, context), + if (!resumeRun) { + await this.invokeRequest(requestContext, (context) => + resolveApproval!(input, context), + ); + this.assertActive(); + return; + } + const result = await this.invokeRequest(requestContext, (context) => + resumeRun( + { + threadId: input.threadId, + runId: input.runId, + resume: [resumeEntryFromApproval(input)], + }, + context, + ), ); this.assertActive(); + if (result.runId !== input.runId) { + this.retireInterruptedRun(input.threadId, input.runId); + } + // A runtime that suspends rather than terminates answers with the run that + // was already streaming, so adopting it blindly would open a second reader + // on one stream. + const key = this.runKey(input.threadId, result.runId); + if (this.consumers.has(key)) return; + this.markRunStarted(input.threadId, result.runId); + this.trackConsumer( + input.threadId, + result.runId, + this.consume(input.threadId, result.runId), + ); } public async resolveConnectionRequest( @@ -1191,7 +1221,7 @@ export class AgentKitClient implements AgentKitController { } this.markRunStarted(threadId, result.runId); const completed = this.consume(threadId, result.runId); - this.consumers.set(this.runKey(threadId, result.runId), completed); + this.trackConsumer(threadId, result.runId, completed); return { runId: result.runId, completed, @@ -1363,11 +1393,26 @@ export class AgentKitClient implements AgentKitController { return this.shutdown(); } + private trackConsumer( + threadId: ThreadId, + runId: RunId, + completed: Promise, + ): void { + this.consumers.set(this.runKey(threadId, runId), completed); + // `consume` reports failures through snapshot state and `onError` before it + // rejects. Observe internally-owned continuations here so a caller that + // cannot receive their completion promise never triggers an unhandled + // rejection; APIs that return `completed` still expose the original promise. + void completed.catch(() => undefined); + } + private async consume(threadId: ThreadId, runId: RunId): Promise { const key = this.runKey(threadId, runId); const abortController = new AbortController(); this.consumerAbortControllers.set(key, abortController); let attempt = 0; + let interruptedForApproval = + this.getThread(threadId).runs[runId]?.status === "awaiting_approval"; try { while (true) { const afterSequence = @@ -1401,6 +1446,11 @@ export class AgentKitClient implements AgentKitController { ); } this.applyEvent(event); + if (event.type === "approval.requested") { + interruptedForApproval = true; + } else if (event.type === "approval.resolved") { + interruptedForApproval = false; + } if ( event.type === "run.completed" || event.type === "run.failed" || @@ -1410,6 +1460,10 @@ export class AgentKitClient implements AgentKitController { } } if (!terminalEvent) { + if (interruptedForApproval) { + this.setConnection("connected"); + return; + } this.reportIntegrity({ code: "run_missing_terminal", threadId, @@ -2103,6 +2157,14 @@ export class AgentKitClient implements AgentKitController { }); } + private retireInterruptedRun(threadId: ThreadId, runId: RunId): void { + const thread = this.getThread(threadId); + this.setThread(threadId, { + ...thread, + activeRunIds: thread.activeRunIds.filter((id) => id !== runId), + }); + } + private runKey(threadId: ThreadId, runId: RunId): string { return `${threadId}\u0000${runId}`; } diff --git a/packages/agentkit/src/client/state.ts b/packages/agentkit/src/client/state.ts index 5c5b44318a9..2e8d3043e46 100644 --- a/packages/agentkit/src/client/state.ts +++ b/packages/agentkit/src/client/state.ts @@ -422,7 +422,7 @@ export function reduceAgentEvent( } case "approval.requested": return { - ...next, + ...updateRun(next, event.runId, { status: "awaiting_approval" }), approvals: { ...next.approvals, [event.request.id]: event.request }, approvalRunIds: { ...next.approvalRunIds, diff --git a/packages/agentkit/src/conformance/index.spec.ts b/packages/agentkit/src/conformance/index.spec.ts index 63adf211b0c..34cc27bd68f 100644 --- a/packages/agentkit/src/conformance/index.spec.ts +++ b/packages/agentkit/src/conformance/index.spec.ts @@ -18,7 +18,12 @@ import type { AgentRunSnapshot, AgentTransport, } from "../protocol/index.js"; -import { AGENTKIT_PROTOCOL_VERSION } from "../protocol/index.js"; +import { + AGENTKIT_PROTOCOL_VERSION, + approvalResponseFromResume, + resumeEntryFromApproval, + resumeOptionId, +} from "../protocol/index.js"; import { assertAgentTransportConformance, type AgentKitConformanceScenario, @@ -109,6 +114,7 @@ function waitForAbort(signal: AbortSignal | undefined): Promise { function createFixtureTransport( preparedScenario: AgentKitConformanceScenario, defect?: FixtureDefect, + resumeWithNewRun = false, ): AgentTransport { const capabilities: AgentCapabilities = { protocolVersion: @@ -322,7 +328,11 @@ function createFixtureTransport( await run.cancelled.promise; yield* run.events.slice(1); } - if (run.scenario === "approval" && run.events.length === 2) { + if ( + run.scenario === "approval" && + run.events.length === 2 && + !resumeWithNewRun + ) { await run.approved.promise; yield* run.events.slice(2); } @@ -451,20 +461,44 @@ function createFixtureTransport( }; if (defect !== "capabilities") { - transport.resolveApproval = async (input) => { + transport.resumeRun = async (input) => { const run = runs.get(input.runId); if (!run || run.threadId !== input.threadId) throw new Error("Unknown run"); + const entry = input.resume[0]!; + if (resumeWithNewRun) { + const resumed: FixtureRun = { + id: `approval-resumed-run-${++runNumber}`, + threadId: run.threadId, + scenario: "approval", + events: [], + subscriptions: 0, + cancelled: deferred(), + approved: deferred(), + }; + resumed.events.push( + event(resumed, 1, { type: "run.started" }), + event(resumed, 2, { + type: "approval.resolved", + approvalId: entry.interruptId, + optionId: resumeOptionId(entry), + response: approvalResponseFromResume(entry), + }), + event(resumed, 3, { type: "run.completed" }), + ); + runs.set(resumed.id, resumed); + return { runId: resumed.id }; + } if (defect !== "approval") { run.events.push( event(run, 3, { type: "approval.resolved", - approvalId: input.approvalId, - optionId: input.optionId, + approvalId: entry.interruptId, + optionId: resumeOptionId(entry), response: defect === "approval-decision" ? { decision: "allow" as "approve" } - : input.response, + : approvalResponseFromResume(entry), }), ); } @@ -472,6 +506,7 @@ function createFixtureTransport( event(run, run.events.length + 1, { type: "run.completed" }), ); run.approved.resolve(); + return { runId: run.id }; }; } @@ -658,6 +693,50 @@ describe("assertAgentTransportConformance", () => { ); }); + it("follows approval continuation onto a replacement run", async () => { + const report = await assertAgentTransportConformance({ + createTransport: (scenario) => + createFixtureTransport(scenario, undefined, scenario === "approval"), + timeoutMs: 250, + }); + + expect(report.checks).toContain("approval continuation"); + }); + + it("accepts the legacy approval bridge for compatibility transports", async () => { + const report = await assertAgentTransportConformance({ + createTransport(scenario) { + const transport = createFixtureTransport(scenario); + const resumeRun = transport.resumeRun!; + transport.resolveApproval = async (input) => { + await resumeRun({ + threadId: input.threadId, + runId: input.runId, + resume: [resumeEntryFromApproval(input)], + }); + }; + delete transport.resumeRun; + const { protocolVersion: _protocolVersion, ...legacyCapabilities } = + transport.capabilities!; + transport.capabilities = legacyCapabilities; + return transport; + }, + timeoutMs: 250, + }); + + expect(report.checks).toContain("approval continuation"); + }); + + it("requires resumeRun for protocol-v2 approval transports", async () => { + const transport = createFixtureTransport("completed"); + transport.resolveApproval = async () => undefined; + delete transport.resumeRun; + + await expect( + assertAgentTransportConformance({ transport }), + ).rejects.toThrow(/approvals is advertised/iu); + }); + it("keeps a baseline profile for transports without lifecycle fixtures", async () => { const transport = createFixtureTransport("completed"); let disposeCalls = 0; diff --git a/packages/agentkit/src/conformance/index.ts b/packages/agentkit/src/conformance/index.ts index 4bb5a2209be..3c2e86fb5fe 100644 --- a/packages/agentkit/src/conformance/index.ts +++ b/packages/agentkit/src/conformance/index.ts @@ -11,8 +11,12 @@ import { AGENTKIT_PROTOCOL_VERSION, createAgentKitProtocolVersionOffer, parseAgentCapabilities, + decodeAgUiEvent, + encodeAgentEvent, + isAgUiEvent, parseAgentEvent, parseAgentEventSequence, + resumeEntryFromApproval, parseAgentQueuedMessage, parseAgentThreadSnapshot, projectAgentCapabilities, @@ -307,6 +311,52 @@ function assertLifecycleCompleteness( return ["rich lifecycle replay", "rich lifecycle snapshot"]; } +/** Key order is not part of a JSON value, so comparison must not depend on it. */ +function canonicalJson(value: unknown): string { + return JSON.stringify(value, (_key, nested) => { + if ( + typeof nested !== "object" || + nested === null || + Array.isArray(nested) + ) { + return nested; + } + return Object.fromEntries( + Object.entries(nested as Record).sort(([a], [b]) => + a < b ? -1 : a > b ? 1 : 0, + ), + ); + }); +} + +/** + * Every domain event must survive the AG-UI wire and come back identical, and + * the frame in between must be one a stock AG-UI client can parse. Checking + * both together is the point: a profile that round-trips through frames only + * Builder can read would pass the first half and deliver no interoperability. + */ +function assertAgUiWireCompatibility( + events: readonly AgentEvent[], + context: { threadId: string; runId: string }, +): void { + for (const event of events) { + const frame = encodeAgentEvent(event); + assert( + isAgUiEvent(frame), + `Event ${event.type} did not encode to a valid AG-UI event.`, + ); + const decoded = decodeAgUiEvent(frame, context); + assert( + decoded !== undefined, + `Event ${event.type} did not decode back from its AG-UI frame.`, + ); + assert( + canonicalJson(decoded) === canonicalJson(event), + `Event ${event.type} changed while round-tripping through AG-UI: ${canonicalJson(event)} became ${canonicalJson(decoded)}.`, + ); + } +} + function assertOrderedEvents(input: { events: readonly AgentEvent[]; threadId: string; @@ -330,6 +380,7 @@ function assertOrderedEvents(input: { runId, path: "runEvents", }); + assertAgUiWireCompatibility(events, { threadId, runId }); if (requireStarted) { assert( events[0]?.type === "run.started", @@ -432,11 +483,13 @@ function assertCapabilitiesTruthful( }; requireOperation(capabilities.actions, transport.invokeAction, "actions"); - requireOperation( - capabilities.approvals, - transport.resolveApproval, - "approvals", - ); + if (capabilities.approvals) { + const operation = + capabilities.protocolVersion === AGENTKIT_PROTOCOL_VERSION + ? transport.resumeRun + : (transport.resumeRun ?? transport.resolveApproval); + requireOperation(true, operation, "approvals"); + } requireOperation(capabilities.feedback, transport.submitFeedback, "feedback"); requireOperation( capabilities.threadHistory, @@ -526,15 +579,19 @@ async function assertUnsupportedOperations( } if (!capabilities.approvals) { await probe( - "resolveApproval", - transport.resolveApproval + "resumeRun", + transport.resumeRun ? () => - transport.resolveApproval!({ + transport.resumeRun!({ threadId, runId: "conformance-unsupported-run", - approvalId: "conformance-unsupported-approval", - optionId: "approve", - response: { decision: "approve" }, + resume: [ + { + interruptId: "conformance-unsupported-approval", + status: "resolved", + payload: { decision: "approve" }, + }, + ], }) : undefined, ); @@ -1073,8 +1130,8 @@ async function assertApprovalContinuation(input: { timeoutMs: number; }): Promise { assert( - input.transport.resolveApproval, - "approvals requires resolveApproval.", + input.transport.resumeRun || input.transport.resolveApproval, + "approvals requires resumeRun or the legacy resolveApproval bridge.", ); const started = await startScenarioRun({ ...input, scenario: "approval" }); const iterator = input.transport @@ -1102,42 +1159,87 @@ async function assertApprovalContinuation(input: { ); if (event.type === "approval.requested") requested = event; } - await input.transport.resolveApproval({ - threadId: input.threadId, - runId: started.runId, + const approval = { approvalId: requested.request.id, optionId: "approve", - response: { decision: "approve", optionIds: ["approve"] }, - }); - events.push( - ...(await collectIterator( + response: { decision: "approve" as const, optionIds: ["approve"] }, + }; + const resumed = input.transport.resumeRun + ? await input.transport.resumeRun({ + threadId: input.threadId, + runId: started.runId, + resume: [resumeEntryFromApproval(approval)], + }) + : await input.transport.resolveApproval!({ + threadId: input.threadId, + runId: started.runId, + ...approval, + }).then(() => ({ runId: started.runId })); + let continuationEvents: AgentEvent[]; + if (resumed.runId === started.runId) { + continuationEvents = await collectIterator( iterator, input.timeoutMs, "approval continuation", - )), - ); - assertOrderedEvents({ - events, - threadId: input.threadId, - runId: started.runId, - terminal: "run.completed", - }); - const resolvedIndex = events.findIndex( + ); + events.push(...continuationEvents); + assertOrderedEvents({ + events, + threadId: input.threadId, + runId: started.runId, + terminal: "run.completed", + }); + } else { + const closed = await withTimeout( + iterator.next(), + input.timeoutMs, + "interrupted approval stream closure", + ); + assert( + closed.done, + "An interrupted run must close before its replacement starts.", + ); + assertOrderedEvents({ + events, + threadId: input.threadId, + runId: started.runId, + terminal: "none", + }); + continuationEvents = await collect( + input.transport.subscribeToRun({ + threadId: input.threadId, + runId: resumed.runId, + }), + input.timeoutMs, + "approval replacement run", + ); + assertOrderedEvents({ + events: continuationEvents, + threadId: input.threadId, + runId: resumed.runId, + terminal: "run.completed", + }); + } + const lifecycleEvents = + resumed.runId === started.runId + ? events + : [...events, ...continuationEvents]; + const resolvedIndex = lifecycleEvents.findIndex( (event) => event.type === "approval.resolved" && event.approvalId === requested?.request.id, ); assert( - resolvedIndex > events.indexOf(requested), + resolvedIndex > lifecycleEvents.indexOf(requested), "Approval continuation must emit approval.resolved after approval.requested.", ); - const resolved = events[resolvedIndex]; + const resolved = lifecycleEvents[resolvedIndex]; assert( resolved?.type === "approval.resolved" && resolved.response.decision === "approve", "Approval continuation must preserve the explicit provider-neutral decision.", ); - return events.length; + return lifecycleEvents.length; } async function assertQueueIdentity(input: { diff --git a/packages/agentkit/src/protocol/agui-codec.spec.ts b/packages/agentkit/src/protocol/agui-codec.spec.ts new file mode 100644 index 00000000000..55523983212 --- /dev/null +++ b/packages/agentkit/src/protocol/agui-codec.spec.ts @@ -0,0 +1,468 @@ +import { EventSchemas, EventType } from "@ag-ui/core"; +import { describe, expect, it } from "vitest"; + +import { decodeAgUiEvent, encodeAgentEvent } from "./agui-codec.js"; +import { + AGENTKIT_METADATA_KEY, + AGENTKIT_NATIVE_EVENT_TYPES, + AGENTKIT_PROFILE_EVENT_TYPES, + approvalResponseFromResume, + isProfileEventType, + resumeEntryFromApproval, + resumeOptionId, +} from "./agui.js"; +import type { AgentEvent } from "./index.js"; +import { AgentProtocolValidationError } from "./validation.js"; + +const THREAD_ID = "thread_1"; +const RUN_ID = "run_1"; +const CONTEXT = { threadId: THREAD_ID, runId: RUN_ID }; + +let nextSequence = 0; + +function base(type: string) { + nextSequence += 1; + return { + id: `evt_${nextSequence}`, + type, + threadId: THREAD_ID, + runId: RUN_ID, + sequence: nextSequence, + occurredAt: "2026-09-09T12:00:00.000Z", + }; +} + +const activity = { + id: "activity_1", + kind: "search" as const, + label: "Searching", + status: "running" as const, +}; + +const toolCall = { + id: "tool_1", + name: "search", + status: "running" as const, + messageId: "msg_1", +}; + +const participant = { + id: "agent_1", + name: "Researcher", + status: "working" as const, +}; + +const task = { id: "task_1", title: "Draft", status: "running" as const }; + +const message = { + id: "msg_1", + role: "assistant" as const, + parts: [{ type: "text" as const, text: "hello" }], +}; + +const taskGroup = { id: "group_1", taskIds: ["task_1"] }; + +/** One event per domain type, so a missed mapping fails rather than hides. */ +function everyEventType(): AgentEvent[] { + nextSequence = 0; + return [ + { ...base("run.started"), agentId: "agent_1" }, + { ...base("run.status"), status: "running" }, + { ...base("agent.registered"), agent: participant }, + { ...base("agent.updated"), agent: participant }, + { ...base("agent.unregistered"), agent: participant }, + { + ...base("agent.interaction"), + interaction: { id: "int_1", kind: "handoff", agentId: "agent_1" }, + }, + { ...base("message.created"), message }, + { ...base("message.delta"), messageId: "msg_1", text: "hel" }, + { ...base("message.completed"), message }, + { ...base("reasoning.delta"), messageId: "msg_1", text: "thinking" }, + { ...base("tool.started"), toolCall }, + { ...base("tool.delta"), toolCallId: "tool_1", inputTextDelta: "{}" }, + { ...base("tool.updated"), toolCall: { ...toolCall, status: "completed" } }, + { ...base("activity.started"), activity }, + { ...base("activity.updated"), activity }, + { ...base("activity.completed"), activity }, + { ...base("task.created"), task }, + { ...base("task.updated"), task }, + { ...base("task.completed"), task }, + { ...base("task-group.created"), taskGroup }, + { ...base("task-group.updated"), taskGroup }, + { ...base("task-group.completed"), taskGroup }, + { ...base("task-group.removed"), taskGroupId: "group_1" }, + { + ...base("approval.requested"), + request: { id: "approval_1", title: "Delete the row?" }, + }, + { + ...base("approval.resolved"), + approvalId: "approval_1", + response: { decision: "approve" }, + }, + { + ...base("connection.requested"), + request: { + id: "conn_1", + provider: "github", + reason: "connect", + status: "requested", + }, + }, + { + ...base("connection.updated"), + request: { + id: "conn_1", + provider: "github", + reason: "connect", + status: "connected", + }, + }, + { + ...base("artifact.created"), + artifact: { id: "artifact_1", kind: "file" }, + }, + { + ...base("widget.created"), + widget: { id: "widget_1", kind: "table", data: { rows: [] } }, + }, + { + ...base("widget.updated"), + widget: { id: "widget_1", kind: "table", data: { rows: [] } }, + }, + { ...base("widget.removed"), widgetId: "widget_1" }, + { + ...base("annotation.created"), + annotation: { id: "note_1", kind: "source", label: "Docs" }, + }, + { + ...base("annotation.updated"), + annotation: { id: "note_1", kind: "source", label: "Docs" }, + }, + { ...base("annotation.removed"), annotationId: "note_1" }, + { + ...base("suggestions.updated"), + suggestions: [{ id: "sug_1", label: "Keep going" }], + }, + { + ...base("action.started"), + invocation: { id: "inv_1", action: "refresh", threadId: THREAD_ID }, + }, + { + ...base("action.completed"), + result: { invocationId: "inv_1", status: "completed" }, + }, + { + ...base("action.failed"), + result: { + invocationId: "inv_1", + status: "failed", + error: { code: "boom", message: "failed" }, + }, + }, + { + ...base("upload.progress"), + progress: { uploadId: "upload_1", loaded: 1, total: 2 }, + }, + { ...base("client.effect"), name: "scroll" }, + { ...base("client.deeplink"), name: "open", data: { path: "/" } }, + { + ...base("thread.updated"), + thread: { + id: THREAD_ID, + createdAt: "2026-09-09T11:00:00.000Z", + updatedAt: "2026-09-09T12:00:00.000Z", + }, + }, + { ...base("queue.updated"), messages: [] }, + { ...base("run.completed"), usage: { totalTokens: 10 } }, + { ...base("run.failed"), error: { code: "internal", message: "exploded" } }, + { ...base("run.cancelled") }, + { ...base("x-builder.custom"), payload: { anything: true } }, + ] as AgentEvent[]; +} + +describe("AG-UI codec", () => { + it("has a fixture for every mapped event type", () => { + const encountered = new Set(everyEventType().map((event) => event.type)); + const mapped = [ + ...Object.keys(AGENTKIT_NATIVE_EVENT_TYPES), + ...AGENTKIT_PROFILE_EVENT_TYPES, + ]; + for (const type of mapped) { + expect(encountered.has(type), `missing fixture for ${type}`).toBe(true); + } + }); + + it("round-trips every domain event without loss", () => { + for (const event of everyEventType()) { + const decoded = decodeAgUiEvent(encodeAgentEvent(event), CONTEXT); + expect(decoded, `${event.type} decoded to nothing`).toEqual(event); + } + }); + + it("preserves formatting on message deltas", () => { + const source = { + ...base("message.delta"), + messageId: "msg_1", + text: "**hello**", + format: "markdown", + } as AgentEvent; + + expect(decodeAgUiEvent(encodeAgentEvent(source), CONTEXT)).toEqual(source); + }); + + it("projects tool-role messages onto valid AG-UI without losing the role", () => { + const source = { + ...base("message.created"), + message: { ...message, role: "tool" }, + } as AgentEvent; + const wire = encodeAgentEvent(source) as { role: string }; + + expect(wire.role).toBe("assistant"); + expect(EventSchemas.safeParse(wire).success).toBe(true); + expect(decodeAgUiEvent(wire, CONTEXT)).toEqual(source); + }); + + it("restores static typing for profile events carried over CUSTOM", () => { + const source = everyEventType().find( + (event) => event.type === "task-group.updated", + )!; + const wire = encodeAgentEvent(source) as { type: EventType; name: string }; + expect(wire.type).toBe(EventType.CUSTOM); + expect(wire.name).toBe("builder.agentkit.v1/task-group.updated"); + + const decoded = decodeAgUiEvent(wire, CONTEXT); + if (decoded?.type !== "task-group.updated") { + throw new Error("expected a task-group.updated event"); + } + expect(decoded.taskGroup.taskIds).toEqual(["task_1"]); + }); + + it("emits frames a stock AG-UI client can parse", () => { + for (const event of everyEventType()) { + const parsed = EventSchemas.safeParse(encodeAgentEvent(event)); + expect(parsed.success, `${event.type} is not valid AG-UI`).toBe(true); + } + }); + + it("keeps the Builder sequence and domain type on every frame", () => { + for (const event of everyEventType()) { + const wire = encodeAgentEvent(event) as { + metadata: Record; + }; + expect(wire.metadata[AGENTKIT_METADATA_KEY]?.seq).toBe(event.sequence); + expect(wire.metadata[AGENTKIT_METADATA_KEY]?.t).toBe(event.type); + } + }); + + it("skips frames from another producer instead of inventing a gap", () => { + const foreign = { + type: EventType.CUSTOM, + name: "someone-else/thing", + value: {}, + }; + expect(decodeAgUiEvent(foreign, CONTEXT)).toBeUndefined(); + expect( + decodeAgUiEvent( + { type: EventType.STATE_SNAPSHOT, snapshot: {} }, + CONTEXT, + ), + ).toBeUndefined(); + }); + + it("refuses a Builder frame it cannot read rather than dropping it", () => { + const wire = encodeAgentEvent(everyEventType()[0]!) as { + metadata: Record>; + }; + wire.metadata[AGENTKIT_METADATA_KEY]!.seq = 0; + expect(() => decodeAgUiEvent(wire, CONTEXT)).toThrow( + AgentProtocolValidationError, + ); + }); + + it("rejects a profiled custom frame with a foreign name", () => { + const wire = encodeAgentEvent( + everyEventType().find((event) => event.type === "widget.removed")!, + ) as { name: string }; + wire.name = "someone-else/thing"; + + expect(() => decodeAgUiEvent(wire, CONTEXT)).toThrow( + AgentProtocolValidationError, + ); + }); + + it("rejects a name that disagrees with the declared event type", () => { + const wire = encodeAgentEvent( + everyEventType().find((event) => event.type === "widget.removed")!, + ) as { name: string }; + wire.name = "builder.agentkit.v1/widget.created"; + expect(() => decodeAgUiEvent(wire, CONTEXT)).toThrow( + AgentProtocolValidationError, + ); + }); + + it.each([null, "cancelled", 1, []])( + "rejects a non-object profile payload: %j", + (value) => { + const wire = encodeAgentEvent( + everyEventType().find((event) => event.type === "run.cancelled")!, + ) as { value: unknown }; + wire.value = value; + + expect(() => decodeAgUiEvent(wire, CONTEXT)).toThrow( + AgentProtocolValidationError, + ); + }, + ); + + it("rejects a native frame that disagrees with its profile event type", () => { + const wire = encodeAgentEvent( + everyEventType().find((event) => event.type === "run.started")!, + ) as { type: EventType }; + wire.type = EventType.RUN_ERROR; + + expect(() => decodeAgUiEvent(wire, CONTEXT)).toThrow( + AgentProtocolValidationError, + ); + }); + + it("runs reconstructed residuals through domain validation", () => { + const wire = encodeAgentEvent({ + ...base("message.delta"), + messageId: "msg_1", + text: "hello", + format: "markdown", + } as AgentEvent) as { + metadata: Record }>; + }; + wire.metadata[AGENTKIT_METADATA_KEY]!.residual.format = "html"; + + expect(() => decodeAgUiEvent(wire, CONTEXT)).toThrow( + AgentProtocolValidationError, + ); + }); + + it("maps approvals onto an AG-UI interrupt outcome", () => { + const approval = everyEventType().find( + (event) => event.type === "approval.requested", + )!; + const wire = encodeAgentEvent(approval) as { + type: EventType; + outcome: { type: string; interrupts: { id: string; reason: string }[] }; + }; + expect(wire.type).toBe(EventType.RUN_FINISHED); + expect(wire.outcome.type).toBe("interrupt"); + expect(wire.outcome.interrupts[0]).toMatchObject({ + id: "approval_1", + reason: "Delete the row?", + }); + }); + + it("keeps explicit denial distinct from interrupt cancellation", () => { + const denied = resumeEntryFromApproval({ + approvalId: "approval_1", + response: { decision: "deny" }, + }); + + expect(denied.status).toBe("resolved"); + expect(approvalResponseFromResume(denied)).toEqual({ decision: "deny" }); + expect( + approvalResponseFromResume({ + interruptId: "approval_1", + status: "cancelled", + }), + ).toEqual({ decision: "deny" }); + expect(() => + approvalResponseFromResume({ + interruptId: "approval_1", + status: "cancelled", + payload: { decision: "approve" }, + }), + ).toThrow(AgentProtocolValidationError); + }); + + it("rejects payloads that redefine trusted event identity", () => { + const native = encodeAgentEvent( + everyEventType().find((event) => event.type === "message.created")!, + ) as { + metadata: Record }>; + }; + native.metadata[AGENTKIT_METADATA_KEY]!.residual.sequence = 999; + + expect(() => decodeAgUiEvent(native, CONTEXT)).toThrow( + AgentProtocolValidationError, + ); + + const custom = encodeAgentEvent( + everyEventType().find((event) => event.type === "widget.removed")!, + ) as { value: Record }; + custom.value.runId = "run-forged"; + + expect(() => decodeAgUiEvent(custom, CONTEXT)).toThrow( + AgentProtocolValidationError, + ); + }); + + it("validates approval payload fields before they reach a runtime", () => { + expect(() => + approvalResponseFromResume({ + interruptId: "approval_1", + status: "resolved", + payload: { decision: "approve", optionIds: [1] }, + }), + ).toThrow(AgentProtocolValidationError); + expect(() => + resumeOptionId({ + interruptId: "approval_1", + status: "resolved", + payload: { decision: "approve", optionId: 1 }, + }), + ).toThrow(AgentProtocolValidationError); + }); + + it("classifies every fixture as native or profile", () => { + for (const event of everyEventType()) { + const known = + event.type in AGENTKIT_NATIVE_EVENT_TYPES || + isProfileEventType(event.type); + expect(known, `${event.type} has no mapping`).toBe(true); + } + }); +}); + +describe("AG-UI reconnect cursor", () => { + function stream(): AgentEvent[] { + return everyEventType().slice(0, 10); + } + + function cursorOf(wire: unknown): number { + return (wire as { metadata: Record }).metadata[ + AGENTKIT_METADATA_KEY + ]!.seq; + } + + it("resumes after the last delivered cursor with no gap or duplicate", () => { + const wire = stream().map(encodeAgentEvent); + const delivered = wire.slice(0, 4); + const lastSequence = cursorOf(delivered.at(-1)); + + const resumed = wire.filter((frame) => cursorOf(frame) > lastSequence); + const replayed = [...delivered, ...resumed].map( + (frame) => decodeAgUiEvent(frame, CONTEXT)!, + ); + + expect(replayed).toEqual(stream()); + }); + + it("makes a dropped frame detectable from the cursor alone", () => { + const wire = stream().map(encodeAgentEvent); + const sequences = [...wire.slice(0, 3), ...wire.slice(4)].map(cursorOf); + + const gaps = sequences.filter( + (value, index) => index > 0 && value !== sequences[index - 1]! + 1, + ); + expect(gaps).toHaveLength(1); + }); +}); diff --git a/packages/agentkit/src/protocol/agui-codec.ts b/packages/agentkit/src/protocol/agui-codec.ts new file mode 100644 index 00000000000..6da54d403a1 --- /dev/null +++ b/packages/agentkit/src/protocol/agui-codec.ts @@ -0,0 +1,338 @@ +import { EventSchemas, EventType, type AGUIEvent } from "@ag-ui/core"; + +import { + AGENTKIT_METADATA_KEY, + AGENTKIT_NATIVE_EVENT_TYPES, + isLosslessEventType, + isNativeEventType, + isProfileEventType, + profileEventName, + profileEventTypeFromName, + readProfileExtension, + writeProfileExtension, + type AgentKitStreamContext, +} from "./agui.js"; +import type { AgentEvent, AgentProtocolMetadata } from "./index.js"; +import { AgentProtocolValidationError, parseAgentEvent } from "./validation.js"; + +/** Base fields every domain event carries; everything else is payload. */ +const BASE_EVENT_KEYS = [ + "id", + "type", + "threadId", + "runId", + "sequence", + "occurredAt", + "metadata", +] as const; + +type UnknownRecord = Record; + +function payloadOf(event: AgentEvent): UnknownRecord { + const payload: UnknownRecord = {}; + for (const [key, value] of Object.entries( + event as unknown as UnknownRecord, + )) { + if ((BASE_EVENT_KEYS as readonly string[]).includes(key)) continue; + payload[key] = value; + } + return payload; +} + +function asRecord(value: unknown): UnknownRecord { + return ( + typeof value === "object" && value !== null ? value : {} + ) as UnknownRecord; +} + +function requireRecord(value: unknown, path: string): UnknownRecord { + if (typeof value !== "object" || value === null || Array.isArray(value)) { + throw new AgentProtocolValidationError(path, "must be an object"); + } + return value as UnknownRecord; +} + +function assertPayloadDoesNotRedefineBase( + payload: UnknownRecord, + path: string, +): void { + for (const key of BASE_EVENT_KEYS) { + if (Object.prototype.hasOwnProperty.call(payload, key)) { + throw new AgentProtocolValidationError( + `${path}.${key}`, + "cannot redefine a protocol envelope field", + ); + } + } +} + +function optionalString(value: unknown): string | undefined { + return typeof value === "string" && value.length > 0 ? value : undefined; +} + +/** + * Projects a domain event onto the fields of its native AG-UI event. This is + * what an AG-UI-only consumer reads; an AgentKit consumer rebuilds from the + * profile residual instead, so a lossy projection here is a rendering + * concession for foreign clients, never a loss of protocol state. + */ +function nativeProjection(event: AgentEvent): UnknownRecord { + switch (event.type) { + case "run.started": + return { threadId: event.threadId, runId: event.runId }; + case "run.completed": + return { + threadId: event.threadId, + runId: event.runId, + outcome: { type: "success" }, + }; + case "run.failed": + return { message: event.error.message, code: event.error.code }; + case "run.status": + return { stepName: event.status }; + case "message.created": + return { + messageId: event.message.id, + role: event.message.role === "tool" ? "assistant" : event.message.role, + }; + case "message.delta": + return { messageId: event.messageId, delta: event.text }; + case "message.completed": + return { messageId: event.message.id }; + case "reasoning.delta": + return { messageId: event.messageId, delta: event.text }; + case "tool.started": + return { + toolCallId: event.toolCall.id, + toolCallName: event.toolCall.name, + parentMessageId: event.toolCall.messageId, + }; + case "tool.delta": + return { + toolCallId: event.toolCallId, + delta: event.inputTextDelta ?? "", + }; + case "tool.updated": + return { + messageId: event.toolCall.messageId ?? event.toolCall.id, + toolCallId: event.toolCall.id, + content: + typeof event.toolCall.output === "string" + ? event.toolCall.output + : JSON.stringify(event.toolCall.output ?? null), + }; + case "activity.started": + case "activity.updated": + case "activity.completed": + return { + messageId: event.activity.id, + activityType: event.activity.kind, + content: { ...event.activity } as UnknownRecord, + replace: event.type !== "activity.started", + }; + case "agent.registered": + return { + subagentRunId: event.agent.id, + name: event.agent.name, + description: event.agent.description, + parentSubagentRunId: event.agent.parentAgentId, + }; + case "agent.unregistered": + return { subagentRunId: event.agent.id }; + case "approval.requested": + return { + threadId: event.threadId, + runId: event.runId, + outcome: { + type: "interrupt", + interrupts: [ + { + id: event.request.id, + reason: event.request.title, + message: event.request.description, + expiresAt: event.request.expiresAt, + }, + ], + }, + }; + default: + return {}; + } +} + +/** + * Encodes one domain event as exactly one AG-UI event. + * + * The 1:1 rule is deliberate. The sequence lives on the AG-UI frame, so fanning + * a domain event out into several frames would either duplicate a cursor value + * or consume slots the producer's own numbering knows nothing about, and the + * reducer's gap detection is only as trustworthy as that numbering. + */ +export function encodeAgentEvent(event: AgentEvent): AGUIEvent { + const type = event.type; + const extension = writeProfileExtension({ + t: type, + seq: event.sequence, + id: event.id, + at: event.occurredAt, + meta: event.metadata, + residual: + isLosslessEventType(type) && + !(event.type === "message.delta" && event.format !== undefined) + ? undefined + : payloadOf(event), + }); + + const frame: UnknownRecord = { + timestamp: Date.parse(event.occurredAt), + metadata: extension, + }; + + if (isNativeEventType(type)) { + return { + ...frame, + ...nativeProjection(event), + type: AGENTKIT_NATIVE_EVENT_TYPES[type], + } as AGUIEvent; + } + + if (!isProfileEventType(type)) { + throw new AgentProtocolValidationError( + "event.type", + `no AG-UI mapping is defined for ${JSON.stringify(type)}`, + ); + } + + return { + ...frame, + type: EventType.CUSTOM, + name: profileEventName(type), + value: payloadOf(event), + } as AGUIEvent; +} + +/** + * Rebuilds a domain event from an AG-UI frame. + * + * Returns `undefined` only for frames that carry no Builder extension — another + * producer's `CUSTOM`, or an AG-UI event type Builder never emits. Those are + * outside the profile and outside its numbering, so skipping them cannot open a + * gap. A frame that *does* carry the extension must decode or throw; treating + * an unreadable Builder frame as absent is what would hide a real gap. + */ +export function decodeAgUiEvent( + value: unknown, + context: AgentKitStreamContext, + path = "event", +): AgentEvent | undefined { + const parsed = EventSchemas.safeParse(value); + if (!parsed.success) { + throw new AgentProtocolValidationError( + path, + `is not a valid AG-UI event: ${parsed.error.issues[0]?.message ?? "unknown error"}`, + ); + } + + const frame = parsed.data as UnknownRecord; + const metadata = asRecord(frame.metadata); + if (!(AGENTKIT_METADATA_KEY in metadata)) return undefined; + + const extension = readProfileExtension(metadata, `${path}.metadata`); + + const base = { + id: extension.id, + threadId: context.threadId, + runId: context.runId, + sequence: extension.seq, + occurredAt: extension.at, + ...(extension.meta === undefined + ? {} + : { metadata: extension.meta as AgentProtocolMetadata }), + }; + + if (frame.type === EventType.CUSTOM) { + const name = typeof frame.name === "string" ? frame.name : ""; + const domainType = profileEventTypeFromName(name); + if (domainType === undefined) { + throw new AgentProtocolValidationError( + `${path}.name`, + `is not an AgentKit profile event name`, + ); + } + if (domainType !== extension.t) { + throw new AgentProtocolValidationError( + `${path}.name`, + `disagrees with the profile event type ${JSON.stringify(extension.t)}`, + ); + } + if (isNativeEventType(domainType)) { + throw new AgentProtocolValidationError( + `${path}.type`, + `must use the native AG-UI mapping for ${JSON.stringify(domainType)}`, + ); + } + const payload = requireRecord(frame.value, `${path}.value`); + assertPayloadDoesNotRedefineBase(payload, `${path}.value`); + return parseAgentEvent({ ...payload, ...base, type: domainType }, path); + } + + if (!isNativeEventType(extension.t)) { + throw new AgentProtocolValidationError( + `${path}.type`, + `has no native mapping for profile event type ${JSON.stringify(extension.t)}`, + ); + } + if (frame.type !== AGENTKIT_NATIVE_EVENT_TYPES[extension.t]) { + throw new AgentProtocolValidationError( + `${path}.type`, + `disagrees with the native mapping for ${JSON.stringify(extension.t)}`, + ); + } + + if (extension.residual !== undefined) { + assertPayloadDoesNotRedefineBase( + extension.residual, + `${path}.metadata.${AGENTKIT_METADATA_KEY}.residual`, + ); + return parseAgentEvent( + { ...extension.residual, ...base, type: extension.t }, + path, + ); + } + + return parseAgentEvent( + decodeLosslessNative(frame, base, extension.t, path), + path, + ); +} + +function decodeLosslessNative( + frame: UnknownRecord, + base: UnknownRecord, + type: string, + path: string, +): AgentEvent { + switch (type) { + case "reasoning.delta": + case "message.delta": + return { + ...base, + type, + messageId: optionalString(frame.messageId) ?? "", + text: typeof frame.delta === "string" ? frame.delta : "", + } as AgentEvent; + default: + throw new AgentProtocolValidationError( + `${path}.metadata`, + `event type ${JSON.stringify(type)} requires a residual payload`, + ); + } +} + +/** + * Validates that a frame is legal AG-UI regardless of the profile, so a + * conformance suite can prove Builder streams stay readable by stock clients. + */ +export function isAgUiEvent(value: unknown): value is AGUIEvent { + return EventSchemas.safeParse(value).success; +} diff --git a/packages/agentkit/src/protocol/agui.ts b/packages/agentkit/src/protocol/agui.ts new file mode 100644 index 00000000000..a22beb19156 --- /dev/null +++ b/packages/agentkit/src/protocol/agui.ts @@ -0,0 +1,388 @@ +import { EventType } from "@ag-ui/core"; + +import type { + AgentApprovalResponse, + AgentEvent, + AgentProtocolMetadata, + AgentResumeEntry, + ApprovalId, + EventId, + RunId, + ThreadId, +} from "./index.js"; +import { + AgentProtocolValidationError, + parseAgentApprovalResponse, +} from "./validation.js"; + +/** + * The profile identifier carried on every frame. AG-UI has no protocol version + * of its own, so a reader that cannot recognise this string must not assume the + * Builder extensions below are absent — it must refuse the stream. + */ +export const AGENTKIT_PROFILE_NAME = "builder.agentkit"; +export const AGENTKIT_PROFILE_VERSION = 1; +export const AGENTKIT_PROFILE_ID = `${AGENTKIT_PROFILE_NAME}.v${AGENTKIT_PROFILE_VERSION}`; + +/** + * Namespace for the Builder events that AG-UI has no counterpart for. They ride + * on `CUSTOM` as `/`. + */ +export const AGENTKIT_PROFILE_EVENT_PREFIX = `${AGENTKIT_PROFILE_ID}/`; + +/** + * Single metadata key holding every Builder extension. `@ag-ui/core` reserves + * `AGUI_METADATA_KEY` ("ag-ui") for itself, so nesting under one owned key + * keeps the two from colliding as either side adds fields. + */ +export const AGENTKIT_METADATA_KEY = AGENTKIT_PROFILE_NAME; + +/** HTTP header advertising the profile so negotiation fails before streaming. */ +export const AGENTKIT_PROFILE_HEADER = "x-agentkit-profile"; + +/** + * The Builder extensions to a base protocol that defines none of them. + * + * `seq` is the load-bearing one: `@ag-ui/core@0.0.59` documents a `resumable` + * transport capability but never defines a sequence, event id, or replay cursor + * anywhere on the wire, so two AG-UI implementations that both advertise it + * would have invented incompatible schemes. Everything else here exists because + * AG-UI's native event for the same value cannot carry it losslessly. + */ +export interface AgentKitProfileExtension { + /** Profile version, so a frame remains self-describing outside its stream. */ + v: number; + /** + * Domain event type. Carried even for native mappings because three domain + * types share `ACTIVITY_SNAPSHOT` and two share `RUN_FINISHED`, so the AG-UI + * `type` alone cannot identify the frame on the way back. + */ + t: string; + /** Per-`(threadId, runId)` monotonic cursor. Also emitted as the SSE `id:`. */ + seq: number; + id: EventId; + /** + * The domain timestamp verbatim. AG-UI's `timestamp` is epoch milliseconds, + * which cannot round-trip an offset, so the original string travels too. + */ + at: string; + meta?: AgentProtocolMetadata; + /** + * Fields the native AG-UI event has no slot for. Present only on the six + * events the mapping table marks partial; a direct mapping omits it entirely. + */ + residual?: Record; +} + +/** + * Domain event types that reach a native AG-UI event, and which one. + * + * Reaching a native event is not the same as fitting in it. Only the types in + * {@link AGENTKIT_LOSSLESS_EVENT_TYPES} survive the trip on native fields + * alone; every other entry here also ships its domain payload in + * `residual`, because AG-UI's event carries a strict subset of the fields. + * `activity.updated` is the sharpest case: AG-UI's `ACTIVITY_DELTA` is a JSON + * Patch against prior state, which a stateless encoder cannot produce, so all + * three activity events map to `ACTIVITY_SNAPSHOT` instead. + */ +export const AGENTKIT_NATIVE_EVENT_TYPES = { + "run.started": EventType.RUN_STARTED, + "run.completed": EventType.RUN_FINISHED, + "run.failed": EventType.RUN_ERROR, + "run.status": EventType.STEP_STARTED, + "message.created": EventType.TEXT_MESSAGE_START, + "message.delta": EventType.TEXT_MESSAGE_CONTENT, + "message.completed": EventType.TEXT_MESSAGE_END, + "reasoning.delta": EventType.REASONING_MESSAGE_CONTENT, + "tool.started": EventType.TOOL_CALL_START, + "tool.delta": EventType.TOOL_CALL_ARGS, + "tool.updated": EventType.TOOL_CALL_RESULT, + "activity.started": EventType.ACTIVITY_SNAPSHOT, + "activity.updated": EventType.ACTIVITY_SNAPSHOT, + "activity.completed": EventType.ACTIVITY_SNAPSHOT, + "agent.registered": EventType.SUBAGENT_STARTED, + "agent.unregistered": EventType.SUBAGENT_FINISHED, + "approval.requested": EventType.RUN_FINISHED, +} as const satisfies Partial>; + +/** + * The mappings that need no residual. Two, out of the twelve the evaluation + * table calls "Direct" — the rest carry structure AG-UI's event has no slot + * for, so "direct" describes the event name, not the payload. + */ +export const AGENTKIT_LOSSLESS_EVENT_TYPES = [ + "reasoning.delta", + "message.delta", +] as const satisfies readonly AgentEvent["type"][]; + +export type AgentKitNativeEventType = keyof typeof AGENTKIT_NATIVE_EVENT_TYPES; + +const LOSSLESS_EVENT_TYPE_SET: ReadonlySet = new Set( + AGENTKIT_LOSSLESS_EVENT_TYPES, +); + +export function isNativeEventType( + type: string, +): type is AgentKitNativeEventType { + return type in AGENTKIT_NATIVE_EVENT_TYPES; +} + +/** + * `message.delta` only fits when it does not set `format`; AG-UI's + * `TEXT_MESSAGE_CONTENT` has no formatting field. + */ +export function isLosslessEventType(type: string): boolean { + return LOSSLESS_EVENT_TYPE_SET.has(type); +} + +/** + * Everything else. These have no AG-UI counterpart at all and are the reason + * the profile exists: adopting AG-UI wholesale would delete the capabilities + * behind them, not express them differently. + * + * `approval.resolved` sits here rather than with the native mappings because + * AG-UI models resolution only as a `resume[]` command on the next run input. + * The command exists, but there is no event announcing it to other observers of + * the thread, so the profile keeps one. + */ +export const AGENTKIT_PROFILE_EVENT_TYPES = [ + "run.cancelled", + "approval.resolved", + "agent.updated", + "agent.interaction", + "connection.requested", + "connection.updated", + "task.created", + "task.updated", + "task.completed", + "task-group.created", + "task-group.updated", + "task-group.completed", + "task-group.removed", + "artifact.created", + "widget.created", + "widget.updated", + "widget.removed", + "annotation.created", + "annotation.updated", + "annotation.removed", + "suggestions.updated", + "action.started", + "action.completed", + "action.failed", + "upload.progress", + "client.effect", + "client.deeplink", + "thread.updated", + "queue.updated", +] as const satisfies readonly AgentEvent["type"][]; + +export type AgentKitProfileEventType = + (typeof AGENTKIT_PROFILE_EVENT_TYPES)[number]; + +const PROFILE_EVENT_TYPE_SET: ReadonlySet = new Set( + AGENTKIT_PROFILE_EVENT_TYPES, +); + +/** Host-defined `x-*` events also travel as profile events. */ +export function isProfileEventType(type: string): boolean { + return PROFILE_EVENT_TYPE_SET.has(type) || type.startsWith("x-"); +} + +export function profileEventName(type: string): string { + return `${AGENTKIT_PROFILE_EVENT_PREFIX}${type}`; +} + +/** + * Returns the domain event type for a `CUSTOM` name, or `undefined` when the + * name belongs to another profile. A foreign `CUSTOM` is another producer's + * business, not a malformed frame, so it is skipped rather than rejected. + */ +export function profileEventTypeFromName(name: string): string | undefined { + if (!name.startsWith(AGENTKIT_PROFILE_EVENT_PREFIX)) return undefined; + return name.slice(AGENTKIT_PROFILE_EVENT_PREFIX.length); +} + +export interface AgentKitStreamContext { + threadId: ThreadId; + runId: RunId; +} + +function requireRecord(value: unknown, path: string): Record { + if (typeof value !== "object" || value === null || Array.isArray(value)) { + throw new AgentProtocolValidationError(path, "expected an object"); + } + return value as Record; +} + +/** + * Reads the Builder extension off an AG-UI frame. A frame that reached an + * AgentKit consumer without it is unreadable, not empty: silently defaulting + * `seq` would let the reducer accept a stream it can no longer detect gaps in. + */ +export function readProfileExtension( + metadata: unknown, + path = "metadata", +): AgentKitProfileExtension { + const container = requireRecord(metadata, path); + const extension = requireRecord( + container[AGENTKIT_METADATA_KEY], + `${path}.${AGENTKIT_METADATA_KEY}`, + ); + const extensionPath = `${path}.${AGENTKIT_METADATA_KEY}`; + + const version = extension.v; + if (version !== AGENTKIT_PROFILE_VERSION) { + throw new AgentProtocolValidationError( + `${extensionPath}.v`, + `unsupported profile version ${JSON.stringify(version)}; expected ${AGENTKIT_PROFILE_VERSION}`, + ); + } + + const domainType = extension.t; + if (typeof domainType !== "string" || domainType.length === 0) { + throw new AgentProtocolValidationError( + `${extensionPath}.t`, + "expected a domain event type", + ); + } + + const seq = extension.seq; + if (typeof seq !== "number" || !Number.isSafeInteger(seq) || seq < 1) { + throw new AgentProtocolValidationError( + `${extensionPath}.seq`, + "expected a positive safe integer sequence", + ); + } + + const id = extension.id; + if (typeof id !== "string" || id.length === 0) { + throw new AgentProtocolValidationError( + `${extensionPath}.id`, + "expected a non-empty event id", + ); + } + + const at = extension.at; + if (typeof at !== "string" || Number.isNaN(Date.parse(at))) { + throw new AgentProtocolValidationError( + `${extensionPath}.at`, + "expected an ISO-8601 timestamp", + ); + } + + const meta = + extension.meta === undefined + ? undefined + : (requireRecord( + extension.meta, + `${extensionPath}.meta`, + ) as AgentProtocolMetadata); + + const residual = + extension.residual === undefined + ? undefined + : requireRecord(extension.residual, `${extensionPath}.residual`); + + return { + v: AGENTKIT_PROFILE_VERSION, + t: domainType, + seq, + id, + at, + meta, + residual, + }; +} + +export function writeProfileExtension( + extension: Omit, +): Record { + const value: AgentKitProfileExtension = { + v: AGENTKIT_PROFILE_VERSION, + t: extension.t, + seq: extension.seq, + id: extension.id, + at: extension.at, + }; + if (extension.meta !== undefined) value.meta = extension.meta; + if (extension.residual !== undefined) value.residual = extension.residual; + return { [AGENTKIT_METADATA_KEY]: value }; +} + +/** + * Approval decisions travel as AG-UI resume entries. Both directions live here + * so the client that writes an entry and the runtime that reads one cannot + * drift into two different payload shapes. + */ +export function resumeEntryFromApproval(input: { + approvalId: ApprovalId; + optionId?: string; + response: AgentApprovalResponse; +}): AgentResumeEntry { + return { + interruptId: input.approvalId, + status: "resolved", + payload: { + ...input.response, + ...(input.optionId === undefined ? {} : { optionId: input.optionId }), + }, + }; +} + +function resumePayload(entry: AgentResumeEntry): Record { + return requireRecord(entry.payload, "resume.payload"); +} + +function approvalResponseFromPayload( + payload: Record, +): AgentApprovalResponse { + if (payload.decision !== "approve" && payload.decision !== "deny") { + throw new AgentProtocolValidationError( + "resume.payload.decision", + 'An approval response must include decision "approve" or "deny".', + ); + } + return parseAgentApprovalResponse( + { + decision: payload.decision, + ...(payload.optionIds === undefined + ? {} + : { optionIds: payload.optionIds }), + ...(payload.other === undefined ? {} : { other: payload.other }), + ...(payload.input === undefined ? {} : { input: payload.input }), + }, + "resume.payload", + ); +} + +export function approvalResponseFromResume( + entry: AgentResumeEntry, +): AgentApprovalResponse { + if (entry.status === "cancelled") { + if (entry.payload === undefined) return { decision: "deny" }; + const payload = resumePayload(entry); + if (payload.decision !== undefined && payload.decision !== "deny") { + throw new AgentProtocolValidationError( + "resume.payload.decision", + 'A cancelled interrupt cannot contain an "approve" decision.', + ); + } + return approvalResponseFromPayload({ ...payload, decision: "deny" }); + } + const payload = resumePayload(entry); + return approvalResponseFromPayload(payload); +} + +export function resumeOptionId(entry: AgentResumeEntry): string | undefined { + if (entry.payload === undefined) return undefined; + const optionId = resumePayload(entry).optionId; + if (optionId === undefined) return undefined; + if (typeof optionId !== "string" || optionId.length === 0) { + throw new AgentProtocolValidationError( + "resume.payload.optionId", + "expected a non-empty string", + ); + } + return optionId; +} diff --git a/packages/agentkit/src/protocol/hardening.spec.ts b/packages/agentkit/src/protocol/hardening.spec.ts index 2deeed93217..023a4764e72 100644 --- a/packages/agentkit/src/protocol/hardening.spec.ts +++ b/packages/agentkit/src/protocol/hardening.spec.ts @@ -212,6 +212,16 @@ describe("AgentKit protocol hardening", () => { versions: [1, 1], }), ).toThrow("must be unique"); + expect( + negotiateAgentKitProtocolVersion({ + protocol: "agentkit", + versions: [1], + }), + ).toMatchObject({ + status: "incompatible", + localVersions: [AGENTKIT_PROTOCOL_VERSION], + peerVersions: [1], + }); }); it("validates abortable subscription inputs without treating abort as run cancellation", () => { diff --git a/packages/agentkit/src/protocol/index.ts b/packages/agentkit/src/protocol/index.ts index 3b0b3f3a443..66237152490 100644 --- a/packages/agentkit/src/protocol/index.ts +++ b/packages/agentkit/src/protocol/index.ts @@ -1198,6 +1198,20 @@ export interface AgentTransport extends AgentTransportThreadOperations { input: CancelRunInput, context?: AgentRequestContext, ): Promise; + /** + * Answers the interrupts that ended a run and starts the run that carries the + * work forward. An interrupt terminates its run under AG-UI semantics, so + * resolution produces a new run to subscribe to rather than resuming a + * stream that is already closed. + */ + resumeRun?( + input: ResumeRunInput, + context?: AgentRequestContext, + ): Promise; + /** + * @deprecated Implement {@link resumeRun}. Kept as a source-compatible + * bridge for protocol-v1 transports while consumers migrate to interrupts. + */ resolveApproval?( input: ResolveApprovalInput, context?: AgentRequestContext, @@ -1236,9 +1250,33 @@ export interface StartRunInput { threadId: ThreadId; messages: AgentMessage[]; options?: AgentRunOptions; + /** + * Resolutions for interrupts raised by an earlier run. Mirrors AG-UI's + * `RunAgentInput.resume`, so a resumed run is an ordinary run start rather + * than a second command channel with its own lifecycle. + */ + resume?: AgentResumeEntry[]; + metadata?: AgentProtocolMetadata; +} + +/** + * One interrupt resolution, shaped exactly like AG-UI's `ResumeEntry` so it + * needs no translation on the wire. + */ +export interface AgentResumeEntry { + interruptId: ApprovalId; + status: "resolved" | "cancelled"; + payload?: unknown; metadata?: AgentProtocolMetadata; } +export interface ResumeRunInput { + threadId: ThreadId; + /** The interrupted run being answered, used to scope authorization. */ + runId: RunId; + resume: AgentResumeEntry[]; +} + export interface StartRunResult { runId: RunId; capabilities?: AgentCapabilities; @@ -1260,6 +1298,7 @@ export interface CancelRunInput { runId: RunId; } +/** @deprecated Use {@link ResumeRunInput} with {@link AgentResumeEntry}. */ export interface ResolveApprovalInput { threadId: ThreadId; runId: RunId; @@ -1317,3 +1356,5 @@ export interface AgentProtocolEnvelope { export * from "./validation.js"; export * from "./compatibility.js"; export * from "./errors.js"; +export * from "./agui.js"; +export * from "./agui-codec.js"; diff --git a/packages/agentkit/src/protocol/validation.spec.ts b/packages/agentkit/src/protocol/validation.spec.ts index d582624e55c..05e8787d753 100644 --- a/packages/agentkit/src/protocol/validation.spec.ts +++ b/packages/agentkit/src/protocol/validation.spec.ts @@ -14,7 +14,7 @@ import { parseAgentProtocolEnvelope, parseInvokeActionInput, parseQueueMessageInput, - parseResolveApprovalInput, + parseResumeRunInput, parseAgentThreadSnapshot, parseStartRunInput, } from "./index.js"; @@ -41,13 +41,22 @@ describe("AgentKit protocol validation", () => { expect(() => parseAgentApprovalResponse({ optionIds: ["approve"] }), ).toThrow("approvalResponse.decision"); + }); + + it("rejects a resume entry without a resolved or cancelled status", () => { expect(() => - parseResolveApprovalInput({ + parseResumeRunInput({ threadId: "thread-1", runId: "run-1", - approvalId: "approval-1", + resume: [{ interruptId: "approval-1", status: "maybe" }], }), - ).toThrow("resolveApproval.response"); + ).toThrow("resumeRun.resume[0].status"); + }); + + it("rejects a resume with no interrupt resolutions", () => { + expect(() => + parseResumeRunInput({ threadId: "thread-1", runId: "run-1", resume: [] }), + ).toThrow("resumeRun.resume"); }); it("keeps custom choice responses distinct from predefined option ids", () => { diff --git a/packages/agentkit/src/protocol/validation.ts b/packages/agentkit/src/protocol/validation.ts index bc02fba38cd..4fabbf7feca 100644 --- a/packages/agentkit/src/protocol/validation.ts +++ b/packages/agentkit/src/protocol/validation.ts @@ -21,6 +21,8 @@ import { type QueueMessageInput, type ResolveApprovalInput, type ResolveConnectionRequestInput, + type AgentResumeEntry, + type ResumeRunInput, type SubscribeToRunInput, type SubmitFeedbackInput, type ThreadId, @@ -2481,10 +2483,55 @@ export function parseStartRunInput( if (input.options !== undefined) { parseAgentRunOptions(input.options, `${path}.options`); } + if (input.resume !== undefined) { + parseAgentResumeEntries(input.resume, `${path}.resume`); + } optionalMetadata(input.metadata, `${path}.metadata`); return value as StartRunInput; } +function parseAgentResumeEntries( + value: unknown, + path: string, +): AgentResumeEntry[] { + array(value, path).forEach((entry, index) => { + const resume = record(entry, `${path}[${index}]`); + knownKeys( + resume, + ["interruptId", "status", "payload", "metadata"], + `${path}[${index}]`, + ); + string(resume.interruptId, `${path}[${index}].interruptId`); + const status = string(resume.status, `${path}[${index}].status`); + if (status !== "resolved" && status !== "cancelled") { + throw new AgentProtocolValidationError( + `${path}[${index}].status`, + "expected resolved or cancelled", + ); + } + optionalMetadata(resume.metadata, `${path}[${index}].metadata`); + }); + return value as AgentResumeEntry[]; +} + +export function parseResumeRunInput( + value: unknown, + path = "resumeRun", +): ResumeRunInput { + const input = record(value, path); + knownKeys(input, ["threadId", "runId", "resume"], path); + string(input.threadId, `${path}.threadId`); + string(input.runId, `${path}.runId`); + const entries = parseAgentResumeEntries(input.resume, `${path}.resume`); + if (entries.length === 0) { + throw new AgentProtocolValidationError( + `${path}.resume`, + "expected at least one interrupt resolution", + ); + } + return value as ResumeRunInput; +} + export function parseSubscribeToRunInput( value: unknown, path = "subscribeToRun", @@ -2649,6 +2696,7 @@ export function parseCancelRunInput( return value as CancelRunInput; } +/** @deprecated Validation bridge for the protocol-v1 approval command. */ export function parseResolveApprovalInput( value: unknown, path = "resolveApproval", diff --git a/packages/agentkit/src/protocol/version.ts b/packages/agentkit/src/protocol/version.ts index b9164427d42..bc965e6bee8 100644 --- a/packages/agentkit/src/protocol/version.ts +++ b/packages/agentkit/src/protocol/version.ts @@ -1,5 +1,5 @@ export const AGENTKIT_PROTOCOL_NAME = "agentkit" as const; -export const AGENTKIT_PROTOCOL_VERSION = 1 as const; +export const AGENTKIT_PROTOCOL_VERSION = 2 as const; export const AGENTKIT_SUPPORTED_PROTOCOL_VERSIONS = [ AGENTKIT_PROTOCOL_VERSION, ] as const; diff --git a/packages/core/src/client/chat/agentkit-protocol.spec.ts b/packages/core/src/client/chat/agentkit-protocol.spec.ts index 68af7ced657..6dc1cd8b1c4 100644 --- a/packages/core/src/client/chat/agentkit-protocol.spec.ts +++ b/packages/core/src/client/chat/agentkit-protocol.spec.ts @@ -3,7 +3,10 @@ import type { AgentEvent, AgentMessage, } from "@agent-native/agentkit/protocol"; -import { createAgentKitProtocolVersionOffer } from "@agent-native/agentkit/protocol"; +import { + createAgentKitProtocolVersionOffer, + resumeEntryFromApproval, +} from "@agent-native/agentkit/protocol"; import { describe, expect, it, vi } from "vitest"; import { createAgentKitProtocolAdapter } from "./agentkit-protocol.js"; @@ -123,6 +126,29 @@ function createRuntime( } describe("createAgentKitProtocolAdapter", () => { + it("rejects resume entries on ordinary Core run starts", async () => { + async function* events(): AsyncIterable { + yield { type: "done", reason: "complete" }; + } + const runtime = createRuntime(events); + const createSession = vi.spyOn(runtime, "createSession"); + const transport = createAgentKitProtocolAdapter(runtime); + + await expect( + transport.startRun({ + threadId: "thread-1", + messages: [userMessage("Continue")], + resume: [ + resumeEntryFromApproval({ + approvalId: "approval-1", + response: approvalResponse("approve"), + }), + ], + }), + ).rejects.toThrow("startRun.resume"); + expect(createSession).not.toHaveBeenCalled(); + }); + it("pauses for a typed connection request and resumes the same run", async () => { async function* connectionEvents(): AsyncIterable { yield { @@ -837,21 +863,126 @@ describe("createAgentKitProtocolAdapter", () => { approvalSeen = next.value?.type === "approval.requested"; } - await transport.resolveApproval?.({ + const resumed = await transport.resumeRun?.({ threadId: "thread-1", runId, - approvalId: "approval-1", - response: approvalResponse("approve"), + resume: [ + resumeEntryFromApproval({ + approvalId: "approval-1", + response: approvalResponse("approve"), + }), + ], }); - const remaining: AgentEvent[] = []; + expect(resumed?.runId).not.toBe(runId); + expect(await iterator.next()).toMatchObject({ done: true }); + const remaining = await drain( + transport.subscribeToRun({ + threadId: "thread-1", + runId: resumed!.runId, + }), + ); + + expect(continueTurnCalled).toBe(true); + expect(remaining.map((event) => event.type)).toContain("approval.resolved"); + expect(remaining.map((event) => event.type)).toContain("run.completed"); + }); + + it("cancels a paused Core turn after its approval stream closes", async () => { + const cancel = vi.fn(async () => ({ status: "cancelled" as const })); + async function* approvalEvents(): AsyncIterable { + yield { + type: "approval-request", + approvalId: "approval-1", + toolCallId: "tool-1", + toolName: "publish", + message: "Publish?", + }; + } + const runtime = createRuntime(approvalEvents, { + capabilities: { + messages: { streaming: true }, + tools: { events: true, approvals: true }, + }, + }); + runtime.createSession = async () => ({ + id: "thread-1", + runtimeId: runtime.id, + startTurn: async () => ({ + id: "turn-1", + runId: "runtime-run-1", + sessionId: "thread-1", + events: approvalEvents(), + cancel, + }), + }); + + const transport = createAgentKitProtocolAdapter(runtime); + const { runId } = await transport.startRun({ + threadId: "thread-1", + messages: [userMessage("Publish")], + }); + const iterator = transport + .subscribeToRun({ threadId: "thread-1", runId }) + [Symbol.asyncIterator](); while (true) { const next = await iterator.next(); - if (next.done) break; - remaining.push(next.value); + if (next.value?.type === "approval.requested") break; } + expect(await iterator.next()).toMatchObject({ done: true }); - expect(continueTurnCalled).toBe(true); - expect(remaining.map((event) => event.type)).toContain("run.completed"); + await transport.cancelRun({ threadId: "thread-1", runId }); + + expect(cancel).toHaveBeenCalledWith({ reason: "protocol-cancel" }); + await expect( + transport.getRun?.({ threadId: "thread-1", runId }), + ).resolves.toMatchObject({ status: "cancelled" }); + }); + + it("cancels a paused Core turn when the adapter is disposed", async () => { + const cancel = vi.fn(async () => ({ status: "cancelled" as const })); + const disposeSession = vi.fn(async () => undefined); + async function* approvalEvents(): AsyncIterable { + yield { + type: "approval-request", + approvalId: "approval-1", + toolCallId: "tool-1", + toolName: "publish", + message: "Publish?", + }; + } + const runtime = createRuntime(approvalEvents, { + capabilities: { + messages: { streaming: true }, + tools: { events: true, approvals: true }, + }, + }); + runtime.createSession = async () => ({ + id: "thread-1", + runtimeId: runtime.id, + startTurn: async () => ({ + id: "turn-1", + runId: "runtime-run-1", + sessionId: "thread-1", + events: approvalEvents(), + cancel, + }), + dispose: disposeSession, + }); + + const transport = createAgentKitProtocolAdapter(runtime); + const { runId } = await transport.startRun({ + threadId: "thread-1", + messages: [userMessage("Publish")], + }); + const events = await drain( + transport.subscribeToRun({ threadId: "thread-1", runId }), + ); + expect(events.map((event) => event.type)).toContain("approval.requested"); + + await transport.dispose(); + + expect(cancel).toHaveBeenCalledWith({ reason: "adapter-dispose" }); + expect(disposeSession).toHaveBeenCalledOnce(); }); it("binds approval decisions to the exact pending request and fails closed", async () => { @@ -907,39 +1038,82 @@ describe("createAgentKitProtocolAdapter", () => { } await expect( - transport.resolveApproval?.({ + transport.resumeRun?.({ threadId: "thread-1", runId, - approvalId: "approval-stale", - response: approvalResponse("approve"), + resume: [ + resumeEntryFromApproval({ + approvalId: "approval-stale", + response: approvalResponse("approve"), + }), + ], }), ).rejects.toThrow( "Approval response approval-stale does not match pending approval approval-current", ); await expect( - transport.resolveApproval?.({ + transport.resumeRun?.({ threadId: "thread-1", runId, - approvalId: "approval-current", - response: {} as AgentApprovalResponse, + resume: [ + { + interruptId: "approval-current", + status: "resolved", + payload: {} as AgentApprovalResponse, + }, + ], }), ).rejects.toThrow("must include decision"); + await expect( + transport.resumeRun?.({ + threadId: "thread-1", + runId, + resume: [ + { + interruptId: "approval-current", + status: "resolved", + payload: { decision: "approve", optionIds: [1] }, + }, + ], + }), + ).rejects.toThrow("optionIds[0]"); + await expect( + transport.resumeRun?.({ + threadId: "thread-1", + runId, + resume: [ + { + interruptId: "approval-current", + status: "resolved", + payload: { decision: "approve", optionId: 1 }, + }, + ], + }), + ).rejects.toThrow("optionId"); expect(continueTurn).not.toHaveBeenCalled(); - await transport.resolveApproval?.({ + const resumed = await transport.resumeRun?.({ threadId: "thread-1", runId, - approvalId: "approval-current", - optionId: "approve", - response: { - ...approvalResponse("deny", { message: "Not yet" }), - optionIds: ["approve"], - other: "Use the staging channel", - }, - }); - await drain({ - [Symbol.asyncIterator]: () => iterator, + resume: [ + resumeEntryFromApproval({ + approvalId: "approval-current", + optionId: "approve", + response: { + ...approvalResponse("deny", { message: "Not yet" }), + optionIds: ["approve"], + other: "Use the staging channel", + }, + }), + ], }); + expect(await iterator.next()).toMatchObject({ done: true }); + await drain( + transport.subscribeToRun({ + threadId: "thread-1", + runId: resumed!.runId, + }), + ); expect(continueTurn).toHaveBeenCalledWith({ turnId: "turn-1", @@ -950,11 +1124,15 @@ describe("createAgentKitProtocolAdapter", () => { }, }); await expect( - transport.resolveApproval?.({ + transport.resumeRun?.({ threadId: "thread-1", runId, - approvalId: "approval-current", - response: approvalResponse("approve"), + resume: [ + resumeEntryFromApproval({ + approvalId: "approval-current", + response: approvalResponse("approve"), + }), + ], }), ).rejects.toThrow("already terminal"); }); @@ -1393,18 +1571,24 @@ describe("createAgentKitProtocolAdapter", () => { if (next.value?.type === "approval.requested") break; } - await transport.resolveApproval!({ + const resumed = await transport.resumeRun!({ threadId: "thread-1", runId, - approvalId: "approval-1", - response: approvalResponse("approve"), + resume: [ + resumeEntryFromApproval({ + approvalId: "approval-1", + response: approvalResponse("approve"), + }), + ], }); - const remaining: AgentEvent[] = []; - while (true) { - const next = await iterator.next(); - if (next.done) break; - remaining.push(next.value); - } + expect(resumed.runId).not.toBe(runId); + expect(await iterator.next()).toMatchObject({ done: true }); + const remaining = await drain( + transport.subscribeToRun({ + threadId: "thread-1", + runId: resumed.runId, + }), + ); expect(continueTurn).toHaveBeenCalledOnce(); expect(remaining.some((event) => event.type === "run.completed")).toBe( diff --git a/packages/core/src/client/chat/agentkit-protocol.ts b/packages/core/src/client/chat/agentkit-protocol.ts index 9da8453df19..5ba4914aa63 100644 --- a/packages/core/src/client/chat/agentkit-protocol.ts +++ b/packages/core/src/client/chat/agentkit-protocol.ts @@ -30,10 +30,14 @@ import type { } from "@agent-native/agentkit/protocol"; import { AGENTKIT_PROTOCOL_VERSION, + AgentProtocolValidationError, + approvalResponseFromResume, createCapabilityUnavailableError, createCapabilityUnsupportedError, inferAgentActivityKind, negotiateAgentKitProtocolVersion, + resumeEntryFromApproval, + resumeOptionId, } from "@agent-native/agentkit/protocol"; import type { @@ -121,6 +125,7 @@ interface ProtocolRun { activeActivities: Map; pumpPromise: Promise | null; continuationPromise: Promise | null; + streamClosed: boolean; terminal: boolean; waitingForContinuation: boolean; pendingApprovalId?: string; @@ -1111,7 +1116,12 @@ export function createAgentKitProtocolAdapter( } const retainedCompleted = [...runs.values()] - .filter((run) => run.terminal && run.activeReaders === 0) + .filter( + (run) => + run.terminal && + !run.waitingForContinuation && + run.activeReaders === 0, + ) .sort((left, right) => left.lastAccessedAtMs - right.lastAccessedAtMs); while (retainedCompleted.length > maxRetainedRuns) { const run = retainedCompleted.shift(); @@ -1597,6 +1607,10 @@ export function createAgentKitProtocolAdapter( } run.waitingForContinuation = true; run.pendingApprovalId = event.approvalId; + // AG-UI models an approval interrupt as RUN_FINISHED. Close the wire + // stream while retaining ownership of the paused Core turn so it can + // still be resumed, cancelled, or disposed. + run.streamClosed = true; return [ { type: "run.status", @@ -2072,7 +2086,8 @@ export function createAgentKitProtocolAdapter( } function ensurePump(run: ProtocolRun): void { - if (run.pumpPromise || run.terminal || !run.turn) return; + if (run.pumpPromise || run.streamClosed || run.terminal || !run.turn) + return; run.pumpPromise = (async () => { try { for await (const event of run.turn.events) { @@ -2096,7 +2111,8 @@ export function createAgentKitProtocolAdapter( }); } } catch (error) { - if (run.terminal) return; + if (run.streamClosed || run.terminal) return; + run.streamClosed = true; run.terminal = true; append(run, { type: "run.failed", @@ -2105,7 +2121,9 @@ export function createAgentKitProtocolAdapter( } finally { run.pumpPromise = null; for (const listener of run.listeners) listener(); - if (!run.terminal && !run.waitingForContinuation) ensurePump(run); + if (!run.streamClosed && !run.terminal && !run.waitingForContinuation) { + ensurePump(run); + } } })(); } @@ -2135,7 +2153,7 @@ export function createAgentKitProtocolAdapter( } yield event; } - if (run.terminal) return; + if (run.streamClosed || run.terminal) return; await notifyWhenChanged(run, signal); } } finally { @@ -2185,6 +2203,12 @@ export function createAgentKitProtocolAdapter( }, async startRun(input) { if (disposed) throw new Error("The AgentKit adapter is disposed."); + if (input.resume !== undefined) { + throw new AgentProtocolValidationError( + "startRun.resume", + "Core requires approval continuations to use resumeRun", + ); + } const turnMetadata = mergeTrustedProtocolMetadata( options.metadata, input.metadata, @@ -2245,6 +2269,7 @@ export function createAgentKitProtocolAdapter( activeActivities: new Map(), pumpPromise: null, continuationPromise: null, + streamClosed: false, terminal: false, waitingForContinuation: false, listeners: new Set(), @@ -2302,23 +2327,39 @@ export function createAgentKitProtocolAdapter( } touchRun(run); if (run.terminal) return; + const streamWasClosed = run.streamClosed; const cancellation = cancelCoreTurn(run, "protocol-cancel"); + run.streamClosed = true; run.terminal = true; let result; try { result = await cancellation; } catch (error) { + run.streamClosed = streamWasClosed; run.terminal = false; ensurePump(run); throw error; } if (result.status === "unsupported") { + run.streamClosed = streamWasClosed; run.terminal = false; ensurePump(run); throw new Error("The Core runtime does not support run cancellation."); } - append(run, { type: "run.status", status: "cancelled" }); - append(run, { type: "run.cancelled" }); + run.waitingForContinuation = false; + run.pendingApprovalId = undefined; + if (streamWasClosed) { + const completedAt = now(); + run.status = "cancelled"; + run.completedAt = completedAt; + run.terminalAtMs = timeMs(completedAt); + run.actions.clear(); + run.activeActivities.clear(); + pruneRetainedRuns(run.terminalAtMs); + } else { + append(run, { type: "run.status", status: "cancelled" }); + append(run, { type: "run.cancelled" }); + } }, async dispose() { if (disposed) return; @@ -2328,10 +2369,15 @@ export function createAgentKitProtocolAdapter( for (const run of runs.values()) { if (!run.terminal) { cancellations.push(cancelCoreTurn(run, "adapter-dispose")); - append(run, { type: "run.status", status: "cancelled" }); - append(run, { type: "run.cancelled" }); + if (!run.streamClosed) { + append(run, { type: "run.status", status: "cancelled" }); + append(run, { type: "run.cancelled" }); + } } + run.streamClosed = true; run.terminal = true; + run.waitingForContinuation = false; + run.pendingApprovalId = undefined; for (const listener of run.listeners) listener(); } await Promise.allSettled(cancellations); @@ -2483,13 +2529,23 @@ export function createAgentKitProtocolAdapter( } if (runtime.capabilities.tools?.approvals) { - transport.resolveApproval = async (input) => { + const resumeRun: NonNullable = async ( + input, + ) => { pruneRetainedRuns(); const run = runs.get(input.runId); if (!run || run.threadId !== input.threadId) { throw new Error(`Unknown AgentKit run: ${input.runId}`); } touchRun(run); + const entry = input.resume[0]; + if (input.resume.length !== 1 || !entry) { + throw new Error( + "The Core runtime resolves exactly one interrupt per resume.", + ); + } + const response = approvalResponseFromResume(entry); + const optionId = resumeOptionId(entry); if (!run.session.continueTurn) { throw new Error( "The Core runtime does not support approval continuation.", @@ -2498,44 +2554,99 @@ export function createAgentKitProtocolAdapter( if (run.continuationPromise) { throw new Error("An approval continuation is already in progress."); } + let replacementRunId: string | undefined; const continuation = (async () => { const activePump = run.pumpPromise; if (activePump) await activePump; - if (run.terminal) { + if (run.terminal && !run.waitingForContinuation) { throw new Error("The AgentKit run is already terminal."); } if (!run.waitingForContinuation) { throw new Error("The AgentKit run is not awaiting approval."); } - if (run.pendingApprovalId !== input.approvalId) { + if (run.pendingApprovalId !== entry.interruptId) { throw new Error( - `Approval response ${input.approvalId} does not match pending approval ${run.pendingApprovalId ?? ""}.`, + `Approval response ${entry.interruptId} does not match pending approval ${run.pendingApprovalId ?? ""}.`, ); } - const decision = explicitApprovalDecision(input.response); + const decision = explicitApprovalDecision(response); const approvalInput = - input.response.other !== undefined - ? { ...input.response.input, other: input.response.other } - : input.response.input; + response.other !== undefined + ? { ...response.input, other: response.other } + : response.input; + const nextRunId = createId("run"); + if (nextRunId === run.runId || runs.has(nextRunId)) { + throw new Error( + `The Core runtime generated duplicate AgentKit run id ${nextRunId}.`, + ); + } const nextTurn = await run.session.continueTurn!({ turnId: run.turn.id, approval: { - id: input.approvalId, + id: entry.interruptId, approved: decision === "approve", message: serializeValue(approvalInput), }, }); + const replacementMetadata = mergeTrustedProtocolMetadata( + run.metadata, + nextTurn.metadata, + { + [AGENT_NATIVE_PROTOCOL_METADATA_KEY]: { + observability: { + protocolRunId: nextRunId, + runtimeRunId: nextTurn.runId, + runtimeId: runtime.id, + sessionId: run.session.id, + turnId: nextTurn.id, + threadId: input.threadId, + interruptedRunId: run.runId, + }, + } satisfies AgentNativeProtocolMetadata, + }, + ); + const replacementRun: ProtocolRun = { + runId: nextRunId, + threadId: input.threadId, + session: run.session, + turn: nextTurn, + events: [], + firstRetainedSequence: 1, + sequence: 0, + status: "queued", + lastAccessedAtMs: timeMs(), + activeReaders: 0, + metadata: replacementMetadata, + activeMessageId: run.activeMessageId, + actions: new Map(run.actions), + activeActivities: new Map(run.activeActivities), + pumpPromise: null, + continuationPromise: null, + streamClosed: false, + terminal: false, + waitingForContinuation: false, + listeners: new Set(), + }; run.waitingForContinuation = false; run.pendingApprovalId = undefined; - run.turn = nextTurn; - append(run, { + run.terminal = true; + run.terminalAtMs = timeMs(); + runs.set(nextRunId, replacementRun); + append(replacementRun, { + type: "run.started", + agentId: runtime.id, + metadata: replacementMetadata, + }); + append(replacementRun, { type: "approval.resolved", - approvalId: input.approvalId, - optionId: input.optionId, - response: input.response, + approvalId: entry.interruptId, + optionId, + response, }); - append(run, { type: "run.status", status: "running" }); - ensurePump(run); + append(replacementRun, { type: "run.status", status: "running" }); + ensurePump(replacementRun); + replacementRunId = nextRunId; + pruneRetainedRuns(run.terminalAtMs); })(); run.continuationPromise = continuation; try { @@ -2545,6 +2656,18 @@ export function createAgentKitProtocolAdapter( run.continuationPromise = null; } } + if (!replacementRunId) { + throw new Error("The Core runtime did not create a replacement run."); + } + return { runId: replacementRunId }; + }; + transport.resumeRun = resumeRun; + transport.resolveApproval = async (input) => { + await resumeRun({ + threadId: input.threadId, + runId: input.runId, + resume: [resumeEntryFromApproval(input)], + }); }; } diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index 1c375c2ba8e..cbffe95a444 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -301,6 +301,9 @@ importers: packages/agentkit: dependencies: + '@ag-ui/core': + specifier: 0.0.59 + version: 0.0.59 '@agent-native/toolkit': specifier: workspace:* version: link:../toolkit @@ -5877,6 +5880,9 @@ packages: '@adobe/css-tools@4.5.0': resolution: {integrity: sha512-6OzddxPio9UiWTCemp4N8cYLV2ZN1ncRnV1cVGtve7dhPOtRkleRyx32GQCYSwDYgaHU3USMm84tNsvKzRCa1Q==} + '@ag-ui/core@0.0.59': + resolution: {integrity: sha512-hDgy4ipTqXieT8YG8Mr917Y+FD/f11VK1GefZ5CwTDCuNqS/oTwjJ5l/DZkicThgS8hQW/Y7wPylPBMBJ8BkUg==} + '@ai-sdk/anthropic@3.0.86': resolution: {integrity: sha512-IVUaTOz46dNniRaALEvphV5WOqPfKTJRZlJhJabz3kaXk1QiwdnzWQhKMKLiprD+872E9YHqmmyFNFFH1MfnBQ==} engines: {node: '>=18'} @@ -21331,6 +21337,10 @@ snapshots: '@adobe/css-tools@4.5.0': optional: true + '@ag-ui/core@0.0.59': + dependencies: + zod: 3.25.76 + '@ai-sdk/anthropic@3.0.86(zod@4.4.3)': dependencies: '@ai-sdk/provider': 3.0.10