diff --git a/libs/ag-ui/src/runtime/reconcile-transcript-http.spec.ts b/libs/ag-ui/src/runtime/reconcile-transcript-http.spec.ts new file mode 100644 index 000000000..9a6fd816f --- /dev/null +++ b/libs/ag-ui/src/runtime/reconcile-transcript-http.spec.ts @@ -0,0 +1,369 @@ +import { once } from 'node:events'; +import { createServer } from 'node:http'; +import type { AddressInfo } from 'node:net'; +import { + EventType, + type Message, + type MessagesSnapshotEvent, + type RunAgentInput, +} from '@ag-ui/client'; +import { describe, expect, it } from 'vitest'; +import { createRun, type RunHandle } from './create-run'; +import { reconcileTranscript } from './reconcile-transcript'; +import { ownTranscript, requestMessages } from './transcript'; + +function deferred() { + let resolve!: (value: T) => void; + const promise = new Promise((yes) => { + resolve = yes; + }); + return { promise, resolve }; +} + +async function bounded(promise: Promise): Promise { + let timer: ReturnType | undefined; + try { + return await Promise.race([ + promise, + new Promise((_, reject) => { + timer = setTimeout( + () => reject(new Error('Reconciliation HTTP milestone timed out')), + 1500 + ); + }), + ]); + } finally { + clearTimeout(timer); + } +} + +interface Exchange { + body: RunAgentInput; + closed: Promise; + send: (...events: unknown[]) => void; +} + +async function serve() { + const exchanges: Exchange[] = []; + const arrivals = new Map>>(); + const server = createServer(async (request, response) => { + const closed = deferred(); + response.on('close', () => closed.resolve()); + const chunks: Buffer[] = []; + for await (const chunk of request) chunks.push(Buffer.from(chunk)); + const exchange: Exchange = { + body: JSON.parse(Buffer.concat(chunks).toString()), + closed: closed.promise, + send: (...events) => { + for (const event of events) + response.write(`data: ${JSON.stringify(event)}\n\n`); + }, + }; + response.writeHead(200, { 'content-type': 'text/event-stream' }); + response.flushHeaders(); + exchanges.push(exchange); + arrivals.get(exchanges.length - 1)?.resolve(exchange); + // Held open so physical cancellation, rather than server EOF, closes it. + }); + server.listen(0, '127.0.0.1'); + await once(server, 'listening'); + return { + url: `http://127.0.0.1:${(server.address() as AddressInfo).port}/events`, + exchanges, + next: (index = 0) => { + if (exchanges[index]) return Promise.resolve(exchanges[index]); + const arrival = arrivals.get(index) ?? deferred(); + arrivals.set(index, arrival); + return bounded(arrival.promise); + }, + close: async () => { + const closed = once(server, 'close'); + server.close(); + server.closeAllConnections(); + await bounded(closed); + }, + }; +} + +function input(runId: string, messages: Message[] = []): RunAgentInput { + return { + threadId: 'thread', + runId, + messages, + state: {}, + tools: [], + context: [], + forwardedProps: {}, + }; +} + +function lifecycle( + exchange: Exchange, + type: EventType.RUN_STARTED | EventType.RUN_FINISHED +) { + return { type, threadId: exchange.body.threadId, runId: exchange.body.runId }; +} + +const corrected = { + id: 'assistant', + role: 'assistant', + name: 'assistant-name', + encryptedValue: 'message-encrypted', + subagentRunId: 'child', + metadata: { nested: [{ keep: 'message', subagentRunId: null }] }, + toolCalls: [ + { + id: 'call', + type: 'function', + function: { name: 'weather', arguments: '{ "city": "Paris" }\n' }, + encryptedValue: 'call-encrypted', + metadata: { nested: ['call', null] }, + }, + ], +} satisfies Message; + +const richSnapshot = [ + { + id: 'system', + role: 'system', + content: 'policy', + name: 'system-name', + metadata: { keep: ['system'] }, + }, + { + id: 'developer', + role: 'developer', + content: 'developer policy', + subagentRunId: '', + encryptedValue: 'developer-encrypted', + }, + { + id: 'user', + role: 'user', + content: [ + { type: 'text', text: 'inspect' }, + { + type: 'image', + source: { type: 'data', value: 'aW1hZ2U=', mimeType: 'image/png' }, + metadata: { nested: [{ label: 'image', nullable: null }] }, + }, + { + type: 'audio', + source: { + type: 'url', + value: 'https://example.test/audio', + mimeType: 'audio/wav', + }, + }, + { + type: 'video', + source: { type: 'data', value: 'dmlkZW8=', mimeType: 'video/mp4' }, + }, + { + type: 'document', + source: { type: 'url', value: 'https://example.test/document' }, + }, + { + type: 'binary', + mimeType: 'application/octet-stream', + data: 'YmluYXJ5', + id: 'binary', + filename: 'legacy.bin', + }, + ], + }, + corrected, + { + id: 'tool', + role: 'tool', + toolCallId: 'call', + content: 'result', + error: 'provider error', + encryptedValue: 'tool-encrypted', + metadata: { nested: { keep: true } }, + }, +] satisfies Message[]; + +function seed() { + // Seeded owned state proves snapshot reconciliation, not streaming reduction. + // Null attribution is local data; the SDK rejects it in wire snapshots. + return ownTranscript([ + { id: 'user', role: 'user', content: 'before' }, + { + id: 'assistant', + role: 'assistant', + content: 'obsolete', + metadata: { obsolete: true }, + toolCalls: [ + { + id: 'call', + type: 'function', + function: { name: 'weather', arguments: '{"city":' }, + }, + ], + }, + { + id: 'reasoning', + role: 'reasoning', + content: 'retained reasoning', + subagentRunId: null, + encryptedValue: 'reasoning-encrypted', + metadata: { subagentRunId: null }, + }, + { + id: 'activity', + role: 'activity', + activityType: 'progress', + content: { steps: [{ label: 'working' }] }, + metadata: { keep: 'activity' }, + }, + ] as unknown as Message[]); +} + +describe('private snapshot reconciliation over actual HTTP', () => { + it('corrects seeded tool data, retains omitted event-only records and sends exact reconciled egress in a second run', async () => { + const server = await serve(); + const factory = createRun({ url: server.url }); + const previous = seed(); + let transcript = previous; + const observed: Message[][] = []; + let second: RunHandle | undefined; + const first = factory.start(input('snapshot'), (event) => { + if (event.type !== EventType.MESSAGES_SNAPSHOT) return; + const messages = (event as MessagesSnapshotEvent).messages; + observed.push(messages); + transcript = reconcileTranscript(transcript, messages); + }); + try { + const initial = await server.next(); + // Snapshot order differs from retained order; repeat with reversed input + // to prove the existing transcript's positions survive another snapshot. + initial.send( + lifecycle(initial, EventType.RUN_STARTED), + { type: EventType.MESSAGES_SNAPSHOT, messages: richSnapshot }, + { + type: EventType.MESSAGES_SNAPSHOT, + messages: [...richSnapshot].reverse(), + }, + lifecycle(initial, EventType.RUN_FINISHED) + ); + expect(await bounded(first.done)).toEqual({ outcome: 'success' }); + await bounded(initial.closed); + expect(observed).toStrictEqual([ + richSnapshot, + [...richSnapshot].reverse(), + ]); + const expected = [ + richSnapshot[2], + corrected, + previous[2], + previous[3], + richSnapshot[0], + richSnapshot[1], + richSnapshot[4], + ]; + expect(transcript).toStrictEqual(expected); + expect(transcript[2]).toBe(previous[2]); + expect(transcript[3]).toBe(previous[3]); + expect(transcript[1]).not.toBe(observed[1][1]); + expect(transcript[1]).not.toHaveProperty('content'); + expect(transcript[1].metadata).not.toHaveProperty('obsolete'); + expect(previous).toStrictEqual(seed()); + expect(Object.isFrozen(transcript)).toBe(true); + const outgoing = [ + richSnapshot[2], + corrected, + { + id: 'reasoning', + role: 'reasoning', + content: 'retained reasoning', + encryptedValue: 'reasoning-encrypted', + metadata: { subagentRunId: null }, + }, + richSnapshot[0], + richSnapshot[1], + richSnapshot[4], + ] satisfies Message[]; + second = factory.start( + input('continue', requestMessages(transcript)), + () => undefined + ); + const continuation = await server.next(1); + expect(continuation.body).toStrictEqual(input('continue', outgoing)); + expect(continuation.body.messages.map((message) => message.role)).toEqual( + ['user', 'assistant', 'reasoning', 'system', 'developer', 'tool'] + ); + expect( + continuation.body.messages.find( + (message) => message.role === 'assistant' + )?.toolCalls?.[0].function.arguments + ).toBe('{ "city": "Paris" }\n'); + expect(transcript).toStrictEqual(expected); + continuation.send( + lifecycle(continuation, EventType.RUN_STARTED), + lifecycle(continuation, EventType.RUN_FINISHED) + ); + expect(await bounded(second.done)).toEqual({ outcome: 'success' }); + await bounded(continuation.closed); + expect(server.exchanges).toHaveLength(2); + } finally { + first.abort(); + second?.abort(); + await server.close(); + } + }); + + it('turns a duplicate snapshot callback failure into an error and closes HTTP without changing prior state or starting a follow-up', async () => { + const previous = seed(); + let transcript = previous; + let failure: unknown; + const delivered: string[] = []; + const server = await serve(); + const handle = createRun({ url: server.url }).start( + input('duplicate'), + (event) => { + delivered.push(event.type); + if (event.type !== EventType.MESSAGES_SNAPSHOT) return; + try { + transcript = reconcileTranscript( + transcript, + (event as MessagesSnapshotEvent).messages + ); + } catch (error) { + failure = error; + throw error; + } + } + ); + try { + const exchange = await server.next(); + exchange.send( + lifecycle(exchange, EventType.RUN_STARTED), + { + type: EventType.MESSAGES_SNAPSHOT, + messages: [ + corrected, + { id: corrected.id, role: 'user', content: 'duplicate' }, + ], + }, + lifecycle(exchange, EventType.RUN_FINISHED) + ); + const result = await bounded(handle.done); + expect(failure).toBeInstanceOf(TypeError); + expect(result).toEqual({ outcome: 'error', error: failure }); + if (result.outcome === 'error') expect(result.error).toBe(failure); + await bounded(exchange.closed); + expect(await handle.done).toBe(result); + expect(delivered).toEqual([ + EventType.RUN_STARTED, + EventType.MESSAGES_SNAPSHOT, + ]); + expect(transcript).toBe(previous); + expect(previous).toStrictEqual(seed()); + expect(server.exchanges).toHaveLength(1); + } finally { + handle.abort(); + await server.close(); + } + }); +}); diff --git a/libs/ag-ui/src/runtime/reconcile-transcript.spec.ts b/libs/ag-ui/src/runtime/reconcile-transcript.spec.ts new file mode 100644 index 000000000..6504722e3 --- /dev/null +++ b/libs/ag-ui/src/runtime/reconcile-transcript.spec.ts @@ -0,0 +1,439 @@ +import type { Message } from '@ag-ui/client'; +import { describe, expect, it } from 'vitest'; +import { reconcileTranscript } from './reconcile-transcript'; +import { ownTranscript, type Transcript } from './transcript'; + +const reasoning = { + id: 'reasoning', + role: 'reasoning', + content: 'thinking', +} satisfies Message; +const activity = { + id: 'activity', + role: 'activity', + activityType: 'progress', + content: { steps: ['working'] }, +} satisfies Message; +const user = (id: string, content = id): Message => ({ + id, + role: 'user', + content, +}); + +function richMessages() { + return [ + { + id: 'system', + role: 'system', + content: 'policy', + name: 'policy-name', + metadata: { nested: [null] }, + }, + { + id: 'developer', + role: 'developer', + content: 'developer', + encryptedValue: 'developer-encrypted', + }, + { + id: 'user', + role: 'user', + subagentRunId: '', + content: [ + { type: 'text', text: 'inspect' }, + { + type: 'image', + source: { type: 'data', value: 'aW1hZ2U=', mimeType: 'image/png' }, + metadata: { nested: [null] }, + }, + { + type: 'audio', + source: { + type: 'url', + value: 'https://example.test/audio', + mimeType: 'audio/wav', + }, + }, + { + type: 'video', + source: { type: 'data', value: 'dmlkZW8=', mimeType: 'video/mp4' }, + }, + { + type: 'document', + source: { type: 'url', value: 'https://example.test/document' }, + }, + { + type: 'binary', + data: 'YmluYXJ5', + mimeType: 'application/octet-stream', + filename: 'file.bin', + id: 'file', + }, + ], + }, + { + id: 'assistant', + role: 'assistant', + name: 'assistant-name', + subagentRunId: 'child', + encryptedValue: 'message-encrypted', + metadata: { nested: [{ value: 'corrected', subagentRunId: null }] }, + toolCalls: [ + { + id: 'call', + type: 'function', + function: { name: 'weather', arguments: '{ "city": "Paris" }\n' }, + encryptedValue: 'call-encrypted', + metadata: { nested: [null, 'call'] }, + }, + ], + }, + { + id: 'tool', + role: 'tool', + toolCallId: 'call', + content: 'result', + error: 'provider error', + encryptedValue: 'result-encrypted', + metadata: { nested: [null] }, + }, + { ...reasoning, encryptedValue: 'reasoning-encrypted' }, + activity, + ] satisfies Message[]; +} + +function expectFrozen(value: unknown): void { + if (value === null || typeof value !== 'object') return; + expect(Object.isFrozen(value)).toBe(true); + for (const child of Object.values(value)) expectFrozen(child); +} + +describe('private snapshot reconciliation', () => { + it('retains old positions and omitted event-only references while replacing values', () => { + const previous = ownTranscript([ + user('b', 'old b'), + reasoning, + user('a', 'old a'), + activity, + ]); + const snapshot = [user('a', 'new a'), user('b', 'new b')]; + const result = reconcileTranscript(previous, snapshot); + expect(result).toStrictEqual([ + snapshot[1], + reasoning, + snapshot[0], + activity, + ]); + expect(result[0]).not.toBe(previous[0]); + expect(result[0]).not.toBe(snapshot[1]); + expect(result[1]).toBe(previous[1]); + expect(result[3]).toBe(previous[3]); + expectFrozen(result); + }); + + it('drops missing ordinary records and appends new IDs in snapshot order once', () => { + const previous = ownTranscript([user('gone'), user('b'), reasoning]); + const snapshot = [user('new-2'), user('b', 'corrected'), user('new-1')]; + expect(reconcileTranscript(previous, snapshot)).toStrictEqual([ + snapshot[1], + reasoning, + snapshot[0], + snapshot[2], + ]); + }); + + it.each(['activity', 'reasoning'] as const)( + 'treats a supplied %s role as authoritative only for that role', + (role) => { + const previous = ownTranscript([reasoning, activity]); + const incoming = { + ...(role === 'activity' ? activity : reasoning), + id: 'new', + }; + const retained = role === 'activity' ? previous[0] : previous[1]; + const result = reconcileTranscript(previous, [incoming]); + expect(result).toStrictEqual([retained, incoming]); + expect(result[0]).toBe(retained); + } + ); + + it('drops omitted old event-only records when both role sets are supplied', () => { + const snapshot = [ + { ...activity, id: 'new-activity' }, + { ...reasoning, id: 'new-reasoning' }, + ]; + expect( + reconcileTranscript(ownTranscript([reasoning, activity]), snapshot) + ).toStrictEqual(snapshot); + }); + + it('keeps only event-only roles for an empty snapshot and reuses stable retained-only results', () => { + const previous = ownTranscript([user('gone'), reasoning, activity]); + const result = reconcileTranscript(previous, []); + expect(result).toStrictEqual([reasoning, activity]); + expect(result).not.toBe(previous); + expect(result[0]).toBe(previous[1]); + expect(result[1]).toBe(previous[2]); + expect(reconcileTranscript(result, [])).toBe(result); + const empty = ownTranscript([]); + expect(reconcileTranscript(empty, [])).toBe(empty); + }); + + it('preserves all seven rich roles and exact corrected argument bytes as complete replacement data', () => { + const snapshot = richMessages(); + const previous = ownTranscript( + snapshot.map((message) => + message.role === 'assistant' + ? { + ...message, + metadata: { obsolete: { keep: false } }, + toolCalls: [ + { + ...message.toolCalls[0], + function: { name: 'weather', arguments: '{"city":' }, + obsolete: true, + }, + ], + } + : message + ) + ); + const result = reconcileTranscript(previous, snapshot); + expect(result).toStrictEqual(snapshot); + const assistant = result.find((message) => message.role === 'assistant'); + expect(assistant?.toolCalls?.[0].function.arguments).toBe( + '{ "city": "Paris" }\n' + ); + expect(assistant?.metadata).not.toHaveProperty('obsolete'); + expect(assistant?.toolCalls?.[0]).not.toHaveProperty('obsolete'); + expectFrozen(result); + expect( + previous.find((message) => message.role === 'assistant')?.toolCalls?.[0] + .function.arguments + ).toBe('{"city":'); + }); + + it('lets same-ID replacement win over event-only retention, including role changes', () => { + const previous = ownTranscript([reasoning, activity, user('user')]); + const snapshot = [ + user('reasoning', 'now user'), + { id: 'activity', role: 'assistant' }, + { ...reasoning, id: 'user' }, + ] satisfies Message[]; + expect(reconcileTranscript(previous, snapshot)).toStrictEqual(snapshot); + }); + + it('does not merge omitted top-level or nested fields back into replacements', () => { + const previous = ownTranscript([ + { + id: 'a', + role: 'assistant', + content: 'old', + encryptedValue: 'old', + metadata: { keep: true, nested: { old: true } }, + toolCalls: [ + { + id: 'call', + type: 'function', + function: { name: 'old', arguments: '{' }, + }, + ], + }, + ]); + const snapshot = [ + { id: 'a', role: 'assistant', metadata: { nested: {} } }, + ] satisfies Message[]; + expect(reconcileTranscript(previous, snapshot)).toStrictEqual(snapshot); + }); + + it('preserves null, empty and nested attribution in the transcript', () => { + const snapshot = [ + { + id: 'null', + role: 'assistant', + subagentRunId: null, + metadata: { subagentRunId: null }, + }, + { + id: 'empty', + role: 'assistant', + subagentRunId: '', + metadata: { nested: { subagentRunId: '' } }, + }, + ] as unknown as Message[]; + expect(reconcileTranscript(ownTranscript([]), snapshot)).toStrictEqual( + snapshot + ); + }); + + it.each(['previous', 'snapshot'] as const)( + 'rejects duplicate IDs in %s without exposing payload data', + (side) => { + const duplicates = [ + user('secret-id', 'secret-content'), + { ...reasoning, id: 'secret-id' }, + ]; + const previous = ownTranscript( + side === 'previous' ? duplicates : [reasoning] + ); + const before = structuredClone(previous); + const snapshot = side === 'snapshot' ? duplicates : []; + expect(() => reconcileTranscript(previous, snapshot)).toThrow(TypeError); + expect(() => reconcileTranscript(previous, snapshot)).toThrow( + /^Duplicate message IDs in (previous transcript|snapshot)$/ + ); + expect(previous).toStrictEqual(before); + } + ); + + it('handles special object-property IDs without collisions or inherited entries', () => { + const previous = ownTranscript([ + user('__proto__'), + user('constructor'), + user('toString'), + ]); + const snapshot = [ + user('toString', 'new 3'), + user('__proto__', 'new 1'), + user('constructor', 'new 2'), + ]; + expect(reconcileTranscript(previous, snapshot)).toStrictEqual([ + snapshot[1], + snapshot[2], + snapshot[0], + ]); + }); + + it('detaches frozen snapshots and mutable nested caller data without freezing the caller', () => { + const nested = { value: 'before' }; + const snapshot = Object.freeze([ + Object.freeze({ + id: 'user', + role: 'user', + content: 'hello', + metadata: { nested }, + } satisfies Message), + ]); + const previous = ownTranscript([reasoning]); + const result = reconcileTranscript(previous, snapshot); + nested.value = 'after'; + expect(result).toStrictEqual([ + reasoning, + { ...snapshot[0], metadata: { nested: { value: 'before' } } }, + ]); + expect(Object.isFrozen(nested)).toBe(false); + expectFrozen(result); + expect(previous).toStrictEqual([reasoning]); + }); + + it('recaptures an owned snapshot without promising deep-equivalent reference reuse', () => { + const previous = ownTranscript([user('a')]); + const result = reconcileTranscript(previous, previous); + expect(result).toStrictEqual(previous); + expect(result).not.toBe(previous); + expect(result[0]).not.toBe(previous[0]); + }); + + it.each([ + ['date', new Date()], + ['map', new Map()], + ['set', new Set()], + ['typed array', new Uint8Array([1])], + ])('rejects unsupported %s atomically', (_name, invalid) => { + const previous = ownTranscript([reasoning, activity]); + let current: Transcript = previous; + const nested = { mutable: [] }; + expect(() => { + current = reconcileTranscript(previous, [ + { ...user('a'), metadata: { nested, invalid } }, + ]); + }).toThrow(TypeError); + expect(current).toBe(previous); + expect(previous).toStrictEqual([reasoning, activity]); + expect(Object.isFrozen(nested)).toBe(false); + }); + + it('rejects cycles atomically and propagates throwing getters before inspecting previous IDs', () => { + const previous = ownTranscript([reasoning, activity]); + const cyclic: Record = {}; + cyclic.self = cyclic; + expect(() => + reconcileTranscript(previous, [{ ...user('a'), metadata: cyclic }]) + ).toThrow(TypeError); + const failure = new Error('capture failed'); + const metadata = Object.defineProperty({}, 'failure', { + enumerable: true, + get: () => { + throw failure; + }, + }); + expect(() => + reconcileTranscript(previous, [{ ...user('a'), metadata }]) + ).toThrow(failure); + // Capture must precede even duplicate-ID rejection in the owned previous input. + const duplicates = ownTranscript([reasoning, reasoning]); + expect(() => + reconcileTranscript(duplicates, [{ ...user('a'), metadata }]) + ).toThrow(failure); + expect(previous).toStrictEqual([reasoning, activity]); + expect(Object.isFrozen(cyclic)).toBe(false); + }); + + it('keeps independent negative controls for replacement, keep-old merge and role-first retention', () => { + const oldAssistant = { + id: 'assistant', + role: 'assistant', + toolCalls: [ + { + id: 'call', + type: 'function', + function: { name: 'weather', arguments: '{"city":' }, + }, + ], + } satisfies Message; + const corrected = { + ...oldAssistant, + toolCalls: [ + { + ...oldAssistant.toolCalls[0], + function: { name: 'weather', arguments: '{ "city": "Paris" }' }, + }, + ], + }; + const previous = ownTranscript([oldAssistant, reasoning, activity]); + const snapshot = [corrected]; + const replacement = ownTranscript(snapshot); + expect(replacement.map((message) => message.role)).not.toContain( + 'reasoning' + ); + expect(replacement.map((message) => message.role)).not.toContain( + 'activity' + ); + const keepOld = [ + ...previous, + ...snapshot.filter( + (incoming) => !previous.some((old) => old.id === incoming.id) + ), + ]; + expect(keepOld[0]).toEqual(oldAssistant); + expect(keepOld[0]).not.toEqual(corrected); + const changedRole = user('reasoning', 'corrected role'); + const roleFirst = previous + .map((old) => + old.role === 'reasoning' || old.role === 'activity' + ? old + : [changedRole].find((incoming) => incoming.id === old.id) + ) + .filter(Boolean); + expect(roleFirst).toContain(previous[1]); + expect(roleFirst).not.toContainEqual(changedRole); + expect(reconcileTranscript(previous, snapshot)).toStrictEqual([ + corrected, + reasoning, + activity, + ]); + expect( + reconcileTranscript(ownTranscript([reasoning]), [changedRole]) + ).toStrictEqual([changedRole]); + }); +}); diff --git a/libs/ag-ui/src/runtime/reconcile-transcript.ts b/libs/ag-ui/src/runtime/reconcile-transcript.ts new file mode 100644 index 000000000..3cbab2562 --- /dev/null +++ b/libs/ag-ui/src/runtime/reconcile-transcript.ts @@ -0,0 +1,46 @@ +import type { Message } from '@ag-ui/client'; +import { ownTranscript, type Transcript } from './transcript'; + +/** Mirrors locked @ag-ui/client 0.0.59 edit-based ordering and role policy. + * Deliberate differences: reject duplicate IDs, always replace same-ID records, + * and leave null-attribution sanitation to request egress. */ +export function reconcileTranscript( + previous: Transcript, + snapshot: readonly Message[] | Transcript +): Transcript { + const captured = ownTranscript(snapshot); + const incoming = new Map(); + let hasActivity = false; + let hasReasoning = false; + for (const message of captured) { + if (incoming.has(message.id)) + throw new TypeError('Duplicate message IDs in snapshot'); + incoming.set(message.id, message); + if (message.role === 'activity') hasActivity = true; + if (message.role === 'reasoning') hasReasoning = true; + } + + const previousIds = new Set(); + const result: Transcript[number][] = []; + for (const message of previous) { + if (previousIds.has(message.id)) + throw new TypeError('Duplicate message IDs in previous transcript'); + previousIds.add(message.id); + const replacement = incoming.get(message.id); + if (replacement) result.push(replacement); + else if ( + (message.role === 'activity' && !hasActivity) || + (message.role === 'reasoning' && !hasReasoning) + ) + result.push(message); + } + for (const message of captured) { + if (!previousIds.has(message.id)) result.push(message); + } + if ( + result.length === previous.length && + result.every((message, index) => message === previous[index]) + ) + return previous; + return Object.freeze(result); +} diff --git a/libs/ag-ui/src/runtime/text-messages-http.spec.ts b/libs/ag-ui/src/runtime/text-messages-http.spec.ts new file mode 100644 index 000000000..61abe322b --- /dev/null +++ b/libs/ag-ui/src/runtime/text-messages-http.spec.ts @@ -0,0 +1,374 @@ +import { once } from 'node:events'; +import { createServer } from 'node:http'; +import type { AddressInfo } from 'node:net'; +import { + EventType, + type BaseEvent, + type Message, + type MessagesSnapshotEvent, + type RunAgentInput, +} from '@ag-ui/client'; +import { describe, expect, it } from 'vitest'; +import { createRun, type RunHandle } from './create-run'; +import { reconcileTranscript } from './reconcile-transcript'; +import { applyTextMessage, type TextMessageEvent } from './text-messages'; +import { ownTranscript, requestMessages, type Transcript } from './transcript'; + +function deferred() { + let resolve!: (value: T) => void; + const promise = new Promise((yes) => { + resolve = yes; + }); + return { promise, resolve }; +} + +async function bounded(promise: Promise): Promise { + let timer: ReturnType | undefined; + try { + return await Promise.race([ + promise, + new Promise((_, reject) => { + timer = setTimeout( + () => reject(new Error('Text HTTP milestone timed out')), + 1500 + ); + }), + ]); + } finally { + clearTimeout(timer); + } +} + +interface Exchange { + body: RunAgentInput; + closed: Promise; + send: (...events: unknown[]) => void; +} + +async function serve() { + const exchanges: Exchange[] = []; + const arrivals = new Map>>(); + const server = createServer(async (request, response) => { + const closed = deferred(); + response.on('close', () => closed.resolve()); + const chunks: Buffer[] = []; + for await (const chunk of request) chunks.push(Buffer.from(chunk)); + const exchange: Exchange = { + body: JSON.parse(Buffer.concat(chunks).toString()), + closed: closed.promise, + send: (...events) => { + for (const event of events) + response.write(`data: ${JSON.stringify(event)}\n\n`); + }, + }; + response.writeHead(200, { 'content-type': 'text/event-stream' }); + response.flushHeaders(); + exchanges.push(exchange); + arrivals.get(exchanges.length - 1)?.resolve(exchange); + // Held open: completion must physically cancel the response reader. + }); + server.listen(0, '127.0.0.1'); + await once(server, 'listening'); + return { + url: `http://127.0.0.1:${(server.address() as AddressInfo).port}/events`, + exchanges, + next: (index = 0) => { + if (exchanges[index]) return Promise.resolve(exchanges[index]); + const arrival = arrivals.get(index) ?? deferred(); + arrivals.set(index, arrival); + return bounded(arrival.promise); + }, + close: async () => { + const closed = once(server, 'close'); + server.close(); + server.closeAllConnections(); + await bounded(closed); + }, + }; +} + +function input(runId: string, messages: Message[]): RunAgentInput { + return { + threadId: 'thread', + runId, + messages, + state: {}, + tools: [], + context: [], + forwardedProps: {}, + }; +} + +function lifecycle( + exchange: Exchange, + type: EventType.RUN_STARTED | EventType.RUN_FINISHED +) { + return { type, threadId: exchange.body.threadId, runId: exchange.body.runId }; +} + +// Only the six normalized message events belong to this helper. Snapshot +// authority and lifecycle authority remain with their existing owners. +function project(previous: Transcript, event: BaseEvent): Transcript { + switch (event.type) { + case EventType.TEXT_MESSAGE_START: + case EventType.TEXT_MESSAGE_CONTENT: + case EventType.TEXT_MESSAGE_END: + case EventType.REASONING_MESSAGE_START: + case EventType.REASONING_MESSAGE_CONTENT: + case EventType.REASONING_MESSAGE_END: + return applyTextMessage(previous, event as TextMessageEvent); + case EventType.MESSAGES_SNAPSHOT: + return reconcileTranscript( + previous, + (event as MessagesSnapshotEvent).messages + ); + default: + return previous; + } +} + +describe('private text accumulation over actual HTTP', () => { + it('observes both normalized CHUNK families and metadata-only continuations, reconciles and sends full owned egress', async () => { + const server = await serve(); + const factory = createRun({ url: server.url }); + const user = { + id: 'u', + role: 'user', + content: 'question', + } satisfies Message; + const seed = ownTranscript([user]); + let transcript = seed; + const events: BaseEvent[] = []; + const observations: Transcript[] = []; + const streamed = deferred(); + let second: RunHandle | undefined; + const first = factory.start( + input('stream', requestMessages(seed)), + (event) => { + events.push(event); + transcript = project(transcript, event); + observations.push(transcript); + if ( + event.type === EventType.REASONING_MESSAGE_CONTENT && + event.delta === '' + ) + streamed.resolve(); + } + ); + try { + const exchange = await server.next(); + expect(exchange.body).toStrictEqual(input('stream', [user])); + exchange.send( + lifecycle(exchange, EventType.RUN_STARTED), + { + type: EventType.TEXT_MESSAGE_CHUNK, + messageId: 'a', + name: 'speaker', + delta: 'Hello ', + metadata: { keep: 'text', phase: { start: true } }, + }, + { type: EventType.TEXT_MESSAGE_CHUNK, delta: '🌍\n' }, + { + type: EventType.TEXT_MESSAGE_CHUNK, + metadata: { phase: { final: true }, usage: { tokens: 2 } }, + }, + { + type: EventType.REASONING_MESSAGE_CHUNK, + messageId: 'r', + delta: '考える ', + metadata: { keep: 'reasoning', phase: { start: true } }, + }, + { type: EventType.REASONING_MESSAGE_CHUNK, delta: '🧠\t' }, + { + type: EventType.REASONING_MESSAGE_CHUNK, + metadata: { phase: { final: true }, usage: { tokens: 3 } }, + } + ); + await bounded(streamed.promise); + const streamedRecords = [ + user, + { + id: 'a', + role: 'assistant', + name: 'speaker', + content: 'Hello 🌍\n', + metadata: { + keep: 'text', + phase: { final: true }, + usage: { tokens: 2 }, + }, + }, + { + id: 'r', + role: 'reasoning', + content: '考える 🧠\t', + metadata: { + keep: 'reasoning', + phase: { final: true }, + usage: { tokens: 3 }, + }, + }, + ] satisfies Message[]; + expect(transcript).toStrictEqual(streamedRecords); + expect(transcript[0]).toBe(seed[0]); + expect(observations[1][1].content).toBe(''); + expect(observations[2][1].content).toBe('Hello '); + expect(observations[3][1].content).toBe('Hello 🌍\n'); + expect(observations[7][2].content).toBe('考える '); + const beforeSnapshot = transcript; + const corrected = { + id: 'a', + role: 'assistant', + name: 'corrected', + content: 'Authoritative answer', + encryptedValue: 'opaque', + metadata: { snapshot: ['only'] }, + toolCalls: [ + { + id: 'call', + type: 'function', + function: { name: 'weather', arguments: '{ "city": "Paris" }\n' }, + encryptedValue: 'opaque-call', + metadata: { preserve: true }, + }, + ], + } satisfies Message; + exchange.send( + { type: EventType.MESSAGES_SNAPSHOT, messages: [user, corrected] }, + lifecycle(exchange, EventType.RUN_FINISHED) + ); + expect(await bounded(first.done)).toEqual({ outcome: 'success' }); + await bounded(exchange.closed); + expect(events.map((event) => event.type)).toEqual([ + EventType.RUN_STARTED, + EventType.TEXT_MESSAGE_START, + EventType.TEXT_MESSAGE_CONTENT, + EventType.TEXT_MESSAGE_CONTENT, + EventType.TEXT_MESSAGE_CONTENT, + EventType.TEXT_MESSAGE_END, + EventType.REASONING_MESSAGE_START, + EventType.REASONING_MESSAGE_CONTENT, + EventType.REASONING_MESSAGE_CONTENT, + EventType.REASONING_MESSAGE_CONTENT, + EventType.REASONING_MESSAGE_END, + EventType.MESSAGES_SNAPSHOT, + EventType.RUN_FINISHED, + ]); + expect(events[4]).toMatchObject({ + type: EventType.TEXT_MESSAGE_CONTENT, + delta: '', + metadata: { usage: { tokens: 2 } }, + }); + expect(events[9]).toMatchObject({ + type: EventType.REASONING_MESSAGE_CONTENT, + delta: '', + metadata: { usage: { tokens: 3 } }, + }); + const expected = [user, corrected, streamedRecords[2]]; + expect(transcript).toStrictEqual(expected); + expect(transcript[2]).toBe(beforeSnapshot[2]); + expect(beforeSnapshot).toStrictEqual(streamedRecords); + expect(seed).toStrictEqual([user]); + expect(Object.isFrozen(transcript)).toBe(true); + const assistant = transcript[1]; + if (assistant.role !== 'assistant') throw new Error('Expected assistant'); + expect(Object.isFrozen(assistant.toolCalls?.[0].function)).toBe(true); + expect(Object.isFrozen(transcript[2].metadata?.usage)).toBe(true); + second = factory.start( + input('continue', requestMessages(transcript)), + () => undefined + ); + const continuation = await server.next(1); + expect(continuation.body).toStrictEqual(input('continue', expected)); + continuation.send( + lifecycle(continuation, EventType.RUN_STARTED), + lifecycle(continuation, EventType.RUN_FINISHED) + ); + expect(await bounded(second.done)).toEqual({ outcome: 'success' }); + await bounded(continuation.closed); + expect(server.exchanges).toHaveLength(2); + } finally { + first.abort(); + second?.abort(); + await server.close(); + } + }); + + it('fails in the helper after an admitted snapshot removes an open target and physically closes without a success override', async () => { + const server = await serve(); + const user = { + id: 'u', + role: 'user', + content: 'question', + } satisfies Message; + const seed = ownTranscript([user]); + let transcript = seed; + let latestValid: Transcript | undefined; + let streamed: Transcript | undefined; + let failure: unknown; + const delivered: EventType[] = []; + const handle = createRun({ url: server.url }).start( + input('removed', requestMessages(seed)), + (event) => { + delivered.push(event.type); + try { + transcript = project(transcript, event); + } catch (error) { + failure = error; + throw error; + } + if (event.type === EventType.TEXT_MESSAGE_CONTENT) + streamed = transcript; + if (event.type === EventType.MESSAGES_SNAPSHOT) + latestValid = transcript; + } + ); + try { + const exchange = await server.next(); + exchange.send( + lifecycle(exchange, EventType.RUN_STARTED), + { + type: EventType.TEXT_MESSAGE_START, + messageId: 'a', + role: 'assistant', + }, + { + type: EventType.TEXT_MESSAGE_CONTENT, + messageId: 'a', + delta: 'before', + }, + { type: EventType.MESSAGES_SNAPSHOT, messages: [user] }, + { + type: EventType.TEXT_MESSAGE_CONTENT, + messageId: 'a', + delta: 'after removal', + }, + lifecycle(exchange, EventType.RUN_FINISHED) + ); + const result = await bounded(handle.done); + expect(failure).toBeInstanceOf(TypeError); + expect(result).toEqual({ outcome: 'error', error: failure }); + if (result.outcome === 'error') expect(result.error).toBe(failure); + await bounded(exchange.closed); + expect(delivered).toEqual([ + EventType.RUN_STARTED, + EventType.TEXT_MESSAGE_START, + EventType.TEXT_MESSAGE_CONTENT, + EventType.MESSAGES_SNAPSHOT, + EventType.TEXT_MESSAGE_CONTENT, + ]); + expect(streamed).toStrictEqual([ + user, + { id: 'a', role: 'assistant', content: 'before' }, + ]); + expect(transcript).toBe(latestValid); + expect(transcript).toStrictEqual([user]); + expect(seed).toStrictEqual([user]); + expect(await handle.done).toBe(result); + expect(server.exchanges).toHaveLength(1); + } finally { + handle.abort(); + await server.close(); + } + }); +}); diff --git a/libs/ag-ui/src/runtime/text-messages.spec.ts b/libs/ag-ui/src/runtime/text-messages.spec.ts new file mode 100644 index 000000000..cb6d156c0 --- /dev/null +++ b/libs/ag-ui/src/runtime/text-messages.spec.ts @@ -0,0 +1,552 @@ +import { EventType, type Message } from '@ag-ui/client'; +import { describe, expect, it } from 'vitest'; +import { applyTextMessage, type TextMessageEvent } from './text-messages'; +import { ownTranscript } from './transcript'; + +const families = [ + { + role: 'assistant', + start: EventType.TEXT_MESSAGE_START, + content: EventType.TEXT_MESSAGE_CONTENT, + end: EventType.TEXT_MESSAGE_END, + }, + { + role: 'reasoning', + start: EventType.REASONING_MESSAGE_START, + content: EventType.REASONING_MESSAGE_CONTENT, + end: EventType.REASONING_MESSAGE_END, + }, +] as const; + +function expectFrozen(value: unknown): void { + if (value === null || typeof value !== 'object') return; + expect(Object.isFrozen(value)).toBe(true); + for (const child of Object.values(value)) expectFrozen(child); +} + +describe('private text and reasoning accumulation', () => { + it.each(['developer', 'system', 'assistant', 'user'] as const)( + 'creates and appends exact content for the %s text role', + (role) => { + const previous = ownTranscript([]); + const started = applyTextMessage(previous, { + type: EventType.TEXT_MESSAGE_START, + messageId: role, + role, + name: 'speaker', + }); + expect(started).toStrictEqual([ + { id: role, role, content: '', name: 'speaker' }, + ]); + const content = applyTextMessage(started, { + type: EventType.TEXT_MESSAGE_CONTENT, + messageId: role, + delta: ' \n café 🧭\t', + }); + expect(content[0].content).toBe(' \n café 🧭\t'); + expect(started[0].content).toBe(''); + expect(previous).toEqual([]); + expectFrozen(content); + } + ); + + it('defaults an omitted text role and ignores event-only data', () => { + const event = { + type: EventType.TEXT_MESSAGE_START, + messageId: 'a', + timestamp: 123, + get rawEvent() { + throw new Error('rawEvent must not be read'); + }, + } as TextMessageEvent; + expect(applyTextMessage(ownTranscript([]), event)).toStrictEqual([ + { id: 'a', role: 'assistant', content: '' }, + ]); + }); + + it('interleaves text and reasoning while preserving all old observations', () => { + const user = ownTranscript([ + { id: 'u', role: 'user', content: 'question' }, + ]); + const text = applyTextMessage(user, { + type: EventType.TEXT_MESSAGE_START, + messageId: 'a', + role: 'assistant', + }); + const reasoning = applyTextMessage(text, { + type: EventType.REASONING_MESSAGE_START, + messageId: 'r', + role: 'reasoning', + }); + const thought = applyTextMessage(reasoning, { + type: EventType.REASONING_MESSAGE_CONTENT, + messageId: 'r', + delta: '考える\n', + }); + const answer = applyTextMessage(thought, { + type: EventType.TEXT_MESSAGE_CONTENT, + messageId: 'a', + delta: ' yes ', + }); + const final = applyTextMessage(answer, { + type: EventType.REASONING_MESSAGE_CONTENT, + messageId: 'r', + delta: '🧠', + }); + expect(final.map((message) => message.content)).toEqual([ + 'question', + ' yes ', + '考える\n🧠', + ]); + expect(reasoning.map((message) => message.content)).toEqual([ + 'question', + '', + '', + ]); + expect(thought[2].content).toBe('考える\n'); + expect(final[0]).toBe(user[0]); + expect(final[1]).toBe(answer[1]); + expect(answer[2]).toBe(thought[2]); + expectFrozen(final); + }); + + it.each(families)( + 'merges $role metadata at every stage, including empty content', + (family) => { + const previous = ownTranscript([]); + const start: TextMessageEvent = { + type: family.start, + role: family.role, + messageId: 'm', + metadata: { keep: 1, replace: { old: true } }, + } as TextMessageEvent; + const started = applyTextMessage(previous, start); + expect(started[0].metadata).toEqual(start.metadata); + const repeat = applyTextMessage(started, { + ...start, + metadata: { replace: { start: true } }, + }); + const content = applyTextMessage(repeat, { + type: family.content, + messageId: 'm', + delta: 'text', + metadata: { replace: { content: true }, content: 2 }, + }); + const empty = applyTextMessage(content, { + type: family.content, + messageId: 'm', + delta: '', + metadata: { replace: { empty: true }, empty: 3 }, + }); + expect(empty[0].content).toBe('text'); + expect(empty[0].metadata).toEqual({ + keep: 1, + replace: { empty: true }, + content: 2, + empty: 3, + }); + const ended = applyTextMessage(empty, { + type: family.end, + messageId: 'm', + metadata: { replace: { end: true } }, + }); + expect(ended[0].metadata).toEqual({ + keep: 1, + replace: { end: true }, + content: 2, + empty: 3, + }); + expect(ended[0].content).toBe('text'); + expect(started[0].metadata).toEqual({ keep: 1, replace: { old: true } }); + expectFrozen(ended); + } + ); + + it.each(families)( + 'keeps stable $role no-ops and permits a supplied empty metadata update', + (family) => { + const previous = ownTranscript([ + { id: 'm', role: family.role, content: 'keep' }, + ]); + const start = { + type: family.start, + role: family.role, + messageId: 'm', + } as TextMessageEvent; + const content = { + type: family.content, + messageId: 'm', + delta: '', + } as TextMessageEvent; + const end = { type: family.end, messageId: 'm' } as TextMessageEvent; + for (const event of [start, start, content, end, end]) { + expect(applyTextMessage(previous, event)).toBe(previous); + } + expect( + applyTextMessage(previous, { ...content, metadata: {} })[0] + ).toStrictEqual({ + ...previous[0], + metadata: {}, + }); + } + ); + + it('preserves opaque tool fields, name, content and attribution on later events', () => { + const source = { + id: 'm', + role: 'assistant', + name: 'original', + content: 'before', + subagentRunId: 'original-child', + encryptedValue: 'opaque', + toolCalls: [ + { + id: 'call', + type: 'function', + function: { name: 'tool', arguments: '{ "x":' }, + encryptedValue: 'call-opaque', + metadata: { keep: [1] }, + }, + ], + metadata: { old: true }, + } satisfies Message; + const previous = ownTranscript([source]); + const repeated = applyTextMessage(previous, { + type: EventType.TEXT_MESSAGE_START, + role: 'assistant', + messageId: 'm', + name: 'replacement', + subagentRunId: 'other-child', + metadata: { start: true }, + }); + const content = applyTextMessage(repeated, { + type: EventType.TEXT_MESSAGE_CONTENT, + messageId: 'm', + delta: ' after', + subagentRunId: '', + metadata: { content: true }, + }); + const ended = applyTextMessage(content, { + type: EventType.TEXT_MESSAGE_END, + messageId: 'm', + subagentRunId: 'other', + metadata: { end: true }, + }); + expect(ended[0]).toStrictEqual({ + ...source, + content: 'before after', + metadata: { old: true, start: true, content: true, end: true }, + }); + expect(previous).toStrictEqual([source]); + expectFrozen(ended); + }); + + it.each(families)( + 'treats undefined $role content as empty for a nonempty append', + (family) => { + // Local captured records can predate content even though the SDK requires + // reasoning content in wire snapshots; this helper defines that fallback. + const previous = ownTranscript([ + { id: 'm', role: family.role }, + ] as Message[]); + const next = applyTextMessage(previous, { + type: family.content, + messageId: 'm', + delta: 'first', + }); + expect(next[0].content).toBe('first'); + expect(previous[0]).not.toHaveProperty('content'); + expect( + applyTextMessage(previous, { + type: family.content, + messageId: 'm', + delta: '', + }) + ).toBe(previous); + } + ); + + it.each(families)( + 'captures $role creation attribution, including empty strings', + (family) => { + for (const subagentRunId of [undefined, null, '', 'child']) { + const event = { + type: family.start, + role: family.role, + messageId: 'm', + subagentRunId, + } as TextMessageEvent; + const result = applyTextMessage(ownTranscript([]), event); + if (typeof subagentRunId === 'string') + expect(result[0].subagentRunId).toBe(subagentRunId); + else expect(result[0]).not.toHaveProperty('subagentRunId'); + } + } + ); + + it.each(['__proto__', 'constructor', 'toString', ''])( + 'uses %s as an ordinary ID', + (messageId) => { + const start = applyTextMessage(ownTranscript([]), { + type: EventType.TEXT_MESSAGE_START, + role: 'assistant', + messageId, + }); + const result = applyTextMessage(start, { + type: EventType.TEXT_MESSAGE_CONTENT, + messageId, + delta: 'ok', + }); + expect(result).toStrictEqual([ + { id: messageId, role: 'assistant', content: 'ok' }, + ]); + } + ); + + it.each(families)( + 'fails atomically for missing or duplicate $role targets', + (family) => { + const single = { + id: 'm', + role: family.role, + content: 'keep', + } satisfies Message; + const previous = ownTranscript([single]); + for (const type of [family.content, family.end]) { + expect(() => + applyTextMessage(previous, { + type, + messageId: 'secret-missing', + delta: 'secret-delta', + } as TextMessageEvent) + ).toThrow(TypeError); + } + const duplicate = ownTranscript([single, single]); + for (const type of [family.start, family.content, family.end]) { + expect(() => + applyTextMessage(duplicate, { + type, + role: family.role, + messageId: 'm', + delta: 'secret-delta', + } as TextMessageEvent) + ).toThrow(TypeError); + } + expect(previous).toStrictEqual([single]); + expect(duplicate).toStrictEqual([single, single]); + } + ); + + it('checks only target uniqueness; full transcript admission belongs to the caller', () => { + const previous = ownTranscript([ + { id: 'unrelated', role: 'user', content: 'one' }, + { id: 'unrelated', role: 'user', content: 'two' }, + { id: 'm', role: 'assistant', content: 'yes' }, + ]); + const next = applyTextMessage(previous, { + type: EventType.TEXT_MESSAGE_CONTENT, + messageId: 'm', + delta: '!', + }); + expect(next[2].content).toBe('yes!'); + expect(next[0]).toBe(previous[0]); + expect(next[1]).toBe(previous[1]); + }); + + it('rejects incompatible roles at START, CONTENT and END without exposing payloads', () => { + const records: Message[] = [ + { id: 'secret-id', role: 'user', content: 'keep' }, + { id: 'secret-id', role: 'reasoning', content: 'keep' }, + { id: 'secret-id', role: 'tool', toolCallId: 'call', content: 'keep' }, + { id: 'secret-id', role: 'activity', activityType: 'test', content: {} }, + ]; + for (const record of records) { + for (const family of families) { + for (const type of [family.start, family.content, family.end]) { + if ( + record.role === family.role || + (record.role === 'user' && + family.role === 'assistant' && + type !== family.start) + ) + continue; + const previous = ownTranscript([record]); + let failure: unknown; + try { + applyTextMessage(previous, { + type, + role: family.role, + messageId: 'secret-id', + delta: 'secret-delta', + metadata: { secret: 'secret-metadata' }, + } as TextMessageEvent); + } catch (error) { + failure = error; + } + expect(failure).toBeInstanceOf(TypeError); + expect(String(failure)).not.toMatch(/secret/); + expect(previous).toStrictEqual([record]); + } + } + } + }); + + it('rejects content appended to multimodal data, even an empty delta', () => { + const source = { + id: 'm', + role: 'user', + content: [{ type: 'text', text: 'keep image caption' }], + } satisfies Message; + const previous = ownTranscript([source]); + for (const delta of ['', 'append']) { + expect(() => + applyTextMessage(previous, { + type: EventType.TEXT_MESSAGE_CONTENT, + messageId: 'm', + delta, + }) + ).toThrow(TypeError); + } + expect(previous).toStrictEqual([source]); + }); + + it('captures incoming metadata despite frozen external parents and later caller mutation', () => { + const nested = { values: ['before'] }; + const metadata = Object.freeze({ nested }); + const event = Object.freeze({ + type: EventType.TEXT_MESSAGE_START, + role: 'assistant', + messageId: 'm', + metadata, + } as const); + const next = applyTextMessage(ownTranscript([]), event); + nested.values.push('later'); + expect(next[0].metadata).toStrictEqual({ nested: { values: ['before'] } }); + expect(next[0].metadata).not.toBe(metadata); + expectFrozen(next); + expect(Object.isFrozen(nested)).toBe(false); + }); + + it('rejects unsupported incoming metadata before shallow spreading can hide it', () => { + class Metadata { + keep = true; + } + const cycle: Record = {}; + cycle.self = cycle; + const failure = new Error('getter failure'); + const getter = { + get value(): never { + throw failure; + }, + }; + const invalid = [ + new Date(), + new Metadata(), + { cycle }, + { fn: () => undefined }, + getter, + ]; + const previous = ownTranscript([ + { + id: 'm', + role: 'assistant', + content: 'keep', + metadata: { original: true }, + }, + ]); + for (const metadata of invalid) { + for (const type of [ + EventType.TEXT_MESSAGE_START, + EventType.TEXT_MESSAGE_CONTENT, + EventType.TEXT_MESSAGE_END, + ] as const) { + expect(() => + applyTextMessage(previous, { + type, + role: 'assistant', + messageId: 'm', + delta: 'new', + metadata, + } as TextMessageEvent) + ).toThrow(); + } + expect(() => + applyTextMessage(ownTranscript([]), { + type: EventType.TEXT_MESSAGE_START, + role: 'assistant', + messageId: 'new', + metadata, + } as TextMessageEvent) + ).toThrow(); + } + expect(previous).toStrictEqual([ + { + id: 'm', + role: 'assistant', + content: 'keep', + metadata: { original: true }, + }, + ]); + }); + + it('captures each needed event field once and never reads irrelevant fields', () => { + const reads: Record = {}; + const values = { + type: EventType.TEXT_MESSAGE_CONTENT, + messageId: 'm', + delta: '!', + metadata: { incoming: true }, + }; + const event = Object.fromEntries( + Object.keys(values).map((key) => [key, undefined]) + ); + for (const [key, value] of Object.entries(values)) + Object.defineProperty(event, key, { + get() { + reads[key] = (reads[key] ?? 0) + 1; + return value; + }, + }); + Object.defineProperty(event, 'rawEvent', { + get() { + throw new Error('irrelevant'); + }, + }); + const result = applyTextMessage( + ownTranscript([{ id: 'm', role: 'assistant', content: 'yes' }]), + event as TextMessageEvent + ); + expect(result[0].content).toBe('yes!'); + expect(reads).toEqual({ type: 1, messageId: 1, delta: 1, metadata: 1 }); + }); +}); + +describe('retained negative controls', () => { + it('a naive mutable append rewrites an earlier observation', () => { + const messages = [{ id: 'm', role: 'assistant', content: 'before' }]; + const observation = [...messages]; + messages[0].content += ' after'; + expect(observation[0].content).not.toBe('before'); + expect(observation[0].content).toBe('before after'); + }); + + it('a content-only reducer loses an empty-delta metadata event', () => { + const previous = ownTranscript([ + { id: 'm', role: 'assistant', content: 'text', metadata: { keep: true } }, + ]); + const event = { + type: EventType.TEXT_MESSAGE_CONTENT, + messageId: 'm', + delta: '', + metadata: { arrived: true }, + } as const; + const contentOnly = event.delta + ? [{ ...previous[0], content: previous[0].content + event.delta }] + : previous; + expect(contentOnly[0].metadata).not.toHaveProperty('arrived'); + expect(applyTextMessage(previous, event)[0].metadata).toEqual({ + keep: true, + arrived: true, + }); + }); +}); diff --git a/libs/ag-ui/src/runtime/text-messages.ts b/libs/ag-ui/src/runtime/text-messages.ts new file mode 100644 index 000000000..800916fe3 --- /dev/null +++ b/libs/ag-ui/src/runtime/text-messages.ts @@ -0,0 +1,99 @@ +import { + EventType, + mergeMetadata, + type TextMessageStartEvent, + type TextMessageContentEvent, + type TextMessageEndEvent, + type ReasoningMessageStartEvent, + type ReasoningMessageContentEvent, + type ReasoningMessageEndEvent, +} from '@ag-ui/client'; +import { ownTranscript, type Transcript } from './transcript'; + +export type TextMessageEvent = + | TextMessageStartEvent + | TextMessageContentEvent + | TextMessageEndEvent + | ReasoningMessageStartEvent + | ReasoningMessageContentEvent + | ReasoningMessageEndEvent; + +/** Applies one normalized message event to an already owned transcript. + * The caller owns admission, full ID uniqueness and root/child routing. + * Only the selected record is captured; unaffected owned records are shared. */ +export function applyTextMessage( + previous: Transcript, + event: TextMessageEvent +): Transcript { + const { type, messageId, metadata } = event; + const start = + type === EventType.TEXT_MESSAGE_START || + type === EventType.REASONING_MESSAGE_START; + const reasoning = + type === EventType.REASONING_MESSAGE_START || + type === EventType.REASONING_MESSAGE_CONTENT || + type === EventType.REASONING_MESSAGE_END; + const role = + type === EventType.TEXT_MESSAGE_START + ? event.role ?? 'assistant' + : 'reasoning'; + let index = -1; + for (let position = 0; position < previous.length; position++) { + if (previous[position].id !== messageId) continue; + if (index !== -1) throw new TypeError('Duplicate target message ID'); + index = position; + } + const current = index === -1 ? undefined : previous[index]; + let record: Transcript[number]; + if (!current) { + if (!start) throw new TypeError('Message target is missing'); + const subagentRunId = event.subagentRunId; + const name = type === EventType.TEXT_MESSAGE_START ? event.name : undefined; + record = { + id: messageId, + role, + content: '', + ...(name !== undefined && { name }), + ...(typeof subagentRunId === 'string' && { subagentRunId }), + }; + } else { + if ( + current.role === 'tool' || + current.role === 'activity' || + (current.role === 'reasoning') !== reasoning || + (start && current.role !== role) + ) + throw new TypeError('Message target role is incompatible'); + record = current; + if ( + type === EventType.TEXT_MESSAGE_CONTENT || + type === EventType.REASONING_MESSAGE_CONTENT + ) { + const { delta } = event; + const content = current.content; + if (content !== undefined && typeof content !== 'string') + throw new TypeError('Message target content must be a string'); + if (delta.length > 0) + record = { ...current, content: (content ?? '') + delta }; + } + if (record === current && metadata === undefined) return previous; + } + // Capture incoming metadata before SDK spreading could hide an unsupported + // prototype. Ownership failures happen before any new value is published. + const captured = ownTranscript([ + metadata === undefined ? record : { ...record, metadata }, + ])[0]; + const changed = + metadata === undefined + ? captured + : Object.freeze({ + ...captured, + metadata: Object.freeze( + mergeMetadata(current?.metadata, captured.metadata) + ), + }); + const next = [...previous]; + if (index === -1) next.push(changed); + else next[index] = changed; + return Object.freeze(next); +} diff --git a/scripts/react-parity/baseline.json b/scripts/react-parity/baseline.json index cc9250174..b2054af25 100644 --- a/scripts/react-parity/baseline.json +++ b/scripts/react-parity/baseline.json @@ -1,11 +1,11 @@ { "schemaVersion": 1, - "baselineHead": "9a87f53641064041cae89e1e83fd2da071440e5b", + "baselineHead": "d329076ac466b54b45a4aea44d0c5edd6977661b", "sourceState": { - "modified": [ - "libs/ag-ui/src/lib/internal/apply-patch.ts" - ], - "untracked": [] + "modified": [], + "untracked": [ + "libs/ag-ui/src/runtime/text-messages.ts" + ] }, "scope": { "libraries": [ @@ -10622,6 +10622,18 @@ "path": "libs/ag-ui/src/runtime/create-run.ts", "sha256": "c830bfa2f0832f4c0ea3dfcc5c64360bfbe067227b60b3b86140dd5a686013cb" }, + { + "id": "source:libs/ag-ui/src/runtime/reconcile-transcript.ts", + "kind": "source", + "path": "libs/ag-ui/src/runtime/reconcile-transcript.ts", + "sha256": "c2efbe4561a3d348422b144f7ab09656aa85b85771870fff4d4963880b9c6b74" + }, + { + "id": "source:libs/ag-ui/src/runtime/text-messages.ts", + "kind": "source", + "path": "libs/ag-ui/src/runtime/text-messages.ts", + "sha256": "6751d4b7f53350ccedafa2e1cae808d80c88acc894e8993441e30d60022e625e" + }, { "id": "source:libs/ag-ui/src/runtime/transcript.ts", "kind": "source", diff --git a/scripts/react-parity/dispositions.json b/scripts/react-parity/dispositions.json index 428537d86..bed76157e 100644 --- a/scripts/react-parity/dispositions.json +++ b/scripts/react-parity/dispositions.json @@ -8674,6 +8674,22 @@ "reason": "Private domain run authority only; excluded from public exports and legacy Angular reachability.", "note": "Composes physical HTTP ownership with captured caller identity, synchronous cancellation and first terminal outcome. Framework-neutral session projection, transcript fidelity and public parity remain separate migration work. Child-attributed RUN_* events are unsupported when delivered by the SDK; protobuf decoding can erase attribution." }, + { + "id": "source:libs/ag-ui/src/runtime/reconcile-transcript.ts", + "taskIds": ["T12", "T13"], + "treatment": "internal", + "status": "in-progress", + "reason": "Private immutable AG-UI snapshot reconciliation only; excluded from public exports and legacy Angular reachability.", + "note": "Captures complete replacement records, preserves existing ID order and omitted reasoning/activity roles, and rejects duplicate IDs atomically. Seeded HTTP coverage proves reconciliation and request egress; stream reduction, session integration and public parity remain separate migration work." + }, + { + "id": "source:libs/ag-ui/src/runtime/text-messages.ts", + "taskIds": ["T12", "T13"], + "treatment": "internal", + "status": "in-progress", + "reason": "Private immutable text and reasoning message accumulation only; excluded from public exports and legacy Angular reachability.", + "note": "Consumes six SDK-normalized message events, captures incoming metadata before shallow merging, preserves prior observations and fails on invalid targets. HTTP coverage composes streaming with snapshot reconciliation and request egress. Caller admission, root/child routing, other protocol families, session integration and public parity remain separate migration work." + }, { "id": "source:libs/ag-ui/src/runtime/transcript.ts", "taskIds": ["T12", "T13"],