Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
163 changes: 162 additions & 1 deletion libs/ag-ui/src/lib/reducer.subagent.spec.ts
Original file line number Diff line number Diff line change
@@ -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 {
Expand All @@ -19,6 +19,7 @@ interface TestDeliveryRun {
currentAssistantMessageId?: string;
eligibleBaselineAssistantId?: string;
protocolRunId?: string;
pendingReasoning?: { messageId: string; startedAt: number };
outcome?: 'success' | 'error' | 'aborted' | 'interrupted' | 'paused';
}

Expand Down Expand Up @@ -52,7 +53,167 @@ function makeStore(generation = 'run-generation-1'): TestStore {

const ev = (e: Record<string, unknown>) => e as unknown as BaseEvent;

describe('reduceEvent attributed child reasoning', () => {
const reasoningTypes = ['START', 'CONTENT', 'CHUNK', 'END'];
const reason = (store: TestStore, type: string, fields: Record<string, unknown> = {}) =>
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);
Expand Down
27 changes: 25 additions & 2 deletions libs/ag-ui/src/lib/reducer.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -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<string, unknown>;

Expand All @@ -664,6 +677,16 @@ function routeSubagentContentEvent(subagentRunId: string, event: BaseEvent, stor
const messages = [...((c['messages'] as Array<Record<string, unknown>>) ?? [])];
const toolCalls = [...((c['toolCalls'] as Array<Record<string, unknown>>) ?? [])];
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: '' });
Expand Down
Loading
Loading