From 900f1e60be7edaf242af7dc7f634eac6e54b1f7f Mon Sep 17 00:00:00 2001 From: Brian Love Date: Thu, 24 Sep 2026 12:45:30 -0700 Subject: [PATCH] fix(ag-ui): scope reasoning timing to its delivery run --- libs/ag-ui/src/lib/reducer.spec.ts | 173 +++++++++++++++++- libs/ag-ui/src/lib/reducer.ts | 60 +++--- .../lib/to-agent.reasoning-ownership.spec.ts | 120 ++++++++++++ scripts/react-parity/baseline.json | 14 +- 4 files changed, 320 insertions(+), 47 deletions(-) create mode 100644 libs/ag-ui/src/lib/to-agent.reasoning-ownership.spec.ts diff --git a/libs/ag-ui/src/lib/reducer.spec.ts b/libs/ag-ui/src/lib/reducer.spec.ts index 6a43088d5..0a665e2da 100644 --- a/libs/ag-ui/src/lib/reducer.spec.ts +++ b/libs/ag-ui/src/lib/reducer.spec.ts @@ -1,4 +1,4 @@ -import { describe, it, expect } from 'vitest'; +import { afterEach, beforeEach, describe, it, expect, vi } from 'vitest'; import { signal } from '@angular/core'; import { Subject } from 'rxjs'; import { @@ -11,7 +11,7 @@ import { type ToolCall, type AgentEvent, } from '@threadplane/chat'; -import { reduceEvent, type ReducerStore, type CustomStreamEvent, type ActivityEntry } from './reducer'; +import { finalizeDeliveryRun, reduceEvent, type ReducerStore, type CustomStreamEvent, type ActivityEntry } from './reducer'; interface TestDeliveryRun { generation: string; @@ -625,6 +625,175 @@ describe('reduceEvent — interrupt', () => { }); describe('reduceEvent — REASONING_MESSAGE_*', () => { + let now = 100; + beforeEach(() => { + now = 100; + vi.spyOn(Date, 'now').mockImplementation(() => now); + }); + afterEach(() => vi.restoreAllMocks()); + + function reasoning(store: ReducerStore, type: string, messageId = 'owned', runId?: string) { + reduceEvent({ type: `REASONING_MESSAGE_${type}`, messageId, runId, delta: 'thinking' } as any, store); + } + + it('consumes a start once and leaves duplicate END references unchanged', () => { + const store = makeStore(); + reasoning(store, 'START'); + now = 200; + reasoning(store, 'END'); + const messages = store.messages(); + expect(messages[0].reasoningDurationMs).toBe(100); + now = 400; + reasoning(store, 'END'); + expect(store.messages()).toBe(messages); + expect(store.messages()[0]).toBe(messages[0]); + }); + + it('does not inherit a same-id start from a previous run', () => { + const store = makeStore(); + reasoning(store, 'START'); + finalizeDeliveryRun(store, store.deliveryRun!, 'success'); + store.deliveryRun = makeStore('next-generation').deliveryRun; + reduceEvent({ type: 'TEXT_MESSAGE_START', messageId: 'owned' } as any, store); + const messages = store.messages(); + now = 400; + reasoning(store, 'END'); + expect(store.messages()).toBe(messages); + expect(messages[0].reasoningDurationMs).toBeUndefined(); + reasoning(store, 'START'); + now = 450; + reasoning(store, 'END'); + expect(store.messages()[0].reasoningDurationMs).toBe(50); + }); + + it('restarts the clock and clears stale duration while retaining reasoning text', () => { + const store = makeStore(); + reasoning(store, 'START'); + reasoning(store, 'CONTENT'); + now = 200; + reasoning(store, 'END'); + now = 300; + reasoning(store, 'START'); + expect(store.messages()[0].reasoningDurationMs).toBeUndefined(); + expect(store.messages()[0].reasoning).toBe('thinking'); + now = 350; + reasoning(store, 'START'); + now = 400; + reasoning(store, 'END'); + expect(store.messages()[0].reasoningDurationMs).toBe(50); + }); + + it.each(['success', 'error', 'paused', 'aborted', 'interrupted'] as const)( + 'clears pending timing on %s and ignores late reasoning events', outcome => { + const store = makeStore(); + reasoning(store, 'START'); + const run = store.deliveryRun!; + finalizeDeliveryRun(store, run, outcome); + const messages = store.messages(); + now = 500; + for (const type of ['END', 'START', 'CONTENT', 'CHUNK']) reasoning(store, type); + expect(store.messages()).toBe(messages); + expect(run).not.toHaveProperty('pendingReasoning'); + expect(finalizeDeliveryRun(store, run, outcome)).toBe(false); + }, + ); + + it.each(['TEXT_MESSAGE_START', 'TOOL_CALL_START', 'MESSAGES_SNAPSHOT'])( + 'discards the previous clock when %s moves assistant ownership', type => { + const store = makeStore(); + reasoning(store, 'START'); + reduceEvent({ type, messageId: 'next', parentMessageId: 'next', toolCallId: 'tool', + toolCallName: 'tool', messages: [{ id: 'next', role: 'assistant', content: 'canonical' }] } as any, store); + const messages = store.messages(); + now = 500; + reasoning(store, 'END'); + expect(store.messages()).toBe(messages); + expect(store.deliveryRun).not.toHaveProperty('pendingReasoning'); + expect(store.deliveryRun?.currentAssistantMessageId).toBe('next'); + }, + ); + + it.each(['START', 'CONTENT', 'CHUNK', 'END'])( + 'ignores foreign-run reasoning %s without disturbing the valid clock', type => { + const store = makeStore(); + reduceEvent({ type: 'RUN_STARTED', runId: 'current' } as any, store); + reasoning(store, 'START', 'owned', 'current'); + const messages = store.messages(); + now = 200; + reasoning(store, type, 'owned', 'foreign'); + expect(store.messages()).toBe(messages); + now = 300; + reasoning(store, 'END', 'owned', 'current'); + expect(store.messages()[0].reasoningDurationMs).toBe(200); + }, + ); + + it.each(['static', 'wrong-generation', 'completed', 'non-assistant', 'unowned', 'not-current', 'no-run'])( + 'does not let END mutate a %s message or claim ownership', scenario => { + const store = makeStore(); + reasoning(store, 'START'); + const message = store.messages()[0]; + if (scenario === 'static') store.messages.set([{ ...message, delivery: staticDelivery('owned') }]); + if (scenario === 'wrong-generation') store.messages.set([{ ...message, delivery: streamingDelivery('other') }]); + if (scenario === 'completed') store.messages.set([{ ...message, delivery: completeDelivery('run-generation-1', 'success') }]); + if (scenario === 'non-assistant') store.messages.set([{ ...message, role: 'user' }]); + if (scenario === 'unowned') store.deliveryRun!.ownedMessageIds.clear(); + if (scenario === 'not-current') store.deliveryRun!.currentAssistantMessageId = 'different'; + if (scenario === 'no-run') store.deliveryRun = null; + const messages = store.messages(); + now = 200; + reasoning(store, 'END'); + expect(store.messages()).toBe(messages); + }, + ); + + it('leaves an END without a local start inert even when another store has that ID', () => { + const other = makeStore(); + reasoning(other, 'START'); + const store = makeStore(); + reduceEvent({ type: 'TEXT_MESSAGE_START', messageId: 'owned' } as any, store); + const messages = store.messages(); + reasoning(store, 'END'); + expect(store.messages()).toBe(messages); + expect(store.deliveryRun?.currentAssistantMessageId).toBe('owned'); + }); + + it('keeps pending timing across a same-id canonical snapshot and keeps its replacement policy', () => { + const store = makeStore(); + reasoning(store, 'START'); + const snapshot = { type: 'MESSAGES_SNAPSHOT', messages: [{ id: 'owned', role: 'assistant', content: 'canonical' }] } as any; + reduceEvent(snapshot, store); + now = 200; + reasoning(store, 'END'); + expect(store.messages()[0].reasoningDurationMs).toBe(100); + expect(store.messages()[0].content).toBe('canonical'); + reduceEvent(snapshot, store); + expect(store.messages()[0].reasoningDurationMs).toBeUndefined(); + const messages = store.messages(); + reasoning(store, 'END'); + expect(store.messages()).toBe(messages); + }); + + it('keeps same-id reasoning clocks isolated between simultaneous stores', () => { + const first = makeStore('first'); + const second = makeStore('second'); + const clock = vi.spyOn(Date, 'now'); + try { + clock.mockReturnValue(100); + reduceEvent({ type: 'REASONING_MESSAGE_START', messageId: 'shared' } as any, first); + clock.mockReturnValue(200); + reduceEvent({ type: 'REASONING_MESSAGE_START', messageId: 'shared' } as any, second); + clock.mockReturnValue(300); + reduceEvent({ type: 'REASONING_MESSAGE_END', messageId: 'shared' } as any, first); + expect(first.messages()[0].reasoningDurationMs).toBe(200); + clock.mockReturnValue(450); + reduceEvent({ type: 'REASONING_MESSAGE_END', messageId: 'shared' } as any, second); + expect(second.messages()[0].reasoningDurationMs).toBe(250); + } finally { + clock.mockRestore(); + } + }); + it('REASONING_MESSAGE_START creates an assistant slot with empty reasoning', () => { const store = makeStore(); reduceEvent({ type: 'REASONING_MESSAGE_START', messageId: 'm1', role: 'assistant' } as any, store); diff --git a/libs/ag-ui/src/lib/reducer.ts b/libs/ag-ui/src/lib/reducer.ts index e96baef34..0fd3fd68c 100644 --- a/libs/ag-ui/src/lib/reducer.ts +++ b/libs/ag-ui/src/lib/reducer.ts @@ -67,6 +67,7 @@ export interface ReducerDeliveryRun { snapshotReplacementIds: Set; eligibleBaselineAssistantId?: string; currentAssistantMessageId?: string; + pendingReasoning?: { messageId: string; startedAt: number }; protocolRunId?: string; outcome?: CompleteOutcome; } @@ -101,6 +102,7 @@ export function finalizeDeliveryRun( ): boolean { if (run.outcome !== undefined) return false; run.outcome = outcome; + delete run.pendingReasoning; store.messages.update(messages => messages.map(message => message.delivery.generation === run.generation ? { ...message, delivery: completeDelivery(run.generation, outcome) } @@ -110,28 +112,9 @@ export function finalizeDeliveryRun( } /** - * Per-message reasoning timing. Populated by REASONING_MESSAGE_START / - * REASONING_MESSAGE_END handlers. The map lives on the module — same - * scope as the reducer function. ReducerStore stays free of timing - * state; consumers read it via `Message.reasoningDurationMs` on - * messages that completed reasoning. - * - * Keyed by messageId. We do not need cross-thread isolation here: - * AG-UI's source agent recreates the reducer pipeline per session, and - * messageIds are unique within a session. - */ -const reasoningTimingMap = new Map(); - -function resolveReasoningDurationMs(messageId: string): number | undefined { - const entry = reasoningTimingMap.get(messageId); - if (!entry || entry.endedAt === undefined) return undefined; - return entry.endedAt - entry.startedAt; -} - -/** - * Pure function: applies a single AG-UI BaseEvent to the store. Caller - * subscribes to source.agent() and forwards each event here. Designed - * for testability — no side effects beyond the supplied store. + * Applies a single AG-UI BaseEvent to the supplied signal store. The adapter + * forwards source subscriber events here; reasoning timing reads Date.now() + * at this imperative boundary and retains its pending start on the delivery run. */ export function reduceEvent(event: BaseEvent, store: ReducerStore): void { // A subagentRunId on a content event means the child produced it: route it @@ -208,15 +191,17 @@ export function reduceEvent(event: BaseEvent, store: ReducerStore): void { return; } case 'REASONING_MESSAGE_START': { + const run = currentRunForEvent(event, store); + if (!run || run.outcome !== undefined) return; const id = messageIdFrom(event); const delivery = ownAssistantMessage(store, id); if (!delivery) return; - reasoningTimingMap.set(id, { startedAt: Date.now() }); + run.pendingReasoning = { messageId: id, startedAt: Date.now() }; // Initialize an assistant slot with empty reasoning if it doesn't already exist. store.messages.update((prev) => prev.some((m) => m.id === id) ? prev.map((m) => m.id === id - ? { ...m, reasoning: m.reasoning ?? '', delivery } + ? { ...m, reasoning: m.reasoning ?? '', reasoningDurationMs: undefined, delivery } : m) : [...prev, { id, role: 'assistant', content: '', reasoning: '', delivery }], ); @@ -224,6 +209,8 @@ export function reduceEvent(event: BaseEvent, store: ReducerStore): void { } case 'REASONING_MESSAGE_CONTENT': case 'REASONING_MESSAGE_CHUNK': { + const run = currentRunForEvent(event, store); + if (!run || run.outcome !== undefined) return; const id = messageIdFrom(event); const delivery = ownAssistantMessage(store, id); if (!delivery) return; @@ -236,18 +223,20 @@ export function reduceEvent(event: BaseEvent, store: ReducerStore): void { return; } case 'REASONING_MESSAGE_END': { + const run = currentRunForEvent(event, store); + if (!run || run.outcome !== undefined) return; const id = messageIdFrom(event); - const entry = reasoningTimingMap.get(id); - if (entry) { - entry.endedAt = Date.now(); - reasoningTimingMap.set(id, entry); - const duration = resolveReasoningDurationMs(id); - if (duration !== undefined) { - store.messages.update((prev) => - prev.map((m) => m.id === id ? { ...m, reasoningDurationMs: duration } : m), - ); - } - } + const pending = run.pendingReasoning; + if (!pending || pending.messageId !== id + || run.currentAssistantMessageId !== id || !run.ownedMessageIds.has(id)) return; + const message = store.messages().find(m => m.id === id && m.role === 'assistant' + && m.delivery.generation === run.generation && m.delivery.phase === 'streaming'); + if (!message) return; + const duration = Date.now() - pending.startedAt; + delete run.pendingReasoning; + store.messages.update(messages => messages.map(m => + m === message ? { ...m, reasoningDurationMs: duration } : m, + )); return; } case 'TEXT_MESSAGE_CONTENT': { @@ -801,6 +790,7 @@ function ownAssistantMessage(store: ReducerStore, id: string) { const currentId = run.currentAssistantMessageId; if (currentId && currentId !== id) { if (run.ownedMessageIds.has(id)) return undefined; + delete run.pendingReasoning; store.messages.update(messages => messages.map(message => message.id === currentId && message.delivery.generation === run.generation diff --git a/libs/ag-ui/src/lib/to-agent.reasoning-ownership.spec.ts b/libs/ag-ui/src/lib/to-agent.reasoning-ownership.spec.ts new file mode 100644 index 000000000..7276e44ef --- /dev/null +++ b/libs/ag-ui/src/lib/to-agent.reasoning-ownership.spec.ts @@ -0,0 +1,120 @@ +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]; + +/** Holds real adapter subscriber callbacks while a source run remains in flight. */ +class ReasoningSource { + state = {}; + subscriber!: Subscriber; + release!: () => void; + callbackCount = 0; + subscribe(subscriber: Subscriber) { + this.subscriber = subscriber; + return { unsubscribe: () => undefined }; + } + runAgent() { + return new Promise(resolve => { + this.release = () => resolve({ result: undefined, newMessages: [] }); + }); + } + abortRun() { /* The test releases transport completion independently. */ } + emit(type: string, runId = 'current', messageId = 'shared') { + this.callbackCount++; + this.subscriber.onEvent?.({ + event: { type, runId, messageId, delta: 'thinking' } as BaseEvent, + input: { runId: 'current' }, + } as Parameters>[0]); + } +} + +describe('toAgent reasoning ownership 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 ReasoningSource(); + 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('isolates simultaneous instances with equal IDs and consumes END once', async () => { + const first = begin(); + const second = begin(); + first.source.emit('REASONING_MESSAGE_START'); + now = 200; + second.source.emit('REASONING_MESSAGE_START'); + now = 300; + first.source.emit('REASONING_MESSAGE_END'); + expect(first.agent.messages()[0].reasoningDurationMs).toBe(200); + const messages = first.agent.messages(); + now = 450; + first.source.emit('REASONING_MESSAGE_END'); + expect(first.agent.messages()).toBe(messages); + second.source.emit('REASONING_MESSAGE_END'); + expect(second.agent.messages()[0].reasoningDurationMs).toBe(250); + expect(first.source.callbackCount).toBe(4); + expect(second.source.callbackCount).toBe(3); + first.source.emit('RUN_FINISHED'); + second.source.emit('RUN_FINISHED'); + first.source.release(); + second.source.release(); + await Promise.all([first.completion, second.completion]); + }); + + it.each(['stop', 'dispose'] as const)('keeps late callbacks inert after %s and does not leak a same-ID start', async operation => { + const previous = begin(); + previous.source.emit('REASONING_MESSAGE_START'); + if (operation === 'stop') await previous.agent.stop(); + else previous.agent.dispose(); + const messages = previous.agent.messages(); + const current = begin(); + current.source.emit('TEXT_MESSAGE_START'); + const currentMessages = current.agent.messages(); + now = 300; + previous.source.emit('REASONING_MESSAGE_END'); + current.source.emit('REASONING_MESSAGE_END'); + expect(previous.agent.messages()).toBe(messages); + expect(current.agent.messages()).toBe(currentMessages); + expect(currentMessages[0].reasoningDurationMs).toBeUndefined(); + expect(previous.source.callbackCount).toBe(3); + current.source.emit('REASONING_MESSAGE_START'); + now = 400; + current.source.emit('REASONING_MESSAGE_END'); + expect(current.agent.messages()[0].reasoningDurationMs).toBe(100); + current.source.emit('RUN_FINISHED'); + previous.source.release(); + current.source.release(); + await Promise.all([previous.completion, current.completion]); + }); + + it.each(['START', 'CONTENT', 'CHUNK', 'END'])( + 'rejects a foreign %s body even when the callback envelope identifies the current run', async type => { + const { source, agent, completion } = begin(); + source.emit('REASONING_MESSAGE_START'); + const messages = agent.messages(); + now = 200; + source.emit(`REASONING_MESSAGE_${type}`, 'foreign'); + expect(agent.messages()).toBe(messages); + now = 300; + source.emit('REASONING_MESSAGE_END'); + expect(agent.messages()[0].reasoningDurationMs).toBe(200); + expect(source.callbackCount).toBe(4); + source.emit('RUN_FINISHED'); + source.release(); + await completion; + }, + ); +}); diff --git a/scripts/react-parity/baseline.json b/scripts/react-parity/baseline.json index 027950de4..69d273ecd 100644 --- a/scripts/react-parity/baseline.json +++ b/scripts/react-parity/baseline.json @@ -1,17 +1,11 @@ { "schemaVersion": 1, - "baselineHead": "4d47995333278f9a19491925bc555d30f8490956", + "baselineHead": "a4a0c7b9b5a364dc70336f59818ac486909bdfef", "sourceState": { "modified": [ - "libs/langgraph/src/runtime/README.md", - "libs/langgraph/src/runtime/history-projection.ts", - "libs/langgraph/src/runtime/message-reducer.ts", - "libs/langgraph/src/runtime/ownership.ts", - "libs/langgraph/src/runtime/stream-projection.ts" + "libs/ag-ui/src/lib/reducer.ts" ], - "untracked": [ - "libs/langgraph/src/runtime/citation-projection.ts" - ] + "untracked": [] }, "scope": { "libraries": [ @@ -10572,7 +10566,7 @@ "id": "source:libs/ag-ui/src/lib/reducer.ts", "kind": "source", "path": "libs/ag-ui/src/lib/reducer.ts", - "sha256": "82ed040383d718c9460b4d8040a942a0b5f969f3060c2522cfb4adfd23686379" + "sha256": "b179449025f05f653d44505a491fbe05a04ae2dd7625bc2cd45fe9f2c3924863" }, { "id": "source:libs/ag-ui/src/lib/run-state-transaction.ts",