From 206facf6dae974b96e080e87e47d8f670d8ef3d3 Mon Sep 17 00:00:00 2001 From: Brian Love Date: Thu, 24 Sep 2026 13:10:30 -0700 Subject: [PATCH 1/2] fix(ag-ui): keep child reasoning in its attributed transcript --- libs/ag-ui/src/lib/reducer.subagent.spec.ts | 163 +++++++++++++++++- libs/ag-ui/src/lib/reducer.ts | 27 ++- .../src/lib/to-agent.child-reasoning.spec.ts | 160 +++++++++++++++++ scripts/react-parity/baseline.json | 2 +- 4 files changed, 348 insertions(+), 4 deletions(-) create mode 100644 libs/ag-ui/src/lib/to-agent.child-reasoning.spec.ts diff --git a/libs/ag-ui/src/lib/reducer.subagent.spec.ts b/libs/ag-ui/src/lib/reducer.subagent.spec.ts index e0c5a89ba..f15408b07 100644 --- a/libs/ag-ui/src/lib/reducer.subagent.spec.ts +++ b/libs/ag-ui/src/lib/reducer.subagent.spec.ts @@ -1,4 +1,4 @@ -import { describe, it, expect } from 'vitest'; +import { describe, it, expect, vi } from 'vitest'; import { signal } from '@angular/core'; import { Subject } from 'rxjs'; import { @@ -19,6 +19,7 @@ interface TestDeliveryRun { currentAssistantMessageId?: string; eligibleBaselineAssistantId?: string; protocolRunId?: string; + pendingReasoning?: { messageId: string; startedAt: number }; outcome?: 'success' | 'error' | 'aborted' | 'interrupted' | 'paused'; } @@ -52,7 +53,167 @@ function makeStore(generation = 'run-generation-1'): TestStore { const ev = (e: Record) => e as unknown as BaseEvent; +describe('reduceEvent attributed child reasoning', () => { + const reasoningTypes = ['START', 'CONTENT', 'CHUNK', 'END']; + const reason = (store: TestStore, type: string, fields: Record = {}) => + reduceEvent(ev({ type: `REASONING_MESSAGE_${type}`, subagentRunId: 'child', messageId: 'shared', ...fields }), store); + + it.each(['shared', 'different'])('isolates child message %s from the parent, sibling and root timer', childMessageId => { + const clock = vi.spyOn(Date, 'now').mockReturnValue(100); + try { + const store = makeStore(); + reduceEvent(ev({ type: 'REASONING_MESSAGE_START', messageId: 'shared' }), store); + reduceEvent(ev({ type: 'REASONING_MESSAGE_CONTENT', messageId: 'shared', delta: 'parent' }), store); + reduceEvent(ev({ type: 'TEXT_MESSAGE_CONTENT', subagentRunId: 'sibling', messageId: childMessageId, delta: 'sibling' }), store); + const parent = store.messages(); + const tools = store.toolCalls(); + const pending = store.deliveryRun!.pendingReasoning; + const sibling = store.activities().get('sibling')!; + const siblingContent = sibling.content(); + clock.mockReturnValue(200); + reason(store, 'START', { messageId: childMessageId }); + reason(store, 'CONTENT', { messageId: childMessageId, delta: ' child\n' }); + reason(store, 'CHUNK', { messageId: childMessageId, delta: 'reasoning ' }); + reason(store, 'END', { messageId: childMessageId }); + expect(store.messages()).toBe(parent); + expect(store.messages()[0]).toBe(parent[0]); + expect(store.toolCalls()).toBe(tools); + expect(store.deliveryRun!.pendingReasoning).toBe(pending); + expect(store.deliveryRun!.currentAssistantMessageId).toBe('shared'); + expect([...store.deliveryRun!.ownedMessageIds]).toEqual(['shared']); + expect(store.activities().get('sibling')).toBe(sibling); + expect(sibling.content()).toBe(siblingContent); + expect(store.activities().get('child')!.content()['messages']).toEqual([ + { id: childMessageId, role: 'assistant', content: '', reasoning: ' child\nreasoning ' }, + ]); + expect(store.deliveryRun!.outcome).toBeUndefined(); + clock.mockReturnValue(400); + reduceEvent(ev({ type: 'REASONING_MESSAGE_END', messageId: 'shared' }), store); + expect(store.messages()[0].reasoningDurationMs).toBe(300); + } finally { + clock.mockRestore(); + } + }); + + it.each(['CONTENT', 'CHUNK'])('buffers %s before STARTED and preserves reasoning through identity refresh', type => { + const store = makeStore(); + reason(store, type, { delta: 'early' }); + const buffered = store.activities().get('child'); + expect(buffered?.content()['messages']).toMatchObject([{ reasoning: 'early', content: '' }]); + reduceEvent(ev({ type: 'SUBAGENT_STARTED', subagentRunId: 'child', name: 'researcher', parentToolCallId: 'parent-tool' }), store); + const announced = store.activities().get('child')!; + expect(announced.generation).not.toBe(buffered?.generation); + expect(announced.content()).toMatchObject({ name: 'researcher', toolCallId: 'parent-tool', messages: [{ reasoning: 'early' }] }); + const activities = store.activities(); + reason(store, 'CHUNK', { delta: ' later' }); + expect(store.activities()).toBe(activities); + expect(store.activities().get('child')).toBe(announced); + expect(announced.content()['messages']).toMatchObject([{ reasoning: 'early later' }]); + }); + + it('keeps reasoning, answer and tool links in one slot across empty deltas and duplicate START', () => { + const store = makeStore(); + reason(store, 'START'); + expect(store.activities().get('child')?.content()['messages']).toEqual([{ id: 'shared', role: 'assistant', content: '', reasoning: '' }]); + reason(store, 'CONTENT', { delta: 'thought' }); + reduceEvent(ev({ type: 'TOOL_CALL_START', subagentRunId: 'child', toolCallId: 'tool', toolCallName: 'lookup', parentMessageId: 'shared' }), store); + reduceEvent(ev({ type: 'TEXT_MESSAGE_START', subagentRunId: 'child', messageId: 'shared' }), store); + reduceEvent(ev({ type: 'TEXT_MESSAGE_CONTENT', subagentRunId: 'child', messageId: 'shared', delta: 'answer' }), store); + reason(store, 'START'); + reason(store, 'START'); + reason(store, 'CONTENT', { delta: '' }); + reason(store, 'CHUNK'); + expect(store.activities().get('child')!.content()['messages']).toEqual([ + { id: 'shared', role: 'assistant', content: 'answer', reasoning: 'thought', toolCallIds: ['tool'] }, + ]); + expect(store.messages()).toEqual([]); + expect(store.toolCalls()).toEqual([]); + }); + + it('END alone and duplicate END do not allocate, bind a run, or change existing content', () => { + const store = makeStore(); + const activities = store.activities(); + const messages = store.messages(); + reason(store, 'END', { runId: 'unbound' }); + expect(store.activities()).toBe(activities); + expect(store.messages()).toBe(messages); + expect(store.deliveryRun!.protocolRunId).toBeUndefined(); + reason(store, 'CONTENT', { delta: 'thought' }); + const entry = store.activities().get('child')!; + const content = entry.content(); + reason(store, 'END', { runId: 'unbound' }); + reason(store, 'END', { runId: 'unbound' }); + expect(entry.content()).toBe(content); + expect(store.deliveryRun!.protocolRunId).toBeUndefined(); + expect(store.deliveryRun!.outcome).toBeUndefined(); + }); + + it.each(['SUBAGENT_FINISHED', 'SUBAGENT_ERROR'])('rejects every reasoning event after %s without binding the root run', terminalType => { + const store = makeStore(); + reduceEvent(ev({ type: 'SUBAGENT_STARTED', subagentRunId: 'child', name: 'researcher' }), store); + reduceEvent(ev({ type: terminalType, subagentRunId: 'child', message: 'failed' }), store); + const activities = store.activities(); + const entry = activities.get('child')!; + const content = entry.content(); + const messages = store.messages(); + for (const type of reasoningTypes) reason(store, type, { runId: 'unbound', delta: 'late' }); + expect(store.activities()).toBe(activities); + expect(entry.content()).toBe(content); + expect(store.messages()).toBe(messages); + expect(store.deliveryRun!.protocolRunId).toBeUndefined(); + expect(store.deliveryRun!.pendingReasoning).toBeUndefined(); + }); + + it('accepts reasoning for a suspended child and retains the re-announced identity', () => { + const store = makeStore(); + reason(store, 'CONTENT', { delta: 'before' }); + reduceEvent(ev({ type: 'SUBAGENT_STARTED', subagentRunId: 'child', name: 'researcher' }), store); + const entry = store.activities().get('child')!; + reduceEvent(ev({ type: 'SUBAGENT_FINISHED', subagentRunId: 'child', outcome: { type: 'suspended' } }), store); + reason(store, 'CHUNK', { delta: ' during' }); + reduceEvent(ev({ type: 'SUBAGENT_STARTED', subagentRunId: 'child', name: 'researcher' }), store); + expect(store.activities().get('child')).toBe(entry); + expect(entry.content()).toMatchObject({ status: 'running', messages: [{ reasoning: 'before during' }] }); + }); + + it.each(['missing', 'settled', 'foreign'])('rejects every child reasoning event with a %s root run before mutation', state => { + const store = makeStore(); + if (state === 'missing') store.deliveryRun = null; + else if (state === 'settled') store.deliveryRun!.outcome = 'success'; + else store.deliveryRun!.protocolRunId = 'current'; + const runBefore = store.deliveryRun ? { ...store.deliveryRun } : null; + const activities = store.activities(); + const messages = store.messages(); + for (const type of reasoningTypes) reason(store, type, { runId: 'foreign', delta: 'late' }); + expect(store.activities()).toBe(activities); + expect(store.messages()).toBe(messages); + expect(store.deliveryRun).toEqual(runBefore); + }); + + it('accepts matching run IDs and binds an initially unbound active run', () => { + const store = makeStore(); + reason(store, 'START', { runId: 'current' }); + reason(store, 'CONTENT', { runId: 'current', delta: 'thought' }); + expect(store.deliveryRun!.protocolRunId).toBe('current'); + expect(store.activities().get('child')?.content()['messages']).toMatchObject([{ reasoning: 'thought' }]); + }); +}); + describe('reduceEvent SUBAGENT_* lifecycle', () => { + it('keeps child reasoning with its child even when the parent uses the same message ID', () => { + const store = makeStore(); + reduceEvent(ev({ type: 'REASONING_MESSAGE_START', messageId: 'shared' }), store); + reduceEvent(ev({ type: 'REASONING_MESSAGE_CONTENT', messageId: 'shared', delta: 'parent reasoning' }), store); + const parentMessages = store.messages(); + reduceEvent(ev({ type: 'SUBAGENT_STARTED', subagentRunId: 'child', name: 'researcher' }), store); + reduceEvent(ev({ type: 'REASONING_MESSAGE_START', messageId: 'shared', subagentRunId: 'child' }), store); + reduceEvent(ev({ type: 'REASONING_MESSAGE_CONTENT', messageId: 'shared', subagentRunId: 'child', delta: 'child reasoning' }), store); + expect(store.messages()).toBe(parentMessages); + expect(store.activities().get('child')?.content()['messages']).toMatchObject([ + { id: 'shared', role: 'assistant', content: '', reasoning: 'child reasoning' }, + ]); + }); + it('SUBAGENT_STARTED creates a running subagent activity entry', () => { const store = makeStore(); reduceEvent(ev({ type: 'SUBAGENT_STARTED', subagentRunId: 'sa-1', name: 'researcher', parentToolCallId: 'call-9' }), store); diff --git a/libs/ag-ui/src/lib/reducer.ts b/libs/ag-ui/src/lib/reducer.ts index 0fd3fd68c..7faaeda98 100644 --- a/libs/ag-ui/src/lib/reducer.ts +++ b/libs/ag-ui/src/lib/reducer.ts @@ -120,8 +120,8 @@ export function reduceEvent(event: BaseEvent, store: ReducerStore): void { // A subagentRunId on a content event means the child produced it: route it // into that subagent's activity entry and never into the parent transcript — // the same structural rule @threadplane/langgraph applies to namespaced - // events. Scope: text + tool events (what our emitters produce). Reasoning/ - // step attribution is deliberately not routed yet (YAGNI). + // events. Scope: text, tool and reasoning message events. Step attribution + // is not routed yet. const subagentRunId = (event as { subagentRunId?: string }).subagentRunId; if (subagentRunId && SUBAGENT_ROUTED_TYPES.has(event.type as string)) { routeSubagentContentEvent(subagentRunId, event, store); @@ -599,6 +599,8 @@ function randomId(): string { const SUBAGENT_ROUTED_TYPES = new Set([ 'TEXT_MESSAGE_START', 'TEXT_MESSAGE_CONTENT', 'TEXT_MESSAGE_END', 'TOOL_CALL_START', 'TOOL_CALL_ARGS', 'TOOL_CALL_END', 'TOOL_CALL_RESULT', + 'REASONING_MESSAGE_START', 'REASONING_MESSAGE_CONTENT', + 'REASONING_MESSAGE_CHUNK', 'REASONING_MESSAGE_END', ]); function subagentArgsBufferKey(subagentRunId: string, toolCallId: string): string { @@ -638,6 +640,17 @@ function ensureSubagentEntry(subagentRunId: string, store: ReducerStore): Activi * TOOL_CALL cases above, but written against the entry's content record * instead of store.messages/store.toolCalls. */ function routeSubagentContentEvent(subagentRunId: string, event: BaseEvent, store: ReducerStore): void { + if (event.type === 'REASONING_MESSAGE_START' || event.type === 'REASONING_MESSAGE_CONTENT' + || event.type === 'REASONING_MESSAGE_CHUNK' || event.type === 'REASONING_MESSAGE_END') { + if (!store.deliveryRun || store.deliveryRun.outcome !== undefined) return; + // END carries no child timing or lifecycle effect, even before STARTED. + // Reject inert/terminal events before currentRunForEvent can bind a run ID. + if (event.type === 'REASONING_MESSAGE_END') return; + const existing = store.activities().get(subagentRunId); + const status = existing?.content()['status']; + if (existing?.activityType === 'subagent' && (status === 'complete' || status === 'error')) return; + if (!currentRunForEvent(event, store)) return; + } const entry = ensureSubagentEntry(subagentRunId, store); // buffer-not-drop: creates on first sight const e = event as unknown as Record; @@ -664,6 +677,16 @@ function routeSubagentContentEvent(subagentRunId: string, event: BaseEvent, stor const messages = [...((c['messages'] as Array>) ?? [])]; const toolCalls = [...((c['toolCalls'] as Array>) ?? [])]; switch (event.type as string) { + case 'REASONING_MESSAGE_START': + case 'REASONING_MESSAGE_CONTENT': + case 'REASONING_MESSAGE_CHUNK': { + const id = e['messageId'] as string; + const idx = messages.findIndex((m) => m['id'] === id); + const delta = event.type === 'REASONING_MESSAGE_START' ? '' : (e['delta'] as string) ?? ''; + if (idx < 0) messages.push({ id, role: 'assistant', content: '', reasoning: delta }); + else messages[idx] = { ...messages[idx], reasoning: `${messages[idx]['reasoning'] ?? ''}${delta}` }; + return { ...c, messages }; + } case 'TEXT_MESSAGE_START': { const id = e['messageId'] as string; if (!messages.some((m) => m['id'] === id)) messages.push({ id, role: 'assistant', content: '' }); diff --git a/libs/ag-ui/src/lib/to-agent.child-reasoning.spec.ts b/libs/ag-ui/src/lib/to-agent.child-reasoning.spec.ts new file mode 100644 index 000000000..438d867ef --- /dev/null +++ b/libs/ag-ui/src/lib/to-agent.child-reasoning.spec.ts @@ -0,0 +1,160 @@ +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'; +import type { AbstractAgent, BaseEvent } from '@ag-ui/client'; +import { toAgent } from './to-agent'; + +type Subscriber = Parameters[0]; + +/** Exercises the real adapter with controlled source subscriber callbacks. */ +class ChildReasoningSource { + state = {}; + subscriber!: Subscriber; + release!: () => void; + subscribe(subscriber: Subscriber) { + this.subscriber = subscriber; + return { unsubscribe: () => undefined }; + } + runAgent() { + return new Promise(resolve => { + this.release = () => resolve({ result: undefined, newMessages: [] }); + }); + } + abortRun() { /* Transport completion is released independently. */ } + emit(type: string, fields: Record = {}) { + this.subscriber.onEvent?.({ + event: { type, runId: 'current', messageId: 'shared', ...fields } as BaseEvent, + input: { runId: 'current' }, + } as Parameters>[0]); + } + reason(type: string, fields: Record = {}) { + this.emit(`REASONING_MESSAGE_${type}`, { subagentRunId: 'child', ...fields }); + } +} + +describe('toAgent child reasoning through source callbacks', () => { + let now = 100; + const agents: ReturnType[] = []; + beforeEach(() => { + now = 100; + vi.spyOn(Date, 'now').mockImplementation(() => now); + }); + afterEach(() => { + for (const agent of agents.splice(0)) agent.dispose(); + vi.restoreAllMocks(); + }); + + function begin() { + const source = new ChildReasoningSource(); + const agent = toAgent(source as unknown as AbstractAgent, { telemetry: false }); + agents.push(agent); + const completion = agent.submit({}); + source.emit('RUN_STARTED'); + return { source, agent, completion }; + } + + it.each(['shared', 'different'])('projects child reasoning %s with stable wrappers and isolated parent/sibling messages', async messageId => { + const { source, agent, completion } = begin(); + source.emit('REASONING_MESSAGE_START'); + source.emit('REASONING_MESSAGE_CONTENT', { delta: 'parent' }); + source.emit('SUBAGENT_STARTED', { subagentRunId: 'child', name: 'researcher' }); + source.emit('SUBAGENT_STARTED', { subagentRunId: 'sibling', name: 'reviewer' }); + source.reason('CONTENT', { subagentRunId: 'sibling', messageId, delta: 'sibling' }); + const parent = agent.messages(); + const tools = agent.toolCalls(); + const child = agent.subagents().get('child')!; + const sibling = agent.subagents().get('sibling')!; + const siblingMessages = sibling.messages(); + now = 200; + source.reason('START', { messageId }); + source.reason('CONTENT', { messageId, delta: 'child\n' }); + source.reason('CHUNK', { messageId, delta: ' thought' }); + source.reason('END', { messageId }); + source.emit('TEXT_MESSAGE_START', { subagentRunId: 'child', messageId }); + source.emit('TEXT_MESSAGE_CONTENT', { subagentRunId: 'child', messageId, delta: 'answer' }); + source.emit('TOOL_CALL_START', { subagentRunId: 'child', parentMessageId: messageId, toolCallId: 'child-tool', toolCallName: 'lookup' }); + source.reason('START', { messageId }); + expect(agent.subagents().get('child')).toBe(child); + expect(agent.subagents().get('sibling')).toBe(sibling); + expect(sibling.messages()).toBe(siblingMessages); + expect(siblingMessages).toMatchObject([{ id: messageId, reasoning: 'sibling', content: '' }]); + expect(agent.messages()).toBe(parent); + expect(agent.toolCalls()).toBe(tools); + expect(agent.clientTools.pending()).toEqual([]); + expect(child.messages()).toMatchObject([{ id: messageId, reasoning: 'child\n thought', content: 'answer', toolCallIds: ['child-tool'], delivery: { phase: 'streaming' } }]); + expect(child.messages()).toHaveLength(1); + expect(child.messages()[0].reasoningDurationMs).toBeUndefined(); + expect(child.status()).toBe('running'); + now = 400; + source.emit('REASONING_MESSAGE_END'); + expect(agent.messages()[0].reasoningDurationMs).toBe(300); + source.emit('RUN_FINISHED'); + source.release(); + await completion; + }); + + it('refreshes buffered child identity once and retains reasoning on re-announcement', async () => { + const { source, agent, completion } = begin(); + source.reason('CHUNK', { delta: 'early' }); + const buffered = agent.subagents().get('child'); + expect(buffered?.messages()).toMatchObject([{ reasoning: 'early' }]); + source.emit('SUBAGENT_STARTED', { subagentRunId: 'child', name: 'researcher', parentToolCallId: 'parent-tool' }); + const announced = agent.subagents().get('child')!; + expect(announced).not.toBe(buffered); + expect(announced.name).toBe('researcher'); + expect(announced.toolCallId).toBe('parent-tool'); + source.reason('CONTENT', { delta: ' later' }); + source.emit('SUBAGENT_STARTED', { subagentRunId: 'child', name: 'researcher', parentToolCallId: 'parent-tool' }); + expect(agent.subagents().get('child')).toBe(announced); + expect(announced.messages()).toMatchObject([{ reasoning: 'early later' }]); + expect(agent.messages()).toEqual([]); + source.emit('RUN_FINISHED'); + source.release(); + await completion; + }); + + it.each(['settled', 'stop', 'dispose'])('retains child and parent references after %s despite late callbacks', async operation => { + const { source, agent, completion } = begin(); + source.emit('REASONING_MESSAGE_START'); + source.reason('CONTENT', { delta: 'child' }); + const child = agent.subagents().get('child'); + expect(child?.messages()).toMatchObject([{ reasoning: 'child' }]); + if (operation === 'settled') source.emit('RUN_FINISHED'); + else if (operation === 'stop') await agent.stop(); + else agent.dispose(); + const parent = agent.messages(); + const children = agent.subagents(); + const childMessages = child!.messages(); + for (const type of ['START', 'CONTENT', 'CHUNK', 'END']) { + source.reason(type, { delta: 'late' }); + source.reason(type, { subagentRunId: 'unseen', delta: 'late' }); + } + expect(agent.messages()).toBe(parent); + expect(agent.subagents()).toBe(children); + expect(child!.messages()).toBe(childMessages); + expect(agent.subagents().has('unseen')).toBe(false); + source.release(); + await completion; + }); + + it.each(['START', 'CONTENT', 'CHUNK', 'END'])('rejects a foreign child %s body with a current callback envelope', async type => { + const { source, agent, completion } = begin(); + source.emit('REASONING_MESSAGE_START'); + source.reason('CONTENT', { delta: 'child' }); + const parent = agent.messages(); + const children = agent.subagents(); + const child = children.get('child'); + expect(child?.messages()).toMatchObject([{ reasoning: 'child' }]); + const messages = child!.messages(); + now = 200; + source.reason(type, { runId: 'foreign', delta: 'foreign' }); + source.reason(type, { runId: 'foreign', subagentRunId: 'unseen', delta: 'foreign' }); + expect(agent.messages()).toBe(parent); + expect(agent.subagents()).toBe(children); + expect(child!.messages()).toBe(messages); + now = 400; + source.emit('REASONING_MESSAGE_END'); + expect(agent.messages()[0].reasoningDurationMs).toBe(300); + source.emit('RUN_FINISHED'); + source.release(); + await completion; + }); +}); diff --git a/scripts/react-parity/baseline.json b/scripts/react-parity/baseline.json index 4120e1146..e13a78709 100644 --- a/scripts/react-parity/baseline.json +++ b/scripts/react-parity/baseline.json @@ -10573,7 +10573,7 @@ "id": "source:libs/ag-ui/src/lib/reducer.ts", "kind": "source", "path": "libs/ag-ui/src/lib/reducer.ts", - "sha256": "b179449025f05f653d44505a491fbe05a04ae2dd7625bc2cd45fe9f2c3924863" + "sha256": "d9ac6e09a4faa7e1aa4b69a6e1bf1a543b3a2485a699ab2c74ba398952978fde" }, { "id": "source:libs/ag-ui/src/lib/run-state-transaction.ts", From 3a50f98a33dbc83cda3259840329a6c921e08ced Mon Sep 17 00:00:00 2001 From: Brian Love Date: Thu, 24 Sep 2026 13:24:53 -0700 Subject: [PATCH 2/2] test(ag-ui): verify request ownership through real HTTP streams --- .../src/lib/to-agent.http-lifecycle.spec.ts | 403 ++++++++++++++++++ 1 file changed, 403 insertions(+) create mode 100644 libs/ag-ui/src/lib/to-agent.http-lifecycle.spec.ts diff --git a/libs/ag-ui/src/lib/to-agent.http-lifecycle.spec.ts b/libs/ag-ui/src/lib/to-agent.http-lifecycle.spec.ts new file mode 100644 index 000000000..3f896b5da --- /dev/null +++ b/libs/ag-ui/src/lib/to-agent.http-lifecycle.spec.ts @@ -0,0 +1,403 @@ +import { afterEach, describe, expect, it, vi } from 'vitest'; +import { createServer, type ServerResponse } from 'node:http'; +import type { Socket } from 'node:net'; +import { + EventType, HttpAgent, type RunAgentInput, type RunStartedEvent, type RunFinishedEvent, + type RunErrorEvent, type TextMessageStartEvent, type TextMessageContentEvent, type TextMessageEndEvent, + type ReasoningMessageStartEvent, type ReasoningMessageContentEvent, type ReasoningMessageEndEvent, + type SubagentStartedEvent, type SubagentFinishedEvent, +} from '@ag-ui/client'; +import { completeDelivery, type Message } from '@threadplane/chat'; +import { toAgent } from './to-agent'; + +const DEADLINE_MS = 2_000; + +function deferred() { + let resolve!: (value: T) => void; + const promise = new Promise((done) => { resolve = done; }); + return { promise, resolve }; +} + +async function within(promise: Promise, milestone: string): Promise { + let timer: ReturnType | undefined; + try { + return await Promise.race([ + promise, + new Promise((_, reject) => { + timer = setTimeout(() => reject(new Error(`Timed out: ${milestone}`)), DEADLINE_MS); + }), + ]); + } finally { + clearTimeout(timer); + } +} + +interface HttpRun { + method: string | undefined; + input: RunAgentInput; + response: ServerResponse; + closed: ReturnType>; + isClosed: boolean; +} + +type SubmitResult = { status: 'fulfilled' } | { status: 'rejected'; reason: unknown }; +const cleanups: Array<() => Promise> = []; + +afterEach(async () => { + for (const cleanup of cleanups.splice(0)) await cleanup(); +}); + +/** Real loopback HTTP and the locked SDK decoder; fetch only captures its input. */ +async function httpFixture() { + const errors: unknown[] = []; + const sockets = new Set(); + const responses = new Set(); + const requests: HttpRun[] = []; + const arrivals: Array>> = []; + const contentEvents = new Map>>(); + const fetchCalls: RequestInit[] = []; + const submissions: Promise[] = []; + const cleanupState: { agent?: ReturnType; unsubscribe?: () => void } = {}; + + const server = createServer((request, response) => { + responses.add(response); + const closed = deferred(); + let run: HttpRun | undefined; + response.on('error', (error) => errors.push(error)); + response.once('close', () => { + if (run) run.isClosed = true; + responses.delete(response); + closed.resolve(); + }); + request.on('error', (error) => errors.push(error)); + const chunks: Buffer[] = []; + request.on('data', (chunk: Buffer) => chunks.push(chunk)); + request.once('end', () => { + try { + const input = JSON.parse(Buffer.concat(chunks).toString('utf8')) as RunAgentInput; + run = { method: request.method, input, response, closed, isClosed: false }; + requests.push(run); + response.writeHead(200, { 'content-type': 'text/event-stream' }); + response.flushHeaders(); + arrivals[requests.length - 1]?.resolve(run); + } catch (error) { + errors.push(error); + response.destroy(); + } + }); + }); + server.on('error', (error) => errors.push(error)); + server.on('clientError', (error, socket) => { + errors.push(error); + socket.destroy(); + }); + server.on('connection', (socket) => { + sockets.add(socket); + socket.once('close', () => sockets.delete(socket)); + }); + + // Registered before listening, so assertion/setup failures also release resources. + cleanups.push(async () => { + cleanupState.agent?.dispose(); + cleanupState.unsubscribe?.(); + for (const response of responses) response.destroy(); + const stopped = new Promise((resolve, reject) => { + if (!server.listening) { resolve(); return; } + server.close((error) => error ? reject(error) : resolve()); + }); + for (const socket of sockets) socket.destroy(); + await within(Promise.all([stopped, ...submissions]), 'HTTP fixture teardown'); + expect(errors, 'unexpected local HTTP server errors').toEqual([]); + }); + + await within(new Promise((resolve, reject) => { + server.once('error', reject); + server.listen(0, '127.0.0.1', () => { + server.off('error', reject); + resolve(); + }); + }), 'HTTP server listening'); + const address = server.address(); + if (!address || typeof address === 'string') throw new Error('Expected a loopback TCP address'); + const realFetch = globalThis.fetch; + const source = new HttpAgent({ + url: `http://127.0.0.1:${address.port}/agent`, + threadId: 'http-lifecycle-thread', + fetch: (url, init) => { + fetchCalls.push(init); + return realFetch(url, init); + }, + }); + const subscription = source.subscribe({ + onTextMessageContentEvent({ event }) { + contentEvents.get(event.messageId)?.resolve(); + }, + }); + cleanupState.unsubscribe = () => subscription.unsubscribe(); + const adapter = toAgent(source, { telemetry: false }); + cleanupState.agent = adapter; + + return { + agent: adapter, + fetchCalls, + requests, + submit(message: string) { + // Attach both handlers immediately, including when a later assertion fails. + const result = adapter.submit({ message }).then( + () => ({ status: 'fulfilled' }), + (reason: unknown) => ({ status: 'rejected', reason }), + ); + submissions.push(result); + return result; + }, + async request(index = 0) { + if (requests[index]) return requests[index]; + const arrival = arrivals[index] ??= deferred(); + return within(arrival.promise, `HTTP request ${index + 1}`); + }, + async partial(run: HttpRun, messageId: string, text: string) { + const observed = deferred(); + contentEvents.set(messageId, observed); + write(run, { type: EventType.RUN_STARTED, threadId: run.input.threadId, runId: run.input.runId }); + write(run, { type: EventType.TEXT_MESSAGE_START, messageId, role: 'assistant' }); + write(run, { type: EventType.TEXT_MESSAGE_CONTENT, messageId, delta: text }); + await within(observed.promise, `SDK text event for ${messageId}`); + // onEvent can precede the adapter's asynchronous reduction. Cancellation + // needs the visible partial answer, not just decoder notification. + await vi.waitFor(() => { + expect(answer(adapter, messageId)).toMatchObject({ + content: text, delivery: { phase: 'streaming' }, + }); + }, { timeout: DEADLINE_MS, interval: 10 }); + }, + }; +} + +type StreamEvent = RunStartedEvent | RunFinishedEvent | RunErrorEvent + | TextMessageStartEvent | TextMessageContentEvent | TextMessageEndEvent + | ReasoningMessageStartEvent | ReasoningMessageContentEvent | ReasoningMessageEndEvent + | SubagentStartedEvent | SubagentFinishedEvent; + +function write(run: HttpRun, event: StreamEvent): void { + run.response.write(`data: ${JSON.stringify(event)}\n\n`); +} + +function finish(run: HttpRun, messageId: string): void { + write(run, { type: EventType.TEXT_MESSAGE_END, messageId }); + write(run, { type: EventType.RUN_FINISHED, threadId: run.input.threadId, runId: run.input.runId }); + run.response.end(); +} + +function answer(agent: ReturnType, id: string): Message { + const message = agent.messages().find((message) => message.id === id); + if (!message) throw new Error(`Expected visible answer ${id}`); + return message; +} + +function expectOpen(run: HttpRun, signal: AbortSignal | null | undefined): asserts signal is AbortSignal { + expect(signal).toBeDefined(); + expect(signal?.aborted).toBe(false); + expect(run.isClosed).toBe(false); + expect(run.response.destroyed).toBe(false); + expect(run.response.writableEnded).toBe(false); +} + +describe('toAgent real HTTP lifecycle', () => { + it('starts one POST only on submit and completes the exact streamed answer', async () => { + const fixture = await httpFixture(); + const { agent } = fixture; + expect(agent.messages()).toEqual([]); + expect(agent.status()).toBe('idle'); + expect(agent.isLoading()).toBe(false); + expect(agent.error()).toBeUndefined(); + expect(agent.state()).toEqual({}); + await agent.ready; + expect(fixture.fetchCalls).toHaveLength(0); + expect(fixture.requests).toHaveLength(0); + + const submitted = fixture.submit('Hello'); + const run = await fixture.request(); + expect(run.method).toBe('POST'); + expect(run.input.threadId).toBe('http-lifecycle-thread'); + expect(run.input.runId).toEqual(expect.any(String)); + expect(run.input.runId.length).toBeGreaterThan(0); + expect(run.input.messages).toEqual([expect.objectContaining({ role: 'user', content: 'Hello' })]); + await fixture.partial(run, 'answer', 'Hello '); + const generation = answer(agent, 'answer').delivery.generation; + write(run, { type: EventType.TEXT_MESSAGE_CONTENT, messageId: 'answer', delta: 'world.' }); + finish(run, 'answer'); + + expect(await within(submitted, 'successful submit settlement')).toEqual({ status: 'fulfilled' }); + expect(answer(agent, 'answer')).toMatchObject({ + role: 'assistant', content: 'Hello world.', delivery: completeDelivery(generation, 'success'), + }); + expect(agent.status()).toBe('idle'); + expect(agent.isLoading()).toBe(false); + expect(agent.error()).toBeUndefined(); + expect(fixture.fetchCalls).toHaveLength(1); + expect(fixture.requests).toHaveLength(1); + }); + + it('retains interrupted delivery when HTTP ends gracefully without a terminal event', async () => { + const fixture = await httpFixture(); + const submitted = fixture.submit('Continue'); + const run = await fixture.request(); + await fixture.partial(run, 'partial', 'Unfinished answer'); + const generation = answer(fixture.agent, 'partial').delivery.generation; + expectOpen(run, fixture.fetchCalls[0].signal); + run.response.end(); + + expect(await within(submitted, 'unterminated submit settlement')).toEqual({ status: 'fulfilled' }); + expect(answer(fixture.agent, 'partial')).toMatchObject({ + content: 'Unfinished answer', delivery: completeDelivery(generation, 'interrupted'), + }); + expect(fixture.agent.error()?.kind).toBe('interrupted'); + expect(fixture.agent.status()).toBe('error'); + expect(fixture.agent.isLoading()).toBe(false); + }); + + it('keeps parent and two child reasoning transcripts separate through actual SSE decoding', async () => { + const fixture = await httpFixture(); + const { agent } = fixture; + const submitted = fixture.submit('Delegate to two children'); + const run = await fixture.request(); + const children = [ + { id: 'child-one', reasoning: 'First child reasoning' }, + { id: 'child-two', reasoning: 'Second child reasoning' }, + ]; + write(run, { type: EventType.RUN_STARTED, threadId: run.input.threadId, runId: run.input.runId }); + write(run, { type: EventType.REASONING_MESSAGE_START, messageId: 'parent-message', role: 'reasoning' }); + write(run, { type: EventType.REASONING_MESSAGE_CONTENT, messageId: 'parent-message', delta: 'Parent reasoning' }); + // Distinct wire IDs keep the stream valid for the locked SDK. Each child + // reuses its own message ID for reasoning and the subsequent answer. + for (const { id, reasoning } of children) { + const messageId = `${id}-message`; + write(run, { type: EventType.SUBAGENT_STARTED, subagentRunId: id, name: id }); + write(run, { type: EventType.REASONING_MESSAGE_START, subagentRunId: id, messageId, role: 'reasoning' }); + write(run, { type: EventType.REASONING_MESSAGE_CONTENT, subagentRunId: id, messageId, delta: reasoning }); + write(run, { type: EventType.REASONING_MESSAGE_END, subagentRunId: id, messageId }); + } + for (const { id } of children) { + const messageId = `${id}-message`; + write(run, { type: EventType.TEXT_MESSAGE_START, subagentRunId: id, messageId, role: 'assistant' }); + write(run, { type: EventType.TEXT_MESSAGE_CONTENT, subagentRunId: id, messageId, delta: `${id} answer` }); + write(run, { type: EventType.TEXT_MESSAGE_END, subagentRunId: id, messageId }); + write(run, { type: EventType.SUBAGENT_FINISHED, subagentRunId: id, outcome: { type: 'success' } }); + } + write(run, { type: EventType.REASONING_MESSAGE_END, messageId: 'parent-message' }); + write(run, { type: EventType.TEXT_MESSAGE_START, messageId: 'parent-message', role: 'assistant' }); + write(run, { type: EventType.TEXT_MESSAGE_CONTENT, messageId: 'parent-message', delta: 'Parent answer' }); + finish(run, 'parent-message'); + + expect(await within(submitted, 'parent and child submit settlement')).toEqual({ status: 'fulfilled' }); + expect(agent.messages().filter((message) => message.role === 'assistant')).toEqual([ + expect.objectContaining({ + id: 'parent-message', content: 'Parent answer', reasoning: 'Parent reasoning', + delivery: completeDelivery(answer(agent, 'parent-message').delivery.generation, 'success'), + }), + ]); + expect([...agent.subagents().keys()]).toEqual(['child-one', 'child-two']); + for (const { id, reasoning } of children) { + const child = agent.subagents().get(id); + expect(child?.status()).toBe('complete'); + // Exactly one child slot contains both the reasoning and its answer. + expect(child?.messages()).toEqual([ + expect.objectContaining({ + id: `${id}-message`, role: 'assistant', content: `${id} answer`, reasoning, + delivery: expect.objectContaining({ phase: 'complete', outcome: 'success' }), + }), + ]); + } + expect(agent.status()).toBe('idle'); + expect(agent.isLoading()).toBe(false); + expect(agent.error()).toBeUndefined(); + expect(run.method).toBe('POST'); + expect(fixture.fetchCalls).toHaveLength(1); + expect(fixture.requests).toHaveLength(1); + }); + + it('retains the partial answer as an error after RUN_ERROR', async () => { + const fixture = await httpFixture(); + const submitted = fixture.submit('Continue'); + const run = await fixture.request(); + await fixture.partial(run, 'partial', 'Before failure'); + const generation = answer(fixture.agent, 'partial').delivery.generation; + write(run, { type: EventType.RUN_ERROR, message: 'Server could not finish', code: 'TEST_FAILURE' }); + run.response.end(); + + expect(await within(submitted, 'errored submit settlement')).toEqual({ status: 'fulfilled' }); + expect(answer(fixture.agent, 'partial')).toMatchObject({ + content: 'Before failure', delivery: completeDelivery(generation, 'error'), + }); + expect(fixture.agent.error()?.message).toBe('Server could not finish'); + expect(fixture.agent.status()).toBe('error'); + expect(fixture.agent.isLoading()).toBe(false); + }); + + it.each(['stop', 'dispose'] as const)('%s aborts the actual open HTTP request before teardown', async (action) => { + const fixture = await httpFixture(); + const { agent } = fixture; + const submitted = fixture.submit('Keep streaming'); + const run = await fixture.request(); + await fixture.partial(run, 'partial', 'Kept partial answer'); + const generation = answer(agent, 'partial').delivery.generation; + const signal = fixture.fetchCalls[0].signal; + expectOpen(run, signal); + + const cancelled = agent[action](); + expect(signal.aborted, `${action} must abort the fetch signal synchronously`).toBe(true); + await within(Promise.resolve(cancelled), `${action} settlement`); + // This is observed before cleanup can destroy any response or socket. + await within(run.closed.promise, `${action} must close the server response`); + expect(run.isClosed).toBe(true); + expect(run.response.writableEnded).toBe(false); + expect(await within(submitted, 'aborted submit settlement')).toEqual({ status: 'fulfilled' }); + expect(answer(agent, 'partial')).toMatchObject({ + content: 'Kept partial answer', delivery: completeDelivery(generation, 'aborted'), + }); + expect(agent.error()).toBeUndefined(); + expect(agent.status()).toBe('idle'); + expect(agent.isLoading()).toBe(false); + + if (action === 'dispose') { + agent.dispose(); + const rejected = await within(fixture.submit('After disposal'), 'disposed submit rejection'); + expect(rejected).toMatchObject({ status: 'rejected', reason: new Error('Agent has been disposed') }); + } + expect(fixture.fetchCalls).toHaveLength(1); + expect(fixture.requests).toHaveLength(1); + }); + + it('stopping a second protocol run preserves the first successful answer and delivery', async () => { + const fixture = await httpFixture(); + const { agent } = fixture; + const firstSubmit = fixture.submit('First'); + const firstRun = await fixture.request(); + await fixture.partial(firstRun, 'first-answer', 'Completed first answer'); + finish(firstRun, 'first-answer'); + expect(await within(firstSubmit, 'first submit settlement')).toEqual({ status: 'fulfilled' }); + const firstAnswer = structuredClone(answer(agent, 'first-answer')); + expect(firstAnswer.delivery).toEqual(completeDelivery(firstAnswer.delivery.generation, 'success')); + + const secondSubmit = fixture.submit('Second'); + const secondRun = await fixture.request(1); + expect(secondRun.input.runId).not.toBe(firstRun.input.runId); + expect(secondRun.input.threadId).toBe(firstRun.input.threadId); + await fixture.partial(secondRun, 'second-answer', 'Partial second answer'); + const generation = answer(agent, 'second-answer').delivery.generation; + const signal = fixture.fetchCalls[1].signal; + expectOpen(secondRun, signal); + const stopped = agent.stop(); + expect(signal.aborted).toBe(true); + await within(stopped, 'second stop settlement'); + await within(secondRun.closed.promise, 'second server response closed before teardown'); + expect(await within(secondSubmit, 'second submit settlement')).toEqual({ status: 'fulfilled' }); + expect(answer(agent, 'first-answer')).toEqual(firstAnswer); + expect(answer(agent, 'second-answer')).toMatchObject({ + content: 'Partial second answer', delivery: completeDelivery(generation, 'aborted'), + }); + expect(agent.error()).toBeUndefined(); + expect(fixture.fetchCalls).toHaveLength(2); + expect(fixture.requests).toHaveLength(2); + }); +});