diff --git a/hooks/hook-handler.sh b/hooks/hook-handler.sh index 8cee95c..b43af65 100755 --- a/hooks/hook-handler.sh +++ b/hooks/hook-handler.sh @@ -90,7 +90,7 @@ if ! is_daemon_alive; then # Wait up to 5 s (50 × 100 ms) for the daemon to accept connections. # The daemon unlinks any stale socket file before binding — see - # GlobalDaemon.start in src/daemon.ts — so we don't need to clean up here. + # Daemon.start in src/daemon.ts — so we don't need to clean up here. for i in $(seq 1 50); do is_daemon_alive && break sleep 0.1 diff --git a/src/daemon.ts b/src/daemon.ts index b013b33..e8085f0 100644 --- a/src/daemon.ts +++ b/src/daemon.ts @@ -2,70 +2,33 @@ // SPDX-License-Identifier: MIT // SPDX-PackageName: weave-claude-code -import * as net from 'net'; import * as fs from 'fs'; +import * as net from 'net'; import * as path from 'path'; import { diag, DiagLogLevel } from '@opentelemetry/api'; -import type { Attributes } from '@opentelemetry/api'; -import type { - HookInput, - SessionStartHookInput, - InstructionsLoadedHookInput, - UserPromptSubmitHookInput, - PreToolUseHookInput, - PermissionRequestHookInput, - SubagentStartHookInput, - SubagentStopHookInput, - TeammateIdleHookInput, - PreCompactHookInput, - StopHookInput, - SessionEndHookInput, -} from '@anthropic-ai/claude-agent-sdk'; import * as weave from 'weave'; -import { loadSettings, VERSION } from './setup.js'; -import { appendToLog } from './utils.js'; -import { lastAssistantTextEndsWith, parseSessionFd } from './parser.js'; -import { TranscriptFile, readFirstTranscriptLine } from './transcriptFile.js'; import { - ATTR, - CompactionAttrs, - setCompactionAttrs, - assistantOutputMessages, - snippet, -} from './genaiSpans.js'; -import { resolveDaemonConfig, daemonConfigFingerprint, missingConfig } from './config.js'; + daemonConfigFingerprint, + missingConfig, + resolveDaemonConfig, +} from './config.js'; import type { DaemonConfig } from './config.js'; -import { newSessionState, upsertInstruction } from './sessionState.js'; -import type { - SessionState, - LoadedInstruction, -} from './sessionState.js'; - -// ───────────────────────────────────────────────────────────────────────────── -// Types -// ───────────────────────────────────────────────────────────────────────────── - -/** Inbound control message sent directly to the socket (not a hook event). - * `shutdown` stops the daemon; `config-hash` asks it to reply with the - * fingerprint of the config it loaded (used by `status` for drift detection) - * plus the daemon's runtime identity (pid, version, entry path). */ +import { loadSettings, VERSION } from './setup.js'; +import { HookHandler } from './hookHandler.js'; +import { appendToLog } from './utils.js'; + +/** Inbound control message sent directly to the socket (not a hook event). */ type ControlMessage = { command: 'shutdown' | 'config-hash'; -} - -/** Raw hook-event payload forwarded by hook-handler.sh. */ -type HookPayload = Record; +}; function isControlMessage(payload: unknown): payload is ControlMessage { if (typeof payload !== 'object' || payload === null) return false; - const cmd = (payload as Record).command; - return cmd === 'shutdown' || cmd === 'config-hash'; + const command = (payload as Record).command; + return command === 'shutdown' || command === 'config-hash'; } -/** Absolute real path of the daemon's own entry script, resolving the npm bin - * symlink to the actual dist/cli.js (or src/cli.ts under tsx). Lets `status` - * report which build the running daemon is executing. Falls back to the raw - * argv path if it can't be resolved. */ +/** Resolve the running daemon's entry script for config-drift reporting. */ function daemonEntryPath(): string { const entry = process.argv[1] ?? ''; try { @@ -75,99 +38,105 @@ function daemonEntryPath(): string { } } -// Keep resumed sessions warm across long idle gaps. -const INACTIVITY_TIMEOUT_MS = 120 * 60 * 1_000; // 120 minutes -// Bound how long stuck in-flight work can keep the daemon alive. -const INFLIGHT_HOLD_MAX_MS = 60 * 60 * 1_000; // 60 minutes -const CONNECTION_TIMEOUT_MS = 5_000; // 5 seconds per connection +const INACTIVITY_TIMEOUT_MS = 120 * 60 * 1_000; +const INFLIGHT_HOLD_MAX_MS = 60 * 60 * 1_000; +const CONNECTION_TIMEOUT_MS = 5_000; +const MAX_SOCKET_PAYLOAD_BYTES = 4 * 1024 * 1024; -const MAX_SOCKET_PAYLOAD_BYTES = 4 * 1024 * 1024; // 4 MiB per message - -export class GlobalDaemon { +export class Daemon { private server?: net.Server; private running = false; private lastActivity = Date.now(); - /** Inactivity shutdown threshold. Overridable via WEAVE_INACTIVITY_MS (ms) for - * testing and for ops (e.g. raising it for long-running agent-teams work). */ - private readonly inactivityMs = Number(process.env.WEAVE_INACTIVITY_MS) || INACTIVITY_TIMEOUT_MS; - private sessions = new Map(); - private sessionQueues = new Map>(); - /** InstructionsLoaded files that arrived before their session existed (the - * hook can fire before SessionStart). Keyed by session_id; drained into the - * session at SessionStart / reconstruction and cleared (also on SessionEnd). */ - private pendingInstructions = new Map(); + private readonly inactivityMs = + Number(process.env.WEAVE_INACTIVITY_MS) || INACTIVITY_TIMEOUT_MS; private tracingEnabled = false; + private readonly hookHandler: HookHandler; constructor( private readonly socketPath: string, private readonly logFile: string, private readonly config: DaemonConfig, - ) {} + ) { + this.hookHandler = new HookHandler( + config.agentName, + (level, message) => this.log(level, message), + ); + } async start(): Promise { if (this.config.weaveProject && this.config.apiKey) { try { await this.initWeave(); - this.log('INFO', `OTel tracer initialized — project=${this.config.weaveProject}, endpoint=${this.config.baseUrl}/agents/otel/v1/traces`); - this.log('INFO', `View traces: https://wandb.ai/${this.config.weaveProject}/weave/agents`); + this.log( + 'INFO', + `OTel tracer initialized — project=${this.config.weaveProject}, endpoint=${this.config.baseUrl}/agents/otel/v1/traces`, + ); + this.log( + 'INFO', + `View traces: https://wandb.ai/${this.config.weaveProject}/weave/agents`, + ); } catch (err) { - this.log('ERROR', `Failed to initialize OTel tracer: ${err} — continuing without tracing`); + this.log( + 'ERROR', + `Failed to initialize OTel tracer: ${err} — continuing without tracing`, + ); this.tracingEnabled = false; } } else { this.log('INFO', 'No weave_project / API key configured — tracing disabled'); } - // Bind the socket, exiting cleanly if another daemon already owns it. - // Concurrent hook invocations can each cold-start a daemon, but only one - // can bind; the losers exit (process.exit(0)) and their hook still reaches - // the winner over the socket. See bindSocketWithHerdProtection. await this.bindSocketWithHerdProtection(); - this.running = true; this.log('INFO', `Daemon started — socket: ${this.socketPath}`); process.on('SIGTERM', () => void this.shutdown('SIGTERM')); - process.on('SIGINT', () => void this.shutdown('SIGINT')); - // Route terminal-close cleanup through shutdown to remove the socket inode. - process.on('SIGHUP', () => void this.shutdown('SIGHUP')); - // The next hook's socket probe handles cleanup after SIGKILL or OOM. + process.on('SIGINT', () => void this.shutdown('SIGINT')); + process.on('SIGHUP', () => void this.shutdown('SIGHUP')); process.on('exit', () => { - try { if (fs.existsSync(this.socketPath)) fs.unlinkSync(this.socketPath); } catch { /* nothing more we can do */ } + try { + if (fs.existsSync(this.socketPath)) fs.unlinkSync(this.socketPath); + } catch { + // The next hook's socket probe handles stale-socket cleanup. + } }); - // Check at most every 60s, but more frequently when the timeout is short - // (env-overridden for tests) so a low WEAVE_INACTIVITY_MS is honored promptly. - const checkEveryMs = Math.min(60_000, Math.max(500, Math.floor(this.inactivityMs / 4))); + const checkEveryMs = Math.min( + 60_000, + Math.max(500, Math.floor(this.inactivityMs / 4)), + ); setInterval(() => this.checkInactivity(), checkEveryMs).unref(); } - /** Probe whether a live daemon is accepting connections on the socket. Uses a - * real connect() attempt — the inode existing is not proof of a listener - * (an ungraceful exit leaves a stale inode behind). */ private socketHasLiveListener(): Promise { if (!fs.existsSync(this.socketPath)) return Promise.resolve(false); - return new Promise((resolve) => { + return new Promise(resolve => { const probe = net.createConnection(this.socketPath); - probe.once('connect', () => { probe.destroy(); resolve(true); }); - probe.once('error', () => { probe.destroy(); resolve(false); }); + probe.once('connect', () => { + probe.destroy(); + resolve(true); + }); + probe.once('error', () => { + probe.destroy(); + resolve(false); + }); }); } - /** Create a fresh server and listen once, resolving on success and rejecting - * on the first listen error. A new server per attempt — one that errored on - * listen cannot be reused. Socket is owner-only (umask 0o077). */ private listenOnce(): Promise { return new Promise((resolve, reject) => { - const prevUmask = process.umask(0o077); - // allowHalfOpen lets handleConnection write a reply after the client - // half-closes (the `config-hash` query). Every branch closes the socket - // explicitly so the high-frequency hook-event path still ends promptly. - const server = net.createServer({ allowHalfOpen: true }, (socket) => this.handleConnection(socket)); - const onError = (err: Error) => { process.umask(prevUmask); reject(err); }; + const previousUmask = process.umask(0o077); + const server = net.createServer( + { allowHalfOpen: true }, + socket => this.handleConnection(socket), + ); + const onError = (err: Error) => { + process.umask(previousUmask); + reject(err); + }; server.once('error', onError); server.listen(this.socketPath, () => { - process.umask(prevUmask); + process.umask(previousUmask); server.removeListener('error', onError); this.server = server; resolve(); @@ -175,12 +144,6 @@ export class GlobalDaemon { }); } - /** - * Bind the daemon socket, tolerant of a herd of concurrent starts. Listen; on - * EADDRINUSE/EEXIST, re-probe: a live listener means another daemon won → exit - * 0; a stale inode is unlinked and retried. Only a confirmed-stale socket is - * ever unlinked, so a late starter can't delete the winner's live socket. - */ private async bindSocketWithHerdProtection(): Promise { const MAX_RECLAIM_ATTEMPTS = 5; for (let attempt = 0; ; attempt++) { @@ -191,36 +154,53 @@ export class GlobalDaemon { const code = (err as NodeJS.ErrnoException).code; if (code !== 'EADDRINUSE' && code !== 'EEXIST') throw err; if (await this.socketHasLiveListener()) { - this.log('INFO', 'Another daemon already owns the socket — exiting to avoid a herd'); + this.log( + 'INFO', + 'Another daemon already owns the socket — exiting to avoid a herd', + ); process.exit(0); } - // Stale inode from an ungraceful exit — reclaim it and retry. if (attempt >= MAX_RECLAIM_ATTEMPTS) throw err; - try { fs.unlinkSync(this.socketPath); } catch { /* already cleaned; retry */ } + try { + fs.unlinkSync(this.socketPath); + } catch { + // Another process already cleaned it; retry the bind. + } } } } - // ── tracer initialization ─────────────────────────────────────────────── - private async initWeave(): Promise { - if (!this.config.weaveProject) throw new Error('weaveProject required to init tracer'); - if (!this.config.apiKey) throw new Error('apiKey required to init tracer'); + if (!this.config.weaveProject) { + throw new Error('weaveProject required to init tracer'); + } + if (!this.config.apiKey) { + throw new Error('apiKey required to init tracer'); + } const [entity, project] = this.config.weaveProject.split('/', 2); if (!entity || !project) { - throw new Error(`Invalid weave_project format: '${this.config.weaveProject}' (expected entity/project)`); + throw new Error( + `Invalid weave_project format: '${this.config.weaveProject}' (expected entity/project)`, + ); } - // weave.init reads the exporter endpoint and API key from the environment. process.env['WF_TRACE_SERVER_URL'] = this.config.baseUrl; process.env['WANDB_API_KEY'] = this.config.apiKey; - // Surface exporter failures in the daemon log. const otelDiag = (message: string, ...args: unknown[]) => - this.log('ERROR', `otel: ${message}${args.length ? ` ${args.map(String).join(' ')}` : ''}`); + this.log( + 'ERROR', + `otel: ${message}${args.length ? ` ${args.map(String).join(' ')}` : ''}`, + ); diag.setLogger( - { verbose: otelDiag, debug: otelDiag, info: otelDiag, warn: otelDiag, error: otelDiag }, + { + verbose: otelDiag, + debug: otelDiag, + info: otelDiag, + warn: otelDiag, + error: otelDiag, + }, DiagLogLevel.WARN, ); @@ -228,8 +208,6 @@ export class GlobalDaemon { this.tracingEnabled = true; } - // ── connection handling ─────────────────────────────────────────────────── - private handleConnection(socket: net.Socket): void { const chunks: Buffer[] = []; let totalBytes = 0; @@ -245,11 +223,13 @@ export class GlobalDaemon { if (totalBytes > MAX_SOCKET_PAYLOAD_BYTES) { rejectedForSize = true; clearTimeout(timer); - this.log('ERROR', `Socket payload exceeded ${MAX_SOCKET_PAYLOAD_BYTES} bytes — closing connection`); + this.log( + 'ERROR', + `Socket payload exceeded ${MAX_SOCKET_PAYLOAD_BYTES} bytes — closing connection`, + ); socket.destroy(); return; } - chunks.push(chunk); }); @@ -259,6 +239,7 @@ export class GlobalDaemon { socket.destroy(); return; } + const raw = Buffer.concat(chunks).toString('utf8').trim(); if (!raw) { socket.end(); @@ -275,7 +256,6 @@ export class GlobalDaemon { } this.lastActivity = Date.now(); - if (isControlMessage(payload)) { if (payload.command === 'config-hash') { socket.end(JSON.stringify({ @@ -291,16 +271,8 @@ export class GlobalDaemon { return; } - // Payload is already buffered, so close now rather than holding the - // half-open socket open while routeEvent runs. socket.end(); - const hookPayload = payload as HookPayload; - const sessionId = hookPayload['session_id'] as string | undefined; - if (sessionId) { - this.enqueueForSession(sessionId, () => this.routeEvent(hookPayload)); - } else { - void this.routeEvent(hookPayload); - } + void this.routeEvent(payload); }); socket.on('error', (err: Error) => { @@ -309,439 +281,15 @@ export class GlobalDaemon { }); } - // ── event routing ───────────────────────────────────────────────────────── - - private async routeEvent(payload: HookPayload): Promise { - const input = payload as HookInput; - const sessionId = input.session_id; - if (!sessionId) { - this.log('ERROR', 'Missing session_id in payload'); - return; - } - - this.log('INFO', `${input.hook_event_name} session=${sessionId}${input.agent_id ? ` agent=${input.agent_id}` : ''}`); - - // Isolate SDK active-span state across concurrent sessions. - await weave.runIsolated(() => this.dispatchEvent(input, sessionId)); - } - - private async dispatchEvent(input: HookInput, sessionId: string): Promise { - try { - switch (input.hook_event_name) { - case 'SessionStart': - await this.handleSessionStart(sessionId, input); - break; - case 'InstructionsLoaded': - // Synchronous: reads the instruction file inline; nothing to await. - this.handleInstructionsLoaded(sessionId, input); - break; - case 'UserPromptSubmit': - await this.handleUserPromptSubmit(sessionId, input); - break; - case 'PreToolUse': - await this.handlePreToolUse(sessionId, input); - break; - case 'PermissionRequest': - await this.handlePermissionRequest(sessionId, input); - break; - case 'PostToolUse': - case 'PostToolUseFailure': - break; - case 'SubagentStart': - await this.handleSubagentStart(sessionId, input); - break; - case 'SubagentStop': - await this.handleSubagentStop(sessionId, input); - break; - case 'TeammateIdle': - await this.handleTeammateIdle(sessionId, input); - break; - case 'PreCompact': - await this.handlePreCompact(sessionId, input); - break; - case 'Stop': - await this.handleStop(sessionId, input); - break; - case 'SessionEnd': - await this.handleSessionEnd(sessionId, input); - break; - default: - break; - } - } catch (err) { - this.log('ERROR', `Error handling ${input.hook_event_name}: ${err}`); - } - } - - // ── event handlers ──────────────────────────────────────────────────────── - - private async handleSessionStart(sessionId: string, input: SessionStartHookInput): Promise { - if (this.sessions.has(sessionId)) return; // idempotent - - const rawPath = input.transcript_path; - if (!rawPath) { - this.log('ERROR', `Missing transcript_path for session ${sessionId}`); - return; - } - - let transcript: TranscriptFile; - try { - transcript = new TranscriptFile(rawPath); - } catch (err) { - this.log('ERROR', `Invalid transcript_path for session ${sessionId}: ${err}`); - return; - } - - const source = input.source; - const initialRequestModel = input.model; - const cwd = input.cwd; - - const conversationId = await this.resolveConversationId(sessionId, transcript.resolvedPath, source); - - const session = newSessionState({ - sessionId, - conversationId, - transcript, - cwd, - source, - initialRequestModel, - agentName: this.config.agentName, - tracingEnabled: this.tracingEnabled, - }); - this.sessions.set(sessionId, session); - this.drainPendingInstructions(session); - - const resumed = conversationId !== sessionId; - this.log('INFO', `Session created: ${sessionId}${resumed ? ` (resumed; conversation=${conversationId})` : ''}`); - this.log( - 'DEBUG', - `SessionStart details: session=${sessionId} conversation=${conversationId} source=${source} model=${initialRequestModel ?? 'unknown'} cwd=${cwd || '(empty)'} transcript_path=${transcript.resolvedPath} transcript_file=${path.basename(transcript.resolvedPath)} active_sessions=${this.sessions.size}`, - ); - } - - private async resolveConversationId( - sessionId: string, - transcriptPath: string, - source: string, - ): Promise { - const MAX_CHAIN_DEPTH = 32; - const MAX_HEAD_READ_ATTEMPTS = 4; - const HEAD_READ_RETRY_MS = 100; - - const transcriptDir = path.dirname(transcriptPath); - const seen = new Set([sessionId]); - let current = sessionId; - let currentPath = transcriptPath; - - for (let depth = 0; depth < MAX_CHAIN_DEPTH; depth++) { - let parent: string | undefined; - // Only the FIRST hop needs retry — ancestor transcripts are static. - const attempts = depth === 0 ? MAX_HEAD_READ_ATTEMPTS : 1; - for (let i = 0; i < attempts; i++) { - const head = readFirstTranscriptLine(currentPath); - const ff = head?.['forkedFrom'] as Record | undefined; - const ffId = ff?.['sessionId']; - if (typeof ffId === 'string' && ffId) { - parent = ffId; - break; - } - if (head !== undefined) break; // head parseable but no fork — root - if (i < attempts - 1) await new Promise((r) => setTimeout(r, HEAD_READ_RETRY_MS)); - } - if (!parent || seen.has(parent)) break; - seen.add(parent); - - const parentPath = path.join(transcriptDir, `${parent}.jsonl`); - current = parent; - if (!fs.existsSync(parentPath)) { - // Parent transcript not on disk (e.g., resumed across machines). - // Stop here — the recorded parent id is still the best stitching - // key we have, even though we can't verify if IT was a fork too. - this.log( - 'DEBUG', - `resolveConversationId: parent transcript not on disk: ${parentPath} — stopping chain walk at ${parent}`, - ); - break; - } - currentPath = parentPath; - } - - if (current !== sessionId && source !== 'resume') { - // Fork detected but `source` doesn't say resume — log so the mismatch - // is visible. We still stitch by the chain root because that's the - // correct behavior; this just surfaces an unexpected hook payload. - this.log( - 'DEBUG', - `resolveConversationId: forkedFrom chain found but source='${source}' (expected 'resume') session=${sessionId} root=${current}`, - ); - } - return current; - } - - private async getOrReconstructSession( - sessionId: string, - input: HookInput, - ): Promise { - const existing = this.sessions.get(sessionId); - if (existing) return existing; - - const rawPath = input.transcript_path; - if (!rawPath) return undefined; - - let transcript: TranscriptFile; - try { - transcript = new TranscriptFile(rawPath); - } catch (err) { - this.log('ERROR', `Cannot reconstruct session ${sessionId}: invalid transcript_path: ${err}`); - return undefined; - } - - // source/model aren't on every hook variant (this reconstructs from a - // UserPromptSubmit), so read them best-effort off the raw record. - const raw = input as Record; - const source = (raw['source'] as string | undefined) ?? 'reconstructed'; - const cwd = input.cwd; - const initialRequestModel = raw['model'] as string | undefined; - const conversationId = await this.resolveConversationId(sessionId, transcript.resolvedPath, source); - - const session = newSessionState({ - sessionId, - conversationId, - transcript, - cwd, - source, - initialRequestModel, - agentName: this.config.agentName, - tracingEnabled: this.tracingEnabled, - }); - this.sessions.set(sessionId, session); - this.drainPendingInstructions(session); - this.log( - 'INFO', - `Session reconstructed after restart: ${sessionId} (conversation=${conversationId})`, - ); - return session; - } - - private handleInstructionsLoaded(sessionId: string, input: InstructionsLoadedHookInput): void { - if (!this.tracingEnabled) return; - const filePath = input.file_path; - let content: string; - try { - content = fs.readFileSync(filePath, 'utf8'); - } catch (err) { - this.log('DEBUG', `InstructionsLoaded: unreadable ${filePath}: ${err}`); - return; - } - - const instruction: LoadedInstruction = { filePath, content }; - const session = this.sessions.get(sessionId); - if (session) { - upsertInstruction(session.systemInstructions, instruction); - } else { - // Session not set up yet; buffer until SessionStart / reconstruct drains it. - const pending = this.pendingInstructions.get(sessionId) ?? []; - upsertInstruction(pending, instruction); - this.pendingInstructions.set(sessionId, pending); - } - this.log( - 'DEBUG', - `InstructionsLoaded: session=${sessionId} reason=${input.load_reason} file=${path.basename(filePath)} bytes=${content.length}${session ? '' : ' (buffered)'}`, - ); - } - - /** Move any instructions buffered before this session existed into its state, - * then discard the buffer. */ - private drainPendingInstructions(session: SessionState): void { - const pending = this.pendingInstructions.get(session.sessionId); - this.pendingInstructions.delete(session.sessionId); - if (!pending?.length) return; - for (const instruction of pending) upsertInstruction(session.systemInstructions, instruction); - this.log('DEBUG', `Drained ${pending.length} buffered instruction file(s) into session ${session.sessionId}`); - } - - private startSessionTurn(session: SessionState, userMessage?: string): weave.Turn | undefined { - if (!session.conversation) return undefined; - const turn = session.conversation.startTurn({ - agentVersion: VERSION, - model: session.initialRequestModel, - userMessage, - systemInstructions: session.systemInstructions.map((i) => i.content), - startTime: new Date(), - }); - turn.setAttributes({ - [ATTR.WEAVE_CWD]: session.cwd, - [ATTR.WEAVE_SOURCE]: session.source, - }); - session.currentTurn = turn; - return turn; - } - - private async handleUserPromptSubmit(sessionId: string, input: UserPromptSubmitHookInput): Promise { - const session = await this.getOrReconstructSession(sessionId, input); - if (!session) { - this.log('ERROR', `Unknown session (no transcript_path to reconstruct): ${sessionId}`); - return; - } - if (!this.tracingEnabled) return; - - const prompt = input.prompt; - this.log( - 'DEBUG', - `UserPromptSubmit: session=${sessionId} current_turn=${session.currentTurn ? 'open' : 'none'} prompt=${snippet(prompt, 120)}`, - ); - - // Close interrupted turns that never received a Stop hook. - this.finalizeOpenTurn(session, 'superseded_by_next_prompt'); - - const turn = this.startSessionTurn(session, prompt); - if (!turn) return; - - // Drain compaction attrs buffered while no turn was open. - if (session.pendingCompaction) { - setCompactionAttrs(turn, session.pendingCompaction); - session.pendingCompaction = undefined; - } - - this.log('INFO', 'Created turn span'); - } - - private async handlePreToolUse(sessionId: string, input: PreToolUseHookInput): Promise { - const session = this.sessions.get(sessionId); - if (!session || !this.tracingEnabled) return; - this.log('DEBUG', `PreToolUse (not yet traced): session=${sessionId} tool=${input.tool_name}`); - } - - private async handlePermissionRequest(sessionId: string, input: PermissionRequestHookInput): Promise { - const session = this.sessions.get(sessionId); - if (!session) return; - this.log('DEBUG', `PermissionRequest (not yet traced): session=${sessionId} tool=${input.tool_name}`); - } - - private async handleSubagentStart(sessionId: string, input: SubagentStartHookInput): Promise { - const session = this.sessions.get(sessionId); - if (!session || !this.tracingEnabled) return; - this.log('DEBUG', `SubagentStart (not yet traced): session=${sessionId} agent=${input.agent_id}`); - } - - private async handleSubagentStop(sessionId: string, input: SubagentStopHookInput): Promise { - const session = await this.getOrReconstructSession(sessionId, input); - if (!session || !this.tracingEnabled) return; - this.log('DEBUG', `SubagentStop (not yet traced): session=${sessionId} agent=${input.agent_id}`); - } - - private async handleTeammateIdle(sessionId: string, input: TeammateIdleHookInput): Promise { - if (!this.tracingEnabled) return; - this.log('DEBUG', `TeammateIdle (not yet traced): session=${sessionId} teammate=${input.teammate_name}`); - } - - private async handlePreCompact(sessionId: string, input: PreCompactHookInput): Promise { - const session = this.sessions.get(sessionId); - if (!session) return; - - // Claude Code sends compaction fields that are absent from the SDK type. - const raw = input as Record; - const summary = raw['summary'] ?? raw['compaction_summary']; - const itemsBefore = raw['items_before']; - const itemsAfter = raw['items_after']; - const attrs: CompactionAttrs = { - summary: typeof summary === 'string' ? summary : undefined, - itemsBefore: typeof itemsBefore === 'number' ? itemsBefore : undefined, - itemsAfter: typeof itemsAfter === 'number' ? itemsAfter : undefined, - }; - - if (session.currentTurn) { - setCompactionAttrs(session.currentTurn, attrs); - this.log('INFO', `PreCompact attached to active turn (session ${sessionId})`); - } else { - // Buffer until the next UserPromptSubmit opens a turn span. - session.pendingCompaction = attrs; - this.log('INFO', `PreCompact buffered; will attach to next turn (session ${sessionId})`); - } - } - - private async handleStop(sessionId: string, input: StopHookInput): Promise { - const session = this.sessions.get(sessionId); - if (!session?.currentTurn) return; - - // Wait for transcript synthesis to flush before reading the final response. - const finalAssistantMessage = input.last_assistant_message; - const parsedSession = await this.parseSessionFileWithRetry( - session.transcript, - finalAssistantMessage, - ); - const currentTurn = parsedSession?.turns.at(-1); - const model = currentTurn?.model; - const transcriptTurns = parsedSession?.turns.length ?? 0; - this.log( - 'DEBUG', - `Stop: session=${sessionId} transcript_path=${session.transcript.resolvedPath} transcript_turns=${transcriptTurns} parsed_model=${model ?? 'unknown'} last_assistant_message_present=${Boolean(input.last_assistant_message)}`, - ); - - const parsedTexts = currentTurn?.text ?? []; - const lastMessage = input.last_assistant_message ?? ''; - const assistantMessages = parsedTexts.length > 0 ? parsedTexts : (lastMessage ? [lastMessage] : []); - - const turnAttrs: Attributes = {}; - if (assistantMessages.length) { - turnAttrs[ATTR.OUTPUT_MESSAGES] = assistantOutputMessages(assistantMessages); - } - const finishReasons = currentTurn?.responses.map(response => response.finishReason) - .filter((reason): reason is string => reason !== undefined); - if (finishReasons?.length) { - turnAttrs[ATTR.RESPONSE_FINISH_REASONS] = finishReasons; - } - if (Object.keys(turnAttrs).length) session.currentTurn.setAttributes(turnAttrs); - // Turn.end() re-emits its request model, so update it through record(). - if (model) { - session.currentTurn.record({ model }); - } - session.currentTurn.end(); - session.currentTurn = undefined; - - this.log('INFO', 'Finished turn'); - } - - private async handleSessionEnd(sessionId: string, input: SessionEndHookInput): Promise { - // Discard any never-drained instruction buffer (e.g. a session that emitted - // InstructionsLoaded but never SessionStart) so the map can't leak. - this.pendingInstructions.delete(sessionId); - const session = this.sessions.get(sessionId); - if (!session) return; - - this.log( - 'DEBUG', - `SessionEnd: session=${sessionId} reason=${input.reason} transcript_path=${session.transcript.resolvedPath} pending_tools=${session.pendingToolCalls.size} open_subagents=${session.subagents.size()}`, - ); - - this.finalizeSession(session, 'session_ended'); - - this.log('INFO', `Finished session ${sessionId}`); - - this.sessions.delete(sessionId); - this.sessionQueues.delete(sessionId); - session.transcript.close(); - } - - private finalizeSession(session: SessionState, orphanReason: string): void { - this.finalizeOpenTurn(session, orphanReason); + private routeEvent(payload: unknown): Promise { + return this.tracingEnabled + ? this.hookHandler.handle(payload) + : Promise.resolve(); } - private finalizeOpenTurn(session: SessionState, orphanReason: string): void { - if (session.currentTurn) { - session.currentTurn.setAttributes({ [ATTR.WEAVE_ORPHAN_REASON]: orphanReason }); - session.currentTurn.end(); - session.currentTurn = undefined; - this.log('DEBUG', `Closed orphaned turn span (${orphanReason})`); - } - } - - // ── lifecycle ───────────────────────────────────────────────────────────── - private checkInactivity(): void { const idle = Date.now() - this.lastActivity; if (idle <= this.inactivityMs) return; - // Keep in-flight work alive up to the hard hold limit. if (idle < INFLIGHT_HOLD_MAX_MS && this.hasInFlightWork()) { this.log('DEBUG', 'Inactivity timeout reached but work in flight — staying up'); return; @@ -750,16 +298,8 @@ export class GlobalDaemon { void this.shutdown('inactivity'); } - /** True if any session has work in flight: an open turn span, a pending tool - * call, or a tracked subagent. Keeps the daemon alive across the inactivity - * timeout so in-flight work isn't cut off mid-flight (see checkInactivity). */ private hasInFlightWork(): boolean { - for (const s of this.sessions.values()) { - if (s.currentTurn) return true; - if (s.pendingToolCalls.size > 0) return true; - if (s.subagents.size() > 0) return true; - } - return false; + return this.hookHandler.hasInFlightWork(); } private async shutdown(reason: string): Promise { @@ -772,13 +312,7 @@ export class GlobalDaemon { private async drain(reason: string): Promise { this.log('INFO', `Shutdown: ${reason}`); this.server?.close(); - for (const session of this.sessions.values()) { - try { - this.finalizeSession(session, 'daemon_shutdown'); - } catch (err) { - this.log('ERROR', `Error finalizing session ${session.sessionId} at shutdown: ${err}`); - } - } + this.hookHandler.finalizeForShutdown(); if (this.tracingEnabled) { try { await weave.flushOTel(); @@ -786,81 +320,39 @@ export class GlobalDaemon { this.log('ERROR', `Error flushing Weave SDK: ${err}`); } } - for (const session of this.sessions.values()) { - session.transcript.close(); - } + this.hookHandler.closeTranscripts(); if (fs.existsSync(this.socketPath)) { fs.unlinkSync(this.socketPath); } } - // ── helpers ─────────────────────────────────────────────────────────────── - - /** Retry parseSessionFile while the transcript writer catches up to Stop. - * If `finalAssistantMessage` is set, require the last assistant call's - * text to end with it (mod trailing whitespace) — guards against reading - * before the synthesis line lands. Default budget: 5 × 200ms = 1s. */ - private async parseSessionFileWithRetry( - transcript: TranscriptFile, - finalAssistantMessage?: string, - attempts = 5, - delayMs = 200, - ): Promise> { - let fd: number; - try { - fd = transcript.getFd(); - } catch (err) { - this.log('ERROR', `Cannot open transcript for parsing: ${err}`); - return null; - } - const expected = (finalAssistantMessage ?? '').trimEnd(); - let result: ReturnType = null; - for (let i = 0; i < attempts; i++) { - result = parseSessionFd(fd); - // Writer caught up: parsed at least one turn AND (no synthesis to verify, - // OR the last assistant call ends with it). - if (result?.turns.length && (!expected || lastAssistantTextEndsWith(result, expected))) { - return result; - } - // No next parse to wait for on the last iteration, so skip the sleep. - if (i < attempts - 1) await new Promise(r => setTimeout(r, delayMs)); - } - return result; - } - - private enqueueForSession(sessionId: string, fn: () => Promise): void { - const prev = this.sessionQueues.get(sessionId) ?? Promise.resolve(); - const next = prev.then(fn).catch((err) => this.log('ERROR', `Queue error for session ${sessionId}: ${err}`)); - this.sessionQueues.set(sessionId, next); - } - - private log(level: 'DEBUG' | 'INFO' | 'ERROR', msg: string): void { + private log(level: 'DEBUG' | 'INFO' | 'ERROR', message: string): void { if (level === 'DEBUG' && !this.config.debug) return; - appendToLog(this.logFile, level, msg); + appendToLog(this.logFile, level, message); } } - -// ───────────────────────────────────────────────────────────────────────────── -// Entry point (invoked by `weave-claude-code daemon`) -// ───────────────────────────────────────────────────────────────────────────── - export async function runDaemon(): Promise { const settings = loadSettings(); const { daemon_socket: socketPath, log_file: logFile } = settings; - fs.mkdirSync(path.dirname(logFile), { recursive: true }); const config = resolveDaemonConfig(settings, process.env); - if (!config.weaveProject || !config.apiKey) { - const missing = missingConfig(!!config.weaveProject, !!config.apiKey, 'WANDB_API_KEY'); - appendToLog(logFile, 'INFO', `Daemon not started — missing configuration: ${missing}`); + const missing = missingConfig( + Boolean(config.weaveProject), + Boolean(config.apiKey), + 'WANDB_API_KEY', + ); + appendToLog( + logFile, + 'INFO', + `Daemon not started — missing configuration: ${missing}`, + ); process.exit(0); } - const daemon = new GlobalDaemon(socketPath, logFile, config); - + const daemon = new Daemon(socketPath, logFile, config); try { await daemon.start(); } catch (err) { diff --git a/src/hookHandler.ts b/src/hookHandler.ts new file mode 100644 index 0000000..95309fd --- /dev/null +++ b/src/hookHandler.ts @@ -0,0 +1,313 @@ +// SPDX-FileCopyrightText: 2026 CoreWeave, Inc. +// SPDX-License-Identifier: MIT +// SPDX-PackageName: weave-claude-code + +import * as fs from 'fs'; +import * as path from 'path'; +import type { + HookInput, + InstructionsLoadedHookInput, + PreCompactHookInput, + SessionEndHookInput, + SessionStartHookInput, + StopHookInput, + UserPromptSubmitHookInput, +} from '@anthropic-ai/claude-agent-sdk'; +import * as weave from 'weave'; +import type { CompactionAttrs } from './genaiSpans.js'; +import { snippet } from './genaiSpans.js'; +import { Session } from './session.js'; +import { TranscriptFile } from './transcriptFile.js'; + +type TraceLog = (level: 'DEBUG' | 'INFO' | 'ERROR', message: string) => void; + +function parseHookInput(payload: unknown): HookInput | undefined { + if (!payload || typeof payload !== 'object') return undefined; + const input = payload as Record; + return typeof input['hook_event_name'] === 'string' + && typeof input['session_id'] === 'string' + ? input as HookInput + : undefined; +} + +export class HookHandler { + private readonly sessions = new Map(); + private readonly sessionQueues = new Map>(); + /** InstructionsLoaded can arrive before SessionStart. */ + private readonly pendingInstructions = new Map>(); + + constructor( + private readonly agentName: string, + private readonly log: TraceLog, + ) {} + + async handle(payload: unknown): Promise { + const input = parseHookInput(payload); + if (!input) { + this.log('ERROR', 'Invalid hook payload'); + 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); + } + } + } + + private async route(input: HookInput): 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)); + } catch (err) { + this.log('ERROR', `Error handling ${input.hook_event_name}: ${err}`); + } + } + + private async dispatchEvent(input: HookInput): Promise { + const sessionId = input.session_id; + switch (input.hook_event_name) { + case 'SessionStart': + await this.handleSessionStart(sessionId, input); + return; + case 'InstructionsLoaded': + this.handleInstructionsLoaded(sessionId, input); + return; + case 'UserPromptSubmit': + await this.handleUserPromptSubmit(sessionId, input); + return; + case 'PreCompact': + this.handlePreCompact(sessionId, input); + return; + case 'Stop': + await this.handleStop(sessionId, input); + return; + case 'SessionEnd': + await this.handleSessionEnd(sessionId, input); + return; + default: + return; + } + } + + private async createSession( + sessionId: string, + transcript: TranscriptFile, + options: { source: string; cwd: string; initialRequestModel?: string }, + ): Promise { + const session = await Session.create({ + sessionId, + transcript, + cwd: options.cwd, + source: options.source, + initialRequestModel: options.initialRequestModel, + agentName: this.agentName, + log: this.log, + }); + this.sessions.set(sessionId, session); + this.drainPendingInstructions(session); + return session; + } + + private async handleSessionStart( + sessionId: string, + input: SessionStartHookInput, + ): Promise { + if (this.sessions.has(sessionId)) return; + + let transcript: TranscriptFile; + try { + transcript = new TranscriptFile(input.transcript_path); + } catch (err) { + this.log('ERROR', `Invalid transcript_path for session ${sessionId}: ${err}`); + return; + } + + const session = await this.createSession(sessionId, transcript, { + source: input.source, + cwd: input.cwd, + initialRequestModel: input.model, + }); + const resumed = session.conversationId !== sessionId; + this.log( + 'INFO', + `Session created: ${sessionId}${resumed ? ` (resumed; conversation=${session.conversationId})` : ''}`, + ); + this.log( + 'DEBUG', + `SessionStart details: session=${sessionId} conversation=${session.conversationId} source=${session.source} model=${session.initialRequestModel ?? 'unknown'} cwd=${session.cwd || '(empty)'} transcript_path=${session.transcriptPath} transcript_file=${path.basename(session.transcriptPath)} active_sessions=${this.sessions.size}`, + ); + } + + private async getOrReconstructSession( + sessionId: string, + input: HookInput, + ): Promise { + const existing = this.sessions.get(sessionId); + if (existing) return existing; + if (!input.transcript_path) return undefined; + + let transcript: TranscriptFile; + try { + transcript = new TranscriptFile(input.transcript_path); + } catch (err) { + this.log('ERROR', `Cannot reconstruct session ${sessionId}: invalid transcript_path: ${err}`); + return undefined; + } + + const raw = input as Record; + const session = await this.createSession(sessionId, transcript, { + source: raw['source'] as string | undefined ?? 'reconstructed', + cwd: input.cwd, + initialRequestModel: raw['model'] as string | undefined, + }); + this.log( + 'INFO', + `Session reconstructed after restart: ${sessionId} (conversation=${session.conversationId})`, + ); + return session; + } + + private handleInstructionsLoaded( + sessionId: string, + input: InstructionsLoadedHookInput, + ): void { + let content: string; + try { + content = fs.readFileSync(input.file_path, 'utf8'); + } catch (err) { + this.log('DEBUG', `InstructionsLoaded: unreadable ${input.file_path}: ${err}`); + return; + } + + const session = this.sessions.get(sessionId); + if (session) { + session.setInstruction(input.file_path, content); + } else { + const pending = this.pendingInstructions.get(sessionId) ?? new Map(); + pending.set(input.file_path, content); + this.pendingInstructions.set(sessionId, pending); + } + this.log( + 'DEBUG', + `InstructionsLoaded: session=${sessionId} reason=${input.load_reason} file=${path.basename(input.file_path)} bytes=${content.length}${session ? '' : ' (buffered)'}`, + ); + } + + private drainPendingInstructions(session: Session): void { + const pending = this.pendingInstructions.get(session.sessionId); + this.pendingInstructions.delete(session.sessionId); + if (!pending) return; + for (const [filePath, content] of pending) { + session.setInstruction(filePath, content); + } + this.log( + 'DEBUG', + `Drained ${pending.size} buffered instruction file(s) into session ${session.sessionId}`, + ); + } + + private async handleUserPromptSubmit( + sessionId: string, + input: UserPromptSubmitHookInput, + ): Promise { + const session = await this.getOrReconstructSession(sessionId, input); + if (!session) { + this.log('ERROR', `Unknown session (no transcript_path to reconstruct): ${sessionId}`); + return; + } + + const result = session.submitPrompt(input.prompt_id, input.prompt); + if (!result.created) return; + this.log( + 'DEBUG', + `UserPromptSubmit: session=${sessionId} current_turn=${result.replacedOpenTurn ? 'open' : 'none'} prompt=${snippet(input.prompt, 120)}`, + ); + this.log('INFO', 'Created turn span'); + } + + private handlePreCompact(sessionId: string, input: PreCompactHookInput): void { + const session = this.sessions.get(sessionId); + if (!session) return; + + // Claude Code sends compaction fields that are absent from the SDK type. + const raw = input as Record; + const attrs: CompactionAttrs = { + summary: (raw['summary'] ?? raw['compaction_summary']) as string | undefined, + itemsBefore: raw['items_before'] as number | undefined, + itemsAfter: raw['items_after'] as number | undefined, + }; + if (session.setCompaction(input.prompt_id, attrs)) { + this.log('INFO', `PreCompact attached to active turn (session ${sessionId})`); + } else { + this.log('INFO', `PreCompact buffered; will attach to next turn (session ${sessionId})`); + } + } + + private async handleStop(sessionId: string, input: StopHookInput): Promise { + const session = await this.getOrReconstructSession(sessionId, input); + if (!session) return; + + const snapshot = await session.snapshotStop( + input.prompt_id, + input.last_assistant_message, + ); + this.log( + 'DEBUG', + `Stop: session=${sessionId} transcript_path=${session.transcriptPath} responses=${snapshot.responseCount} model=${snapshot.model ?? 'unknown'} last_assistant_message_present=${Boolean(input.last_assistant_message)}`, + ); + this.log('INFO', 'Recorded turn stop snapshot'); + } + + private async handleSessionEnd( + sessionId: string, + input: SessionEndHookInput, + ): Promise { + this.pendingInstructions.delete(sessionId); + const session = this.sessions.get(sessionId) + ?? await this.getOrReconstructSession(sessionId, input); + if (!session) return; + + const turnCount = session.finishAtSessionEnd(input.prompt_id); + this.log( + 'DEBUG', + `SessionEnd: session=${sessionId} reason=${input.reason} transcript_path=${session.transcriptPath} turns=${turnCount}`, + ); + this.sessions.delete(sessionId); + session.close(); + this.log('INFO', `Finished session ${sessionId}`); + } + + hasInFlightWork(): boolean { + for (const session of this.sessions.values()) { + if (session.hasInFlightWork()) return true; + } + return false; + } + + finalizeForShutdown(): void { + for (const session of this.sessions.values()) { + try { + session.finishOpenTurns('daemon_shutdown'); + } catch (err) { + this.log('ERROR', `Error finalizing session ${session.sessionId} at shutdown: ${err}`); + } + } + } + + closeTranscripts(): void { + for (const session of this.sessions.values()) { + session.close(); + } + } +} diff --git a/src/session.ts b/src/session.ts new file mode 100644 index 0000000..179a8d7 --- /dev/null +++ b/src/session.ts @@ -0,0 +1,442 @@ +// SPDX-FileCopyrightText: 2026 CoreWeave, Inc. +// SPDX-License-Identifier: MIT +// SPDX-PackageName: weave-claude-code + +import * as fs from 'fs'; +import * as path from 'path'; +import type { Attributes } from '@opentelemetry/api'; +import * as weave from 'weave'; +import { emitChatSpans } from './chatSpans.js'; +import { + ATTR, + assistantOutputMessages, + buildIntegrationAttrs, + parseTimestamp, + setCompactionAttrs, +} from './genaiSpans.js'; +import type { CompactionAttrs } from './genaiSpans.js'; +import { + assistantResponses, + extractAssistantTextBlocks, + lastAssistantTextEndsWith, + parseSessionFd, +} from './parser.js'; +import type { AssistantResponse, ParsedSession } from './parser.js'; +import { VERSION } from './setup.js'; +import { TranscriptFile, readFirstTranscriptLine } from './transcriptFile.js'; + +type TraceLog = (level: 'DEBUG' | 'INFO' | 'ERROR', message: string) => void; + +type TurnTrace = { + span: weave.Turn; + promptId?: string; + userText?: string; + /** A Stop snapshot is quiescent but remains reopenable because hooks block. */ + phase: 'active' | 'stopped'; + /** Number of provider responses already present when this prompt began. */ + responseOffset: number; + /** Frozen when a newer prompt starts, preventing cross-turn replay. */ + responseLimit?: number; + /** Supports repeated/blockable Stop hooks without duplicate chat spans. */ + seenResponses: Set; +}; + +type NewSessionOptions = { + sessionId: string; + transcript: TranscriptFile; + cwd: string; + source: string; + initialRequestModel?: string; + agentName: string; + log: TraceLog; +}; + +type StartTurnOptions = { + promptId?: string; + userMessage?: string; + recoverCurrentTurn?: boolean; + responseOffsetFloor?: number; + makeCurrent?: boolean; +}; + +export class Session { + readonly sessionId: string; + readonly conversationId: string; + readonly transcript: TranscriptFile; + readonly cwd: string; + readonly source: string; + readonly initialRequestModel?: string; + + /** File path → latest loaded contents, preserving first-load order. */ + private readonly systemInstructions = new Map(); + private readonly turns = new Set(); + private currentTurn?: TurnTrace; + private pendingCompaction?: CompactionAttrs; + + private constructor( + options: NewSessionOptions, + conversationId: string, + private readonly log: TraceLog, + ) { + this.sessionId = options.sessionId; + this.conversationId = conversationId; + this.transcript = options.transcript; + this.cwd = options.cwd; + this.source = options.source; + this.initialRequestModel = options.initialRequestModel; + + const version = readFirstTranscriptLine(options.transcript.resolvedPath)?.version; + const integrationAttrs = buildIntegrationAttrs({ + version: VERSION, + meta: { claude_code_app_version: version }, + }); + this.conversation = weave.startConversation({ + conversationId, + agentName: options.agentName, + attributes: { ...integrationAttrs, [ATTR.WEAVE_PLUGIN_VERSION]: VERSION }, + }); + } + + private readonly conversation: weave.Conversation; + + static async create(options: NewSessionOptions): Promise { + const conversationId = await resolveConversationId( + options.sessionId, + options.transcript.resolvedPath, + options.log, + ); + return new Session(options, conversationId, options.log); + } + + get transcriptPath(): string { + return this.transcript.resolvedPath; + } + + setInstruction(filePath: string, content: string): void { + this.systemInstructions.set(filePath, content); + } + + submitPrompt( + promptId: string | undefined, + prompt: string, + ): { created: boolean; replacedOpenTurn: boolean } { + if (promptId !== undefined && this.turnForPrompt(promptId)) { + return { created: false, replacedOpenTurn: false }; + } + + const previous = this.currentTurn; + let responseOffsetFloor: number | undefined; + if (previous) { + previous.responseLimit ??= assistantResponses( + parseSessionFd(this.transcript.getFd()) ?? { turns: [] }, + ).length; + responseOffsetFloor = previous.responseLimit; + this.finalizeTurn(previous, 'superseded_by_next_prompt'); + } + + const turn = this.startTurn({ + promptId, + userMessage: prompt, + responseOffsetFloor, + }); + if (this.pendingCompaction) { + setCompactionAttrs(turn.span, this.pendingCompaction); + this.pendingCompaction = undefined; + } + return { created: true, replacedOpenTurn: Boolean(previous) }; + } + + setCompaction(promptId: string | undefined, attrs: CompactionAttrs): boolean { + const turn = this.turnForPrompt(promptId); + if (!turn) { + this.pendingCompaction = attrs; + return false; + } + setCompactionAttrs(turn.span, attrs); + return true; + } + + async snapshotStop( + promptId: string | undefined, + lastAssistantMessage?: string, + ): Promise<{ responseCount: number; model?: string }> { + const turn = this.turnForPrompt(promptId) ?? this.ensureTurn(promptId); + const parsed = await this.parseTranscriptWithRetry(lastAssistantMessage); + const responses = parsed ? this.responsesForTurn(parsed, turn) : []; + this.recordTurnOutput(turn, responses, { lastMessage: lastAssistantMessage }); + turn.phase = 'stopped'; + return { + responseCount: responses.length, + model: responses.filter(response => response.model).at(-1)?.model, + }; + } + + finishAtSessionEnd(promptId: string | undefined): number { + const parsed = this.parseTranscript(); + this.reconcileFinalTurn(promptId, parsed); + return this.finishTurns('session_ended', parsed); + } + + finishOpenTurns(orphanReason: string): number { + return this.finishTurns(orphanReason, this.parseTranscript()); + } + + private finishTurns( + orphanReason: string, + parsed: ParsedSession | null, + ): number { + const turnCount = this.turns.size; + for (const turn of [...this.turns]) { + this.recordFinalTurnOutput(turn, orphanReason, parsed); + this.endTurn(turn); + } + return turnCount; + } + + hasInFlightWork(): boolean { + return [...this.turns].some(turn => turn.phase === 'active'); + } + + close(): void { + this.transcript.close(); + } + + private turnForPrompt(promptId: string | undefined): TurnTrace | undefined { + return promptId === undefined + ? this.currentTurn + : [...this.turns].find(turn => turn.promptId === promptId); + } + + private transcriptCursor( + options: StartTurnOptions, + ): { responseOffset: number; startTime?: Date; userText?: string } { + const parsed = parseSessionFd(this.transcript.getFd()); + if (!parsed) return { responseOffset: 0, userText: options.userMessage }; + + const responses = assistantResponses(parsed); + const current = parsed.turns.at(-1); + const transcriptHasPrompt = options.userMessage !== undefined + && current?.userText === options.userMessage; + const includeCurrent = options.recoverCurrentTurn || transcriptHasPrompt; + const responseOffset = includeCurrent + ? responses.length - (current?.responses.length ?? 0) + : responses.length; + return { + responseOffset: Math.max(responseOffset, options.responseOffsetFloor ?? 0), + startTime: includeCurrent ? parseTimestamp(current?.startTime) : undefined, + userText: options.userMessage ?? (includeCurrent ? current?.userText : undefined), + }; + } + + private startTurn(options: StartTurnOptions = {}): TurnTrace { + const cursor = this.transcriptCursor(options); + const span = this.conversation.startTurn({ + agentVersion: VERSION, + model: this.initialRequestModel, + userMessage: cursor.userText, + systemInstructions: [...this.systemInstructions.values()], + startTime: cursor.startTime, + }); + span.setAttributes({ + [ATTR.WEAVE_CWD]: this.cwd, + [ATTR.WEAVE_SOURCE]: this.source, + }); + const turn: TurnTrace = { + span, + promptId: options.promptId, + userText: cursor.userText, + phase: 'active', + responseOffset: cursor.responseOffset, + seenResponses: new Set(), + }; + this.turns.add(turn); + if (options.makeCurrent !== false) this.currentTurn = turn; + return turn; + } + + private ensureTurn(promptId: string | undefined): TurnTrace { + return this.turnForPrompt(promptId) ?? this.startTurn({ + promptId, + // An exact prompt_id must not claim the last transcript turn. A legacy + // hook has no competing identity, so it can safely recover that turn. + recoverCurrentTurn: promptId === undefined, + makeCurrent: !this.currentTurn || this.currentTurn.promptId === promptId, + }); + } + + private reconcileFinalTurn( + promptId: string | undefined, + parsed: ParsedSession | null, + ): void { + const finalTranscriptTurn = parsed?.turns.at(-1); + if (!parsed || !finalTranscriptTurn) return; + + let turn = promptId === undefined + ? [...this.turns].find(candidate => + candidate.userText !== undefined + && candidate.userText === finalTranscriptTurn.userText) + ?? (this.currentTurn?.promptId === undefined ? this.currentTurn : undefined) + : this.turnForPrompt(promptId); + const legacyTurn = this.currentTurn; + if (!turn && promptId !== undefined && legacyTurn && legacyTurn.promptId === undefined) { + turn = legacyTurn; + turn.promptId = promptId; + } + const stoppedUnknownPrompt = promptId === undefined + && this.currentTurn?.promptId !== undefined + && this.currentTurn.phase === 'stopped'; + if (!turn && !stoppedUnknownPrompt) { + turn = this.startTurn({ + promptId, + userMessage: finalTranscriptTurn.userText, + recoverCurrentTurn: true, + }); + } + if (!turn || turn.responseLimit !== undefined) return; + + const finalResponseOffset = assistantResponses(parsed).length + - finalTranscriptTurn.responses.length; + if (turn.userText === undefined) { + // An exact prompt_id binds a root reconstructed from an earlier terminal + // hook to the matching final transcript turn. + turn.responseOffset = finalResponseOffset; + turn.userText = finalTranscriptTurn.userText; + if (turn.userText !== undefined) { + turn.span.record({ + messages: [{ role: 'user', parts: [{ type: 'text', content: turn.userText }] }], + }); + } + } else { + turn.responseOffset = Math.max(turn.responseOffset, finalResponseOffset); + } + } + + private responsesForTurn( + parsed: ParsedSession, + turn: TurnTrace, + ): AssistantResponse[] { + return assistantResponses(parsed).slice(turn.responseOffset, turn.responseLimit); + } + + private recordTurnOutput( + turn: TurnTrace, + responses: AssistantResponse[], + options: { lastMessage?: string; orphanReason?: string } = {}, + ): void { + emitChatSpans(turn.span, responses, { seen: turn.seenResponses }); + + const text = responses.flatMap(response => extractAssistantTextBlocks(response.content)); + if (!text.length && options.lastMessage) text.push(options.lastMessage); + const attributes: Attributes = {}; + if (text.length) attributes[ATTR.OUTPUT_MESSAGES] = assistantOutputMessages(text); + const finishReasons = responses + .map(response => response.finishReason) + .filter((reason): reason is string => Boolean(reason)); + if (finishReasons.length) attributes[ATTR.RESPONSE_FINISH_REASONS] = finishReasons; + if (options.orphanReason) attributes[ATTR.WEAVE_ORPHAN_REASON] = options.orphanReason; + if (Object.keys(attributes).length) turn.span.setAttributes(attributes); + + const model = responses.filter(response => response.model).at(-1)?.model; + if (model) turn.span.record({ model }); + } + + private recordFinalTurnOutput( + turn: TurnTrace, + orphanReason: string, + parsed: ParsedSession | null, + ): void { + const responses = parsed ? this.responsesForTurn(parsed, turn) : []; + this.recordTurnOutput(turn, responses, { + orphanReason: turn.phase === 'active' ? orphanReason : undefined, + }); + } + + private finalizeTurn(turn: TurnTrace, orphanReason: string): void { + this.recordFinalTurnOutput(turn, orphanReason, this.parseTranscript()); + this.endTurn(turn); + } + + private endTurn(turn: TurnTrace): void { + turn.span.end(); + this.turns.delete(turn); + if (this.currentTurn === turn) this.currentTurn = undefined; + } + + private parseTranscript(): ParsedSession | null { + try { + return parseSessionFd(this.transcript.getFd()); + } catch (error) { + this.log('DEBUG', `Could not recover chat spans while closing turn: ${error}`); + return null; + } + } + + /** Retry parsing while the transcript writer catches up to Stop. */ + private async parseTranscriptWithRetry( + finalAssistantMessage?: string, + attempts = 5, + delayMs = 200, + ): Promise { + let fd: number; + try { + fd = this.transcript.getFd(); + } catch (err) { + this.log('ERROR', `Cannot open transcript for parsing: ${err}`); + return null; + } + const expected = (finalAssistantMessage ?? '').trimEnd(); + let result: ParsedSession | null = null; + for (let i = 0; i < attempts; i++) { + result = parseSessionFd(fd); + if (result?.turns.length && (!expected || lastAssistantTextEndsWith(result, expected))) { + return result; + } + if (i < attempts - 1) await new Promise(resolve => setTimeout(resolve, delayMs)); + } + return result; + } +} + +async function resolveConversationId( + sessionId: string, + transcriptPath: string, + log: TraceLog, +): Promise { + const MAX_CHAIN_DEPTH = 32; + const MAX_HEAD_READ_ATTEMPTS = 4; + const HEAD_READ_RETRY_MS = 100; + + const transcriptDir = path.dirname(transcriptPath); + const seen = new Set([sessionId]); + let current = sessionId; + let currentPath = transcriptPath; + + for (let depth = 0; depth < MAX_CHAIN_DEPTH; depth++) { + let parent: string | undefined; + const attempts = depth === 0 ? MAX_HEAD_READ_ATTEMPTS : 1; + for (let i = 0; i < attempts; i++) { + const head = readFirstTranscriptLine(currentPath); + if (head?.forkedFrom?.sessionId) { + parent = head.forkedFrom.sessionId; + break; + } + if (head !== undefined) break; + if (i < attempts - 1) await new Promise(resolve => setTimeout(resolve, HEAD_READ_RETRY_MS)); + } + if (!parent || seen.has(parent)) break; + seen.add(parent); + + const parentPath = path.join(transcriptDir, `${parent}.jsonl`); + current = parent; + if (!fs.existsSync(parentPath)) { + log( + 'DEBUG', + `resolveConversationId: parent transcript not on disk: ${parentPath} — stopping chain walk at ${parent}`, + ); + break; + } + currentPath = parentPath; + } + + return current; +} diff --git a/src/sessionState.ts b/src/sessionState.ts deleted file mode 100644 index a9b0ada..0000000 --- a/src/sessionState.ts +++ /dev/null @@ -1,222 +0,0 @@ -// SPDX-FileCopyrightText: 2026 CoreWeave, Inc. -// SPDX-License-Identifier: MIT -// SPDX-PackageName: weave-claude-code - -import * as path from 'path'; -import * as weave from 'weave'; -import { VERSION } from './setup.js'; -import { isTextBlock } from './parser.js'; -import { TranscriptFile, readFirstTranscriptLine } from './transcriptFile.js'; -import { sha256Hex } from './utils.js'; -import { ATTR, buildIntegrationAttrs, addPermissionResolvedEvent } from './genaiSpans.js'; -import type { CompactionAttrs } from './genaiSpans.js'; - -/** Stores the tool span opened at PreToolUse so PostToolUse can close it. */ -export type PendingToolCall = { - tool: weave.Tool; - toolName: string; - toolInput: Record; - /** True once a PermissionRequest event has been emitted for this tool. */ - permissionRequested?: boolean; -} - -type ActiveChat = { - responseKey: string; - llm: weave.LLM; -} - -/** Emit `weave.permission_resolved` on a pending tool call's span, if one was requested. */ -export function resolvePermissionIfPending(pending: PendingToolCall, approved: boolean): void { - if (!pending.permissionRequested) return; - addPermissionResolvedEvent(pending.tool, { - approved, - timestamp: new Date(), - }); -} - -export function hashPrompt(prompt: string): string { - return sha256Hex(prompt); -} - -export function subagentsDirFor(sessionTranscriptPath: string): string { - const projectDir = path.dirname(sessionTranscriptPath); - const sessionDirName = path.basename(sessionTranscriptPath, '.jsonl'); - return path.join(projectDir, sessionDirName, 'subagents'); -} - -export function computeSubagentTranscriptPath(parentTranscriptPath: string, agentId: string): string { - return path.join(subagentsDirFor(parentTranscriptPath), `agent-${agentId}.jsonl`); -} - -export function extractUserMessageContent(line: Record | undefined): string | undefined { - if (!line || line['type'] !== 'user') return undefined; - const msg = line['message']; - if (!msg || typeof msg !== 'object') return undefined; - const content = (msg as Record)['content']; - if (typeof content === 'string') return content; - if (Array.isArray(content)) { - const parts = content.filter(isTextBlock).map(block => block.text); - return parts.length > 0 ? parts.join('') : undefined; - } - return undefined; -} - -export type LoadedInstruction = { filePath: string; content: string }; - -export function upsertInstruction(list: LoadedInstruction[], item: LoadedInstruction): void { - const idx = list.findIndex((i) => i.filePath === item.filePath); - if (idx >= 0) list[idx] = item; - else list.push(item); -} - -const SUBAGENT_TRANSCRIPT_RETRY_DELAYS_MS = [0, 50, 100, 150]; -export async function readSubagentFirstLineWithRetry( - transcriptPath: string, -): Promise | undefined> { - for (const delay of SUBAGENT_TRANSCRIPT_RETRY_DELAYS_MS) { - if (delay > 0) await new Promise(r => setTimeout(r, delay)); - const line = readFirstTranscriptLine(transcriptPath); - if (line && line['type'] === 'user') return line; - } - return undefined; -} - -export type SubagentTracker = { - subagentType: string; - detectedAt: Date; - toolUseId?: string; - subAgent?: weave.SubAgent; - agentId?: string; - promptHash?: string; - ended?: boolean; - transcriptPath?: string; - pendingTeammateIdle?: boolean; - teamName?: string; -} - -export type TeamMember = { - subAgent: weave.SubAgent; - conversation: weave.Conversation; - coordinatorTranscriptPath: string; - emitted: boolean; -} - -export type SessionState = { - sessionId: string; - conversationId: string; - transcript: TranscriptFile; - cwd: string; - source: string; - initialRequestModel?: string; - - conversation?: weave.Conversation; - - currentTurn?: weave.Turn; - - pendingToolCalls: Map; - subagents: SubagentTracking; - - activeChat?: ActiveChat; - emittedChatSpanResponseKeys: Set; - - /** Compaction attrs buffered while no turn span is open. Drained on next UserPromptSubmit. */ - pendingCompaction?: CompactionAttrs; - - systemInstructions: LoadedInstruction[]; -} - -export class SubagentTracking { - private trackers: SubagentTracker[] = []; - - /** Add a pending tracker at PreToolUse, before SubagentStart correlates an agent_id. */ - add(tracker: SubagentTracker): void { - this.trackers.push(tracker); - } - - findUnmatchedByContent(promptHash: string, subagentType: string): SubagentTracker | undefined { - let best: SubagentTracker | undefined; - for (const t of this.trackers) { - if (t.agentId) continue; - if (t.promptHash !== promptHash) continue; - if (t.subagentType !== subagentType) continue; - if (!best || t.detectedAt.getTime() < best.detectedAt.getTime()) best = t; - } - return best; - } - - byAgentId(agentId: string): SubagentTracker | undefined { - return this.trackers.find(t => t.agentId === agentId); - } - - findPendingTeammateIdle(subagentType: string): SubagentTracker | undefined { - let best: SubagentTracker | undefined; - for (const t of this.trackers) { - if (!t.pendingTeammateIdle) continue; - if (t.subagentType !== subagentType) continue; - if (!best || t.detectedAt.getTime() < best.detectedAt.getTime()) best = t; - } - return best; - } - - byToolUseId(toolUseId: string): SubagentTracker | undefined { - return this.trackers.find(t => t.toolUseId === toolUseId); - } - - remove(tracker: SubagentTracker): void { - const idx = this.trackers.indexOf(tracker); - if (idx >= 0) this.trackers.splice(idx, 1); - } - - size(): number { - return this.trackers.length; - } - - all(): SubagentTracker[] { - return [...this.trackers]; - } -} - -type NewSessionStateOptions = { - sessionId: string; - conversationId: string; - transcript: TranscriptFile; - cwd: string; - source: string; - initialRequestModel: string | undefined; - agentName: string; - tracingEnabled: boolean; -}; - -export function newSessionState(options: NewSessionStateOptions): SessionState { - const { sessionId, conversationId, transcript, cwd, source, initialRequestModel } = - options; - // Preserve the Claude Code version when reconstructing a session. - const headLine = readFirstTranscriptLine(transcript.resolvedPath); - const version = headLine?.['version']; - const claudeCodeAppVersion = typeof version === 'string' ? version : undefined; - const integrationAttrs = buildIntegrationAttrs({ - version: VERSION, - meta: { claude_code_app_version: claudeCodeAppVersion }, - }); - const conversation = options.tracingEnabled - ? weave.startConversation({ - conversationId, - agentName: options.agentName, - attributes: { ...integrationAttrs, [ATTR.WEAVE_PLUGIN_VERSION]: VERSION }, - }) - : undefined; - - return { - sessionId, - conversationId, - transcript, - cwd, - source, - initialRequestModel, - conversation, - pendingToolCalls: new Map(), - subagents: new SubagentTracking(), - emittedChatSpanResponseKeys: new Set(), - systemInstructions: [], - }; -} diff --git a/src/transcriptFile.ts b/src/transcriptFile.ts index 8de363f..3a65a67 100644 --- a/src/transcriptFile.ts +++ b/src/transcriptFile.ts @@ -9,6 +9,11 @@ import { isPathWithinBase } from './utils.js'; const O_RDONLY_NOFOLLOW = fs.constants.O_RDONLY | fs.constants.O_NOFOLLOW; +export type TranscriptHead = Record & { + version?: string; + forkedFrom?: { sessionId: string }; +}; + /** * Represents a Claude Code transcript file. * @@ -70,7 +75,7 @@ export class TranscriptFile { * ancestor transcripts in the fork chain: opens its own fd and closes it * before returning. */ -export function readFirstTranscriptLine(transcriptPath: string): Record | undefined { +export function readFirstTranscriptLine(transcriptPath: string): TranscriptHead | undefined { const resolved = path.resolve(transcriptPath); if (!isPathWithinBase(resolved, os.homedir())) return undefined; @@ -95,7 +100,7 @@ export function readFirstTranscriptLine(transcriptPath: string): Record; + return JSON.parse(line) as TranscriptHead; } catch { return undefined; } finally { diff --git a/tests/daemon-idle-inflight.test.ts b/tests/daemon-idle-inflight.test.ts index c52fbb7..7b15d63 100644 --- a/tests/daemon-idle-inflight.test.ts +++ b/tests/daemon-idle-inflight.test.ts @@ -2,14 +2,9 @@ // SPDX-License-Identifier: MIT // SPDX-PackageName: weave-claude-code -// The daemon idles out after a quiet window, but the inactivity check only held -// it open for in-flight cross-session *team* work. A plain long-running tool or -// turn (longer than the timeout, with no other session active) tripped the -// timeout mid-flight: the daemon exited, dropped the still-open turn/tool spans, -// and the resumed work landed on a fresh, amnesiac daemon. -// -// The fix: also hold the daemon open while any session has an open turn span, a -// pending tool call, or a tracked subagent. +// Active root work pins the daemon across its idle window. A blockable Stop +// keeps that root reopenable but makes it quiescent; later call-state slices +// extend the same predicate for tools and subagents. import { test } from 'node:test'; import assert from 'node:assert/strict'; @@ -56,7 +51,7 @@ test('daemon stays up past the inactivity timeout while a turn span is open', as } }); -test('daemon still idles out once the turn closes and nothing is in flight', async () => { +test('daemon idles out once a stopped turn is quiescent', async () => { const d = await startTestDaemon({ env: { WEAVE_INACTIVITY_MS: '1000' } }); try { const sessionId = 'inflight-002'; @@ -65,10 +60,10 @@ test('daemon still idles out once the turn closes and nothing is in flight', asy await d.send({ hook_event_name: 'UserPromptSubmit', session_id: sessionId, transcript_path: transcript, prompt: 'a quick task' }); await d.send({ hook_event_name: 'Stop', session_id: sessionId, transcript_path: transcript }); - // Turn span closed → nothing in flight → the daemon must still decide to - // idle out (the in-flight hold must not pin it open forever). + // Stop is blockable and retains the root for a continuation, but without + // active work it must not pin the daemon open indefinitely. const shuttingDown = await d.waitForLog(/Inactivity timeout — shutting down/, 3500); - assert.ok(shuttingDown, `daemon should idle out after the turn closes; log was:\n${d.readLog()}`); + assert.ok(shuttingDown, `daemon should idle out after the turn becomes quiescent; log was:\n${d.readLog()}`); } finally { await d.stop(); } diff --git a/tests/helpers.ts b/tests/helpers.ts index 8541f43..3bd02cf 100644 --- a/tests/helpers.ts +++ b/tests/helpers.ts @@ -12,7 +12,7 @@ import { InMemorySpanExporter, SimpleSpanProcessor, type ReadableSpan } from '@o import * as weave from 'weave'; import { MARKETPLACE_NAME, type Settings } from '../src/setup.ts'; -import { GlobalDaemon } from '../src/daemon.ts'; +import { Daemon } from '../src/daemon.ts'; import { resolveApiKey, resolveProject } from '../src/config.ts'; const HERE = path.dirname(fileURLToPath(import.meta.url)); @@ -119,7 +119,7 @@ export type DaemonDriver = { export function makeGenaiDaemon(agentName = 'claude-code'): DaemonDriver { const logFile = path.join(os.tmpdir(), `wcp-genai-${process.pid}.log`); - const d = new GlobalDaemon('/tmp/unused.sock', logFile, { + const d = new Daemon('/tmp/unused.sock', logFile, { weaveProject: 'e/p', apiKey: 'k', baseUrl: 'https://x', agentName, debug: false, }); (d as unknown as { tracingEnabled: boolean }).tracingEnabled = true; diff --git a/tests/system-instructions-integration.test.ts b/tests/system-instructions-integration.test.ts index 09c84e9..0a7cc13 100644 --- a/tests/system-instructions-integration.test.ts +++ b/tests/system-instructions-integration.test.ts @@ -62,6 +62,7 @@ test('buffers InstructionsLoaded fired before SessionStart, then accumulates in await d.routeEvent(loadInstr(sid, '/x/CLAUDE.md', 'PROJECT', 'session_start')); await d.routeEvent({ hook_event_name: 'UserPromptSubmit', session_id: sid, prompt: 'do it' }); await d.routeEvent({ hook_event_name: 'Stop', session_id: sid }); + await d.routeEvent({ hook_event_name: 'SessionEnd', session_id: sid, reason: 'clear' }); await flushWeave(); const [turn] = turnRoots(exporter.getFinishedSpans()); @@ -92,6 +93,7 @@ test('re-loading the same file replaces its content rather than duplicating', as await d.routeEvent(loadInstr(sid, '/x/CLAUDE.md', 'V2', 'compact')); await d.routeEvent({ hook_event_name: 'UserPromptSubmit', session_id: sid, prompt: 'do it' }); await d.routeEvent({ hook_event_name: 'Stop', session_id: sid }); + await d.routeEvent({ hook_event_name: 'SessionEnd', session_id: sid, reason: 'clear' }); await flushWeave(); const [turn] = turnRoots(exporter.getFinishedSpans()); @@ -119,6 +121,7 @@ test('stamps system instructions on every turn root (no session span to hang the await d.routeEvent({ hook_event_name: 'Stop', session_id: sid }); await d.routeEvent({ hook_event_name: 'UserPromptSubmit', session_id: sid, prompt: 'turn two' }); await d.routeEvent({ hook_event_name: 'Stop', session_id: sid }); + await d.routeEvent({ hook_event_name: 'SessionEnd', session_id: sid, reason: 'clear' }); await flushWeave(); const turns = turnRoots(exporter.getFinishedSpans()); @@ -142,6 +145,7 @@ test('omits gen_ai.system_instructions when no instructions were loaded', async await d.routeEvent({ hook_event_name: 'SessionStart', session_id: sid, transcript_path: file, source: 'startup', cwd: '/x' }); await d.routeEvent({ hook_event_name: 'UserPromptSubmit', session_id: sid, prompt: 'do it' }); await d.routeEvent({ hook_event_name: 'Stop', session_id: sid }); + await d.routeEvent({ hook_event_name: 'SessionEnd', session_id: sid, reason: 'clear' }); await flushWeave(); const [turn] = turnRoots(exporter.getFinishedSpans()); diff --git a/tests/turn-lifecycle.test.ts b/tests/turn-lifecycle.test.ts new file mode 100644 index 0000000..fefc0ca --- /dev/null +++ b/tests/turn-lifecycle.test.ts @@ -0,0 +1,403 @@ +// 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 * as fs from 'node:fs'; +import * as os from 'node:os'; +import * as path from 'node:path'; +import type { ReadableSpan } from '@opentelemetry/sdk-trace-base'; +import { ATTR } from '../src/genaiSpans.ts'; +import { VERSION } from '../src/setup.ts'; +import { + flushWeave, + initWeaveInMemory, + makeGenaiDaemon, + spanParentId, +} from './helpers.ts'; + +type Transcript = { + file: string; + append(...entries: Record[]): void; +}; + +function makeTranscript(t: TestContext, sessionId: string): Transcript { + const dir = fs.mkdtempSync(path.join(os.homedir(), '.weave-turn-lifecycle-')); + const file = path.join(dir, `${sessionId}.jsonl`); + fs.writeFileSync(file, ''); + t.after(() => fs.rmSync(dir, { recursive: true, force: true })); + return { + file, + append(...entries) { + fs.appendFileSync(file, entries.map(entry => JSON.stringify(entry)).join('\n') + '\n'); + }, + }; +} + +function userEntry( + text: string, + options: { timestamp?: string; version?: string } = {}, +): Record { + return { + type: 'user', + ...options, + message: { role: 'user', content: text }, + }; +} + +function assistantEntry( + id: string, + text: string, + options: { + timestamp?: string; + usage?: Record; + finishReason?: string; + } = {}, +): Record { + return { + type: 'assistant', + ...(options.timestamp ? { timestamp: options.timestamp } : {}), + message: { + role: 'assistant', + id, + model: 'claude-opus-4-8', + usage: options.usage ?? { input_tokens: 100, output_tokens: 50 }, + content: [{ type: 'text', text }], + ...(options.finishReason ? { stop_reason: options.finishReason } : {}), + }, + }; +} + +function turns(spans: ReadableSpan[]): ReadableSpan[] { + return spans.filter(span => span.attributes[ATTR.OPERATION_NAME] === 'invoke_agent'); +} + +function chats(spans: ReadableSpan[]): ReadableSpan[] { + return spans.filter(span => span.attributes[ATTR.OPERATION_NAME] === 'chat'); +} + +test('Stop snapshots only new normalized responses and SessionEnd closes the root', async (t) => { + const exporter = await initWeaveInMemory(); + exporter.reset(); + const sessionId = 'root-stop-snapshots'; + const transcript = makeTranscript(t, sessionId); + transcript.append(userEntry('do it', { + timestamp: '2026-01-01T00:00:00.000Z', + version: '1.2.3', + })); + const daemon = makeGenaiDaemon(); + + await daemon.routeEvent({ + hook_event_name: 'SessionStart', session_id: sessionId, + transcript_path: transcript.file, source: 'startup', cwd: '/x', + }); + await daemon.routeEvent({ + hook_event_name: 'UserPromptSubmit', session_id: sessionId, prompt: 'do it', + }); + transcript.append(assistantEntry('response-a', 'working', { + timestamp: '2026-01-01T00:00:01.000Z', + usage: { input_tokens: 10, output_tokens: 4, cache_read_input_tokens: 20 }, + })); + await daemon.routeEvent({ hook_event_name: 'Stop', session_id: sessionId }); + await flushWeave(); + + assert.equal(turns(exporter.getFinishedSpans()).length, 0, 'blockable Stop retains the root'); + assert.deepEqual(chats(exporter.getFinishedSpans()).map(span => span.attributes[ATTR.RESPONSE_ID]), [ + 'response-a', + ]); + + transcript.append(assistantEntry('response-b', 'done', { + timestamp: '2026-01-01T00:00:02.000Z', + finishReason: 'end_turn', + })); + await daemon.routeEvent({ hook_event_name: 'Stop', session_id: sessionId }); + await daemon.routeEvent({ + hook_event_name: 'SessionEnd', session_id: sessionId, reason: 'clear', + }); + await flushWeave(); + + const spans = exporter.getFinishedSpans(); + const [turn] = turns(spans); + const responseSpans = chats(spans); + assert.ok(turn); + assert.equal(responseSpans.length, 2, 'repeated Stop does not replay response-a'); + assert.deepEqual(responseSpans.map(span => span.attributes[ATTR.RESPONSE_ID]), [ + 'response-a', + 'response-b', + ]); + assert.ok(responseSpans.every(span => spanParentId(span) === turn.spanContext().spanId)); + assert.equal(turn.attributes[ATTR.WEAVE_ORPHAN_REASON], undefined); + assert.deepEqual(turn.attributes[ATTR.RESPONSE_FINISH_REASONS], ['end_turn']); + assert.equal(turn.attributes[ATTR.WEAVE_INTEGRATION_NAME], 'weave-claude-code'); + assert.equal(turn.attributes[ATTR.WEAVE_INTEGRATION_VERSION], VERSION); + assert.equal(turn.attributes['weave.integration.meta.claude_code_app_version'], '1.2.3'); + assert.equal(responseSpans[0].attributes[ATTR.USAGE_INPUT_TOKENS], 30); +}); + +test('a newer prompt closes an interrupted root without replaying its response', async (t) => { + const exporter = await initWeaveInMemory(); + exporter.reset(); + const sessionId = 'root-interrupted'; + const transcript = makeTranscript(t, sessionId); + transcript.append(userEntry('first')); + const daemon = makeGenaiDaemon(); + + await daemon.routeEvent({ + hook_event_name: 'SessionStart', session_id: sessionId, + transcript_path: transcript.file, source: 'startup', cwd: '/x', + }); + await daemon.routeEvent({ hook_event_name: 'UserPromptSubmit', session_id: sessionId, prompt: 'first' }); + transcript.append( + assistantEntry('only-once', 'first answer'), + userEntry('second'), + ); + await daemon.routeEvent({ hook_event_name: 'UserPromptSubmit', session_id: sessionId, prompt: 'second' }); + transcript.append(userEntry('third')); + await daemon.routeEvent({ hook_event_name: 'UserPromptSubmit', session_id: sessionId, prompt: 'third' }); + await daemon.routeEvent({ hook_event_name: 'SessionEnd', session_id: sessionId, reason: 'clear' }); + await flushWeave(); + + const spans = exporter.getFinishedSpans(); + assert.equal(spans.filter(span => span.attributes[ATTR.RESPONSE_ID] === 'only-once').length, 1); + const rootSpans = turns(spans); + assert.equal(rootSpans.length, 3); + const first = rootSpans.find(span => String(span.attributes[ATTR.INPUT_MESSAGES]).includes('first')); + const second = rootSpans.find(span => String(span.attributes[ATTR.INPUT_MESSAGES]).includes('second')); + assert.ok(first && second); + assert.equal(first.attributes[ATTR.WEAVE_ORPHAN_REASON], 'superseded_by_next_prompt'); + assert.equal(second.attributes[ATTR.WEAVE_ORPHAN_REASON], 'superseded_by_next_prompt'); +}); + +test('an identical prompt submitted during transcript lag does not replay prior output', async (t) => { + const exporter = await initWeaveInMemory(); + exporter.reset(); + const sessionId = 'root-repeated-prompt-race'; + const transcript = makeTranscript(t, sessionId); + transcript.append(userEntry('same prompt')); + const daemon = makeGenaiDaemon(); + + await daemon.routeEvent({ + hook_event_name: 'SessionStart', session_id: sessionId, + transcript_path: transcript.file, source: 'startup', cwd: '/x', + }); + await daemon.routeEvent({ + hook_event_name: 'UserPromptSubmit', session_id: sessionId, + prompt_id: 'prompt-a', prompt: 'same prompt', + }); + transcript.append(assistantEntry('response-a', 'first answer')); + await daemon.routeEvent({ + hook_event_name: 'Stop', session_id: sessionId, prompt_id: 'prompt-a', + }); + + // The second hook can arrive before its identical user line reaches JSONL. + await daemon.routeEvent({ + hook_event_name: 'UserPromptSubmit', session_id: sessionId, + prompt_id: 'prompt-b', prompt: 'same prompt', + }); + await daemon.routeEvent({ + hook_event_name: 'Stop', session_id: sessionId, prompt_id: 'prompt-b', + }); + await flushWeave(); + + const laggingSpans = exporter.getFinishedSpans(); + assert.equal(turns(laggingSpans).length, 1, 'only the completed first root exports during lag'); + assert.deepEqual( + chats(laggingSpans).map(span => span.attributes[ATTR.RESPONSE_ID]), + ['response-a'], + 'the lagging second root does not replay the first response', + ); + + transcript.append( + userEntry('same prompt'), + assistantEntry('response-b', 'second answer'), + ); + await daemon.routeEvent({ + hook_event_name: 'Stop', session_id: sessionId, prompt_id: 'prompt-b', + }); + await daemon.routeEvent({ + hook_event_name: 'SessionEnd', session_id: sessionId, + prompt_id: 'prompt-b', reason: 'clear', + }); + await flushWeave(); + + const spans = exporter.getFinishedSpans(); + assert.equal(turns(spans).length, 2); + assert.deepEqual( + chats(spans).map(span => span.attributes[ATTR.RESPONSE_ID]), + ['response-a', 'response-b'], + ); +}); + +test('duplicate prompt_id is idempotent', async (t) => { + const exporter = await initWeaveInMemory(); + exporter.reset(); + const sessionId = 'root-prompt-id'; + const transcript = makeTranscript(t, sessionId); + transcript.append(userEntry('once')); + const daemon = makeGenaiDaemon(); + + await daemon.routeEvent({ + hook_event_name: 'SessionStart', session_id: sessionId, + transcript_path: transcript.file, source: 'startup', cwd: '/x', + }); + const prompt = { + hook_event_name: 'UserPromptSubmit', session_id: sessionId, + prompt_id: 'prompt-1', prompt: 'once', + }; + await daemon.routeEvent(prompt); + await daemon.routeEvent(prompt); + await daemon.routeEvent({ + hook_event_name: 'SessionEnd', session_id: sessionId, + prompt_id: 'prompt-1', reason: 'clear', + }); + await flushWeave(); + + assert.equal(turns(exporter.getFinishedSpans()).length, 1); +}); + +test('an out-of-order Stop does not replace the foreground prompt', async (t) => { + const exporter = await initWeaveInMemory(); + exporter.reset(); + const sessionId = 'root-out-of-order-stop'; + const transcript = makeTranscript(t, sessionId); + transcript.append(userEntry('foreground')); + const daemon = makeGenaiDaemon(); + + await daemon.routeEvent({ + hook_event_name: 'SessionStart', session_id: sessionId, + transcript_path: transcript.file, source: 'startup', cwd: '/x', + }); + await daemon.routeEvent({ + hook_event_name: 'UserPromptSubmit', session_id: sessionId, + prompt_id: 'foreground-id', prompt: 'foreground', + }); + await daemon.routeEvent({ + hook_event_name: 'Stop', session_id: sessionId, + prompt_id: 'background-id', transcript_path: transcript.file, + }); + transcript.append(userEntry('next')); + await daemon.routeEvent({ + hook_event_name: 'UserPromptSubmit', session_id: sessionId, + prompt_id: 'next-id', prompt: 'next', + }); + await daemon.routeEvent({ hook_event_name: 'SessionEnd', session_id: sessionId, reason: 'clear' }); + await flushWeave(); + + const foreground = turns(exporter.getFinishedSpans()).find(span => + String(span.attributes[ATTR.INPUT_MESSAGES]).includes('foreground')); + assert.ok(foreground); + assert.equal(foreground.attributes[ATTR.WEAVE_ORPHAN_REASON], 'superseded_by_next_prompt'); +}); + +test('SessionEnd alone reconstructs the final turn, including its input', async (t) => { + const exporter = await initWeaveInMemory(); + exporter.reset(); + const sessionId = 'root-session-end-restart'; + const transcript = makeTranscript(t, sessionId); + transcript.append( + userEntry('finish it', { timestamp: '2026-01-01T00:00:00.000Z' }), + assistantEntry('restart-final', 'finished', { + timestamp: '2026-01-01T00:00:01.000Z', + finishReason: 'end_turn', + }), + ); + const daemon = makeGenaiDaemon(); + + await daemon.routeEvent({ + hook_event_name: 'SessionEnd', session_id: sessionId, + prompt_id: 'final-prompt', transcript_path: transcript.file, + cwd: '/x', reason: 'clear', + }); + await flushWeave(); + + const spans = exporter.getFinishedSpans(); + const [turn] = turns(spans); + const chat = spans.find(span => span.attributes[ATTR.RESPONSE_ID] === 'restart-final'); + assert.ok(turn && chat); + assert.equal( + turn.attributes[ATTR.INPUT_MESSAGES], + JSON.stringify([{ role: 'user', parts: [{ type: 'text', content: 'finish it' }] }]), + ); + assert.equal(spanParentId(chat), turn.spanContext().spanId); +}); + +test('shutdown does not synthesize a historical turn for an idle resumed session', async (t) => { + const exporter = await initWeaveInMemory(); + exporter.reset(); + const sessionId = 'root-idle-resume-shutdown'; + const transcript = makeTranscript(t, sessionId); + transcript.append( + userEntry('historical prompt'), + assistantEntry('historical-response', 'historical answer'), + ); + const daemon = makeGenaiDaemon(); + + await daemon.routeEvent({ + hook_event_name: 'SessionStart', session_id: sessionId, + transcript_path: transcript.file, source: 'resume', cwd: '/x', + }); + await daemon.drain('test shutdown'); + await flushWeave(); + + const spans = exporter.getFinishedSpans(); + assert.equal(turns(spans).length, 0); + assert.equal(chats(spans).length, 0); +}); + +test('restart-first Stop does not claim another transcript prompt', async (t) => { + const exporter = await initWeaveInMemory(); + exporter.reset(); + const sessionId = 'root-stop-restart'; + const transcript = makeTranscript(t, sessionId); + transcript.append( + userEntry('older'), + assistantEntry('older-response', 'old'), + userEntry('newer'), + assistantEntry('newer-response', 'new'), + ); + const daemon = makeGenaiDaemon(); + + await daemon.routeEvent({ + hook_event_name: 'Stop', session_id: sessionId, prompt_id: 'older-prompt', + transcript_path: transcript.file, cwd: '/x', + }); + await daemon.routeEvent({ hook_event_name: 'SessionEnd', session_id: sessionId, reason: 'clear' }); + await flushWeave(); + + const spans = exporter.getFinishedSpans(); + assert.equal(turns(spans).length, 1); + assert.equal(chats(spans).length, 0); +}); + +test('SessionEnd binds a restart-first root with the same prompt_id', async (t) => { + const exporter = await initWeaveInMemory(); + exporter.reset(); + const sessionId = 'root-stop-same-prompt-restart'; + const transcript = makeTranscript(t, sessionId); + transcript.append( + userEntry('final prompt'), + assistantEntry('final-response', 'finished'), + ); + const daemon = makeGenaiDaemon(); + + await daemon.routeEvent({ + hook_event_name: 'Stop', session_id: sessionId, prompt_id: 'prompt-a', + transcript_path: transcript.file, cwd: '/x', + }); + await daemon.routeEvent({ + hook_event_name: 'SessionEnd', session_id: sessionId, prompt_id: 'prompt-a', + transcript_path: transcript.file, reason: 'clear', + }); + await flushWeave(); + + const spans = exporter.getFinishedSpans(); + const [turn] = turns(spans); + const chat = spans.find(span => span.attributes[ATTR.RESPONSE_ID] === 'final-response'); + assert.ok(turn && chat); + assert.equal( + turn.attributes[ATTR.INPUT_MESSAGES], + JSON.stringify([{ role: 'user', parts: [{ type: 'text', content: 'final prompt' }] }]), + ); + assert.equal(spanParentId(chat), turn.spanContext().spanId); +});