diff --git a/src/callLifecycle.ts b/src/callLifecycle.ts index 2ff5ccf..8188105 100644 --- a/src/callLifecycle.ts +++ b/src/callLifecycle.ts @@ -66,6 +66,9 @@ export type TracedAgent = CallScope & PermissionContext & { agentId?: string; outcome?: ToolResult; stopSeen: boolean; + /** A protocol-specific terminal event can decide the result before nested + * calls have drained. */ + completion?: AgentCompletion; children: Set; /** Chat responses already emitted from this Agent's Stop snapshots. */ seenResponses: Set; @@ -73,6 +76,11 @@ export type TracedAgent = CallScope & PermissionContext & { export type TracedCall = TracedTool | TracedAgent; +type AgentCompletion = ( + | { outcome: ToolResult; failureType?: string } + | { orphanReason: string } +) & { endTime?: Date }; + /** Secondary indexes for the identities exposed by Claude's hooks. */ export type CallState = { byToolUseId: Map; @@ -107,7 +115,8 @@ export function beginCall( parent: CallParent, args: ToolCall, ): TracedCall | undefined { - if (parent.kind === 'agent' && parent.stopSeen && parent.outcome) return undefined; + if (parent.kind === 'agent' + && (parent.completion || (parent.stopSeen && parent.outcome))) return undefined; if (state.byToolUseId.has(args.toolUseId) || state.toolUseTombstones.has(args.toolUseId)) return undefined; @@ -295,10 +304,13 @@ export function denyCall(state: CallState, toolUseId: string, reason: string): v call.span.result = reason; call.span.setAttributes({ [ATTR.ERROR_TYPE]: 'permission_denied' }); call.span.end({ error: new Error(reason) }); + completeCall(state, call); } else { - finishAgentSpan(call, { ok: false, error: reason }, 'permission_denied'); + finishAgentCall(state, call, { + outcome: { ok: false, error: reason }, + failureType: 'permission_denied', + }); } - completeCall(state, call); } /** Apply the exact tool_use_id terminal event once. */ @@ -320,6 +332,24 @@ export function recordCallOutcome( finishAgentIfReady(state, call); } +/** Record an Agent result whose protocol has a later terminal event. */ +export function deferAgentOutcome(call: TracedAgent, outcome: ToolResult): void { + resolvePermission(call, true); + call.outcome ??= outcome; +} + +/** Complete an Agent from a protocol-specific terminal event. */ +export function finishAgentCall( + state: CallState, + call: TracedAgent, + completion: { outcome: ToolResult; failureType?: string } | { orphanReason: string }, + endTime?: Date, +): void { + if ('outcome' in completion) resolvePermission(call, true); + call.completion ??= { ...completion, ...(endTime ? { endTime } : {}) }; + finishAgentIfReady(state, call); +} + function finishToolCall(call: TracedTool, outcome: ToolResult): void { resolvePermission(call, true); if (outcome.ok) { @@ -355,6 +385,23 @@ function finishAgentSpan( }); } +function endAgent(call: TracedAgent, completion: AgentCompletion): void { + if ('outcome' in completion) { + finishAgentSpan( + call, + completion.outcome, + completion.failureType, + completion.endTime, + ); + return; + } + call.span.setAttributes({ [ATTR.WEAVE_ORPHAN_REASON]: completion.orphanReason }); + call.span.end({ + error: new Error(`call did not complete (${completion.orphanReason})`), + ...(completion.endTime ? { endTime: completion.endTime } : {}), + }); +} + function errorType(error: string): string { return error.trim().match(/^[A-Z][A-Za-z_]*Error/)?.[0] ?? 'tool_error'; } @@ -408,8 +455,11 @@ export function recordAgentStop( } function finishAgentIfReady(state: CallState, call: TracedAgent): void { - if (!call.stopSeen || !call.outcome || call.children.size) return; - finishAgentSpan(call, call.outcome); + if (call.children.size) return; + const completion = call.completion + ?? (call.stopSeen && call.outcome ? { outcome: call.outcome } : undefined); + if (!completion) return; + endAgent(call, completion); completeCall(state, call); } @@ -439,15 +489,43 @@ export function finalizeOpenCalls( roots: Iterable, reason: string, endTime?: Date, + defer: (call: TracedCall) => boolean = () => false, ): string[] { const closed: string[] = []; + const finalCompletion = (call: TracedAgent, at?: Date): AgentCompletion => { + if (call.completion) { + return at && !call.completion.endTime + ? { ...call.completion, endTime: at } + : call.completion; + } + if (call.outcome) { + return { outcome: call.outcome, ...(at ? { endTime: at } : {}) }; + } + if (!call.toolUseId && call.stopSeen) { + return { + outcome: { ok: true, output: null }, + ...(at ? { endTime: at } : {}), + }; + } + return { orphanReason: reason, ...(at ? { endTime: at } : {}) }; + }; + const awaitChildren = (completion: AgentCompletion): AgentCompletion => + 'outcome' in completion + ? { + outcome: completion.outcome, + ...(completion.failureType ? { failureType: completion.failureType } : {}), + } + : { orphanReason: completion.orphanReason }; const closeChildren = (parent: CallParent) => { for (const call of [...parent.children].reverse()) { if (call.kind === 'agent') closeChildren(call); - if (call.kind === 'agent' && call.outcome) { - finishAgentSpan(call, call.outcome, undefined, endTime); - } else if (call.kind === 'agent' && !call.toolUseId && call.stopSeen) { - call.span.end(endTime ? { endTime } : undefined); + if (defer(call)) continue; + if (call.kind === 'agent' && call.children.size) { + call.completion = awaitChildren(finalCompletion(call)); + continue; + } + if (call.kind === 'agent') { + endAgent(call, finalCompletion(call, endTime)); } else { call.span.setAttributes({ [ATTR.WEAVE_ORPHAN_REASON]: reason }); call.span.end({ diff --git a/src/hookHandler.ts b/src/hookHandler.ts index f7fff54..8c9760e 100644 --- a/src/hookHandler.ts +++ b/src/hookHandler.ts @@ -17,6 +17,7 @@ import type { StopHookInput, SubagentStartHookInput, SubagentStopHookInput, + TeammateIdleHookInput, UserPromptSubmitHookInput, } from '@anthropic-ai/claude-agent-sdk'; import * as weave from 'weave'; @@ -48,6 +49,7 @@ import { ATTR, assistantOutputMessages, snippet } from './genaiSpans.js'; import type { SpanParent } from './genaiSpans.js'; import { parseSessionFd } from './parser.js'; import { Session } from './session.js'; +import { TeamCoordinator } from './teamCoordinator.js'; import { TranscriptFile, readSubagentPrompt, @@ -118,9 +120,11 @@ function parseHookInput(payload: unknown): HookInput | undefined { export class HookHandler { private readonly sessions = new Map(); - private readonly sessionQueues = new Map>(); + private eventQueue = Promise.resolve(); + private eventSequence = 0; /** InstructionsLoaded can arrive before SessionStart. */ private readonly pendingInstructions = new Map>(); + private readonly teams = new TeamCoordinator(); constructor( private readonly agentName: string, @@ -134,33 +138,26 @@ export class HookHandler { return; } - const sessionId = input.session_id; - const previous = this.sessionQueues.get(sessionId) ?? Promise.resolve(); - const next = previous.then(() => this.route(input)); - this.sessionQueues.set(sessionId, next); - try { - await next; - } finally { - if (this.sessionQueues.get(sessionId) === next) { - this.sessionQueues.delete(sessionId); - } - } + const sequence = ++this.eventSequence; + const next = this.eventQueue.then(() => this.route(input, sequence)); + this.eventQueue = next; + await next; } - private async route(input: HookInput): Promise { + private async route(input: HookInput, sequence: number): Promise { const sessionId = input.session_id; this.log( 'INFO', `${input.hook_event_name} session=${sessionId}${input.agent_id ? ` agent=${input.agent_id}` : ''}`, ); try { - await weave.runIsolated(() => this.dispatchEvent(input)); + await weave.runIsolated(() => this.dispatchEvent(input, sequence)); } catch (err) { this.log('ERROR', `Error handling ${input.hook_event_name}: ${err}`); } } - private async dispatchEvent(input: HookInput): Promise { + private async dispatchEvent(input: HookInput, sequence: number): Promise { const sessionId = input.session_id; switch (input.hook_event_name) { case 'SessionStart': @@ -173,7 +170,7 @@ export class HookHandler { await this.handleUserPromptSubmit(sessionId, input); return; case 'PreToolUse': - await this.handlePreToolUse(sessionId, input); + await this.handlePreToolUse(sessionId, input, sequence); return; case 'PermissionRequest': await this.handlePermissionRequest(sessionId, input); @@ -183,7 +180,7 @@ export class HookHandler { return; case 'PostToolUse': case 'PostToolUseFailure': - await this.handlePostToolResult(sessionId, input); + await this.handlePostToolResult(sessionId, input, sequence); return; case 'SubagentStart': await this.handleSubagentStart(sessionId, input); @@ -191,6 +188,9 @@ export class HookHandler { case 'SubagentStop': await this.handleSubagentStop(sessionId, input); return; + case 'TeammateIdle': + await this.handleTeammateIdle(input, sequence); + return; case 'PreCompact': this.handlePreCompact(sessionId, input); return; @@ -332,7 +332,11 @@ export class HookHandler { return; } - const result = session.submitPrompt(input.prompt_id, input.prompt); + const result = session.submitPrompt( + input.prompt_id, + input.prompt, + call => call.kind === 'agent' && this.teams.has(call), + ); if (!result.created) return; this.log( 'DEBUG', @@ -362,6 +366,7 @@ export class HookHandler { private async handlePreToolUse( sessionId: string, input: PreToolUseHookInput, + sequence: number, ): Promise { const descriptor = this.toolCall(input); if (!descriptor) return; @@ -377,6 +382,9 @@ export class HookHandler { return; } const call = beginCall(session.calls, parent, descriptor); + if (call?.kind === 'agent') { + this.teams.registerDispatch(session, call, sequence); + } if (call && !input.agent_id) call.root.phase = 'active'; } @@ -435,6 +443,7 @@ export class HookHandler { private async handlePostToolResult( sessionId: string, input: PostToolResultHookInput, + sequence: number, ): Promise { const result = this.toolResult(input); if (typeof input.tool_use_id !== 'string' || !result) return; @@ -449,9 +458,20 @@ export class HookHandler { ?? await this.getOrReconstructSession(sessionId, input); if (!session) return; - if (!existingCall) await this.recoverCall(session, input, descriptor!); + const call = existingCall + ?? await this.recoverCall(session, input, descriptor!); + if (!existingCall && call?.kind === 'agent') { + this.teams.registerDispatch(session, call, sequence); + } + if (call?.kind === 'agent') { + const update = this.teams.postOutcome(call, result); + if (update.handled) { + this.settleTeamCompletions(update.completions, session); + return; + } + } recordCallOutcome(session.calls, input.tool_use_id, result); - session.finishSupersededTurns(); + this.settleSession(session); } private toolCall( @@ -525,9 +545,13 @@ export class HookHandler { ?? await this.getOrReconstructSession(sessionId, input); if (!session) return; - if (!existingCall) await this.recoverCall(session, input, descriptor!); - denyCall(session.calls, input.tool_use_id, input.reason); - session.finishSupersededTurns(); + const call = existingCall + ?? await this.recoverCall(session, input, descriptor!); + const completions = call?.kind === 'agent' + ? this.teams.deny(call, input.reason) + : undefined; + if (!completions) denyCall(session.calls, input.tool_use_id, input.reason); + this.settleTeamCompletions(completions ?? [], session); } private async handleSubagentStart( @@ -643,6 +667,13 @@ export class HookHandler { ) : undefined; const lifecycle = match.kind === 'found' ? match.call : recovered; + if (lifecycle && this.teams.has(lifecycle)) { + this.log( + 'DEBUG', + `Subagent stopped: agentId=${input.agent_id} awaiting TeammateIdle`, + ); + return; + } const parent = match.kind === 'found' ? match.call.span : recovered?.span ?? turn?.span; @@ -688,7 +719,31 @@ export class HookHandler { 'DEBUG', `Subagent stopped: agentId=${input.agent_id} type=${input.agent_type} model=${transcript.model ?? 'unknown'} match=${match.kind}`, ); - session.finishSupersededTurns(); + this.settleSession(session); + } + + private async handleTeammateIdle( + input: TeammateIdleHookInput, + sequence: number, + ): Promise { + const result = await this.teams.recordIdle({ + sequence, + teamName: input.team_name, + memberName: input.teammate_name, + transcriptPath: input.transcript_path, + }); + if (!result.completions.length) { + this.log( + 'DEBUG', + `TeammateIdle: ${result.status} ${input.team_name}::${input.teammate_name}`, + ); + return; + } + this.log( + 'DEBUG', + `TeammateIdle: completed ${input.team_name}::${input.teammate_name}`, + ); + this.settleTeamCompletions(result.completions); } private async handleStop(sessionId: string, input: StopHookInput): Promise { @@ -715,14 +770,46 @@ export class HookHandler { ?? await this.getOrReconstructSession(sessionId, input); if (!session) return; - const turnCount = session.finishAtSessionEnd(input.prompt_id); + const result = session.finishAtSessionEnd( + input.prompt_id, + call => call.kind === 'agent' + && (this.teams.has(call) || this.isNamedBackgroundAgent(call)), + ); this.log( 'DEBUG', - `SessionEnd: session=${sessionId} reason=${input.reason} transcript_path=${session.transcriptPath} turns=${turnCount}`, + `SessionEnd: session=${sessionId} reason=${input.reason} transcript_path=${session.transcriptPath} turns=${result.turnCount}`, ); - this.sessions.delete(sessionId); + if (result.deferred) { + this.log('INFO', `SessionEnd: deferred session ${sessionId} until background work completes`); + return; + } + this.closeSession(session); + } + + private isNamedBackgroundAgent(call: TracedAgent): boolean { + const name = call.input['name']; + return typeof name === 'string' && name.trim().length > 0; + } + + private settleSession(session: Session): void { + if (session.completeDeferredEnd()) { + this.closeSession(session); + } + } + + private settleTeamCompletions( + completions: Array<{ owner: Session }>, + fallback?: Session, + ): void { + const owners = new Set(completions.map(completion => completion.owner)); + if (!owners.size && fallback) owners.add(fallback); + for (const owner of owners) this.settleSession(owner); + } + + private closeSession(session: Session): void { + this.sessions.delete(session.sessionId); session.close(); - this.log('INFO', `Finished session ${sessionId}`); + this.log('INFO', `Finished session ${session.sessionId}`); } hasInFlightWork(): boolean { @@ -734,13 +821,15 @@ export class HookHandler { /** Admission must be stopped before taking this snapshot. */ async waitForPendingEvents(): Promise { - await Promise.all([...this.sessionQueues.values()]); + await this.eventQueue; } finalizeForShutdown(): void { for (const session of this.sessions.values()) { try { - session.finishOpenTurns('daemon_shutdown'); + const endTime = new Date(Date.now() + 1); + this.teams.orphanSession(session.sessionId, 'daemon_shutdown', endTime); + session.finishOpenTurns('daemon_shutdown', endTime); } catch (err) { this.log('ERROR', `Error finalizing session ${session.sessionId} at shutdown: ${err}`); } diff --git a/src/session.ts b/src/session.ts index 7bbd27a..899d89a 100644 --- a/src/session.ts +++ b/src/session.ts @@ -84,6 +84,7 @@ export class Session { readonly cwd: string; readonly source: string; readonly initialRequestModel?: string; + readonly conversation: weave.Conversation; /** File path → latest loaded contents, preserving first-load order. */ private readonly systemInstructions = new Map(); @@ -91,6 +92,7 @@ export class Session { readonly calls = newCallState(); private currentTurn?: TurnTrace; private pendingCompaction?: CompactionAttrs; + private endRequested = false; private constructor( options: NewSessionOptions, @@ -116,8 +118,6 @@ export class Session { }); } - private readonly conversation: weave.Conversation; - static async create(options: NewSessionOptions): Promise { const conversationId = await resolveConversationId( options.sessionId, @@ -138,10 +138,12 @@ export class Session { submitPrompt( promptId: string | undefined, prompt: string, + defer: (call: TracedCall) => boolean = () => false, ): { created: boolean; replacedOpenTurn: boolean } { if (promptId !== undefined && this.turnForPrompt(promptId)) { return { created: false, replacedOpenTurn: false }; } + this.endRequested = false; const previous = this.currentTurn; let responseOffsetFloor: number | undefined; @@ -150,7 +152,9 @@ export class Session { parseSessionFd(this.transcript.getFd()) ?? { turns: [] }, ).length; responseOffsetFloor = previous.responseLimit; - if (promptId === undefined || previous.children.size === 0) { + if (promptId === undefined) { + this.finalizeTurn(previous, 'superseded_by_next_prompt', defer); + } else if (previous.children.size === 0) { this.finalizeTurn(previous, 'superseded_by_next_prompt'); } } @@ -192,34 +196,56 @@ export class Session { }; } - finishAtSessionEnd(promptId: string | undefined): number { + finishAtSessionEnd( + promptId: string | undefined, + defer: (call: TracedCall) => boolean = () => false, + ): { turnCount: number; deferred: boolean } { + this.endRequested = true; const parsed = this.parseTranscript(); this.reconcileFinalTurn(promptId, parsed); - return this.finishTurns('session_ended', parsed, spanCloseTime()); + const turnCount = this.finishTurns( + 'session_ended', + parsed, + spanCloseTime(), + defer, + ); + return { turnCount, deferred: this.hasOpenCalls() }; } - finishOpenTurns(orphanReason: string): number { - return this.finishTurns(orphanReason, this.parseTranscript(), spanCloseTime()); + finishOpenTurns( + orphanReason: string, + endTime = spanCloseTime(), + ): number { + return this.finishTurns(orphanReason, this.parseTranscript(), endTime); } private finishTurns( orphanReason: string, parsed: ParsedSession | null, endTime: Date, + defer: (call: TracedCall) => boolean = () => false, ): number { const turnCount = this.turns.size; for (const turn of [...this.turns]) { this.recordFinalTurnOutput(turn, orphanReason, parsed); - this.closeTurn(turn, orphanReason, endTime); + this.closeTurn(turn, orphanReason, endTime, defer); } return turnCount; } hasInFlightWork(): boolean { - return [...this.turns].some(turn => turn.children.size > 0) + return this.hasOpenCalls() || [...this.turns].some(turn => turn.phase === 'active'); } + completeDeferredEnd(): boolean { + this.finishSupersededTurns(); + if (!this.endRequested || this.hasOpenCalls()) return false; + const endTime = spanCloseTime(); + for (const turn of [...this.turns]) this.endTurn(turn, endTime); + return true; + } + close(): void { this.transcript.close(); } @@ -418,12 +444,17 @@ export class Session { }); } - private finalizeTurn(turn: TurnTrace, orphanReason: string): void { + private finalizeTurn( + turn: TurnTrace, + orphanReason: string, + defer: (call: TracedCall) => boolean = () => false, + ): void { this.recordFinalTurnOutput(turn, orphanReason, this.parseTranscript()); this.closeTurn( turn, orphanReason, turn.children.size ? spanCloseTime() : undefined, + defer, ); } @@ -431,15 +462,22 @@ export class Session { turn: TurnTrace, orphanReason: string, endTime?: Date, + defer: (call: TracedCall) => boolean = () => false, ): void { for (const toolUseId of finalizeOpenCalls( this.calls, [turn], orphanReason, endTime, + defer, )) { this.log('DEBUG', `Closed pending call: ${toolUseId}`); } + if (turn.children.size) return; + this.endTurn(turn, endTime); + } + + private endTurn(turn: TurnTrace, endTime?: Date): void { turn.span.end(endTime ? { endTime } : undefined); this.turns.delete(turn); if (this.currentTurn === turn) this.currentTurn = undefined; @@ -453,6 +491,10 @@ export class Session { } } + private hasOpenCalls(): boolean { + return [...this.turns].some(turn => turn.children.size > 0); + } + private parseTranscript(): ParsedSession | null { try { return parseSessionFd(this.transcript.getFd()); diff --git a/src/teamCoordinator.ts b/src/teamCoordinator.ts new file mode 100644 index 0000000..bcdbf5e --- /dev/null +++ b/src/teamCoordinator.ts @@ -0,0 +1,187 @@ +// SPDX-FileCopyrightText: 2026 CoreWeave, Inc. +// SPDX-License-Identifier: MIT +// SPDX-PackageName: weave-claude-code + +import type * as weave from 'weave'; +import { emitChatSpans } from './chatSpans.js'; +import { + deferAgentOutcome, + denyCall, + finishAgentCall, +} from './callLifecycle.js'; +import type { ToolResult, TracedAgent } from './callLifecycle.js'; +import { ATTR, assistantOutputMessages, parseTimestamp } from './genaiSpans.js'; +import type { ParsedTurn } from './parser.js'; +import { VERSION } from './setup.js'; +import { readTeammateTurns } from './teamTranscripts.js'; +import type { Session } from './session.js'; + +type PendingTeam = { + session: Session; + call: TracedAgent; + teamName: string; + memberName: string; + sequence: number; +}; + +export type TeamCompletion = { + owner: Session; + teamName: string; + memberName: string; +}; + +export type TeamUpdate = { + handled: boolean; + completions: TeamCompletion[]; +}; + +export type TeamIdleUpdate = { + status: 'missing' | 'ambiguous' | 'completed'; + completions: TeamCompletion[]; +}; + +const text = (value: unknown) => + typeof value === 'string' && value.trim() ? value.trim() : undefined; + +function emitTeammate( + conversation: weave.Conversation, + memberName: string, + turns: ParsedTurn[], +): { model?: string; text?: string } { + const responses = turns.flatMap(turn => turn.responses); + const model = turns.filter(turn => turn.model).at(-1)?.model; + const span = conversation.startTurn({ + agentName: memberName, + agentVersion: VERSION, + model, + userMessage: turns[0]?.userText, + startTime: parseTimestamp(turns[0]?.startTime ?? responses[0]?.startTime), + }); + try { + emitChatSpans(span, responses, { agentName: memberName }); + const output = turns.flatMap(turn => turn.text); + if (output.length) { + span.setAttributes({ [ATTR.OUTPUT_MESSAGES]: assistantOutputMessages(output) }); + } + if (model) span.record({ model }); + return { model, text: turns.at(-1)?.text.join('\n') || undefined }; + } finally { + span.end({ endTime: parseTimestamp(responses.at(-1)?.endTime) ?? new Date() }); + } +} + +/** Correlates explicit Agent Team dispatches with their cross-session idle + * event. Recovery and weak-evidence matching are deliberately separate. */ +export class TeamCoordinator { + private readonly pending = new Set(); + + registerDispatch( + session: Session, + call: TracedAgent, + sequence: number, + ): void { + if (!call.toolUseId || this.has(call)) return; + const teamName = text(call.input['team_name']); + if (!teamName) return; + this.pending.add({ + session, + call, + teamName, + memberName: text(call.input['name']) ?? call.agentType, + sequence, + }); + } + + has(call: TracedAgent): boolean { + return [...this.pending].some(candidate => candidate.call === call); + } + + postOutcome(call: TracedAgent, outcome: ToolResult): TeamUpdate { + const pending = this.find(call); + if (!pending) return { handled: false, completions: [] }; + if (outcome.ok) { + deferAgentOutcome(call, outcome); + return { handled: true, completions: [] }; + } + + finishAgentCall(pending.session.calls, call, { outcome }); + this.pending.delete(pending); + return { handled: true, completions: [this.completed(pending)] }; + } + + async recordIdle(input: { + sequence: number; + teamName: string; + memberName: string; + transcriptPath?: string; + }): Promise { + const candidates = [...this.pending].filter(candidate => + candidate.sequence < input.sequence + && candidate.teamName === input.teamName + && candidate.memberName === input.memberName + && candidate.call.outcome?.ok); + if (candidates.length !== 1) { + return { + status: candidates.length ? 'ambiguous' : 'missing', + completions: [], + }; + } + + const [pending] = candidates; + const turns = input.transcriptPath + ? readTeammateTurns(input.transcriptPath) + : []; + if (!turns.length) return { status: 'missing', completions: [] }; + + const emitted = emitTeammate( + pending.session.conversation, + pending.memberName, + turns, + ); + if (emitted.model) { + pending.call.span.setAttributes({ [ATTR.RESPONSE_MODEL]: emitted.model }); + } + const original = pending.call.outcome; + finishAgentCall(pending.session.calls, pending.call, { + outcome: { + ok: true, + output: emitted.text ?? (original?.ok ? original.output : null), + }, + }); + this.pending.delete(pending); + return { status: 'completed', completions: [this.completed(pending)] }; + } + + deny(call: TracedAgent, reason: string): TeamCompletion[] | undefined { + const pending = this.find(call); + if (!pending || !call.toolUseId) return undefined; + denyCall(pending.session.calls, call.toolUseId, reason); + this.pending.delete(pending); + return [this.completed(pending)]; + } + + orphanSession(sessionId: string, reason: string, endTime: Date): void { + for (const pending of [...this.pending]) { + if (pending.session.sessionId !== sessionId) continue; + finishAgentCall( + pending.session.calls, + pending.call, + { orphanReason: reason }, + endTime, + ); + this.pending.delete(pending); + } + } + + private find(call: TracedAgent): PendingTeam | undefined { + return [...this.pending].find(candidate => candidate.call === call); + } + + private completed(pending: PendingTeam): TeamCompletion { + return { + owner: pending.session, + teamName: pending.teamName, + memberName: pending.memberName, + }; + } +} diff --git a/src/teamTranscripts.ts b/src/teamTranscripts.ts new file mode 100644 index 0000000..9949e98 --- /dev/null +++ b/src/teamTranscripts.ts @@ -0,0 +1,22 @@ +// SPDX-FileCopyrightText: 2026 CoreWeave, Inc. +// SPDX-License-Identifier: MIT +// SPDX-PackageName: weave-claude-code + +import { parseSessionFd } from './parser.js'; +import type { ParsedTurn } from './parser.js'; +import { TranscriptFile } from './transcriptFile.js'; + +/** Read the transcript explicitly named by TeammateIdle. Correlation through + * metadata and neighboring transcript discovery belongs to the recovery layer. */ +export function readTeammateTurns(transcriptPath: string): ParsedTurn[] { + let transcript: TranscriptFile | undefined; + try { + transcript = new TranscriptFile(transcriptPath); + return parseSessionFd(transcript.getFd())?.turns + .filter(turn => turn.responses.length) ?? []; + } catch { + return []; + } finally { + transcript?.close(); + } +} diff --git a/tests/agent-teams-core.test.ts b/tests/agent-teams-core.test.ts new file mode 100644 index 0000000..1ac5548 --- /dev/null +++ b/tests/agent-teams-core.test.ts @@ -0,0 +1,439 @@ +// SPDX-FileCopyrightText: 2026 CoreWeave, Inc. +// SPDX-License-Identifier: MIT +// SPDX-PackageName: weave-claude-code + +import { test, type TestContext } from 'node:test'; +import assert from 'node:assert/strict'; +import { ATTR } from '../src/genaiSpans.ts'; +import { + assistantEntry, + flushWeave, + initWeaveInMemory, + makeGenaiDaemon, + makeTranscript, + spanParentId, + userEntry, +} from './helpers.ts'; + +const TEAM = 'review-team'; +const MEMBER = 'reviewer'; + +function teammateEntries(sessionId: string, text: string, responseId: string) { + return [ + { type: 'agent-setting', agentSetting: MEMBER, sessionId }, + { type: 'user', teamName: TEAM, message: { role: 'user', content: `task: ${text}` } }, + assistantEntry(responseId, { type: 'text', text }), + ]; +} + +async function coordinator(t: TestContext, label: string, promptId?: string) { + const exporter = await initWeaveInMemory(); + exporter.reset(); + const sid = `team-${label}`; + const transcript = makeTranscript(t, sid, label); + transcript.append(userEntry('delegate reviews')); + const daemon = makeGenaiDaemon(); + await daemon.routeEvent({ + hook_event_name: 'SessionStart', + session_id: sid, + transcript_path: transcript.file, + source: 'startup', + cwd: '/x', + }); + await daemon.routeEvent({ + hook_event_name: 'UserPromptSubmit', + session_id: sid, + transcript_path: transcript.file, + prompt: 'delegate reviews', + ...(promptId ? { prompt_id: promptId } : {}), + }); + return { exporter, daemon, sid, transcript }; +} + +const teamInput = (prompt: string): Record => ({ + subagent_type: MEMBER, + prompt, + team_name: TEAM, + name: MEMBER, +}); + +async function preDispatch( + daemon: ReturnType, + sid: string, + toolUseId: string, + input: Record, +) { + await daemon.routeEvent({ + hook_event_name: 'PreToolUse', + session_id: sid, + tool_use_id: toolUseId, + tool_name: 'Agent', + tool_input: input, + }); +} + +async function postDispatch( + daemon: ReturnType, + sid: string, + toolUseId: string, + input: Record, +) { + await daemon.routeEvent({ + hook_event_name: 'PostToolUse', + session_id: sid, + tool_use_id: toolUseId, + tool_name: 'Agent', + tool_input: input, + tool_response: 'dispatched', + }); +} + +async function dispatch( + daemon: ReturnType, + sid: string, + toolUseId: string, + prompt: string, +) { + const input = teamInput(prompt); + await preDispatch(daemon, sid, toolUseId, input); + await postDispatch(daemon, sid, toolUseId, input); +} + +async function idle( + t: TestContext, + daemon: ReturnType, + label: string, + text: string, + responseId: string, +) { + const sessionId = `${label}-member`; + const transcript = makeTranscript(t, sessionId, label); + transcript.append(...teammateEntries(sessionId, text, responseId)); + await daemon.routeEvent({ + hook_event_name: 'TeammateIdle', + session_id: sessionId, + transcript_path: transcript.file, + team_name: TEAM, + teammate_name: MEMBER, + }); +} + +test('explicit team dispatch completes on TeammateIdle', async (t) => { + const { exporter, daemon, sid } = await coordinator(t, 'direct-idle'); + await dispatch(daemon, sid, 'team-call', 'review'); + await idle(t, daemon, 'direct-idle', 'review result', 'team-msg'); + await daemon.routeEvent({ + hook_event_name: 'SessionEnd', + session_id: sid, + reason: 'clear', + }); + await flushWeave(); + + const spans = exporter.getFinishedSpans(); + const agent = spans.find(span => + span.attributes[ATTR.WEAVE_SUBAGENT_SPAWNING_TOOL_CALL_ID] === 'team-call'); + assert.ok(agent); + assert.equal(agent.attributes[ATTR.OUTPUT_MESSAGES], JSON.stringify([ + { role: 'assistant', content: 'review result' }, + ])); + assert.ok(spans.some(span => span.attributes[ATTR.RESPONSE_ID] === 'team-msg')); + assert.ok(spans.some(span => + span.attributes[ATTR.AGENT_NAME] === MEMBER + && span.attributes[ATTR.WEAVE_SUBAGENT_SPAWNING_TOOL_CALL_ID] === undefined)); +}); + +test('an ordinary named background Agent completes through SubagentStop', async (t) => { + const { exporter, daemon, sid, transcript } = await coordinator(t, 'named-background'); + const input = { description: 'background work', prompt: 'ordinary task', name: 'worker' }; + const agentId = 'ordinary-named-agent'; + const subPath = transcript.subagent( + agentId, + userEntry('ordinary task'), + assistantEntry('ordinary-named-msg', { type: 'text', text: 'ordinary result' }), + ); + await preDispatch(daemon, sid, 'ordinary-named-call', input); + await daemon.routeEvent({ + hook_event_name: 'SubagentStart', + session_id: sid, + agent_id: agentId, + agent_type: 'general-purpose', + }); + await postDispatch(daemon, sid, 'ordinary-named-call', input); + await daemon.routeEvent({ + hook_event_name: 'SubagentStop', + session_id: sid, + agent_id: agentId, + agent_type: 'general-purpose', + agent_transcript_path: subPath, + }); + await daemon.routeEvent({ + hook_event_name: 'SessionEnd', + session_id: sid, + reason: 'clear', + }); + await flushWeave(); + + const agents = exporter.getFinishedSpans().filter(span => + span.attributes[ATTR.WEAVE_SUBAGENT_SPAWNING_TOOL_CALL_ID] === 'ordinary-named-call'); + assert.equal(agents.length, 1); + assert.equal(agents[0].attributes[ATTR.WEAVE_ORPHAN_REASON], undefined); + assert.ok(exporter.getFinishedSpans().some(span => + span.attributes[ATTR.RESPONSE_ID] === 'ordinary-named-msg')); +}); + +test('an explicit Team call survives SessionEnd', async (t) => { + const { exporter, daemon, sid } = await coordinator(t, 'session-end'); + await dispatch(daemon, sid, 'deferred-team-call', 'inspect'); + await daemon.routeEvent({ + hook_event_name: 'SessionEnd', + session_id: sid, + reason: 'clear', + }); + + const internals = daemon as unknown as { hasInFlightWork(): boolean }; + assert.equal(internals.hasInFlightWork(), true); + assert.equal(exporter.getFinishedSpans().some(span => + span.attributes[ATTR.WEAVE_SUBAGENT_SPAWNING_TOOL_CALL_ID] === 'deferred-team-call'), false); + + await idle(t, daemon, 'session-end', 'inspection result', 'session-end-msg'); + await flushWeave(); + + const spans = exporter.getFinishedSpans(); + const agent = spans.find(span => + span.attributes[ATTR.WEAVE_SUBAGENT_SPAWNING_TOOL_CALL_ID] === 'deferred-team-call'); + assert.ok(agent); + assert.equal(agent.attributes[ATTR.WEAVE_ORPHAN_REASON], undefined); + assert.equal(internals.hasInFlightWork(), false); +}); + +test('SessionEnd closes nested tools before retaining their Team Agent', async (t) => { + const { exporter, daemon, sid, transcript } = await coordinator(t, 'nested-session-end'); + const input = teamInput('inspect'); + await preDispatch(daemon, sid, 'nested-team-call', input); + const agentId = 'nested-team-agent'; + transcript.subagent(agentId, userEntry('inspect')); + await daemon.routeEvent({ + hook_event_name: 'SubagentStart', + session_id: sid, + agent_id: agentId, + agent_type: MEMBER, + }); + await daemon.routeEvent({ + hook_event_name: 'PreToolUse', + session_id: sid, + agent_id: agentId, + tool_use_id: 'nested-team-read', + tool_name: 'Read', + tool_input: { file_path: '/still-open.ts' }, + }); + await postDispatch(daemon, sid, 'nested-team-call', input); + await daemon.routeEvent({ + hook_event_name: 'SessionEnd', + session_id: sid, + reason: 'clear', + }); + await flushWeave(); + + const nestedTool = exporter.getFinishedSpans().find(span => + span.attributes['gen_ai.tool.call.id'] === 'nested-team-read'); + assert.ok(nestedTool); + assert.equal(nestedTool.attributes[ATTR.WEAVE_ORPHAN_REASON], 'session_ended'); + assert.equal(exporter.getFinishedSpans().some(span => + span.attributes[ATTR.WEAVE_SUBAGENT_SPAWNING_TOOL_CALL_ID] === 'nested-team-call'), false); + + await idle(t, daemon, 'nested-session-end', 'nested result', 'nested-team-msg'); + await flushWeave(); + const agent = exporter.getFinishedSpans().find(span => + span.attributes[ATTR.WEAVE_SUBAGENT_SPAWNING_TOOL_CALL_ID] === 'nested-team-call'); + assert.ok(agent); + assert.equal(agent.attributes[ATTR.WEAVE_ORPHAN_REASON], undefined); + assert.equal(spanParentId(nestedTool), agent.spanContext().spanId); +}); + +test('SessionEnd waits for an ordinary named Agent Stop and Post', async (t) => { + const { exporter, daemon, sid, transcript } = await coordinator(t, 'named-session-end'); + const input = { description: 'background work', prompt: 'ordinary task', name: 'worker' }; + const agentId = 'ordinary-ended-agent'; + const subPath = transcript.subagent( + agentId, + userEntry('ordinary task'), + assistantEntry('ordinary-ended-msg', { type: 'text', text: 'ordinary result' }), + ); + await preDispatch(daemon, sid, 'ordinary-ended-call', input); + await daemon.routeEvent({ + hook_event_name: 'SubagentStart', + session_id: sid, + agent_id: agentId, + agent_type: 'general-purpose', + }); + await daemon.routeEvent({ + hook_event_name: 'SessionEnd', + session_id: sid, + reason: 'clear', + }); + await daemon.routeEvent({ + hook_event_name: 'SubagentStop', + session_id: sid, + agent_id: agentId, + agent_type: 'general-purpose', + agent_transcript_path: subPath, + }); + await postDispatch(daemon, sid, 'ordinary-ended-call', input); + await flushWeave(); + + const agents = exporter.getFinishedSpans().filter(span => + span.attributes[ATTR.WEAVE_SUBAGENT_SPAWNING_TOOL_CALL_ID] === 'ordinary-ended-call'); + assert.equal(agents.length, 1); + assert.equal(agents[0].attributes[ATTR.WEAVE_ORPHAN_REASON], undefined); + assert.ok(exporter.getFinishedSpans().some(span => + span.attributes[ATTR.RESPONSE_ID] === 'ordinary-ended-msg')); +}); + +test('new work cancels a deferred SessionEnd', async (t) => { + const { exporter, daemon, sid, transcript } = await coordinator(t, 'resume-after-end'); + await dispatch(daemon, sid, 'resume-team-call', 'late review'); + await daemon.routeEvent({ + hook_event_name: 'SessionEnd', + session_id: sid, + reason: 'clear', + }); + + transcript.append(userEntry('continue after resume')); + await daemon.routeEvent({ + hook_event_name: 'UserPromptSubmit', + session_id: sid, + transcript_path: transcript.file, + prompt_id: 'resumed-prompt', + prompt: 'continue after resume', + }); + await idle(t, daemon, 'resume-after-end', 'late result', 'resume-after-end-msg'); + await flushWeave(); + + assert.equal(exporter.getFinishedSpans().some(span => + span.attributes[ATTR.AGENT_NAME] === 'claude-code' + && String(span.attributes[ATTR.INPUT_MESSAGES]).includes('continue after resume')), false); + + await daemon.routeEvent({ + hook_event_name: 'SessionEnd', + session_id: sid, + prompt_id: 'resumed-prompt', + reason: 'clear', + }); + await flushWeave(); + assert.equal(exporter.getFinishedSpans().some(span => + span.attributes[ATTR.AGENT_NAME] === 'claude-code' + && String(span.attributes[ATTR.INPUT_MESSAGES]).includes('continue after resume')), true); +}); + +test('a failed team spawn closes its deferred owner', async (t) => { + const { exporter, daemon, sid } = await coordinator(t, 'failure'); + const input = teamInput('fail review'); + await preDispatch(daemon, sid, 'failed-team-call', input); + await daemon.routeEvent({ + hook_event_name: 'SessionEnd', + session_id: sid, + reason: 'clear', + }); + assert.equal(exporter.getFinishedSpans().some(span => + span.attributes[ATTR.WEAVE_SUBAGENT_SPAWNING_TOOL_CALL_ID] === 'failed-team-call'), false); + + await daemon.routeEvent({ + hook_event_name: 'PostToolUseFailure', + session_id: sid, + tool_use_id: 'failed-team-call', + tool_name: 'Agent', + tool_input: input, + error: 'SpawnError: unavailable', + }); + await flushWeave(); + + const spans = exporter.getFinishedSpans(); + const agent = spans.find(span => + span.attributes[ATTR.WEAVE_SUBAGENT_SPAWNING_TOOL_CALL_ID] === 'failed-team-call'); + assert.ok(agent); + assert.equal(agent.attributes[ATTR.ERROR_TYPE], 'SpawnError'); + assert.equal(agent.attributes[ATTR.WEAVE_ORPHAN_REASON], undefined); + assert.ok(spans.some(span => span.attributes[ATTR.AGENT_NAME] === 'claude-code')); +}); + +test('legacy team work retains only its exact owning turn', async (t) => { + const { exporter, daemon, sid, transcript } = await coordinator(t, 'legacy-owner-turn'); + await dispatch(daemon, sid, 'legacy-owner-call', 'background review'); + for (const prompt of ['second prompt', 'third prompt']) { + transcript.append(userEntry(prompt)); + await daemon.routeEvent({ + hook_event_name: 'UserPromptSubmit', + session_id: sid, + transcript_path: transcript.file, + prompt, + }); + } + + const finishedRoots = exporter.getFinishedSpans().filter(span => + span.attributes[ATTR.AGENT_NAME] === 'claude-code'); + assert.equal(finishedRoots.length, 1); + + await idle(t, daemon, 'legacy-owner-turn', 'background result', 'legacy-owner-msg'); + await daemon.routeEvent({ + hook_event_name: 'SessionEnd', + session_id: sid, + reason: 'clear', + }); + await flushWeave(); + assert.ok(exporter.getFinishedSpans().some(span => + span.attributes[ATTR.RESPONSE_ID] === 'legacy-owner-msg')); +}); + +test('legacy team work does not retain a later turn before an explicit prompt', async (t) => { + const { exporter, daemon, sid, transcript } = await coordinator(t, 'mixed-owner-turn'); + await dispatch(daemon, sid, 'mixed-owner-call', 'background review'); + + transcript.append(userEntry('second legacy prompt')); + await daemon.routeEvent({ + hook_event_name: 'UserPromptSubmit', + session_id: sid, + transcript_path: transcript.file, + prompt: 'second legacy prompt', + }); + transcript.append(userEntry('third explicit prompt')); + await daemon.routeEvent({ + hook_event_name: 'UserPromptSubmit', + session_id: sid, + transcript_path: transcript.file, + prompt_id: 'third-prompt-id', + prompt: 'third explicit prompt', + }); + + const finishedRoots = exporter.getFinishedSpans().filter(span => + span.attributes[ATTR.AGENT_NAME] === 'claude-code'); + assert.equal(finishedRoots.length, 1); + assert.ok(String(finishedRoots[0].attributes[ATTR.INPUT_MESSAGES]) + .includes('second legacy prompt')); + + await idle(t, daemon, 'mixed-owner-turn', 'background result', 'mixed-owner-msg'); + await daemon.routeEvent({ + hook_event_name: 'SessionEnd', + session_id: sid, + prompt_id: 'third-prompt-id', + reason: 'clear', + }); + await flushWeave(); + assert.ok(exporter.getFinishedSpans().some(span => + span.attributes[ATTR.RESPONSE_ID] === 'mixed-owner-msg')); +}); + +test('shutdown orphans an explicit Team Agent', async (t) => { + const { exporter, daemon, sid } = await coordinator(t, 'shutdown'); + await dispatch(daemon, sid, 'shutdown-team-call', 'remote task'); + await daemon.drain('SIGTERM'); + await flushWeave(); + + const spans = exporter.getFinishedSpans(); + const root = spans.find(span => span.attributes[ATTR.AGENT_NAME] === 'claude-code'); + const agent = spans.find(span => + span.attributes[ATTR.WEAVE_SUBAGENT_SPAWNING_TOOL_CALL_ID] === 'shutdown-team-call'); + assert.ok(root && agent); + assert.equal(agent.attributes[ATTR.WEAVE_ORPHAN_REASON], 'daemon_shutdown'); + assert.equal(spanParentId(agent), root.spanContext().spanId); + assert.deepEqual(agent.endTime, root.endTime); +});