diff --git a/app/src/providers/ChatRuntimeProvider.tsx b/app/src/providers/ChatRuntimeProvider.tsx index fceb293f4..9eb8328a8 100644 --- a/app/src/providers/ChatRuntimeProvider.tsx +++ b/app/src/providers/ChatRuntimeProvider.tsx @@ -51,10 +51,15 @@ import { setToolTimelineForThread, setWorkflowProposalForThread, streamDeltaReceived, + subagentAwaitingUser, + subagentDone, + subagentIterationStarted, + subagentSpawned, + subagentToolCallReceived, + subagentToolResultReceived, + toolArgsDeltaReceived, toolCallReceived, toolResultReceived, - type ToolTimelineEntry, - type ToolTimelineEntryStatus, upsertArtifactFailedForThread, upsertArtifactInProgressForThread, upsertArtifactReadyForThread, @@ -72,11 +77,6 @@ import { setSelectedThread, } from '../store/threadSlice'; import { IS_PROD } from '../utils/config'; -import { - formatTimelineEntry, - isKnownClientTool, - promptFromArgsBuffer, -} from '../utils/toolTimelineFormatting'; const logChatRuntime = debug('openhuman:chat-runtime'); const USER_FACING_AGENT_ERROR_MESSAGE = @@ -332,27 +332,6 @@ function parseWorkflowProposal(output: string): WorkflowProposal | null { }; } -export function findPendingDelegationContext( - entries: ToolTimelineEntry[], - round: number -): { sourceToolName?: string; prompt?: string; spawnEntryId?: string } { - for (let i = entries.length - 1; i >= 0; i -= 1) { - const entry = entries[i]; - if (entry.status !== 'running' || entry.round !== round) continue; - if ( - ['spawn_subagent', 'spawn_async_subagent'].includes(entry.name) || - entry.name.startsWith('delegate_') - ) { - return { - sourceToolName: entry.name, - prompt: entry.detail ?? promptFromArgsBuffer(entry.argsBuffer), - spawnEntryId: entry.id, - }; - } - } - return {}; -} - const ChatRuntimeProvider = ({ children }: { children: React.ReactNode }) => { const dispatch = useAppDispatch(); const { refetch: refetchSnapshot } = useRefetchSnapshotOnTurnEnd(); @@ -485,22 +464,6 @@ const ChatRuntimeProvider = ({ children }: { children: React.ReactNode }) => { useEffect(() => { if (socketStatus !== 'connected') return; - const decorateEntry = (entry: ToolTimelineEntry): ToolTimelineEntry => { - const formatted = formatTimelineEntry(entry); - // The server now attaches a human label/detail for dynamic - // Composio/MCP/integration tools the client can't know. Trust it for - // those; for the fixed set of built-ins the client formatter labels - // well (with args-aware detail), the client label stays authoritative. - if (entry.displayName && !isKnownClientTool(entry.name)) { - return { - ...entry, - displayName: entry.displayName, - detail: entry.detail ?? formatted.detail, - }; - } - return { ...entry, displayName: formatted.title, detail: formatted.detail ?? entry.detail }; - }; - // When a turn ends, any follow-ups the user queued behind it are about to be // dispatched by the backend as fresh turns. Nothing else persists their // prompt — the web channel never writes user messages; the composer does @@ -699,88 +662,47 @@ const ChatRuntimeProvider = ({ children }: { children: React.ReactNode }) => { }) ); - const existing = store.getState().chatRuntime.toolTimelineByThread[event.thread_id] ?? []; - const pendingContext = findPendingDelegationContext(existing, event.round); - // Collapse the parent's `spawn_subagent`/`spawn_async_subagent`/`delegate_*` tool-call row into - // the subagent row so the timeline shows ONE entry per delegation - // instead of "Research" (the tool call) + "Researching" (the child). - // The tool call's prompt is carried onto the subagent as the parent's - // delegation message, which the drawer renders as the opening turn. - const base = pendingContext.spawnEntryId - ? existing.filter(e => e.id !== pendingContext.spawnEntryId) - : existing; + // Collapse the parent spawn/delegate row into the subagent row (one + // entry per delegation) — merge now lives in the reducer (Phase 3). dispatch( - setToolTimelineForThread({ + subagentSpawned({ threadId: event.thread_id, - entries: [ - ...base, - decorateEntry({ - id: `${event.thread_id}:subagent:${event.skill_id}:${event.tool_name}`, - name: `subagent:${event.tool_name}`, - round: event.round, - status: 'running', - detail: pendingContext.prompt, - sourceToolName: pendingContext.sourceToolName, - subagent: { - taskId: event.skill_id, - agentId: event.tool_name, - displayName: event.subagent?.display_name, - workerThreadId: event.subagent?.worker_thread_id, - mode: event.subagent?.mode, - dedicatedThread: event.subagent?.dedicated_thread, - prompt: pendingContext.prompt, - toolCalls: [], - transcript: [], - }, - }), - ], + round: event.round, + rowId: `${event.thread_id}:subagent:${event.skill_id}:${event.tool_name}`, + taskId: event.skill_id, + agentId: event.tool_name, + displayName: event.subagent?.display_name, + workerThreadId: event.subagent?.worker_thread_id, + mode: event.subagent?.mode, + dedicatedThread: event.subagent?.dedicated_thread, }) ); }, onSubagentAwaitingUser: (event: ChatSubagentDoneEvent) => { - const subagentRowId = `${event.thread_id}:subagent:${event.skill_id}:${event.tool_name}`; - const existing = store.getState().chatRuntime.toolTimelineByThread[event.thread_id] ?? []; - if (existing.length > 0) { - const entries = existing.map(entry => { - if (entry.id !== subagentRowId || entry.status !== 'running') return entry; - return decorateEntry({ - ...entry, - status: 'awaiting_user' as ToolTimelineEntryStatus, - subagent: entry.subagent - ? { ...entry.subagent, status: 'awaiting_user' } - : entry.subagent, - }); - }); - dispatch(setToolTimelineForThread({ threadId: event.thread_id, entries })); - } + dispatch( + subagentAwaitingUser({ + threadId: event.thread_id, + rowId: `${event.thread_id}:subagent:${event.skill_id}:${event.tool_name}`, + }) + ); }, onSubagentDone: (event: ChatSubagentDoneEvent) => { - const subagentRowId = `${event.thread_id}:subagent:${event.skill_id}:${event.tool_name}`; - const existing = store.getState().chatRuntime.toolTimelineByThread[event.thread_id] ?? []; - if (existing.length > 0) { - const entries = existing.map(entry => { - if (entry.id !== subagentRowId || entry.status !== 'running') return entry; - return decorateEntry({ - ...entry, - status: (event.success ? 'success' : 'error') as ToolTimelineEntryStatus, - subagent: entry.subagent - ? { - ...entry.subagent, - iterations: event.subagent?.iterations ?? entry.subagent.iterations, - elapsedMs: event.subagent?.elapsed_ms ?? entry.subagent.elapsedMs, - outputChars: event.subagent?.output_chars ?? entry.subagent.outputChars, - // Worktree isolation metadata (#3376) — present only for - // workers that ran with `isolation = "worktree"`. Drives the - // inline worktree row's open/diff/remove affordances. - worktreePath: event.subagent?.worktree_path ?? entry.subagent.worktreePath, - changedFiles: event.subagent?.changed_files ?? entry.subagent.changedFiles, - isDirty: event.subagent?.dirty_status ?? entry.subagent.isDirty, - } - : entry.subagent, - }); - }); - dispatch(setToolTimelineForThread({ threadId: event.thread_id, entries })); - } + // Worktree isolation metadata (#3376) — present only for workers that ran + // with `isolation = "worktree"`; drives the inline worktree row's + // open/diff/remove affordances. Undefined fields leave the row untouched. + dispatch( + subagentDone({ + threadId: event.thread_id, + rowId: `${event.thread_id}:subagent:${event.skill_id}:${event.tool_name}`, + success: event.success, + iterations: event.subagent?.iterations, + elapsedMs: event.subagent?.elapsed_ms, + outputChars: event.subagent?.output_chars, + worktreePath: event.subagent?.worktree_path, + changedFiles: event.subagent?.changed_files, + isDirty: event.subagent?.dirty_status, + }) + ); const current = store.getState().chatRuntime.inferenceStatusByThread[event.thread_id]; if (!current) return; @@ -794,59 +716,35 @@ const ChatRuntimeProvider = ({ children }: { children: React.ReactNode }) => { onSubagentIterationStart: event => { const taskId = event.subagent?.task_id ?? event.skill_id; const agentId = event.subagent?.agent_id ?? event.tool_name; - const rowId = `${event.thread_id}:subagent:${taskId}:${agentId}`; - const existing = store.getState().chatRuntime.toolTimelineByThread[event.thread_id] ?? []; - const idx = existing.findIndex(entry => entry.id === rowId); - if (idx < 0) return; - const entry = existing[idx]; - if (!entry.subagent) return; - const next = [...existing]; - next[idx] = { - ...entry, - subagent: { - ...entry.subagent, - childIteration: event.subagent?.child_iteration ?? entry.subagent.childIteration, - childMaxIterations: - event.subagent?.child_max_iterations ?? entry.subagent.childMaxIterations, - }, - }; - dispatch(setToolTimelineForThread({ threadId: event.thread_id, entries: next })); + dispatch( + subagentIterationStarted({ + threadId: event.thread_id, + rowId: `${event.thread_id}:subagent:${taskId}:${agentId}`, + childIteration: event.subagent?.child_iteration, + childMaxIterations: event.subagent?.child_max_iterations, + }) + ); }, onSubagentToolCall: event => { const taskId = event.subagent?.task_id ?? event.skill_id; const agentId = event.subagent?.agent_id; if (!agentId) return; const rowId = `${event.thread_id}:subagent:${taskId}:${agentId}`; - const existing = store.getState().chatRuntime.toolTimelineByThread[event.thread_id] ?? []; - const idx = existing.findIndex(entry => entry.id === rowId); - if (idx < 0) return; - const entry = existing[idx]; - if (!entry.subagent) return; - // De-dupe on call_id — the same call should not append twice if - // the socket layer redelivers (e.g. on reconnect during a run). - if (entry.subagent.toolCalls.some(c => c.callId === event.tool_call_id)) return; - const next = [...existing]; - next[idx] = { - ...entry, - subagent: { - ...entry.subagent, - toolCalls: [ - ...entry.subagent.toolCalls, - { - callId: event.tool_call_id, - toolName: event.tool_name, - status: 'running', - iteration: event.subagent?.child_iteration, - args: event.args, - displayName: event.tool_display_label, - detail: event.tool_display_detail, - }, - ], - }, - }; - dispatch(setToolTimelineForThread({ threadId: event.thread_id, entries: next })); - // Mirror the call into the ordered transcript so the drawer renders - // it right after the text that triggered it (chronological view). + // Reducer owns the toolCalls upsert (dedup on call_id) — no getState(). + dispatch( + subagentToolCallReceived({ + threadId: event.thread_id, + rowId, + callId: event.tool_call_id, + toolName: event.tool_name, + iteration: event.subagent?.child_iteration, + args: event.args, + displayName: event.tool_display_label, + detail: event.tool_display_detail, + }) + ); + // Mirror the call into the ordered transcript so the drawer renders it + // right after the text that triggered it (self-guarded / self-deduped). dispatch( recordSubagentTranscriptTool({ threadId: event.thread_id, @@ -898,28 +796,19 @@ const ChatRuntimeProvider = ({ children }: { children: React.ReactNode }) => { const agentId = event.subagent?.agent_id; if (!agentId) return; const rowId = `${event.thread_id}:subagent:${taskId}:${agentId}`; - const existing = store.getState().chatRuntime.toolTimelineByThread[event.thread_id] ?? []; - const idx = existing.findIndex(entry => entry.id === rowId); - if (idx < 0) return; - const entry = existing[idx]; - if (!entry.subagent) return; - const callIdx = entry.subagent.toolCalls.findIndex(c => c.callId === event.tool_call_id); - if (callIdx < 0) return; - const updatedCalls = [...entry.subagent.toolCalls]; - updatedCalls[callIdx] = { - ...updatedCalls[callIdx], - status: event.success ? 'success' : 'error', - elapsedMs: event.subagent?.elapsed_ms ?? updatedCalls[callIdx].elapsedMs, - outputChars: event.subagent?.output_chars ?? updatedCalls[callIdx].outputChars, - result: event.output ?? updatedCalls[callIdx].result, - // Carry the structured failure so the child row keeps its "why / next" - // copy live instead of losing it until a snapshot reload (#4459). A - // successful result clears any stale failure on the row. - failure: event.success ? undefined : parseToolFailure(event.failure), - }; - const next = [...existing]; - next[idx] = { ...entry, subagent: { ...entry.subagent, toolCalls: updatedCalls } }; - dispatch(setToolTimelineForThread({ threadId: event.thread_id, entries: next })); + // Reducer owns the nested toolCall settle (no-op if the call is absent). + dispatch( + subagentToolResultReceived({ + threadId: event.thread_id, + rowId, + callId: event.tool_call_id, + success: event.success, + elapsedMs: event.subagent?.elapsed_ms, + outputChars: event.subagent?.output_chars, + result: event.output, + failure: event.failure, + }) + ); dispatch( resolveSubagentTranscriptTool({ threadId: event.thread_id, @@ -1051,45 +940,16 @@ const ChatRuntimeProvider = ({ children }: { children: React.ReactNode }) => { ); }, onToolArgsDelta: event => { - const cr = store.getState().chatRuntime; - const existing = cr.toolTimelineByThread[event.thread_id] ?? []; - let matchIdx = -1; - if (event.tool_call_id) { - matchIdx = existing.findIndex(entry => entry.id === event.tool_call_id); - } - if (matchIdx < 0 && event.tool_name) { - matchIdx = existing.findIndex( - entry => - entry.status === 'running' && - entry.name === event.tool_name && - entry.round === event.round - ); - } - - let entries: ToolTimelineEntry[]; - if (matchIdx >= 0) { - entries = [...existing]; - entries[matchIdx] = decorateEntry({ - ...entries[matchIdx], - argsBuffer: `${entries[matchIdx].argsBuffer ?? ''}${event.delta}`, - name: - entries[matchIdx].name.length === 0 && event.tool_name - ? event.tool_name - : entries[matchIdx].name, - }); - } else { - entries = [ - ...existing, - decorateEntry({ - id: event.tool_call_id, - name: event.tool_name ?? '', - round: event.round, - status: 'running', - argsBuffer: event.delta, - }), - ]; - } - dispatch(setToolTimelineForThread({ threadId: event.thread_id, entries })); + // Match + append + decorate now live in the reducer (Phase 3). + dispatch( + toolArgsDeltaReceived({ + threadId: event.thread_id, + round: event.round, + delta: event.delta, + toolName: event.tool_name, + toolCallId: event.tool_call_id, + }) + ); }, onTaskBoardUpdated: (event: ChatTaskBoardUpdatedEvent) => { if (!event.task_board) return; diff --git a/app/src/providers/__tests__/ChatRuntimeProvider.test.tsx b/app/src/providers/__tests__/ChatRuntimeProvider.test.tsx index 77d4c6311..353a349a2 100644 --- a/app/src/providers/__tests__/ChatRuntimeProvider.test.tsx +++ b/app/src/providers/__tests__/ChatRuntimeProvider.test.tsx @@ -9,13 +9,14 @@ import { store } from '../../store'; import { clearAllChatRuntime, enqueueFollowup, + findPendingDelegationContext, registerParallelRequest, resetSessionTokenUsage, setPendingPlanReviewForThread, } from '../../store/chatRuntimeSlice'; import { setStatusForUser } from '../../store/socketSlice'; import { clearAllThreads, loadThreads, setSelectedThread } from '../../store/threadSlice'; -import ChatRuntimeProvider, { findPendingDelegationContext } from '../ChatRuntimeProvider'; +import ChatRuntimeProvider from '../ChatRuntimeProvider'; vi.mock('../../services/chatService', async () => { const actual = await vi.importActual('../../services/chatService'); diff --git a/app/src/store/__tests__/chatRuntimeSlice.test.ts b/app/src/store/__tests__/chatRuntimeSlice.test.ts index 4e1ab4830..36cf57ee3 100644 --- a/app/src/store/__tests__/chatRuntimeSlice.test.ts +++ b/app/src/store/__tests__/chatRuntimeSlice.test.ts @@ -26,6 +26,13 @@ import reducer, { setTaskBoardForThread, setToolTimelineForThread, streamDeltaReceived, + subagentAwaitingUser, + subagentDone, + subagentIterationStarted, + subagentSpawned, + subagentToolCallReceived, + subagentToolResultReceived, + toolArgsDeltaReceived, toolCallReceived, toolResultReceived, upsertArtifactFailedForThread, @@ -1072,3 +1079,193 @@ describe('streamDeltaReceived (Phase 3 reducer-side merge)', () => { expect(state.processingByThread['t1']).toBeUndefined(); }); }); + +describe('subagent event reducers (Phase 3)', () => { + const spawn = (threadId = 't1') => + reducer( + undefined, + subagentSpawned({ + threadId, + round: 0, + rowId: 't1:subagent:task-1:researcher', + taskId: 'task-1', + agentId: 'researcher', + displayName: 'Researcher', + }) + ); + + it('subagentSpawned collapses the parent spawn row into the subagent row', () => { + // Seed a running parent delegate row for round 0. + let state = reducer( + undefined, + setToolTimelineForThread({ + threadId: 't1', + entries: [ + { + id: 'spawn-1', + name: 'spawn_subagent', + round: 0, + status: 'running', + detail: 'go research', + }, + ], + }) + ); + state = reducer( + state, + subagentSpawned({ + threadId: 't1', + round: 0, + rowId: 't1:subagent:task-1:researcher', + taskId: 'task-1', + agentId: 'researcher', + }) + ); + const rows = state.toolTimelineByThread['t1']; + expect(rows).toHaveLength(1); + expect(rows[0]).toMatchObject({ + id: 't1:subagent:task-1:researcher', + name: 'subagent:researcher', + status: 'running', + detail: 'go research', // carried from the collapsed spawn row's prompt + }); + expect(rows[0].subagent).toMatchObject({ taskId: 'task-1', agentId: 'researcher' }); + }); + + it('subagentToolCallReceived appends and de-dupes on callId; result settles it', () => { + let state = spawn(); + const row = 't1:subagent:task-1:researcher'; + state = reducer( + state, + subagentToolCallReceived({ threadId: 't1', rowId: row, callId: 'c1', toolName: 'grep' }) + ); + // Redelivery is a no-op. + state = reducer( + state, + subagentToolCallReceived({ threadId: 't1', rowId: row, callId: 'c1', toolName: 'grep' }) + ); + let sub = state.toolTimelineByThread['t1'][0].subagent!; + expect(sub.toolCalls).toHaveLength(1); + expect(sub.toolCalls[0].status).toBe('running'); + + state = reducer( + state, + subagentToolResultReceived({ + threadId: 't1', + rowId: row, + callId: 'c1', + success: true, + result: 'ok', + }) + ); + sub = state.toolTimelineByThread['t1'][0].subagent!; + expect(sub.toolCalls[0]).toMatchObject({ status: 'success', result: 'ok' }); + }); + + it('subagentDone settles the row + metadata; awaiting/iteration update in place', () => { + let state = spawn(); + const row = 't1:subagent:task-1:researcher'; + state = reducer( + state, + subagentIterationStarted({ + threadId: 't1', + rowId: row, + childIteration: 2, + childMaxIterations: 5, + }) + ); + expect(state.toolTimelineByThread['t1'][0].subagent).toMatchObject({ + childIteration: 2, + childMaxIterations: 5, + }); + + const awaiting = reducer(state, subagentAwaitingUser({ threadId: 't1', rowId: row })); + expect(awaiting.toolTimelineByThread['t1'][0].status).toBe('awaiting_user'); + + const done = reducer( + state, + subagentDone({ threadId: 't1', rowId: row, success: true, iterations: 3, elapsedMs: 42 }) + ); + expect(done.toolTimelineByThread['t1'][0]).toMatchObject({ status: 'success' }); + expect(done.toolTimelineByThread['t1'][0].subagent).toMatchObject({ + iterations: 3, + elapsedMs: 42, + }); + }); + + it('settles an awaiting_user row on done (it must not stay stuck)', () => { + let state = spawn(); + const row = 't1:subagent:task-1:researcher'; + state = reducer(state, subagentAwaitingUser({ threadId: 't1', rowId: row })); + expect(state.toolTimelineByThread['t1'][0].status).toBe('awaiting_user'); + state = reducer(state, subagentDone({ threadId: 't1', rowId: row, success: true })); + expect(state.toolTimelineByThread['t1'][0].status).toBe('success'); + }); + + it('done is a no-op once the row is terminal (already settled)', () => { + let state = spawn(); + const row = 't1:subagent:task-1:researcher'; + state = reducer(state, subagentDone({ threadId: 't1', rowId: row, success: true })); + // A second done cannot re-settle a terminal row. + const again = reducer(state, subagentDone({ threadId: 't1', rowId: row, success: false })); + expect(again.toolTimelineByThread['t1'][0].status).toBe('success'); + }); + + it('subagentSpawned is idempotent — a redelivered event does not duplicate the row', () => { + let state = spawn(); + const before = state.toolTimelineByThread['t1'].length; + state = reducer( + state, + subagentSpawned({ + threadId: 't1', + round: 0, + rowId: 't1:subagent:task-1:researcher', + taskId: 'task-1', + agentId: 'researcher', + }) + ); + expect(state.toolTimelineByThread['t1']).toHaveLength(before); + }); +}); + +describe('toolArgsDeltaReceived (Phase 3 reducer-side merge)', () => { + it('creates a running row when args arrive before the tool_call, then appends', () => { + let state = reducer( + undefined, + toolArgsDeltaReceived({ + threadId: 't1', + round: 0, + delta: '{"q":', + toolName: 'search', + toolCallId: 'c1', + }) + ); + state = reducer( + state, + toolArgsDeltaReceived({ threadId: 't1', round: 0, delta: '"hi"}', toolCallId: 'c1' }) + ); + const rows = state.toolTimelineByThread['t1']; + expect(rows).toHaveLength(1); + expect(rows[0]).toMatchObject({ + id: 'c1', + name: 'search', + status: 'running', + argsBuffer: '{"q":"hi"}', + }); + }); + + it('falls back to the newest running row of the same name+round when no id matches', () => { + let state = reducer( + undefined, + setToolTimelineForThread({ + threadId: 't1', + entries: [{ id: 'r1', name: 'search', round: 0, status: 'running', argsBuffer: '{' }], + }) + ); + state = reducer( + state, + toolArgsDeltaReceived({ threadId: 't1', round: 0, delta: '}', toolName: 'search' }) + ); + expect(state.toolTimelineByThread['t1'][0].argsBuffer).toBe('{}'); + }); +}); diff --git a/app/src/store/chatRuntimeSlice.ts b/app/src/store/chatRuntimeSlice.ts index 2b33f59dd..549daf451 100644 --- a/app/src/store/chatRuntimeSlice.ts +++ b/app/src/store/chatRuntimeSlice.ts @@ -13,7 +13,11 @@ import type { PersistedTurnState, TaskBoard, } from '../types/turnState'; -import { formatTimelineEntry, isKnownClientTool } from '../utils/toolTimelineFormatting'; +import { + formatTimelineEntry, + isKnownClientTool, + promptFromArgsBuffer, +} from '../utils/toolTimelineFormatting'; import { resetUserScopedState } from './resetActions'; const turnStateLog = debug('chatRuntime.turnState'); @@ -252,6 +256,33 @@ function decorateEntry(entry: ToolTimelineEntry): ToolTimelineEntry { return { ...entry, displayName: formatted.title, detail: formatted.detail ?? entry.detail }; } +/** + * Find the parent `spawn_*`/`delegate_*` tool row a just-spawned subagent should + * collapse into, so the timeline shows one entry per delegation. Returns the + * source tool name, the delegation prompt, and the spawn row's id to remove. + * Pure — searches the round's running rows newest-first. + */ +export function findPendingDelegationContext( + entries: ToolTimelineEntry[], + round: number +): { sourceToolName?: string; prompt?: string; spawnEntryId?: string } { + for (let i = entries.length - 1; i >= 0; i -= 1) { + const entry = entries[i]; + if (entry.status !== 'running' || entry.round !== round) continue; + if ( + ['spawn_subagent', 'spawn_async_subagent'].includes(entry.name) || + entry.name.startsWith('delegate_') + ) { + return { + sourceToolName: entry.name, + prompt: entry.detail ?? promptFromArgsBuffer(entry.argsBuffer), + spawnEntryId: entry.id, + }; + } + } + return {}; +} + export interface ToolTimelineEntry { id: string; name: string; @@ -1208,6 +1239,236 @@ const chatRuntimeSlice = createSlice({ list.push({ kind, round, seq: list.length, text: delta }); } }, + /** + * Reducer-side merge for a `tool_args_delta` socket event (Phase 3). + * Appends the streamed args to the matching row (by `toolCallId`, else the + * newest running row of the same name+round), or creates a running row when + * the args arrive before the tool-call event. Re-decorates each time. + */ + toolArgsDeltaReceived: ( + state, + action: PayloadAction<{ + threadId: string; + round: number; + delta: string; + toolName?: string; + toolCallId?: string; + }> + ) => { + const { threadId, round, delta, toolName, toolCallId } = action.payload; + const entries = (state.toolTimelineByThread[threadId] ??= []); + let matchIdx = -1; + if (toolCallId) matchIdx = entries.findIndex(e => e.id === toolCallId); + if (matchIdx < 0 && toolName) { + matchIdx = entries.findIndex( + e => e.status === 'running' && e.name === toolName && e.round === round + ); + } + if (matchIdx >= 0) { + const prev = entries[matchIdx]; + entries[matchIdx] = decorateEntry({ + ...prev, + argsBuffer: `${prev.argsBuffer ?? ''}${delta}`, + name: prev.name.length === 0 && toolName ? toolName : prev.name, + }); + } else { + entries.push( + decorateEntry({ + id: toolCallId ?? '', + name: toolName ?? '', + round, + status: 'running', + argsBuffer: delta, + }) + ); + } + }, + /** + * Reducer-side merges for the sub-agent event family (Phase 3). Each locates + * the delegation's timeline row by its precomputed `rowId` and updates the + * nested `subagent` activity in place — no `getState()` / full-array rebuild + * in the provider. + */ + subagentSpawned: ( + state, + action: PayloadAction<{ + threadId: string; + round: number; + rowId: string; + taskId: string; + agentId: string; + displayName?: string; + workerThreadId?: string; + mode?: string; + dedicatedThread?: boolean; + }> + ) => { + const { + threadId, + round, + rowId, + taskId, + agentId, + displayName, + workerThreadId, + mode, + dedicatedThread, + } = action.payload; + const entries = (state.toolTimelineByThread[threadId] ??= []); + // Idempotent: a socket redelivery must not append a second row with the + // same id (later updates find only the first). Not gated by the provider's + // event-seen map, so guard here. + if (entries.some(e => e.id === rowId)) return; + const pending = findPendingDelegationContext(entries, round); + // Collapse the parent spawn/delegate row into the subagent row so the + // timeline shows one entry per delegation. + if (pending.spawnEntryId) { + const spawnIdx = entries.findIndex(e => e.id === pending.spawnEntryId); + if (spawnIdx >= 0) entries.splice(spawnIdx, 1); + } + entries.push( + decorateEntry({ + id: rowId, + name: `subagent:${agentId}`, + round, + status: 'running', + detail: pending.prompt, + sourceToolName: pending.sourceToolName, + subagent: { + taskId, + agentId, + displayName, + workerThreadId, + mode, + dedicatedThread, + prompt: pending.prompt, + toolCalls: [], + transcript: [], + }, + }) + ); + }, + subagentAwaitingUser: (state, action: PayloadAction<{ threadId: string; rowId: string }>) => { + const entry = state.toolTimelineByThread[action.payload.threadId]?.find( + e => e.id === action.payload.rowId && e.status === 'running' + ); + if (!entry) return; + entry.status = 'awaiting_user'; + if (entry.subagent) entry.subagent.status = 'awaiting_user'; + }, + subagentDone: ( + state, + action: PayloadAction<{ + threadId: string; + rowId: string; + success: boolean; + iterations?: number; + elapsedMs?: number; + outputChars?: number; + worktreePath?: string; + changedFiles?: string[]; + isDirty?: boolean; + }> + ) => { + const { + threadId, + rowId, + success, + iterations, + elapsedMs, + outputChars, + worktreePath, + changedFiles, + isDirty, + } = action.payload; + // Settle a still-in-flight row: `running`, or `awaiting_user` (a subagent + // paused for input that then completes must not stay stuck at + // awaiting_user). Already-terminal rows are left as-is. + const entry = state.toolTimelineByThread[threadId]?.find( + e => e.id === rowId && (e.status === 'running' || e.status === 'awaiting_user') + ); + if (!entry) return; + entry.status = success ? 'success' : 'error'; + if (entry.subagent) { + const s = entry.subagent; + if (iterations !== undefined) s.iterations = iterations; + if (elapsedMs !== undefined) s.elapsedMs = elapsedMs; + if (outputChars !== undefined) s.outputChars = outputChars; + if (worktreePath !== undefined) s.worktreePath = worktreePath; + if (changedFiles !== undefined) s.changedFiles = changedFiles; + if (isDirty !== undefined) s.isDirty = isDirty; + } + }, + subagentIterationStarted: ( + state, + action: PayloadAction<{ + threadId: string; + rowId: string; + childIteration?: number; + childMaxIterations?: number; + }> + ) => { + const { threadId, rowId, childIteration, childMaxIterations } = action.payload; + const entry = state.toolTimelineByThread[threadId]?.find(e => e.id === rowId); + if (!entry?.subagent) return; + if (childIteration !== undefined) entry.subagent.childIteration = childIteration; + if (childMaxIterations !== undefined) entry.subagent.childMaxIterations = childMaxIterations; + }, + subagentToolCallReceived: ( + state, + action: PayloadAction<{ + threadId: string; + rowId: string; + callId: string; + toolName: string; + iteration?: number; + args?: unknown; + displayName?: string; + detail?: string; + }> + ) => { + const { threadId, rowId, callId, toolName, iteration, args, displayName, detail } = + action.payload; + const entry = state.toolTimelineByThread[threadId]?.find(e => e.id === rowId); + if (!entry?.subagent) return; + // De-dupe on call_id — a redelivered event must not append twice. + if (entry.subagent.toolCalls.some(c => c.callId === callId)) return; + entry.subagent.toolCalls.push({ + callId, + toolName, + status: 'running', + iteration, + args, + displayName, + detail, + }); + }, + subagentToolResultReceived: ( + state, + action: PayloadAction<{ + threadId: string; + rowId: string; + callId: string; + success: boolean; + elapsedMs?: number; + outputChars?: number; + result?: string; + failure?: unknown; + }> + ) => { + const { threadId, rowId, callId, success, elapsedMs, outputChars, result, failure } = + action.payload; + const entry = state.toolTimelineByThread[threadId]?.find(e => e.id === rowId); + if (!entry?.subagent) return; + const call = entry.subagent.toolCalls.find(c => c.callId === callId); + if (!call) return; + call.status = success ? 'success' : 'error'; + if (elapsedMs !== undefined) call.elapsedMs = elapsedMs; + if (outputChars !== undefined) call.outputChars = outputChars; + if (result !== undefined) call.result = result; + // A successful result clears any stale failure on the row. + call.failure = success ? undefined : parseToolFailure(failure); + }, /** * Optimistically mark a detached background sub-agent as cancelled after the * user confirms a cancel via `openhuman.subagent_cancel`. The aborted run @@ -1808,6 +2069,13 @@ export const { clearToolTimelineForThread, setTurnTimelinesForThread, streamDeltaReceived, + subagentAwaitingUser, + subagentDone, + subagentIterationStarted, + subagentSpawned, + subagentToolCallReceived, + subagentToolResultReceived, + toolArgsDeltaReceived, toolCallReceived, toolResultReceived, clearProcessingForThread, diff --git a/app/test/playwright/specs/tinyplace-messaging.spec.ts b/app/test/playwright/specs/tinyplace-messaging.spec.ts index 7ff9e3b54..14ebbb7b3 100644 --- a/app/test/playwright/specs/tinyplace-messaging.spec.ts +++ b/app/test/playwright/specs/tinyplace-messaging.spec.ts @@ -16,11 +16,11 @@ // Requires the web session harness (app/scripts/e2e-web-session.sh) with // TINYPLACE_API_BASE_URL exported so the core hits a real backend. The // messaging e2e runner (e2e/tinyplace-messaging/run-ui.sh) wires this up. -import { test, expect } from '@playwright/test'; +import { expect, test } from '@playwright/test'; -import { bootAuthenticatedPage } from '../helpers/core-rpc'; // The core-launch helper is shared with the core-level suite (plain ESM). import { launchAgent, receiveMessage } from '../../../../e2e/tinyplace-messaging/lib/core.mjs'; +import { bootAuthenticatedPage } from '../helpers/core-rpc'; const CORE_RPC_URL = process.env.PW_CORE_RPC_URL || 'http://127.0.0.1:17788/rpc'; const CORE_RPC_TOKEN = process.env.PW_CORE_RPC_TOKEN || 'openhuman-playwright-token'; @@ -37,10 +37,26 @@ async function freshMnemonic(): Promise { } const PLACEHOLDER_ACCOUNTS = [ - { chain: 'evm', address: '0x0000000000000000000000000000000000000001', derivationPath: "m/44'/60'/0'/0/0" }, - { chain: 'btc', address: 'bc1qplaceholderplaceholderplaceholderplac0000', derivationPath: "m/84'/0'/0'/0/0" }, - { chain: 'solana', address: '11111111111111111111111111111111', derivationPath: "m/44'/501'/0'/0'" }, - { chain: 'tron', address: 'T0000000000000000000000000000000001', derivationPath: "m/44'/195'/0'/0/0" }, + { + chain: 'evm', + address: '0x0000000000000000000000000000000000000001', + derivationPath: "m/44'/60'/0'/0/0", + }, + { + chain: 'btc', + address: 'bc1qplaceholderplaceholderplaceholderplac0000', + derivationPath: "m/84'/0'/0'/0/0", + }, + { + chain: 'solana', + address: '11111111111111111111111111111111', + derivationPath: "m/44'/501'/0'/0'", + }, + { + chain: 'tron', + address: 'T0000000000000000000000000000000001', + derivationPath: "m/44'/195'/0'/0/0", + }, ]; /** Call Alice's (the app's) core over JSON-RPC and unwrap the {logs,result}. */ @@ -68,7 +84,9 @@ test.describe('tiny.place direct messaging (UI)', () => { test.beforeAll(async () => { // 1) Give the app's core a fresh tiny.place identity + published Signal keys. const mnemonic = await freshMnemonic(); - const encryptedMnemonic = await aliceRpc('openhuman.encrypt_secret', { plaintext: mnemonic }); + const encryptedMnemonic = await aliceRpc('openhuman.encrypt_secret', { + plaintext: mnemonic, + }); await aliceRpc('openhuman.wallet_setup', { consentGranted: true, source: 'imported', @@ -94,7 +112,9 @@ test.describe('tiny.place direct messaging (UI)', () => { bob?.stop(); }); - test('sends an encrypted DM from the UI that the peer decrypts, and renders the peer reply', async ({ page }) => { + test('sends an encrypted DM from the UI that the peer decrypts, and renders the peer reply', async ({ + page, + }) => { await bootAuthenticatedPage(page, 'pw-messaging-user', '/agent-world/messaging'); // The DM composer: enter the peer's cryptoId and open the thread. @@ -122,7 +142,10 @@ test.describe('tiny.place direct messaging (UI)', () => { // Now the peer replies; the plaintext must render in the UI thread. const reply = `peer → ui @ ${Date.now()}`; - await bob.rpc('openhuman.tinyplace_signal_send_message', { recipient: aliceCryptoId, plaintext: reply }); + await bob.rpc('openhuman.tinyplace_signal_send_message', { + recipient: aliceCryptoId, + plaintext: reply, + }); // Wait until the reply envelope is actually in Alice's mailbox — inspected, // not decrypted (decrypt advances the ratchet and can only run once, so we @@ -131,10 +154,10 @@ test.describe('tiny.place direct messaging (UI)', () => { .poll( async () => { const list = await aliceRpc('openhuman.tinyplace_messages_list', { limit: 50 }); - const envelopes = Array.isArray(list) ? list : list?.messages ?? []; + const envelopes = Array.isArray(list) ? list : (list?.messages ?? []); return envelopes.some((e: any) => e.from === bob.cryptoId); }, - { timeout: 15_000, intervals: [500, 1000] }, + { timeout: 15_000, intervals: [500, 1000] } ) .toBe(true);