fix(agents): stream child runs in shared UI

This commit is contained in:
程序员阿江(Relakkes)
2026-08-11 14:44:05 +08:00
parent 3cb9e24faa
commit ab0d89e16e
24 changed files with 2129 additions and 39 deletions
+270 -2
View File
@@ -57,7 +57,7 @@ vi.mock('../api/workflows', async (importOriginal) => {
import { subagentsApi } from '../api/subagents'
import { createDefaultSessionState, useChatStore } from '../stores/chatStore'
import { useActivityPanelStore } from '../stores/activityPanelStore'
import { useTabStore } from '../stores/tabStore'
import { SUBAGENT_TAB_PREFIX, useTabStore } from '../stores/tabStore'
import { memberSessionId, useTeamStore } from '../stores/teamStore'
import { useWorkflowStore } from '../stores/workflowStore'
import { SubagentRunPage, TeamMemberRunPage } from './SubagentRunPage'
@@ -203,6 +203,161 @@ describe('SubagentRunPage', () => {
expectSharedSessionSurface('subagent')
})
it.each([
['foreground SubAgent', 'tool-direct', 'direct-agent', false],
['background Agent', 'tool-background', 'background-agent', false],
['workflow Agent', 'agent:workflow-agent', 'workflow-agent', true],
] as const)('renders %s text, thinking and tools before the next transcript poll', async (
_kind,
toolUseId,
agentId,
byAgentId,
) => {
const response = subagentRun({
toolUseId,
agentId,
status: 'running',
messages: [],
prompt: `Prompt for ${agentId}`,
})
if (byAgentId) {
vi.mocked(subagentsApi.getRunByAgent).mockResolvedValue(response)
} else {
vi.mocked(subagentsApi.getRunByTool).mockResolvedValue(response)
}
render(
<SubagentRunPage
sourceSessionId="session-1"
toolUseId={toolUseId}
title={agentId}
/>,
)
const conversation = await screen.findByTestId('subagent-conversation')
const send = (
event: Extract<import('../types/chat').ServerMessage, { type: 'agent_run_event' }>['event'],
) => useChatStore.getState().handleServerMessage('session-1', {
type: 'agent_run_event',
runAgentId: agentId,
streamId: `stream-${agentId}`,
targetAgentId: agentId,
event,
})
act(() => {
send({ type: 'thinking', text: `Live thinking from ${agentId}` })
send({ type: 'content_start', blockType: 'text' })
send({ type: 'content_delta', text: `Live answer from ${agentId}` })
send({
type: 'content_start',
blockType: 'tool_use',
toolName: 'Read',
toolUseId: 'live-read',
})
send({ type: 'content_delta', toolInput: '{"file_path":"src/live.ts"}' })
})
const runSessionId = `${SUBAGENT_TAB_PREFIX}session-1__${toolUseId}`
expect(useChatStore.getState().sessions[runSessionId]?.activeToolName).toBe('Read')
await waitFor(() => {
expect(useChatStore.getState().sessions[runSessionId]?.streamingToolInput)
.toContain('src/live.ts')
})
act(() => {
send({
type: 'tool_use_complete',
toolName: 'Read',
toolUseId: 'live-read',
input: { file_path: 'src/live.ts' },
})
send({
type: 'tool_result',
toolUseId: 'live-read',
content: 'live file contents',
isError: false,
})
})
expect(conversation).toHaveTextContent(`Live thinking from ${agentId}`)
expect(conversation).toHaveTextContent(`Live answer from ${agentId}`)
expect(conversation).toHaveTextContent('live.ts')
expect(conversation).toHaveTextContent('live file contents')
expect(useChatStore.getState().sessions['session-1']?.messages ?? []).toEqual([])
expect(byAgentId ? subagentsApi.getRunByAgent : subagentsApi.getRunByTool)
.toHaveBeenCalledTimes(1)
fireEvent.click(screen.getByRole('button', { name: 'Refresh SubAgent run' }))
await waitFor(() => {
expect(byAgentId ? subagentsApi.getRunByAgent : subagentsApi.getRunByTool)
.toHaveBeenCalledTimes(2)
})
expect(conversation).toHaveTextContent(`Live answer from ${agentId}`)
expect(conversation).toHaveTextContent('live file contents')
})
it('keeps a live SubAgent turn across a stale poll, then reconciles the next durable poll', async () => {
const stalePoll = deferred<SubagentRunResponse>()
vi.mocked(subagentsApi.getRunByTool)
.mockResolvedValueOnce(subagentRun({
status: 'running',
messages: [],
}))
.mockReturnValueOnce(stalePoll.promise)
.mockResolvedValueOnce(subagentRun({
status: 'completed',
messages: [{
id: 'durable-live-turn',
type: 'assistant',
content: [{ type: 'text', text: 'Durable answer after the live turn' }],
timestamp: TRANSCRIPT_TIMESTAMP,
}],
}))
render(
<SubagentRunPage
sourceSessionId="session-1"
toolUseId="tool-1"
title="SubAgent"
/>,
)
const conversation = await screen.findByTestId('subagent-conversation')
fireEvent.click(screen.getByRole('button', { name: 'Refresh SubAgent run' }))
await waitFor(() => expect(subagentsApi.getRunByTool).toHaveBeenCalledTimes(2))
act(() => {
const send = (
event: Extract<import('../types/chat').ServerMessage, { type: 'agent_run_event' }>['event'],
) => useChatStore.getState().handleServerMessage('session-1', {
type: 'agent_run_event',
runAgentId: 'abc123',
streamId: 'stale-boundary-stream',
targetAgentId: 'abc123',
event,
})
send({ type: 'content_start', blockType: 'text' })
send({ type: 'content_delta', text: 'Live answer that the stale poll must not erase' })
send({ type: 'status', state: 'idle' })
})
expect(conversation).toHaveTextContent('Live answer that the stale poll must not erase')
await act(async () => {
stalePoll.resolve(subagentRun({ status: 'completed', messages: [] }))
await stalePoll.promise
})
expect(conversation).toHaveTextContent('Live answer that the stale poll must not erase')
const runSessionId = `${SUBAGENT_TAB_PREFIX}session-1__tool-1`
expect(useChatStore.getState().sessions[runSessionId]?.agentStreamRevision)
.toBe(2)
fireEvent.click(screen.getByRole('button', { name: 'Refresh SubAgent run' }))
await waitFor(() => expect(subagentsApi.getRunByTool).toHaveBeenCalledTimes(3))
expect(await screen.findByText('Durable answer after the live turn')).toBeInTheDocument()
expect(conversation).not.toHaveTextContent('Live answer that the stale poll must not erase')
expect(useChatStore.getState().sessions[runSessionId]?.agentStreamRevision)
.toBe(0)
})
it('uses live workflow progress instead of sealing a running agent as completed', async () => {
vi.mocked(subagentsApi.getRunByAgent).mockResolvedValue(subagentRun({
agentId: 'abc123',
@@ -923,6 +1078,58 @@ describe('SubagentRunPage', () => {
expect(within(panel).queryByText('Do not show the whole shared list')).not.toBeInTheDocument()
})
it('routes a scoped nested teammate Agent into its shared run UI', async () => {
vi.mocked(subagentsApi.getRunByTool).mockResolvedValue(subagentRun({
agentId: 'nested-team-agent',
status: 'running',
messages: [],
}))
const createdAt = Date.parse('2026-08-09T00:00:00.000Z')
useChatStore.setState({ sessions: { 'session-1': createDefaultSessionState() } })
useTeamStore.setState({
workbenchesBySession: {
'session-1': {
teamName: 'review-team',
loading: false,
error: null,
snapshots: [{
version: 'review-team-v1',
generatedAt: TRANSCRIPT_TIMESTAMP,
team: {
name: 'review-team',
leadSessionId: 'session-1',
createdAt: String(createdAt),
members: [{
agentId: 'reviewer@review-team',
name: 'reviewer',
role: 'security-reviewer',
status: 'running',
}],
},
tasks: [],
messages: [],
}],
},
},
})
render(<SubagentRunPage sourceSessionId="session-1" toolUseId="nested-Agent:0" title="Nested review" />)
const conversation = await screen.findByTestId('subagent-conversation')
act(() => {
useChatStore.getState().handleServerMessage('session-1', {
type: 'agent_run_event',
runAgentId: 'nested-team-agent',
streamId: 'nested-team-stream',
targetAgentId: 'nested-team-agent',
targetAgentScopeId: JSON.stringify(['review-team', 'session-1', createdAt]),
event: { type: 'content_delta', text: 'Scoped nested answer' },
})
})
await waitFor(() => expect(conversation).toHaveTextContent('Scoped nested answer'))
expect(useChatStore.getState().sessions['session-1']?.streamingText).toBe('')
})
it('isolates a teammate shared task list before its workbench snapshot arrives', async () => {
vi.mocked(subagentsApi.getRunByTool).mockResolvedValue(subagentRun({
agentId: 'reviewer@review-team',
@@ -970,7 +1177,7 @@ describe('SubagentRunPage', () => {
expect(within(panel).queryByText('Never leak the early shared task')).not.toBeInTheDocument()
})
it('renders an Agent Teams member in the shared run desktop and returns to the workbench', async () => {
it('renders and streams an Agent Teams member in the shared run desktop', async () => {
const member = {
agentId: 'reviewer@review-team',
name: 'reviewer',
@@ -992,6 +1199,7 @@ describe('SubagentRunPage', () => {
incarnationId: 'review-team:2026-08-09:lead-session',
leadAgentId: 'lead@review-team',
leadSessionId: 'lead-session',
createdAt: String(Date.parse('2026-08-09T00:00:00.000Z')),
members: [member, peer],
},
tasks: [],
@@ -1067,6 +1275,66 @@ describe('SubagentRunPage', () => {
expect(await screen.findByTestId('team-member-conversation')).toHaveTextContent('Auth review is in progress.')
const transcript = screen.getByTestId('team-member-conversation')
act(() => {
const targetAgentScopeId = JSON.stringify([
snapshot.team.name,
snapshot.team.leadSessionId,
Number(snapshot.team.createdAt),
])
const send = (
event: Extract<import('../types/chat').ServerMessage, { type: 'agent_run_event' }>['event'],
) => useChatStore.getState().handleServerMessage('lead-session', {
type: 'agent_run_event',
runAgentId: 'reviewer-fragment-uuid',
streamId: 'teammate-live-stream',
targetAgentId: member.agentId,
targetAgentScopeId,
event,
})
send({ type: 'thinking', text: 'Live teammate thinking' })
send({ type: 'content_start', blockType: 'text' })
send({ type: 'content_delta', text: 'Live teammate answer' })
send({
type: 'content_start',
blockType: 'tool_use',
toolName: 'Read',
toolUseId: 'member-live-read',
})
send({
type: 'tool_use_complete',
toolName: 'Read',
toolUseId: 'member-live-read',
input: { file_path: 'src/member-live.ts' },
})
send({
type: 'tool_result',
toolUseId: 'member-live-read',
content: 'member live file contents',
isError: false,
})
})
expect(transcript).toHaveTextContent('Live teammate thinking')
expect(transcript).toHaveTextContent('Live teammate answer')
expect(transcript).toHaveTextContent('member-live.ts')
expect(transcript).toHaveTextContent('member live file contents')
expect(useChatStore.getState().sessions['lead-session']?.messages ?? []).toEqual([])
act(() => {
useChatStore.getState().handleServerMessage('lead-session', {
type: 'agent_run_event',
runAgentId: 'reviewer-fragment-uuid',
streamId: 'teammate-live-stream',
targetAgentId: member.agentId,
targetAgentScopeId: JSON.stringify([
snapshot.team.name,
snapshot.team.leadSessionId,
Number(snapshot.team.createdAt),
]),
event: { type: 'status', state: 'idle' },
})
})
expect(transcript).toHaveTextContent('Live teammate answer')
expect(transcript).toHaveTextContent('member live file contents')
expect(transcript).toHaveTextContent('Prioritize the auth flow.')
expect(transcript).toHaveTextContent('Check src/auth.ts before merge.')
// Drive the real transcript adapter: both lead-to-member and member-to-member
+74 -14
View File
@@ -45,6 +45,17 @@ const EMPTY_DISMISSED_BACKGROUND_TASK_KEYS: string[] = []
const EMPTY_RUN_TASKS: CLITask[] = []
const EMPTY_OWNER_AGENT_IDS: string[] = []
function teamStreamScopeId(team: {
name: string
leadSessionId?: string
createdAt?: string
} | undefined): string | undefined {
if (!team?.createdAt) return undefined
const createdAt = Number(team.createdAt)
if (!Number.isFinite(createdAt)) return undefined
return JSON.stringify([team.name, team.leadSessionId ?? '', createdAt])
}
export function SubagentRunPage({
sourceSessionId,
toolUseId,
@@ -59,6 +70,7 @@ export function SubagentRunPage({
const t = useTranslation()
const [data, setData] = useState<SubagentRunResponse | null>(null)
const [responseActivityEpoch, setResponseActivityEpoch] = useState(0)
const [responseStreamRevision, setResponseStreamRevision] = useState(0)
const [loading, setLoading] = useState(true)
const [error, setError] = useState<string | null>(null)
const requestIdRef = useRef(0)
@@ -156,7 +168,9 @@ export function SubagentRunPage({
const load = useCallback(async (options?: { resetData?: boolean }) => {
const requestId = requestIdRef.current + 1
requestIdRef.current = requestId
const requestedActivityEpoch = useChatStore.getState().sessions[tabId]?.historyMutationEpoch ?? 0
const requestedSession = useChatStore.getState().sessions[tabId]
const requestedActivityEpoch = requestedSession?.historyMutationEpoch ?? 0
const requestedStreamRevision = requestedSession?.agentStreamRevision ?? 0
setLoading(true)
setError(null)
if (options?.resetData) setData(null)
@@ -169,6 +183,7 @@ export function SubagentRunPage({
if (requestIdRef.current !== requestId) return
setData(nextData)
setResponseActivityEpoch(requestedActivityEpoch)
setResponseStreamRevision(requestedStreamRevision)
} catch (err) {
if (requestIdRef.current !== requestId) return
setError(err instanceof Error ? err.message : String(err))
@@ -214,32 +229,67 @@ export function SubagentRunPage({
}, runActivity, {
preferCurrent: (existing.historyMutationEpoch ?? 0) !== responseActivityEpoch,
})
const currentStreamRevision = existing.agentStreamRevision ?? 0
const preserveLiveConversation = currentStreamRevision > 0 && (
existing.chatState !== 'idle' ||
currentStreamRevision !== responseStreamRevision
)
return {
sessions: {
...state.sessions,
[tabId]: {
...existing,
messages: [...transcriptMessages, ...localMessages],
messages: preserveLiveConversation
? existing.messages
: [...transcriptMessages, ...localMessages],
agentTaskNotifications: mergedActivity.agentTaskNotifications,
backgroundAgentTasks: mergedActivity.backgroundAgentTasks,
connectionState: 'connected',
chatState: effectiveStatus === 'running' || hasPendingMessage
? 'thinking'
: 'idle',
chatState: preserveLiveConversation
? existing.chatState
: effectiveStatus === 'running' || hasPendingMessage
? 'thinking'
: 'idle',
...(!preserveLiveConversation ? {
agentStreamRevision: 0,
streamingText: '',
streamingToolInput: '',
activeToolUseId: null,
activeToolName: null,
activeThinkingId: null,
} : {}),
},
},
}
})
}, [data, effectiveStatus, responseActivityEpoch, tabId])
}, [data, effectiveStatus, responseActivityEpoch, responseStreamRevision, tabId])
useEffect(() => {
if (!data?.agentId) return
return registerAgentRunSession(sourceSessionId, tabId, [data.agentId], {
...(teamMember || teamAgentName
? { eventIdPrefix: data.agentId }
: {}),
})
}, [data?.agentId, sourceSessionId, tabId, teamAgentName, teamMember])
const streamScopeId = teamStreamScopeId(teamSnapshot?.team)
const usesTeamFragmentIds = Boolean(teamMember || teamAgentName)
const unregisterUnscoped = registerAgentRunSession(
sourceSessionId,
tabId,
[data.agentId],
usesTeamFragmentIds ? { eventIdPrefix: data.agentId } : {},
)
const unregisterScoped = streamScopeId
? registerAgentRunSession(sourceSessionId, tabId, [data.agentId], {
streamScopeId,
...(usesTeamFragmentIds
? {
eventIdPrefix: data.agentId,
streamEventIdPrefix: data.agentId,
}
: {}),
})
: () => undefined
return () => {
unregisterScoped()
unregisterUnscoped()
}
}, [data?.agentId, sourceSessionId, tabId, teamAgentName, teamMember, teamSnapshot?.team])
useEffect(() => {
if (workflowOwnerAliases.length === 0) return
@@ -338,10 +388,15 @@ export function TeamMemberRunPage({
useEffect(() => {
if (!member) return
const memberTeam = snapshot?.team ?? useTeamStore.getState().getTeamByMemberSessionId(tabId)
const streamScopeId = teamStreamScopeId(memberTeam ?? undefined)
const unregisterLogicalOwners = registerAgentRunSession(leadSessionId, runSessionId, [
member.agentId,
member.name,
], { ownerScopeId: memberTeam ? teamIdentityKey(memberTeam) : undefined })
], {
ownerScopeId: memberTeam ? teamIdentityKey(memberTeam) : undefined,
...(streamScopeId ? { streamScopeId } : {}),
streamEventIdPrefix: 'runAgentId',
})
// Until the first transcript response identifies the concrete fragments,
// a configured session id might itself be that physical owner. Leave its
// events pending so the first replay uses the same fragment namespace as
@@ -352,7 +407,11 @@ export function TeamMemberRunPage({
leadSessionId,
runSessionId,
[ownerAgentId],
{ eventIdPrefix: ownerAgentId },
{
eventIdPrefix: ownerAgentId,
...(streamScopeId ? { streamScopeId } : {}),
streamEventIdPrefix: ownerAgentId,
},
))
: []
const unregisterMemberSession = memberOwnerAgentIdsKnown &&
@@ -361,6 +420,7 @@ export function TeamMemberRunPage({
leadSessionId,
runSessionId,
[member.sessionId],
streamScopeId ? { streamScopeId } : undefined,
)
: () => undefined
return () => {
+175
View File
@@ -158,6 +158,7 @@ vi.mock('./cliTaskStore', () => ({
}))
import { sessionsApi } from '../api/sessions'
import type { ServerMessage } from '../types/chat'
import { useSettingsStore } from './settingsStore'
import { runsForOwner, runsForSession, useWorkflowStore } from './workflowStore'
import {
@@ -5014,6 +5015,180 @@ describe('chatStore history mapping', () => {
})
})
it('replays live agent output that arrives before the run page registers its target id', () => {
vi.useFakeTimers()
const runSessionId = '__subagent__test-session-1__early-workflow'
useChatStore.setState({
sessions: {
[TEST_SESSION_ID]: makeSession(),
[runSessionId]: makeSession(),
},
})
useChatStore.getState().handleServerMessage(TEST_SESSION_ID, {
type: 'agent_run_event',
runAgentId: 'early-workflow-agent',
streamId: 'early-workflow-stream',
targetAgentId: 'early-workflow-agent',
event: { type: 'content_delta', text: 'arrived before REST identity' },
})
expect(useChatStore.getState().sessions[runSessionId]?.streamingText).toBe('')
const unregister = registerAgentRunSession(
TEST_SESSION_ID,
runSessionId,
['early-workflow-agent'],
)
try {
vi.advanceTimersByTime(60)
expect(useChatStore.getState().sessions[runSessionId]?.streamingText).toBe(
'arrived before REST identity',
)
expect(useChatStore.getState().sessions[TEST_SESSION_ID]?.streamingText).toBe('')
} finally {
unregister()
vi.useRealTimers()
}
})
it('does not replay buffered output after an unregistered run has completed', () => {
vi.useFakeTimers()
const runSessionId = '__subagent__test-session-1__completed-before-open'
useChatStore.setState({
sessions: {
[TEST_SESSION_ID]: makeSession(),
[runSessionId]: makeSession({
messages: [{
id: 'durable-answer',
type: 'assistant_text',
content: 'durable completed answer',
timestamp: 1,
}],
}),
},
})
for (const event of [
{ type: 'content_start', blockType: 'text' },
{ type: 'content_delta', text: 'transient answer' },
{ type: 'status', state: 'idle' },
] as const) {
useChatStore.getState().handleServerMessage(TEST_SESSION_ID, {
type: 'agent_run_event',
runAgentId: 'completed-agent',
streamId: 'completed-stream',
targetAgentId: 'completed-agent',
event,
})
}
const unregister = registerAgentRunSession(
TEST_SESSION_ID,
runSessionId,
['completed-agent'],
)
try {
vi.advanceTimersByTime(60)
expect(useChatStore.getState().sessions[runSessionId]?.messages).toEqual([
expect.objectContaining({ id: 'durable-answer', content: 'durable completed answer' }),
])
expect(useChatStore.getState().sessions[runSessionId]?.streamingText).toBe('')
} finally {
unregister()
vi.useRealTimers()
}
})
it('isolates a reused Team member target by creation scope', () => {
vi.useFakeTimers()
const oldSessionId = 'team-member:old-scope:worker'
const newSessionId = 'team-member:new-scope:worker'
useChatStore.setState({
sessions: {
[TEST_SESSION_ID]: makeSession(),
[oldSessionId]: makeSession(),
[newSessionId]: makeSession(),
},
})
const unregisterOld = registerAgentRunSession(
TEST_SESSION_ID,
oldSessionId,
['worker@reused-team'],
{ streamScopeId: 'old-team-scope' },
)
const unregisterNew = registerAgentRunSession(
TEST_SESSION_ID,
newSessionId,
['worker@reused-team'],
{ streamScopeId: 'new-team-scope' },
)
try {
useChatStore.getState().handleServerMessage(TEST_SESSION_ID, {
type: 'agent_run_event',
runAgentId: 'physical-new-worker',
streamId: 'new-worker-stream',
targetAgentId: 'worker@reused-team',
targetAgentScopeId: 'new-team-scope',
event: { type: 'content_delta', text: 'new incarnation output' },
})
vi.advanceTimersByTime(60)
expect(useChatStore.getState().sessions[newSessionId]?.streamingText).toBe(
'new incarnation output',
)
expect(useChatStore.getState().sessions[oldSessionId]?.streamingText).toBe('')
} finally {
unregisterNew()
unregisterOld()
vi.useRealTimers()
}
})
it('ignores a superseded foreground stream after its background continuation starts', () => {
vi.useFakeTimers()
const runSessionId = '__subagent__test-session-1__background-handoff'
useChatStore.setState({
sessions: {
[TEST_SESSION_ID]: makeSession(),
[runSessionId]: makeSession(),
},
})
const unregister = registerAgentRunSession(
TEST_SESSION_ID,
runSessionId,
['handoff-agent'],
)
const send = (streamId: string, event: Extract<ServerMessage, { type: 'agent_run_event' }>['event']) => {
useChatStore.getState().handleServerMessage(TEST_SESSION_ID, {
type: 'agent_run_event',
runAgentId: 'handoff-agent',
streamId,
targetAgentId: 'handoff-agent',
event,
})
}
try {
send('foreground-stream', { type: 'content_delta', text: 'abandoned partial' })
vi.advanceTimersByTime(60)
send('foreground-stream', { type: 'streaming_fallback', cause: 'stream_retry' })
send('foreground-stream', { type: 'status', state: 'idle' })
send('background-stream', { type: 'content_start', blockType: 'text' })
send('background-stream', { type: 'content_delta', text: 'background answer' })
send('foreground-stream', { type: 'content_delta', text: 'late foreground text' })
vi.advanceTimersByTime(60)
const session = useChatStore.getState().sessions[runSessionId]
expect(session?.streamingText).toBe('background answer')
expect(session?.messages).not.toEqual(expect.arrayContaining([
expect.objectContaining({ content: expect.stringContaining('abandoned partial') }),
]))
expect(session?.chatState).toBe('streaming')
} finally {
unregister()
vi.useRealTimers()
}
})
it('replays owner task events that arrive before the run identity is registered', () => {
const runSessionId = '__subagent__test-session-1__slow-agent'
useWorkflowStore.setState({ runs: {} })
+227 -6
View File
@@ -162,6 +162,8 @@ export type PerSessionState = {
pendingBackgroundTaskStopFailures?: Record<string, string>
stopAllSubagentsRequested?: boolean
historyMutationEpoch?: number
/** Changes when a directed child stream starts or settles. */
agentStreamRevision?: number
suppressNextTaskNotificationResponse?: boolean
replaceHistoryOnCompletion?: boolean
activeGoal?: ActiveGoalState | null
@@ -210,6 +212,7 @@ const DEFAULT_SESSION_STATE: PerSessionState = {
pendingBackgroundTaskStopFailures: {},
stopAllSubagentsRequested: false,
historyMutationEpoch: 0,
agentStreamRevision: 0,
suppressNextTaskNotificationResponse: false,
replaceHistoryOnCompletion: false,
activeGoal: null,
@@ -393,11 +396,32 @@ type PendingOwnedTaskEvent = {
}
const pendingOwnedTaskEvents = new Map<string, Map<string, PendingOwnedTaskEvent[]>>()
const MAX_PENDING_OWNED_TASK_EVENTS = 200
type AgentStreamRouteRegistration = {
count: number
eventIdPrefix?: string | 'runAgentId'
}
type AgentStreamEventTarget = {
sessionId: string
eventIdPrefix?: string
}
const runSessionIdsByStreamAlias = new Map<
string,
Map<string, Map<string, AgentStreamRouteRegistration>>
>()
const pendingAgentRunEvents = new Map<string, Extract<ServerMessage, { type: 'agent_run_event' }>[]>()
const activeAgentStreamIdBySession = new Map<string, string>()
const retiredAgentStreamIdsBySession = new Map<string, Set<string>>()
const MAX_PENDING_AGENT_RUN_EVENTS = 1000
const MAX_RETIRED_AGENT_STREAM_IDS = 32
function ownedTaskRouteKey(ownerAgentId: string, ownerScopeId?: string): string {
return JSON.stringify([ownerScopeId ?? null, ownerAgentId])
}
function agentStreamRouteKey(targetAgentId: string, targetAgentScopeId?: string): string {
return JSON.stringify([targetAgentScopeId ?? null, targetAgentId])
}
/**
* Bind an owning runtime agent to the synthetic chat session that renders its
* transcript. Owner-scoped task events arrive on the parent session socket;
@@ -408,7 +432,12 @@ export function registerAgentRunSession(
sourceSessionId: string,
runSessionId: string,
ownerAliases: Array<string | null | undefined>,
options: { ownerScopeId?: string; eventIdPrefix?: string } = {},
options: {
ownerScopeId?: string
eventIdPrefix?: string
streamScopeId?: string
streamEventIdPrefix?: string | 'runAgentId'
} = {},
): () => void {
const aliases = [...new Set(ownerAliases
.map((alias) => alias?.trim())
@@ -433,6 +462,24 @@ export function registerAgentRunSession(
}
runSessionIdsByOwner.set(sourceSessionId, byOwner)
const byStreamAlias = runSessionIdsByStreamAlias.get(sourceSessionId) ??
new Map<string, Map<string, AgentStreamRouteRegistration>>()
for (const alias of aliases) {
const routeKey = agentStreamRouteKey(alias, options.streamScopeId)
const sessions = byStreamAlias.get(routeKey) ?? new Map<string, AgentStreamRouteRegistration>()
const existing = sessions.get(runSessionId)
sessions.set(runSessionId, {
count: (existing?.count ?? 0) + 1,
...(options.streamEventIdPrefix
? { eventIdPrefix: options.streamEventIdPrefix }
: existing?.eventIdPrefix
? { eventIdPrefix: existing.eventIdPrefix }
: {}),
})
byStreamAlias.set(routeKey, sessions)
}
runSessionIdsByStreamAlias.set(sourceSessionId, byStreamAlias)
// A run page learns its concrete agent id from the first API response, but
// lifecycle events can beat that response. Replay the bounded owner queue as
// soon as the explicit join exists so fast workflows do not disappear.
@@ -449,24 +496,189 @@ export function registerAgentRunSession(
})
}
replayPendingAgentRunEvents(sourceSessionId)
return () => {
const current = runSessionIdsByOwner.get(sourceSessionId)
if (!current) return
if (current) {
for (const alias of aliases) {
const routeKey = ownedTaskRouteKey(alias, options.ownerScopeId)
const sessions = current.get(routeKey)
const registration = sessions?.get(runSessionId)
if (registration && registration.count > 1) {
sessions?.set(runSessionId, { ...registration, count: registration.count - 1 })
} else {
sessions?.delete(runSessionId)
}
if (sessions?.size === 0) current.delete(routeKey)
}
if (current.size === 0) runSessionIdsByOwner.delete(sourceSessionId)
}
const currentStreams = runSessionIdsByStreamAlias.get(sourceSessionId)
if (!currentStreams) return
for (const alias of aliases) {
const routeKey = ownedTaskRouteKey(alias, options.ownerScopeId)
const sessions = current.get(routeKey)
const routeKey = agentStreamRouteKey(alias, options.streamScopeId)
const sessions = currentStreams.get(routeKey)
const registration = sessions?.get(runSessionId)
if (registration && registration.count > 1) {
sessions?.set(runSessionId, { ...registration, count: registration.count - 1 })
} else {
sessions?.delete(runSessionId)
}
if (sessions?.size === 0) current.delete(routeKey)
if (sessions?.size === 0) currentStreams.delete(routeKey)
}
if (current.size === 0) runSessionIdsByOwner.delete(sourceSessionId)
if (currentStreams.size === 0) runSessionIdsByStreamAlias.delete(sourceSessionId)
}
}
function agentStreamTargets(
sourceSessionId: string,
event: Extract<ServerMessage, { type: 'agent_run_event' }>,
): AgentStreamEventTarget[] {
const registrations = runSessionIdsByStreamAlias
.get(sourceSessionId)
?.get(agentStreamRouteKey(event.targetAgentId, event.targetAgentScopeId))
if (!registrations) return []
return [...registrations].map(([sessionId, registration]) => ({
sessionId,
...(registration.eventIdPrefix
? {
eventIdPrefix: registration.eventIdPrefix === 'runAgentId'
? event.runAgentId
: registration.eventIdPrefix,
}
: {}),
}))
}
function prefixAgentStreamId(prefix: string, value: string): string {
return value.startsWith(`${prefix}/`) ? value : `${prefix}/${value}`
}
function projectAgentRunStreamEvent(
event: Extract<ServerMessage, { type: 'agent_run_event' }>['event'],
eventIdPrefix: string | undefined,
): Extract<ServerMessage, { type: 'agent_run_event' }>['event'] {
if (!eventIdPrefix) return event
if (event.type === 'content_start' && event.toolUseId) {
return {
...event,
toolUseId: prefixAgentStreamId(eventIdPrefix, event.toolUseId),
...(event.parentToolUseId
? { parentToolUseId: prefixAgentStreamId(eventIdPrefix, event.parentToolUseId) }
: {}),
}
}
if (event.type === 'tool_use_complete' || event.type === 'tool_result') {
return {
...event,
toolUseId: prefixAgentStreamId(eventIdPrefix, event.toolUseId),
...(event.parentToolUseId
? { parentToolUseId: prefixAgentStreamId(eventIdPrefix, event.parentToolUseId) }
: {}),
}
}
return event
}
function retireAgentStream(sessionId: string, streamId: string): void {
const retired = retiredAgentStreamIdsBySession.get(sessionId) ?? new Set<string>()
retired.add(streamId)
while (retired.size > MAX_RETIRED_AGENT_STREAM_IDS) {
const oldest = retired.values().next().value
if (!oldest) break
retired.delete(oldest)
}
retiredAgentStreamIdsBySession.set(sessionId, retired)
}
function advanceAgentStreamRevision(sessionId: string): void {
useChatStore.setState((state) => ({
sessions: {
...state.sessions,
[sessionId]: {
...(state.sessions[sessionId] ?? createDefaultSessionState()),
agentStreamRevision: (state.sessions[sessionId]?.agentStreamRevision ?? 0) + 1,
},
},
}))
}
function dispatchAgentRunEvent(
sourceSessionId: string,
event: Extract<ServerMessage, { type: 'agent_run_event' }>,
): boolean {
const targets = agentStreamTargets(sourceSessionId, event)
if (targets.length === 0) return false
for (const target of targets) {
const isTerminal = event.event.type === 'error' ||
(event.event.type === 'status' && event.event.state === 'idle')
const retired = retiredAgentStreamIdsBySession.get(target.sessionId)
if (retired?.has(event.streamId)) continue
const activeStreamId = activeAgentStreamIdBySession.get(target.sessionId)
if (isTerminal) {
// A foreground run can be superseded by its background continuation.
// Its delayed terminal must not settle the newer stream.
if (activeStreamId && activeStreamId !== event.streamId) {
retireAgentStream(target.sessionId, event.streamId)
continue
}
activeAgentStreamIdBySession.delete(target.sessionId)
retireAgentStream(target.sessionId, event.streamId)
if (activeStreamId === event.streamId) advanceAgentStreamRevision(target.sessionId)
} else if (activeStreamId !== event.streamId) {
if (activeStreamId) {
retireAgentStream(target.sessionId, activeStreamId)
useChatStore.getState().handleServerMessage(target.sessionId, {
type: 'streaming_fallback',
cause: 'stream_retry',
})
}
activeAgentStreamIdBySession.set(target.sessionId, event.streamId)
advanceAgentStreamRevision(target.sessionId)
}
useChatStore.getState().handleServerMessage(
target.sessionId,
projectAgentRunStreamEvent(event.event, target.eventIdPrefix),
)
}
return true
}
function replayPendingAgentRunEvents(sourceSessionId: string): void {
const pending = pendingAgentRunEvents.get(sourceSessionId)
if (!pending?.length) return
const remaining = pending.filter(event => !dispatchAgentRunEvent(sourceSessionId, event))
if (remaining.length > 0) pendingAgentRunEvents.set(sourceSessionId, remaining)
else pendingAgentRunEvents.delete(sourceSessionId)
}
function bufferAgentRunEvent(
sourceSessionId: string,
event: Extract<ServerMessage, { type: 'agent_run_event' }>,
): void {
const pending = pendingAgentRunEvents.get(sourceSessionId) ?? []
const isTerminal = event.event.type === 'error' ||
(event.event.type === 'status' && event.event.state === 'idle')
if (isTerminal) {
const remaining = pending.filter(candidate => !(
candidate.streamId === event.streamId &&
candidate.targetAgentId === event.targetAgentId &&
candidate.targetAgentScopeId === event.targetAgentScopeId
))
if (remaining.length > 0) pendingAgentRunEvents.set(sourceSessionId, remaining)
else pendingAgentRunEvents.delete(sourceSessionId)
return
}
pending.push(event)
if (pending.length > MAX_PENDING_AGENT_RUN_EVENTS) {
pending.splice(0, pending.length - MAX_PENDING_AGENT_RUN_EVENTS)
}
pendingAgentRunEvents.set(sourceSessionId, pending)
}
function taskEventTargets(
sourceSessionId: string,
data: Record<string, unknown>,
@@ -2490,6 +2702,9 @@ export const useChatStore = create<ChatStore>((set, get) => ({
clearPendingToolParentUseIds(sessionId)
clearPendingToolInputDelta(sessionId)
pendingOwnedTaskEvents.delete(sessionId)
pendingAgentRunEvents.delete(sessionId)
activeAgentStreamIdBySession.delete(sessionId)
retiredAgentStreamIdsBySession.delete(sessionId)
set((s) => ({ sessions: updateSessionIn(s.sessions, sessionId, () => ({
messages: [],
activeGoal: null,
@@ -2504,6 +2719,12 @@ export const useChatStore = create<ChatStore>((set, get) => ({
},
handleServerMessage: (sessionId, msg) => {
if (msg.type === 'agent_run_event') {
if (!dispatchAgentRunEvent(sessionId, msg)) {
bufferAgentRunEvent(sessionId, msg)
}
return
}
const update = (updater: (session: PerSessionState) => Partial<PerSessionState>) => {
set((s) => ({ sessions: updateSessionIn(s.sessions, sessionId, updater) }))
}
+79
View File
@@ -943,6 +943,85 @@ describe('teamStore incremental transcript polling', () => {
}
})
it('keeps a settled live member turn when an older terminal poll resolves', async () => {
const staleTerminalPoll = deferred<any>()
getMemberTranscriptMock
.mockResolvedValueOnce({ messages: [] })
.mockReturnValueOnce(staleTerminalPoll.promise)
.mockResolvedValueOnce({
messages: [{
id: 'durable-live-answer',
type: 'assistant',
content: [{ type: 'text', text: 'Durable teammate answer' }],
timestamp: '2026-08-10T00:00:02.000Z',
}],
})
const member = {
agentId: 'worker@terminal-race-team',
name: 'worker',
role: 'worker',
status: 'running' as const,
}
const team = {
name: 'terminal-race-team',
leadSessionId: 'terminal-race-lead',
incarnationId: 'terminal-race-incarnation',
createdAt: '2026-08-10T00:00:00.000Z',
members: [member],
}
const sessionId = memberSessionId(member.agentId, team.incarnationId)
useTeamStore.setState({
activeTeam: team,
memberTeamBySession: { [sessionId]: team },
})
await useTeamStore.getState().refreshMemberSession(sessionId)
const unregister = registerAgentRunSession(
team.leadSessionId,
sessionId,
[member.agentId],
)
try {
const pendingPoll = useTeamStore.getState().refreshMemberSession(sessionId)
const send = (event: Extract<import('../types/chat').ServerMessage, { type: 'agent_run_event' }>['event']) => {
useChatStore.getState().handleServerMessage(team.leadSessionId, {
type: 'agent_run_event',
runAgentId: 'terminal-race-worker-run',
streamId: 'terminal-race-stream',
targetAgentId: member.agentId,
event,
})
}
send({ type: 'content_start', blockType: 'text' })
send({ type: 'content_delta', text: 'Live teammate answer' })
send({ type: 'status', state: 'idle' })
const completedTeam = {
...team,
members: [{ ...member, status: 'completed' as const }],
}
useTeamStore.setState({
activeTeam: completedTeam,
memberTeamBySession: { [sessionId]: completedTeam },
})
staleTerminalPoll.resolve({ messages: [] })
await pendingPoll
expect(useChatStore.getState().sessions[sessionId]?.messages).toEqual([
expect.objectContaining({ content: 'Live teammate answer' }),
])
expect(useChatStore.getState().sessions[sessionId]?.agentStreamRevision).toBe(2)
await useTeamStore.getState().refreshMemberSession(sessionId)
expect(useChatStore.getState().sessions[sessionId]?.messages).toEqual([
expect.objectContaining({ content: 'Durable teammate answer' }),
])
expect(useChatStore.getState().sessions[sessionId]?.agentStreamRevision).toBe(0)
} finally {
unregister()
}
})
it('joins physical-owner live activity with its fragment-scoped transcript identity', async () => {
getMemberTranscriptMock
.mockResolvedValueOnce({
+30 -2
View File
@@ -106,6 +106,7 @@ function createMemberSessionState() {
agentTaskNotifications: {},
backgroundAgentTasks: {},
historyMutationEpoch: 0,
agentStreamRevision: 0,
elapsedTimer: null,
}
}
@@ -552,6 +553,7 @@ function syncMemberSessionMessages(
activity?: ReturnType<typeof reconstructRunActivityFromTranscript>,
requestedMutationEpoch?: number,
requestedTaskUpdatedAt?: Map<string, number>,
requestedStreamRevision?: number,
) {
const isTerminal = member.status === 'completed' ||
member.status === 'error' ||
@@ -581,6 +583,14 @@ function syncMemberSessionMessages(
useChatStore.setState((state) => {
const existing = state.sessions[sessionId]
const nextState = existing ?? createMemberSessionState()
const currentStreamRevision = nextState.agentStreamRevision ?? 0
const preserveLiveConversation = currentStreamRevision > 0 && (
nextState.chatState !== 'idle' ||
(
requestedStreamRevision !== undefined &&
currentStreamRevision !== requestedStreamRevision
)
)
const mutationEpochChanged = requestedMutationEpoch !== undefined &&
(nextState.historyMutationEpoch ?? 0) !== requestedMutationEpoch
const taskFreshnessChanged = existing && Object.values(
@@ -627,13 +637,25 @@ function syncMemberSessionMessages(
...state.sessions,
[sessionId]: {
...nextState,
messages,
messages: preserveLiveConversation ? nextState.messages : messages,
...(activity ? {
agentTaskNotifications: mergedActivity!.agentTaskNotifications,
backgroundAgentTasks: mergedActivity!.backgroundAgentTasks,
} : {}),
connectionState: 'connected',
chatState: isActive ? 'thinking' : 'idle',
chatState: preserveLiveConversation
? nextState.chatState
: isActive
? 'thinking'
: 'idle',
...(!preserveLiveConversation ? {
agentStreamRevision: 0,
streamingText: '',
streamingToolInput: '',
activeToolUseId: null,
activeToolName: null,
activeThinkingId: null,
} : {}),
},
},
}
@@ -1059,6 +1081,7 @@ export const useTeamStore = create<TeamStore>((set, get) => ({
const incarnationKey = memberIncarnationKey(team, member)
const requestedSession = useChatStore.getState().sessions[sessionId]
const requestedMutationEpoch = requestedSession?.historyMutationEpoch ?? 0
const requestedStreamRevision = requestedSession?.agentStreamRevision ?? 0
const requestedTaskUpdatedAt = new Map<string, number>()
for (const task of Object.values(requestedSession?.backgroundAgentTasks ?? {})) {
requestedTaskUpdatedAt.set(task.taskId, task.updatedAt)
@@ -1175,6 +1198,7 @@ export const useTeamStore = create<TeamStore>((set, get) => ({
activity,
requestedMutationEpoch,
requestedTaskUpdatedAt,
requestedStreamRevision,
)
} catch {
const currentTeam = get().getTeamByMemberSessionId(sessionId)
@@ -1197,6 +1221,10 @@ export const useTeamStore = create<TeamStore>((set, get) => ({
currentMember,
snapshot,
existingMessages,
undefined,
undefined,
undefined,
requestedStreamRevision,
)
}
})()
+19
View File
@@ -84,6 +84,14 @@ export type UIAttachment = {
export type ServerMessage =
| { type: 'connected'; sessionId: string }
| { type: 'session_state'; turnState: 'running' | 'idle' }
| {
type: 'agent_run_event'
runAgentId: string
streamId: string
targetAgentId: string
targetAgentScopeId?: string
event: AgentRunStreamMessage
}
| { type: 'content_start'; blockType: 'text' | 'tool_use'; toolName?: string; toolUseId?: string; originalToolUseId?: string; parentToolUseId?: string }
| { type: 'content_delta'; text?: string; toolInput?: string }
| { type: 'tool_use_complete'; toolName: string; toolUseId: string; originalToolUseId?: string; input: unknown; parentToolUseId?: string }
@@ -149,6 +157,17 @@ export type ServerMessage =
| { type: 'task_update'; taskId: string; status: string; progress?: string }
| { type: 'session_title_updated'; sessionId: string; title: string }
export type AgentRunStreamMessage =
| { type: 'content_start'; blockType: 'text' | 'tool_use'; toolName?: string; toolUseId?: string; originalToolUseId?: string; parentToolUseId?: string }
| { type: 'content_delta'; text?: string; toolInput?: string }
| { type: 'tool_use_complete'; toolName: string; toolUseId: string; originalToolUseId?: string; input: unknown; parentToolUseId?: string }
| { type: 'tool_result'; toolUseId: string; originalToolUseId?: string; content: unknown; isError: boolean; parentToolUseId?: string }
| { type: 'thinking'; text: string; complete?: boolean }
| { type: 'status'; state: ChatState; verb?: string; attemptStart?: boolean }
| { type: 'api_retry'; attempt: number; maxRetries: number; retryDelayMs: number; errorStatus: number | null; errorType?: string; errorMessage?: string }
| { type: 'streaming_fallback'; cause: StreamingFallbackCause }
| { type: 'error'; message: string; code: string; retryable?: boolean; businessErrorCode?: string }
export type TokenUsage = {
input_tokens: number
output_tokens: number
+134
View File
@@ -0,0 +1,134 @@
import { afterEach, beforeEach, describe, expect, test } from 'bun:test'
import { getIsInteractive, setIsInteractive } from '../bootstrap/state.js'
import type { CanUseToolFn } from '../hooks/useCanUseTool.js'
import { getDefaultAppState } from '../state/AppStateStore.js'
import { emitAgentRunMessage, setAgentRunMessageSink } from '../utils/sdkEventQueue.js'
import { Stream } from '../utils/stream.js'
import {
__runHeadlessStreamingForTests,
bindAgentRunMessageSink,
} from './print.js'
import { RemoteIO } from './remoteIO.js'
import { StructuredIO } from './structuredIO.js'
async function* emptyInput(): AsyncGenerator<string> {}
describe('headless Agent stream outbound binding', () => {
let wasInteractive = true
beforeEach(() => {
wasInteractive = getIsInteractive()
setIsInteractive(false)
setAgentRunMessageSink(undefined)
})
afterEach(() => {
setAgentRunMessageSink(undefined)
setIsInteractive(wasInteractive)
})
test('sends the private child stream directly through RemoteIO only', async () => {
const regular = new StructuredIO(emptyInput())
const removeRegular = bindAgentRunMessageSink(regular)
emitAgentRunMessage({
runAgentId: 'regular-agent',
streamId: 'regular-stream',
targetAgentId: 'regular-agent',
}, { kind: 'complete' })
removeRegular()
regular.outbound.done()
expect(await regular.outbound.next()).toEqual({ done: true, value: undefined })
const remote = new StructuredIO(emptyInput())
Object.setPrototypeOf(remote, RemoteIO.prototype)
const removeRemote = bindAgentRunMessageSink(remote)
emitAgentRunMessage({
runAgentId: 'desktop-agent',
streamId: 'desktop-stream',
targetAgentId: 'desktop-agent',
}, {
kind: 'message',
message: { type: 'stream_event', event: { type: 'content_block_delta' } },
})
expect((await remote.outbound.next()).value).toMatchObject({
type: 'system',
subtype: 'agent_run_message',
run_agent_id: 'desktop-agent',
stream_id: 'desktop-stream',
target_agent_id: 'desktop-agent',
event_kind: 'message',
})
removeRemote()
})
test('binds the real headless RemoteIO stream to the same outbound FIFO', async () => {
const input = new Stream<string>()
const remote = new StructuredIO(input)
Object.setPrototypeOf(remote, RemoteIO.prototype)
let appState = getDefaultAppState()
const sigintListenersBefore = new Set(process.listeners('SIGINT'))
const output = __runHeadlessStreamingForTests(
remote,
[],
[],
[],
[],
(() => undefined) as unknown as CanUseToolFn,
{},
() => appState,
update => {
appState = update(appState)
},
[],
{
verbose: undefined,
jsonSchema: undefined,
permissionPromptToolName: undefined,
allowedTools: undefined,
thinkingConfig: undefined,
maxTurns: undefined,
maxBudgetUsd: undefined,
taskBudget: undefined,
systemPrompt: undefined,
appendSystemPrompt: undefined,
userSpecifiedModel: undefined,
fallbackModel: undefined,
outputFormat: 'stream-json',
},
)
const iterator = output[Symbol.asyncIterator]()
emitAgentRunMessage({
runAgentId: 'joined-agent',
streamId: 'joined-stream',
targetAgentId: 'joined-agent',
}, { kind: 'complete' })
const sentinel = { type: 'agent-run-stream-test-sentinel' } as never
remote.outbound.enqueue(sentinel)
const beforeSentinel = []
for (;;) {
const next = await iterator.next()
if (next.value === sentinel) break
beforeSentinel.push(next.value)
}
expect(beforeSentinel).toContainEqual(
expect.objectContaining({
type: 'system',
subtype: 'agent_run_message',
run_agent_id: 'joined-agent',
stream_id: 'joined-stream',
}),
)
input.done()
while (!(await iterator.next()).done) {
// Drain until the real headless input loop removes its sink and closes.
}
for (const listener of process.listeners('SIGINT')) {
if (!sigintListenersBefore.has(listener)) process.off('SIGINT', listener)
}
})
})
+16 -1
View File
@@ -361,7 +361,10 @@ import { unassignTeammateTasks } from '../utils/tasks.js'
import { getRunningTasks } from '../utils/task/framework.js'
import { isBackgroundTask } from '../tasks/types.js'
import { stopTask } from '../tasks/stopTask.js'
import { drainSdkEvents } from '../utils/sdkEventQueue.js'
import {
drainSdkEvents,
setAgentRunMessageSink,
} from '../utils/sdkEventQueue.js'
import { initializeGrowthBook } from '../services/analytics/growthbook.js'
import { errorMessage, toError } from '../utils/errors.js'
import { sleep } from '../utils/sleep.js'
@@ -991,6 +994,13 @@ export async function runHeadless(
)
}
export function bindAgentRunMessageSink(structuredIO: StructuredIO): () => void {
if (!(structuredIO instanceof RemoteIO)) return () => undefined
return setAgentRunMessageSink(event => {
structuredIO.outbound.enqueue(event)
})
}
function runHeadlessStreaming(
structuredIO: StructuredIO,
mcpClients: MCPServerConnection[],
@@ -1040,6 +1050,7 @@ function runHeadlessStreaming(
let abortController: AbortController | undefined
// Same queue sendRequest() enqueues to — one FIFO for everything.
const output = structuredIO.outbound
const removeAgentRunMessageSink = bindAgentRunMessageSink(structuredIO)
// Ctrl+C in -p mode: abort the in-flight query, then shut down gracefully.
// gracefulShutdown persists session state and flushes analytics, with a
@@ -2658,6 +2669,7 @@ function runHeadlessStreaming(
unsubscribeSkillChanges()
unsubscribeAuthStatus?.()
statusListeners.delete(rateLimitListener)
removeAgentRunMessageSink()
output.done()
}
}
@@ -4177,6 +4189,7 @@ function runHeadlessStreaming(
unsubscribeSkillChanges()
unsubscribeAuthStatus?.()
statusListeners.delete(rateLimitListener)
removeAgentRunMessageSink()
output.done()
}
})()
@@ -4184,6 +4197,8 @@ function runHeadlessStreaming(
return output
}
export { runHeadlessStreaming as __runHeadlessStreamingForTests }
/**
* Creates a CanUseToolFn that incorporates a custom permission prompt tool.
* This function converts the permissionPromptTool into a CanUseToolFn that can be used in ask.tsx
@@ -520,6 +520,23 @@
]
}
],
"agent-run-stream": [
{
"in": "system/agent_run_message",
"out": [
{
"type": "agent_run_event",
"runAgentId": "agent_golden_physical",
"streamId": "agent_golden_invocation",
"targetAgentId": "agent_golden_target",
"event": {
"type": "status",
"state": "idle"
}
}
]
}
],
"task-lifecycle": [
{
"in": "system/task_started",
@@ -0,0 +1,328 @@
import { afterEach, describe, expect, test } from 'bun:test'
import {
__resetWebSocketHandlerStateForTests,
translateCliMessage,
} from '../ws/handler.js'
type AgentRunRoute = {
runAgentId: string
streamId: string
targetAgentId: string
targetAgentScopeId?: string
}
function agentRunMessage(
route: AgentRunRoute,
message: unknown,
) {
return {
type: 'system',
subtype: 'agent_run_message',
run_agent_id: route.runAgentId,
stream_id: route.streamId,
target_agent_id: route.targetAgentId,
...(route.targetAgentScopeId
? { target_agent_scope_id: route.targetAgentScopeId }
: {}),
event_kind: 'message',
message,
}
}
function terminalAgentRunMessage(
route: AgentRunRoute,
eventKind: 'complete' | 'cancelled' | 'error',
error?: string,
) {
return {
type: 'system',
subtype: 'agent_run_message',
run_agent_id: route.runAgentId,
stream_id: route.streamId,
target_agent_id: route.targetAgentId,
...(route.targetAgentScopeId
? { target_agent_scope_id: route.targetAgentScopeId }
: {}),
event_kind: eventKind,
...(error ? { error } : {}),
}
}
function stream(event: unknown) {
return { type: 'stream_event', event }
}
afterEach(() => {
__resetWebSocketHandlerStateForTests()
})
describe('translateCliMessage: agent_run_message', () => {
test('keeps text, thinking, tool input and tool result inside the targeted scoped run', () => {
const sessionId = 'agent-run-stream-session'
const route: AgentRunRoute = {
runAgentId: 'physical-agent',
streamId: 'physical-agent-invocation-1',
targetAgentId: 'logical-agent',
targetAgentScopeId: '["team-a","lead-session",123]',
}
const translated = [
...translateCliMessage(agentRunMessage(
route,
stream({ type: 'message_start', message: { id: 'message-1', usage: {} } }),
), sessionId),
...translateCliMessage(agentRunMessage(
route,
stream({ type: 'content_block_start', index: 0, content_block: { type: 'thinking' } }),
), sessionId),
...translateCliMessage(agentRunMessage(
route,
stream({ type: 'content_block_delta', index: 0, delta: { type: 'thinking_delta', thinking: 'checking' } }),
), sessionId),
...translateCliMessage(agentRunMessage(
route,
stream({ type: 'content_block_start', index: 1, content_block: { type: 'text' } }),
), sessionId),
...translateCliMessage(agentRunMessage(
route,
stream({ type: 'content_block_delta', index: 1, delta: { type: 'text_delta', text: 'live answer' } }),
), sessionId),
...translateCliMessage(agentRunMessage(
route,
stream({
type: 'content_block_start',
index: 2,
content_block: { type: 'tool_use', id: 'tool-1', name: 'Read' },
}),
), sessionId),
...translateCliMessage(agentRunMessage(
route,
stream({ type: 'content_block_delta', index: 2, delta: { type: 'input_json_delta', partial_json: '{"file_path":"a.ts"}' } }),
), sessionId),
...translateCliMessage(agentRunMessage(
route,
stream({ type: 'content_block_stop', index: 2 }),
), sessionId),
...translateCliMessage(agentRunMessage(
route,
{
type: 'user',
message: {
content: [{ type: 'tool_result', tool_use_id: 'tool-1', content: 'contents' }],
},
},
), sessionId),
]
expect(translated).toEqual(expect.arrayContaining([
{ type: 'agent_run_event', ...route, event: { type: 'thinking', text: 'checking' } },
{ type: 'agent_run_event', ...route, event: { type: 'content_delta', text: 'live answer' } },
{
type: 'agent_run_event',
...route,
event: { type: 'content_delta', toolInput: '{"file_path":"a.ts"}' },
},
{
type: 'agent_run_event',
...route,
event: {
type: 'tool_use_complete',
toolName: 'Read',
toolUseId: 'tool-1',
input: { file_path: 'a.ts' },
parentToolUseId: undefined,
},
},
{
type: 'agent_run_event',
...route,
event: {
type: 'tool_result',
toolUseId: 'tool-1',
content: 'contents',
isError: false,
parentToolUseId: undefined,
},
},
]))
expect(translated.every(message => message.type === 'agent_run_event')).toBe(true)
})
test('isolates interleaved tool JSON by invocation stream, even for the same physical agent', () => {
const sessionId = 'parallel-agent-run-streams'
const routeFor = (streamId: string): AgentRunRoute => ({
runAgentId: 'same-physical-agent',
streamId,
targetAgentId: 'same-logical-agent',
})
const startTool = (route: AgentRunRoute, toolName: string) => translateCliMessage(
agentRunMessage(route, stream({
type: 'content_block_start',
index: 0,
content_block: { type: 'tool_use', id: 'same-id', name: toolName },
})),
sessionId,
)
const delta = (route: AgentRunRoute, partialJson: string) => translateCliMessage(
agentRunMessage(route, stream({
type: 'content_block_delta',
index: 0,
delta: { type: 'input_json_delta', partial_json: partialJson },
})),
sessionId,
)
const stop = (route: AgentRunRoute) => translateCliMessage(
agentRunMessage(route, stream({ type: 'content_block_stop', index: 0 })),
sessionId,
)
const first = routeFor('invocation-a')
const second = routeFor('invocation-b')
startTool(first, 'Read')
startTool(second, 'Write')
delta(first, '{"a":1}')
delta(second, '{"b":2}')
expect(stop(first)).toEqual([expect.objectContaining({
streamId: first.streamId,
event: expect.objectContaining({ toolName: 'Read', input: { a: 1 } }),
})])
expect(stop(second)).toEqual([expect.objectContaining({
streamId: second.streamId,
event: expect.objectContaining({ toolName: 'Write', input: { b: 2 } }),
})])
})
test('wraps retry and fallback status without leaking them into the root conversation', () => {
const sessionId = 'agent-run-retry-fallback'
const route: AgentRunRoute = {
runAgentId: 'retrying-agent',
streamId: 'retrying-agent-invocation',
targetAgentId: 'retrying-agent',
}
expect(translateCliMessage(agentRunMessage(route, {
type: 'system',
subtype: 'api_retry',
attempt: 2,
max_retries: 4,
retry_delay_ms: 750,
error_status: 529,
error: 'overloaded_error',
}), sessionId)).toEqual([{
type: 'agent_run_event',
...route,
event: {
type: 'api_retry',
attempt: 2,
maxRetries: 4,
retryDelayMs: 750,
errorStatus: 529,
errorType: 'overloaded_error',
},
}])
expect(translateCliMessage(agentRunMessage(route, {
type: 'system',
subtype: 'streaming_fallback',
cause: 'watchdog',
}), sessionId)).toEqual([{
type: 'agent_run_event',
...route,
event: { type: 'streaming_fallback', cause: 'watchdog' },
}])
})
test('keeps child stream deduplication state out of the root conversation', () => {
const sessionId = 'agent-run-root-isolation'
const route: AgentRunRoute = {
runAgentId: 'child-agent',
streamId: 'child-stream',
targetAgentId: 'child-agent',
}
const messageId = 'same-provider-message-id'
translateCliMessage(agentRunMessage(
route,
stream({ type: 'message_start', message: { id: messageId, usage: {} } }),
), sessionId)
translateCliMessage(agentRunMessage(
route,
stream({ type: 'content_block_start', index: 0, content_block: { type: 'text' } }),
), sessionId)
translateCliMessage(agentRunMessage(
route,
stream({ type: 'content_block_delta', index: 0, delta: { type: 'text_delta', text: 'child text' } }),
), sessionId)
expect(translateCliMessage({
type: 'assistant',
message: {
id: messageId,
content: [{ type: 'text', text: 'root buffered text' }],
},
}, sessionId)).toEqual([
{ type: 'content_start', blockType: 'text' },
{ type: 'content_delta', text: 'root buffered text' },
])
})
test('settles complete, cancelled and error terminals on the targeted stream only', () => {
const route: AgentRunRoute = {
runAgentId: 'physical-agent',
streamId: 'physical-agent-invocation',
targetAgentId: 'logical-agent',
targetAgentScopeId: '["team-a","lead-session",123]',
}
expect(translateCliMessage(
terminalAgentRunMessage(route, 'complete'),
'agent-run-complete-session',
)).toEqual([{
type: 'agent_run_event',
...route,
event: { type: 'status', state: 'idle' },
}])
expect(translateCliMessage(
terminalAgentRunMessage(route, 'cancelled'),
'agent-run-cancelled-session',
)).toEqual([
{
type: 'agent_run_event',
...route,
event: { type: 'streaming_fallback', cause: 'stream_retry' },
},
{
type: 'agent_run_event',
...route,
event: { type: 'status', state: 'idle' },
},
])
expect(translateCliMessage(
terminalAgentRunMessage(route, 'error', 'child failed'),
'agent-run-error-session',
)).toEqual([{
type: 'agent_run_event',
...route,
event: {
type: 'error',
message: 'child failed',
code: 'AGENT_RUN_ERROR',
},
}])
})
test('drops frames that do not identify an exact run, stream and target', () => {
const base = {
type: 'system',
subtype: 'agent_run_message',
run_agent_id: 'physical-agent',
stream_id: 'invocation',
target_agent_id: 'logical-agent',
event_kind: 'complete',
}
for (const field of ['run_agent_id', 'stream_id', 'target_agent_id'] as const) {
expect(translateCliMessage({ ...base, [field]: ' ' }, `missing-${field}`)).toEqual([])
}
})
})
@@ -271,6 +271,18 @@ export const goldenScenarios: GoldenScenario[] = [
},
],
},
{
id: 'agent-run-stream',
description: 'A child run completion stays wrapped for its synthetic conversation.',
messages: [{
type: 'system',
subtype: 'agent_run_message',
run_agent_id: 'agent_golden_physical',
stream_id: 'agent_golden_invocation',
target_agent_id: 'agent_golden_target',
event_kind: 'complete',
}],
},
{
id: 'task-lifecycle',
description: 'Detached task frames update Activity without reviving an idle foreground turn.',
@@ -663,6 +663,136 @@ describe('WebSocket handler session isolation', () => {
])
})
it('forwards directed Agent output while a foreground admission awaits send acknowledgement', async () => {
const sessionId = `agent-run-during-admission-${crypto.randomUUID()}`
const ws = makeClientSocket(sessionId)
const outputCallbacks = new Set<(cliMsg: any) => void>()
let resolveSend!: (sent: boolean) => void
spyOn(conversationService, 'hasSession').mockReturnValue(true)
spyOn(conversationService, 'getPendingPermissionRequests').mockReturnValue([])
spyOn(conversationService, 'onOutput').mockImplementation((_sid, callback) => {
outputCallbacks.add(callback)
})
spyOn(conversationService, 'removeOutputCallback').mockImplementation((_sid, callback) => {
outputCallbacks.delete(callback)
})
spyOn(sessionService, 'getCustomTitle').mockResolvedValue('Existing title')
spyOn(conversationService, 'sendMessage').mockImplementation(
() => new Promise<boolean>((resolve) => {
resolveSend = resolve
}),
)
handleWebSocket.open(ws)
ws.sent.length = 0
handleWebSocket.message(ws, JSON.stringify({
type: 'user_message',
content: 'Start another foreground turn',
}))
await flushMicrotasks(30)
const directedDelta = {
type: 'system',
subtype: 'agent_run_message',
run_agent_id: 'background-physical-agent',
stream_id: 'background-invocation',
target_agent_id: 'background-logical-agent',
event_kind: 'message',
message: {
type: 'stream_event',
event: {
type: 'content_block_delta',
index: 0,
delta: { type: 'text_delta', text: 'background is still live' },
},
},
}
for (const callback of [...outputCallbacks]) callback(directedDelta)
expect(ws.sent.map((payload) => JSON.parse(payload))).toContainEqual({
type: 'agent_run_event',
runAgentId: 'background-physical-agent',
streamId: 'background-invocation',
targetAgentId: 'background-logical-agent',
event: { type: 'content_delta', text: 'background is still live' },
})
resolveSend(true)
await flushMicrotasks(30)
})
it('lets directed Agent terminals through the stop fence after suppressing late content', () => {
const sessionId = `agent-run-terminal-after-stop-${crypto.randomUUID()}`
const ws = makeClientSocket(sessionId)
let outputCallback: ((cliMsg: any) => void) | null = null
spyOn(globalThis, 'setTimeout').mockImplementation(() => 1 as any)
spyOn(conversationService, 'hasSession').mockReturnValue(true)
spyOn(conversationService, 'sendInterrupt').mockReturnValue(true)
spyOn(conversationService, 'onOutput').mockImplementation((_sid, callback) => {
outputCallback = callback
})
handleWebSocket.open(ws)
__markActiveTurnForTests(sessionId)
ws.sent.length = 0
handleWebSocket.message(ws, JSON.stringify({ type: 'stop_generation' }))
const route = {
type: 'system',
subtype: 'agent_run_message',
run_agent_id: 'stopped-physical-agent',
target_agent_id: 'stopped-logical-agent',
}
outputCallback?.({
...route,
stream_id: 'stopped-content-invocation',
event_kind: 'message',
message: {
type: 'stream_event',
event: {
type: 'content_block_delta',
index: 0,
delta: { type: 'text_delta', text: 'must stay hidden after stop' },
},
},
})
outputCallback?.({
...route,
stream_id: 'stopped-content-invocation',
event_kind: 'complete',
})
outputCallback?.({
...route,
stream_id: 'stopped-error-invocation',
event_kind: 'error',
error: 'stopped child settled with an error',
})
const sent = ws.sent.map((payload) => JSON.parse(payload))
expect(sent).not.toContainEqual(expect.objectContaining({
type: 'agent_run_event',
event: { type: 'content_delta', text: 'must stay hidden after stop' },
}))
expect(sent).toContainEqual({
type: 'agent_run_event',
runAgentId: 'stopped-physical-agent',
streamId: 'stopped-content-invocation',
targetAgentId: 'stopped-logical-agent',
event: { type: 'status', state: 'idle' },
})
expect(sent).toContainEqual({
type: 'agent_run_event',
runAgentId: 'stopped-physical-agent',
streamId: 'stopped-error-invocation',
targetAgentId: 'stopped-logical-agent',
event: {
type: 'error',
message: 'stopped child settled with an error',
code: 'AGENT_RUN_ERROR',
},
})
})
it('suppresses an interrupted result for every client bound to the stopped session', () => {
const sessionId = `stopped-turn-multi-client-${crypto.randomUUID()}`
const first = makeClientSocket(sessionId)
+19
View File
@@ -58,6 +58,14 @@ export const RUNTIME_CONFIG_APPLIED_EVENT = 'runtime_config_applied' as const
export type ServerMessage =
| { type: 'connected'; sessionId: string }
| { type: 'session_state'; turnState: 'running' | 'idle' }
| {
type: 'agent_run_event'
runAgentId: string
streamId: string
targetAgentId: string
targetAgentScopeId?: string
event: AgentRunStreamMessage
}
| { type: 'content_start'; blockType: 'text' | 'tool_use'; toolName?: string; toolUseId?: string; originalToolUseId?: string; parentToolUseId?: string }
| { type: 'content_delta'; text?: string; toolInput?: string }
| { type: 'tool_use_complete'; toolName: string; toolUseId: string; originalToolUseId?: string; input: unknown; parentToolUseId?: string }
@@ -130,6 +138,17 @@ export type ServerMessage =
| { type: 'task_update'; taskId: string; status: string; progress?: string }
| { type: 'session_title_updated'; sessionId: string; title: string }
export type AgentRunStreamMessage =
| { type: 'content_start'; blockType: 'text' | 'tool_use'; toolName?: string; toolUseId?: string; originalToolUseId?: string; parentToolUseId?: string }
| { type: 'content_delta'; text?: string; toolInput?: string }
| { type: 'tool_use_complete'; toolName: string; toolUseId: string; originalToolUseId?: string; input: unknown; parentToolUseId?: string }
| { type: 'tool_result'; toolUseId: string; originalToolUseId?: string; content: unknown; isError: boolean; parentToolUseId?: string }
| { type: 'thinking'; text: string; complete?: boolean }
| { type: 'status'; state: ChatState; verb?: string; attemptStart?: boolean }
| { type: 'api_retry'; attempt: number; maxRetries: number; retryDelayMs: number; errorStatus: number | null; errorType?: string; errorMessage?: string }
| { type: 'streaming_fallback'; cause: StreamingFallbackCause }
| { type: 'error'; message: string; code: string; retryable?: boolean; businessErrorCode?: string }
export type TokenUsage = {
input_tokens: number
output_tokens: number
+97 -8
View File
@@ -8,6 +8,7 @@
import type { ServerWebSocket } from 'bun'
import type {
AgentRunStreamMessage,
ClientMessage,
PermissionMode,
ServerMessage,
@@ -889,7 +890,11 @@ async function handleUserMessage(
bindAllClientSessionOutputs(sessionId, {
shouldForward: (cliMsg) => {
if (userMessageSent || (cliMsg.type === 'result' && cliMsg.is_error)) {
if (
userMessageSent ||
(cliMsg.type === 'result' && cliMsg.is_error) ||
isAgentRunMessageFrame(cliMsg)
) {
return true
}
return shouldForwardCurrentTurnLocalCommand(cliMsg)
@@ -2523,21 +2528,95 @@ function getStreamState(sessionId: string): SessionStreamState {
return state
}
function isAgentRunStreamMessage(
message: ServerMessage,
): message is AgentRunStreamMessage {
return message.type === 'content_start' ||
message.type === 'content_delta' ||
message.type === 'tool_use_complete' ||
message.type === 'tool_result' ||
message.type === 'thinking' ||
message.type === 'status' ||
message.type === 'api_retry' ||
message.type === 'streaming_fallback' ||
message.type === 'error'
}
function translateAgentRunMessage(
cliMsg: any,
sessionId: string,
): ServerMessage[] {
const runAgentId = typeof cliMsg.run_agent_id === 'string'
? cliMsg.run_agent_id.trim()
: ''
const streamId = typeof cliMsg.stream_id === 'string'
? cliMsg.stream_id.trim()
: ''
const targetAgentId = typeof cliMsg.target_agent_id === 'string'
? cliMsg.target_agent_id.trim()
: ''
const targetAgentScopeId = typeof cliMsg.target_agent_scope_id === 'string'
? cliMsg.target_agent_scope_id.trim()
: ''
if (!runAgentId || !streamId || !targetAgentId) return []
const route = {
runAgentId,
streamId,
targetAgentId,
...(targetAgentScopeId ? { targetAgentScopeId } : {}),
}
const streamSessionId = `${sessionId}\u0000agent-run:${streamId}`
if (cliMsg.event_kind === 'complete') {
sessionStreamStates.delete(streamSessionId)
return [{
type: 'agent_run_event',
...route,
event: { type: 'status', state: 'idle' },
}]
}
if (cliMsg.event_kind === 'cancelled') {
sessionStreamStates.delete(streamSessionId)
return [
{
type: 'agent_run_event',
...route,
event: { type: 'streaming_fallback', cause: 'stream_retry' },
},
{
type: 'agent_run_event',
...route,
event: { type: 'status', state: 'idle' },
},
]
}
if (cliMsg.event_kind === 'error') {
sessionStreamStates.delete(streamSessionId)
return [{
type: 'agent_run_event',
...route,
event: {
type: 'error',
message: typeof cliMsg.error === 'string' ? cliMsg.error : 'Agent run failed',
code: 'AGENT_RUN_ERROR',
},
}]
}
if (cliMsg.event_kind !== 'message' || !cliMsg.message) return []
return translateCliMessage(cliMsg.message, streamSessionId)
.filter(isAgentRunStreamMessage)
.map(event => ({ type: 'agent_run_event', ...route, event }))
}
/** Clean up stream state when session disconnects */
function cleanupStreamState(sessionId: string) {
sessionStreamStates.delete(sessionId)
const agentPrefix = `${sessionId}\u0000agent-run:`
for (const key of sessionStreamStates.keys()) {
if (key.startsWith(agentPrefix)) sessionStreamStates.delete(key)
}
}
function cleanupSessionRuntimeState(
@@ -2771,6 +2850,9 @@ export async function ensureCliSessionStartedForControl(
}
export function translateCliMessage(cliMsg: any, sessionId: string): ServerMessage[] {
if (isAgentRunMessageFrame(cliMsg)) {
return translateAgentRunMessage(cliMsg, sessionId)
}
const streamState = getStreamState(sessionId)
switch (cliMsg.type) {
case 'assistant': {
@@ -3932,6 +4014,10 @@ function hasStoppedTurnBoundary(sessionId: string): boolean {
activeUserTurns.get(sessionId)?.replacementAfterStop === true
}
function isAgentRunMessageFrame(cliMsg: any): boolean {
return cliMsg?.type === 'system' && cliMsg.subtype === 'agent_run_message'
}
function isAgentScopedPermissionRequest(cliMsg: any): boolean {
return cliMsg?.type === 'control_request' &&
cliMsg.request?.subtype === 'can_use_tool' &&
@@ -3951,6 +4037,9 @@ function shouldSuppressCliOutputDuringStop(
taskLifecycle: CliBackgroundTaskLifecycle | null,
): boolean {
if (taskLifecycle !== null) return false
if (isAgentRunMessageFrame(cliMsg)) {
return agentStopRequestedSessions.has(sessionId) && cliMsg.event_kind === 'message'
}
if (cliMsg?.type === 'control_cancel_request' || cliMsg?.type === 'control_response') {
return false
}
@@ -128,6 +128,7 @@ describe('resumed Agent ownership', () => {
expect(runAgentSpy).toHaveBeenCalledWith(expect.objectContaining({
spawningToolUseId: 'toolu_original_agent',
ownerAgentId: 'current-parent',
streamTargetAgentId: agentId,
}))
})
})
+1
View File
@@ -241,6 +241,7 @@ export async function resumeAgentBackground({
description: meta?.description,
spawningToolUseId,
ownerAgentId,
streamTargetAgentId: agentId,
persistedAgentType: resumedAgentType,
alreadyPersistedMessageCount: resumedMessages.length,
contentReplacementState: resumedReplacementState,
@@ -0,0 +1,330 @@
import { afterEach, beforeEach, describe, expect, mock, test } from 'bun:test'
import { mkdtemp, rm } from 'node:fs/promises'
import { tmpdir } from 'node:os'
import { join } from 'node:path'
import type { AgentRunMessageEvent } from '../../utils/sdkEventQueue.js'
let queryMessages: unknown[] = []
mock.module('../../query.js', () => ({
query: async function* () {
for (const message of queryMessages) yield message
},
}))
type RouteOptions =
| { spawningToolUseId: string; isAsync?: boolean }
| { querySource: 'workflow_agent' }
| { streamTargetAgentId: string }
describe('runAgent live stream bridge', () => {
let configDir: string
let previousConfigDir: string | undefined
beforeEach(async () => {
previousConfigDir = process.env.CLAUDE_CONFIG_DIR
configDir = await mkdtemp(join(tmpdir(), 'agent-stream-test-'))
process.env.CLAUDE_CONFIG_DIR = configDir
queryMessages = [
{
type: 'stream_event',
event: { type: 'message_start', message: { id: 'agent-message', usage: {} } },
},
{
type: 'stream_event',
event: {
type: 'content_block_start',
index: 0,
content_block: { type: 'text', text: '' },
},
},
{
type: 'stream_event',
event: {
type: 'content_block_delta',
index: 0,
delta: { type: 'text_delta', text: 'live token' },
},
},
{
type: 'stream_event',
event: {
type: 'content_block_start',
index: 1,
content_block: { type: 'thinking', thinking: '' },
},
},
{
type: 'stream_event',
event: {
type: 'content_block_delta',
index: 1,
delta: { type: 'thinking_delta', thinking: 'live thought' },
},
},
{
type: 'stream_event',
event: {
type: 'content_block_start',
index: 2,
content_block: {
type: 'tool_use',
id: 'tool-call',
name: 'Read',
input: {},
},
},
},
{
type: 'stream_event',
event: {
type: 'content_block_delta',
index: 2,
delta: { type: 'input_json_delta', partial_json: '{"file_path":"/tmp/a"}' },
},
},
{
type: 'stream_event',
event: { type: 'content_block_stop', index: 2 },
},
{
type: 'system',
subtype: 'streaming_fallback',
cause: 'stream_retry',
},
{
type: 'system',
subtype: 'api_retry',
attempt: 2,
max_retries: 3,
retry_delay_ms: 50,
error_status: 503,
error: 'server_error',
},
{
type: 'user',
uuid: '11111111-1111-4111-8111-111111111111',
timestamp: '2026-08-11T00:00:00.000Z',
message: {
role: 'user',
content: [{
type: 'tool_result',
tool_use_id: 'tool-call',
content: 'tool output',
is_error: false,
}],
},
},
]
})
afterEach(async () => {
const queue = await import('../../utils/sdkEventQueue.js')
queue.setAgentRunMessageSink(undefined)
queue.drainSdkEvents()
if (previousConfigDir === undefined) delete process.env.CLAUDE_CONFIG_DIR
else process.env.CLAUDE_CONFIG_DIR = previousConfigDir
await rm(configDir, { recursive: true, force: true })
})
async function createRun(routeOptions: RouteOptions) {
const [
{ getDefaultAppState },
{ createFileStateCacheWithSizeLimit },
{ asSystemPrompt },
agentModule,
] = await Promise.all([
import('../../state/AppStateStore.js'),
import('../../utils/fileStateCache.js'),
import('../../utils/systemPromptType.js'),
import('./runAgent.js'),
])
const agentDefinition = {
agentType: 'stream-reviewer',
whenToUse: 'Verify streaming',
rawSystemPrompt: 'Stream.',
getSystemPrompt: () => 'Stream.',
source: 'projectSettings',
} as const
const parentState = getDefaultAppState()
const toolUseContext = {
options: {
commands: [],
debug: false,
mainLoopModel: 'sonnet',
tools: [],
verbose: false,
thinkingConfig: { type: 'disabled' as const },
mcpClients: [],
mcpResources: {},
isNonInteractiveSession: true,
agentDefinitions: { activeAgents: [agentDefinition], allAgents: [agentDefinition] },
},
abortController: new AbortController(),
readFileState: createFileStateCacheWithSizeLimit(),
getAppState: () => parentState,
setAppState: () => {},
setResponseLength: () => {},
messages: [],
toolUseId: 'parent-agent-tool',
}
return agentModule.runAgent({
agentDefinition,
promptMessages: [],
toolUseContext: toolUseContext as never,
canUseTool: (async () => ({ behavior: 'allow' })) as never,
isAsync: 'isAsync' in routeOptions ? routeOptions.isAsync ?? false : false,
querySource: 'querySource' in routeOptions
? routeOptions.querySource
: 'agent:custom',
override: {
userContext: {},
systemContext: {},
systemPrompt: asSystemPrompt([]),
agentId: 'physical-run-agent' as never,
},
availableTools: [],
...('spawningToolUseId' in routeOptions
? { spawningToolUseId: routeOptions.spawningToolUseId }
: {}),
...('streamTargetAgentId' in routeOptions
? { streamTargetAgentId: routeOptions.streamTargetAgentId }
: {}),
})
}
test.each([
['foreground SubAgent', { spawningToolUseId: 'agent-tool' }],
['background Agent', { spawningToolUseId: 'background-tool', isAsync: true }],
['workflow Agent', { querySource: 'workflow_agent' }],
['teammate Agent', { streamTargetAgentId: 'worker@team' }],
] as const)('forwards text, thinking, tool, result, and retry events for %s', async (_kind, routeOptions) => {
const [print, remoteIoModule, structuredIoModule, teammate] = await Promise.all([
import('../../cli/print.js'),
import('../../cli/remoteIO.js'),
import('../../cli/structuredIO.js'),
import('../../utils/teammateContext.js'),
])
const remote = new structuredIoModule.StructuredIO((async function* () {})())
Object.setPrototypeOf(remote, remoteIoModule.RemoteIO.prototype)
const events: Array<AgentRunMessageEvent & { uuid: string; session_id: string }> = []
const removeSink = print.bindAgentRunMessageSink(remote)
const yielded: unknown[] = []
const consume = async () => {
for await (const message of await createRun(routeOptions)) {
yielded.push(message)
}
}
try {
if ('streamTargetAgentId' in routeOptions) {
const context = {
...teammate.createTeammateContext({
agentId: routeOptions.streamTargetAgentId,
agentName: 'worker',
teamName: 'team',
planModeRequired: false,
parentSessionId: 'leader-session',
abortController: new AbortController(),
}),
streamScopeId: 'team-stream-scope',
}
await teammate.runWithTeammateContext(context, consume)
} else {
await consume()
}
} finally {
removeSink()
}
remote.outbound.done()
while (true) {
const next = await remote.outbound.next()
if (next.done) break
events.push(next.value as AgentRunMessageEvent & { uuid: string; session_id: string })
}
expect(yielded).toEqual([queryMessages.at(-1)])
expect(events.map(event => event.event_kind)).toEqual([
...queryMessages.map(() => 'message' as const),
'complete',
])
expect(events.slice(0, -1).map(event => event.message)).toEqual(queryMessages)
expect(events.every(event => event.run_agent_id === 'physical-run-agent')).toBe(true)
expect(new Set(events.map(event => event.stream_id)).size).toBe(1)
expect(events[0]?.stream_id).toBeString()
expect(events.every(event => event.target_agent_id === (
'streamTargetAgentId' in routeOptions
? routeOptions.streamTargetAgentId
: 'physical-run-agent'
))).toBe(true)
expect(events.every(event => event.target_agent_scope_id === (
'streamTargetAgentId' in routeOptions ? 'team-stream-scope' : undefined
))).toBe(true)
})
test('emits one terminal event when the consumer returns the generator early', async () => {
const queue = await import('../../utils/sdkEventQueue.js')
queryMessages = [
{
type: 'stream_event',
event: {
type: 'content_block_delta',
index: 0,
delta: { type: 'text_delta', text: 'before return' },
},
},
{
type: 'attachment',
attachment: { type: 'structured_output', data: { stop: true } },
},
{
type: 'stream_event',
event: {
type: 'content_block_delta',
index: 0,
delta: { type: 'text_delta', text: 'must not be consumed' },
},
},
]
const events: Array<AgentRunMessageEvent> = []
const removeSink = queue.setAgentRunMessageSink(event => events.push(event))
try {
const generator = await createRun({ spawningToolUseId: 'agent-tool' })
expect((await generator.next()).value).toEqual(queryMessages[1])
expect(events.map(event => event.event_kind)).toEqual(['message'])
await generator.return()
} finally {
removeSink()
}
expect(events.map(event => event.event_kind)).toEqual(['message', 'cancelled'])
})
test('drops directed stream events without occupying the generic SDK queue', async () => {
const queue = await import('../../utils/sdkEventQueue.js')
queue.drainSdkEvents()
queue.emitAgentRunMessage({
runAgentId: 'physical-run-agent',
streamId: 'isolated-stream',
targetAgentId: 'worker@team',
}, {
kind: 'message',
message: queryMessages[2],
})
queue.enqueueSdkEvent({
type: 'system',
subtype: 'session_state_changed',
state: 'running',
})
expect(queue.drainSdkEvents()).toEqual([
expect.objectContaining({
subtype: 'session_state_changed',
state: 'running',
}),
])
})
})
+52
View File
@@ -78,6 +78,8 @@ import {
} from '../../utils/telemetry/perfettoTracing.js'
import type { ContentReplacementState } from '../../utils/toolResultStorage.js'
import { createAgentId } from '../../utils/uuid.js'
import { emitAgentRunMessage } from '../../utils/sdkEventQueue.js'
import { getTeammateContext } from '../../utils/teammateContext.js'
import { resolveAgentTools } from './agentToolUtils.js'
import { type AgentDefinition, isBuiltInAgent } from './loadAgentsDir.js'
@@ -309,6 +311,7 @@ export async function* runAgent({
description,
spawningToolUseId,
ownerAgentId,
streamTargetAgentId,
persistedAgentType,
alreadyPersistedMessageCount,
workflow,
@@ -370,6 +373,8 @@ export async function* runAgent({
spawningToolUseId?: string
/** Parent agent that owns this run's task lifecycle. Undefined means root. */
ownerAgentId?: string
/** Logical target used only to route live output to a synthetic desktop run. */
streamTargetAgentId?: string
/** Stable upstream identity for a resumed agent. A named teammate may use
* the general-purpose runtime definition after its original definition is
* no longer active, but its transcript metadata must keep the teammate name
@@ -407,6 +412,25 @@ export async function* runAgent({
)
const agentId = override?.agentId ? override.agentId : createAgentId()
const teammateContext = getTeammateContext()
const agentRunRoute = spawningToolUseId || ownerAgentId || streamTargetAgentId || querySource === 'workflow_agent'
? {
runAgentId: agentId,
streamId: createAgentId(),
targetAgentId: streamTargetAgentId ?? agentId,
...(teammateContext?.streamScopeId
? { targetAgentScopeId: teammateContext.streamScopeId }
: {}),
}
: undefined
let agentRunTerminalEmitted = false
const emitAgentRunTerminal = (
event: { kind: 'complete' } | { kind: 'cancelled' } | { kind: 'error'; error: string },
) => {
if (!agentRunRoute || agentRunTerminalEmitted) return
agentRunTerminalEmitted = true
emitAgentRunMessage(agentRunRoute, event)
}
// Register agent in Perfetto trace for hierarchy visualization
if (isPerfettoTracingEnabled()) {
@@ -842,6 +866,18 @@ export async function* runAgent({
maxTurns: maxTurns ?? agentDefinition.maxTurns,
})) {
onQueryProgress?.()
if (
agentRunRoute &&
(message.type === 'stream_event' ||
message.type === 'assistant' ||
message.type === 'user' ||
(
message.type === 'system' &&
(message.subtype === 'streaming_fallback' || message.subtype === 'api_retry')
))
) {
emitAgentRunMessage(agentRunRoute, { kind: 'message', message })
}
// Forward subagent API request starts to parent's metrics display
// so TTFT/OTPS update during subagent execution.
if (
@@ -899,7 +935,23 @@ export async function* runAgent({
if (isBuiltInAgent(agentDefinition) && agentDefinition.callback) {
agentDefinition.callback()
}
if (agentRunRoute) {
emitAgentRunTerminal({ kind: 'complete' })
}
} catch (error) {
if (error instanceof AbortError) {
emitAgentRunTerminal({ kind: 'cancelled' })
} else {
emitAgentRunTerminal({
kind: 'error',
error: error instanceof Error ? error.message : String(error),
})
}
throw error
} finally {
// A consumer can stop the async generator with return()/break() without
// throwing. Discard its unfinished partial before settling the stream.
emitAgentRunTerminal({ kind: 'cancelled' })
// Clean up agent-specific MCP servers (runs on normal completion, abort, or error)
await mcpCleanup()
// Clean up agent's session hooks
+70
View File
@@ -107,15 +107,47 @@ type AgentToolActivityEvent = {
owner_agent_id?: string
}
export type AgentRunMessageEvent = {
type: 'system'
subtype: 'agent_run_message'
run_agent_id: string
stream_id: string
target_agent_id: string
target_agent_scope_id?: string
event_kind: 'message' | 'complete' | 'cancelled' | 'error'
message?: unknown
error?: string
}
export type SdkEvent =
| TaskStartedEvent
| TaskProgressEvent
| TaskNotificationSdkEvent
| SessionStateChangedEvent
| AgentToolActivityEvent
| AgentRunMessageEvent
const MAX_QUEUE_SIZE = 1000
const queue: SdkEvent[] = []
type EnvelopedAgentRunMessage = AgentRunMessageEvent & {
uuid: UUID
session_id: string
}
let agentRunMessageSink: ((event: EnvelopedAgentRunMessage) => void) | undefined
/**
* Agent runs execute below the main query generator, so their token deltas
* cannot wait for the next drainSdkEvents() call. The headless printer binds
* this sink to its existing outbound FIFO while it is alive.
*/
export function setAgentRunMessageSink(
sink: ((event: EnvelopedAgentRunMessage) => void) | undefined,
): () => void {
agentRunMessageSink = sink
return () => {
if (agentRunMessageSink === sink) agentRunMessageSink = undefined
}
}
export function enqueueSdkEvent(event: SdkEvent): void {
// SDK events are only consumed (drained) in headless/streaming mode.
@@ -123,12 +155,50 @@ export function enqueueSdkEvent(event: SdkEvent): void {
if (!getIsNonInteractiveSession()) {
return
}
if (event.subtype === 'agent_run_message') {
if (agentRunMessageSink) {
agentRunMessageSink({
...event,
uuid: randomUUID(),
session_id: getSessionId(),
})
}
return
}
if (queue.length >= MAX_QUEUE_SIZE) {
queue.shift()
}
queue.push(event)
}
export function emitAgentRunMessage(
route: {
runAgentId: string
streamId: string
targetAgentId?: string
targetAgentScopeId?: string
},
event:
| { kind: 'message'; message: unknown }
| { kind: 'complete' }
| { kind: 'cancelled' }
| { kind: 'error'; error: string },
): void {
enqueueSdkEvent({
type: 'system',
subtype: 'agent_run_message',
run_agent_id: route.runAgentId,
stream_id: route.streamId,
target_agent_id: route.targetAgentId ?? route.runAgentId,
...(route.targetAgentScopeId
? { target_agent_scope_id: route.targetAgentScopeId }
: {}),
event_kind: event.kind,
...(event.kind === 'message' ? { message: event.message } : {}),
...(event.kind === 'error' ? { error: event.error } : {}),
})
}
export function drainSdkEvents(): Array<
SdkEvent & { uuid: UUID; session_id: string }
> {
+15 -2
View File
@@ -31,6 +31,7 @@ import {
} from '../tasks.js'
import {
createTeammateContext,
getTeammateContext,
runWithTeammateContext,
} from '../teammateContext.js'
import {
@@ -788,12 +789,14 @@ describe('in-process teammate activity synchronization', () => {
parentSessionId: 'leader-session',
}
const member = createMember(agentName, teamName)
await writeTeamFileAsync(teamName, {
const teamFile: TeamFile = {
name: teamName,
createdAt: Date.now(),
leadAgentId: `team-lead@${teamName}`,
leadSessionId: identity.parentSessionId,
members: [member],
})
}
await writeTeamFileAsync(teamName, teamFile)
await createTask(teamName, {
subject: 'Earlier explicit assignment',
description: 'Makes an accidental first-turn claim observable',
@@ -848,9 +851,13 @@ describe('in-process teammate activity synchronization', () => {
} as unknown as ToolUseContext
let firstPrompt = ''
let streamTargetAgentId: string | undefined
let streamScopeId: string | undefined
const runAgent = spyOn(runAgentModule, 'runAgent').mockImplementation(
async function* (input: Parameters<typeof runAgentModule.runAgent>[0]) {
firstPrompt = JSON.stringify(input.promptMessages[0])
streamTargetAgentId = input.streamTargetAgentId
streamScopeId = getTeammateContext()?.streamScopeId
expect(readTeamFile(teamName)?.members[0]?.isActive).toBe(true)
abortController.abort()
},
@@ -874,6 +881,12 @@ describe('in-process teammate activity synchronization', () => {
expect(result.success).toBe(true)
expect(firstPrompt).toContain('Review the release')
expect(firstPrompt).not.toContain('Audit first-turn delivery')
expect(streamTargetAgentId).toBe(identity.agentId)
expect(streamScopeId).toBe(JSON.stringify([
teamName,
identity.parentSessionId,
teamFile.createdAt,
]))
expect(readTeamFile(teamName)?.members[0]?.isActive).toBe(false)
const unclaimedTask = (await listTasks(teamName))
.find(task => task.subject === 'Audit first-turn delivery')
+14 -3
View File
@@ -94,7 +94,10 @@ import {
type Task,
} from '../tasks.js'
import type { TeammateContext } from '../teammateContext.js'
import { runWithTeammateContext } from '../teammateContext.js'
import {
createTeamStreamScopeId,
runWithTeammateContext,
} from '../teammateContext.js'
import {
createIdleNotification,
getLastPeerDmSummary,
@@ -115,7 +118,7 @@ import {
createPermissionRequest,
sendPermissionRequestViaMailbox,
} from './permissionSync.js'
import { setMemberActive } from './teamHelpers.js'
import { readTeamFile, setMemberActive } from './teamHelpers.js'
import { TEAMMATE_SYSTEM_PROMPT_ADDENDUM } from './teammatePromptAddendum.js'
type SetAppStateFn = (updater: (prev: AppState) => AppState) => void
@@ -976,6 +979,13 @@ export async function runInProcessTeammate(
invokingRequestId,
} = config
const { setAppState } = toolUseContext
const teamFile = readTeamFile(identity.teamName)
const scopedTeammateContext: TeammateContext = {
...teammateContext,
...(teamFile
? { streamScopeId: createTeamStreamScopeId(teamFile) }
: {}),
}
logForDebugging(
`[inProcessRunner] Starting agent loop for ${identity.agentId}`,
@@ -1207,7 +1217,7 @@ export async function runInProcessTeammate(
let workWasAborted = false
// Run agent within contexts
const runActiveTurn = () => runWithTeammateContext(teammateContext, async () => {
const runActiveTurn = () => runWithTeammateContext(scopedTeammateContext, async () => {
return runWithAgentContext(agentContext, async () => {
// Mark task as running (not idle)
updateTaskState(
@@ -1250,6 +1260,7 @@ export async function runInProcessTeammate(
availableTools: toolUseContext.options.tools,
allowedTools,
contentReplacementState: teammateReplacementState,
streamTargetAgentId: identity.agentId,
})) {
// Check lifecycle abort first (kills whole teammate)
if (abortController.signal.aborted) {
+14
View File
@@ -32,12 +32,26 @@ export type TeammateContext = {
planModeRequired: boolean
/** Leader's session ID (for transcript correlation) */
parentSessionId: string
/** Immutable Team creation tuple used only to isolate live UI routing. */
streamScopeId?: string
/** Discriminator - always true for in-process teammates */
isInProcess: true
/** Abort controller for lifecycle management (linked to parent) */
abortController: AbortController
}
export function createTeamStreamScopeId(team: {
name: string
leadSessionId?: string
createdAt: number
}): string {
return JSON.stringify([
team.name,
team.leadSessionId ?? '',
Number.isFinite(team.createdAt) ? team.createdAt : 0,
])
}
const teammateContextStorage = new AsyncLocalStorage<TeammateContext>()
/**
@@ -6,8 +6,10 @@ import {
} from '../../Tool.js'
const capturedOwners: Array<string | undefined> = []
const runAgentMock = mock((params: { ownerAgentId?: string }) => {
const capturedSources: Array<string | undefined> = []
const runAgentMock = mock((params: { ownerAgentId?: string; querySource?: string }) => {
capturedOwners.push(params.ownerAgentId)
capturedSources.push(params.querySource)
return (async function* () {})()
})
const runAgentModule = await import('../../tools/AgentTool/runAgent.js')
@@ -54,10 +56,12 @@ async function invoke(ownerAgentId?: string): Promise<void> {
describe('runWorkflowAgent ownership', () => {
test('persists nested ownership without adding an owner to root workers', async () => {
capturedOwners.length = 0
capturedSources.length = 0
await invoke('parent-agent')
await invoke()
expect(capturedOwners).toEqual(['parent-agent', undefined])
expect(capturedSources).toEqual(['workflow_agent', 'workflow_agent'])
})
})