From e08486b268f694261dd7f3eaee78d38a231f3160 Mon Sep 17 00:00:00 2001 From: Rick Gao Date: Thu, 23 Jul 2026 15:11:25 -0700 Subject: [PATCH] fix(daemon): drain and finalize traces deterministically --- src/callLifecycle.ts | 14 ++- src/daemon.ts | 69 ++++++++--- src/hookHandler.ts | 5 + src/session.ts | 33 +++-- tests/daemon-idle-inflight.test.ts | 45 +++++++ tests/daemon-shutdown-finalizes-turn.test.ts | 122 +++++++++++++++++++ tests/daemon-socket-ownership.test.ts | 120 ++++++++++++++++++ tests/permission-tracing.test.ts | 5 + 8 files changed, 385 insertions(+), 28 deletions(-) create mode 100644 tests/daemon-shutdown-finalizes-turn.test.ts create mode 100644 tests/daemon-socket-ownership.test.ts diff --git a/src/callLifecycle.ts b/src/callLifecycle.ts index ddab792..2ff5ccf 100644 --- a/src/callLifecycle.ts +++ b/src/callLifecycle.ts @@ -337,6 +337,7 @@ function finishAgentSpan( call: TracedAgent, outcome: ToolResult, failureType?: string, + endTime?: Date, ): void { const output = outcome.ok ? outcome.output : outcome.error; if (output !== undefined && output !== null && output !== '') { @@ -344,11 +345,14 @@ function finishAgentSpan( call.span.setAttributes({ [ATTR.OUTPUT_MESSAGES]: assistantOutputMessages([text]) }); } if (outcome.ok) { - call.span.end(); + call.span.end(endTime ? { endTime } : undefined); return; } call.span.setAttributes({ [ATTR.ERROR_TYPE]: failureType ?? errorType(outcome.error) }); - call.span.end({ error: new Error(outcome.error) }); + call.span.end({ + error: new Error(outcome.error), + ...(endTime ? { endTime } : {}), + }); } function errorType(error: string): string { @@ -434,19 +438,21 @@ export function finalizeOpenCalls( state: CallState, roots: Iterable, reason: string, + endTime?: Date, ): string[] { const closed: string[] = []; 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); + finishAgentSpan(call, call.outcome, undefined, endTime); } else if (call.kind === 'agent' && !call.toolUseId && call.stopSeen) { - call.span.end(); + call.span.end(endTime ? { endTime } : undefined); } else { call.span.setAttributes({ [ATTR.WEAVE_ORPHAN_REASON]: reason }); call.span.end({ error: new Error(`call did not complete (${reason})`), + ...(endTime ? { endTime } : {}), }); } completeCall(state, call, false); diff --git a/src/daemon.ts b/src/daemon.ts index e8085f0..81f408e 100644 --- a/src/daemon.ts +++ b/src/daemon.ts @@ -45,6 +45,9 @@ const MAX_SOCKET_PAYLOAD_BYTES = 4 * 1024 * 1024; export class Daemon { private server?: net.Server; + /** Socket inode bound by this process; a successor may reuse the path while + * this daemon drains. */ + private ownedSocketInode?: number; private running = false; private lastActivity = Date.now(); private readonly inactivityMs = @@ -64,6 +67,13 @@ export class Daemon { } async start(): Promise { + // Install cleanup before any await can expose a partially started daemon. + this.running = true; + process.on('SIGTERM', () => void this.shutdown('SIGTERM')); + process.on('SIGINT', () => void this.shutdown('SIGINT')); + process.on('SIGHUP', () => void this.shutdown('SIGHUP')); + process.on('exit', () => this.releaseOwnedSocket()); + if (this.config.weaveProject && this.config.apiKey) { try { await this.initWeave(); @@ -87,20 +97,8 @@ export class Daemon { } 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')); - process.on('SIGHUP', () => void this.shutdown('SIGHUP')); - process.on('exit', () => { - try { - if (fs.existsSync(this.socketPath)) fs.unlinkSync(this.socketPath); - } catch { - // The next hook's socket probe handles stale-socket cleanup. - } - }); - const checkEveryMs = Math.min( 60_000, Math.max(500, Math.floor(this.inactivityMs / 4)), @@ -135,10 +133,11 @@ export class Daemon { reject(err); }; server.once('error', onError); + this.server = server; server.listen(this.socketPath, () => { process.umask(previousUmask); server.removeListener('error', onError); - this.server = server; + this.captureSocketOwnership(); resolve(); }); }); @@ -170,6 +169,31 @@ export class Daemon { } } + /** Capture an early bind only while this process's server is still live. */ + private captureSocketOwnership(): void { + if (this.ownedSocketInode !== undefined || !this.server?.listening) return; + try { + this.ownedSocketInode = fs.statSync(this.socketPath).ino; + } catch { + // The socket was already removed. + } + } + + /** Remove only the socket inode this daemon created. */ + private releaseOwnedSocket(): void { + this.captureSocketOwnership(); + const owned = this.ownedSocketInode; + this.ownedSocketInode = undefined; + if (owned === undefined) return; + try { + if (fs.statSync(this.socketPath).ino === owned) { + fs.unlinkSync(this.socketPath); + } + } catch { + // The socket was already removed. + } + } + private async initWeave(): Promise { if (!this.config.weaveProject) { throw new Error('weaveProject required to init tracer'); @@ -311,7 +335,21 @@ export class Daemon { private async drain(reason: string): Promise { this.log('INFO', `Shutdown: ${reason}`); - this.server?.close(); + // The path can become visible just before listen's callback records it. + this.captureSocketOwnership(); + let serverClosed: Promise | undefined; + if (this.server?.listening) { + // Node stops listening and unlinks a Unix socket when close() is called; + // the callback only waits for already accepted connections. + serverClosed = new Promise(resolve => + this.server!.close(() => resolve())); + // A successor may now bind this path, so never inspect it again. + this.ownedSocketInode = undefined; + } else { + this.releaseOwnedSocket(); + } + await serverClosed; + await this.hookHandler.waitForPendingEvents(); this.hookHandler.finalizeForShutdown(); if (this.tracingEnabled) { try { @@ -321,9 +359,6 @@ export class Daemon { } } this.hookHandler.closeTranscripts(); - if (fs.existsSync(this.socketPath)) { - fs.unlinkSync(this.socketPath); - } } private log(level: 'DEBUG' | 'INFO' | 'ERROR', message: string): void { diff --git a/src/hookHandler.ts b/src/hookHandler.ts index 36e5271..f7fff54 100644 --- a/src/hookHandler.ts +++ b/src/hookHandler.ts @@ -732,6 +732,11 @@ export class HookHandler { return false; } + /** Admission must be stopped before taking this snapshot. */ + async waitForPendingEvents(): Promise { + await Promise.all([...this.sessionQueues.values()]); + } + finalizeForShutdown(): void { for (const session of this.sessions.values()) { try { diff --git a/src/session.ts b/src/session.ts index 9d9942d..7bbd27a 100644 --- a/src/session.ts +++ b/src/session.ts @@ -30,6 +30,11 @@ import { TranscriptFile, readFirstTranscriptLine } from './transcriptFile.js'; type TraceLog = (level: 'DEBUG' | 'INFO' | 'ERROR', message: string) => void; +/** Keep forced parent closes after children opened during this clock tick. */ +function spanCloseTime(): Date { + return new Date(Date.now() + 1); +} + export type TurnTrace = { kind: 'turn'; span: weave.Turn; @@ -190,21 +195,22 @@ export class Session { finishAtSessionEnd(promptId: string | undefined): number { const parsed = this.parseTranscript(); this.reconcileFinalTurn(promptId, parsed); - return this.finishTurns('session_ended', parsed); + return this.finishTurns('session_ended', parsed, spanCloseTime()); } finishOpenTurns(orphanReason: string): number { - return this.finishTurns(orphanReason, this.parseTranscript()); + return this.finishTurns(orphanReason, this.parseTranscript(), spanCloseTime()); } private finishTurns( orphanReason: string, parsed: ParsedSession | null, + endTime: Date, ): number { const turnCount = this.turns.size; for (const turn of [...this.turns]) { this.recordFinalTurnOutput(turn, orphanReason, parsed); - this.closeTurn(turn, orphanReason); + this.closeTurn(turn, orphanReason, endTime); } return turnCount; } @@ -414,14 +420,27 @@ export class Session { private finalizeTurn(turn: TurnTrace, orphanReason: string): void { this.recordFinalTurnOutput(turn, orphanReason, this.parseTranscript()); - this.closeTurn(turn, orphanReason); + this.closeTurn( + turn, + orphanReason, + turn.children.size ? spanCloseTime() : undefined, + ); } - private closeTurn(turn: TurnTrace, orphanReason: string): void { - for (const toolUseId of finalizeOpenCalls(this.calls, [turn], orphanReason)) { + private closeTurn( + turn: TurnTrace, + orphanReason: string, + endTime?: Date, + ): void { + for (const toolUseId of finalizeOpenCalls( + this.calls, + [turn], + orphanReason, + endTime, + )) { this.log('DEBUG', `Closed pending call: ${toolUseId}`); } - turn.span.end(); + turn.span.end(endTime ? { endTime } : undefined); this.turns.delete(turn); if (this.currentTurn === turn) this.currentTurn = undefined; } diff --git a/tests/daemon-idle-inflight.test.ts b/tests/daemon-idle-inflight.test.ts index 1848159..9498164 100644 --- a/tests/daemon-idle-inflight.test.ts +++ b/tests/daemon-idle-inflight.test.ts @@ -92,3 +92,48 @@ test('an open tool keeps a stopped turn alive', async () => { await d.stop(); } }); + +test('shutdown drains queued hooks before finalizing', async () => { + const d = await startTestDaemon(); + try { + const sessionId = 'shutdown-queued-hooks'; + const transcript = writeTranscript(d.home, sessionId); + await d.send({ + hook_event_name: 'SessionStart', + session_id: sessionId, + transcript_path: transcript, + }); + await d.send({ + hook_event_name: 'UserPromptSubmit', + session_id: sessionId, + transcript_path: transcript, + prompt: 'queue work', + }); + assert.ok(await d.waitForLog(/Created turn span/)); + + // Stop remains in its transcript retry loop while the tool queues behind it. + await d.send({ + hook_event_name: 'Stop', + session_id: sessionId, + transcript_path: transcript, + last_assistant_message: 'not flushed yet', + }); + await d.send({ + hook_event_name: 'PreToolUse', + session_id: sessionId, + transcript_path: transcript, + tool_use_id: 'queued-tool', + tool_name: 'Read', + tool_input: { file_path: '/x' }, + }); + d.proc.kill('SIGTERM'); + + assert.ok( + await d.waitForExit(5000), + `daemon did not exit; log was:\n${d.readLog()}`, + ); + assert.match(d.readLog(), /Closed pending call: queued-tool/); + } finally { + await d.stop(); + } +}); diff --git a/tests/daemon-shutdown-finalizes-turn.test.ts b/tests/daemon-shutdown-finalizes-turn.test.ts new file mode 100644 index 0000000..ff96a9a --- /dev/null +++ b/tests/daemon-shutdown-finalizes-turn.test.ts @@ -0,0 +1,122 @@ +// SPDX-FileCopyrightText: 2026 CoreWeave, Inc. +// SPDX-License-Identifier: MIT +// SPDX-PackageName: weave-claude-code + +import { test } 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'; + +for (const closure of [ + { + name: 'SessionEnd', + reason: 'session_ended', + close: (daemon: ReturnType, sid: string) => + daemon.routeEvent({ + hook_event_name: 'SessionEnd', + session_id: sid, + reason: 'clear', + }), + }, + { + name: 'daemon drain', + reason: 'daemon_shutdown', + close: (daemon: ReturnType) => + daemon.drain('inactivity'), + }, +]) { + test(`${closure.name} exports an open turn and its completed child`, async (t) => { + const exporter = await initWeaveInMemory(); + exporter.reset(); + const sid = `turn-close-${closure.name}`; + const transcript = makeTranscript(t, sid, 'turn-close'); + transcript.append( + userEntry('do it'), + assistantEntry('msg-a', { type: 'text', text: 'working' }), + ); + 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, prompt: 'do it' }); + await daemon.routeEvent({ hook_event_name: 'PreToolUse', session_id: sid, tool_use_id: 'tool-1', tool_name: 'Read', tool_input: { file_path: '/foo' } }); + await daemon.routeEvent({ hook_event_name: 'PostToolUse', session_id: sid, tool_use_id: 'tool-1', tool_response: 'ok' }); + await closure.close(daemon, sid); + await flushWeave(); + + const spans = exporter.getFinishedSpans(); + const turn = spans.find(span => + span.attributes[ATTR.OPERATION_NAME] === 'invoke_agent'); + const chat = spans.find(span => + span.attributes[ATTR.RESPONSE_ID] === 'msg-a'); + const tool = spans.find(span => + span.attributes[ATTR.OPERATION_NAME] === 'execute_tool'); + assert.ok(turn && chat && tool); + assert.equal(turn.attributes[ATTR.WEAVE_ORPHAN_REASON], closure.reason); + assert.equal(spanParentId(chat), turn.spanContext().spanId); + assert.equal(spanParentId(tool), turn.spanContext().spanId); + assert.ok(spans.indexOf(chat) < spans.indexOf(turn)); + }); +} + +test('daemon drain orphans an open Agent under its turn', async (t) => { + const exporter = await initWeaveInMemory(); + exporter.reset(); + const sid = 'turn-close-agent'; + const transcript = makeTranscript(t, sid, 'turn-close-agent'); + transcript.append(userEntry('delegate it')); + 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, prompt: 'delegate it' }); + await daemon.routeEvent({ + hook_event_name: 'PreToolUse', session_id: sid, tool_use_id: 'agent-1', + tool_name: 'Agent', tool_input: { subagent_type: 'reviewer', prompt: 'review' }, + }); + await daemon.drain('SIGTERM'); + await flushWeave(); + + const spans = exporter.getFinishedSpans(); + const turn = spans.find(span => + span.attributes[ATTR.AGENT_NAME] === 'claude-code'); + const agent = spans.find(span => + span.attributes[ATTR.AGENT_NAME] === 'reviewer'); + assert.ok(turn && agent); + assert.equal(agent.attributes[ATTR.WEAVE_ORPHAN_REASON], 'daemon_shutdown'); + assert.equal(spanParentId(agent), turn.spanContext().spanId); + assert.deepEqual(agent.endTime, turn.endTime); +}); + +test('abandoning a permission-pending tool records an orphan, not a denial', async (t) => { + const exporter = await initWeaveInMemory(); + exporter.reset(); + const sid = 'turn-close-permission'; + const transcript = makeTranscript(t, sid, 'turn-close-permission'); + transcript.append(userEntry('run it')); + const daemon = makeGenaiDaemon(); + const toolInput = { command: 'sleep 10' }; + + 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, prompt: 'run it' }); + await daemon.routeEvent({ hook_event_name: 'PreToolUse', session_id: sid, tool_use_id: 'pending-tool', tool_name: 'Bash', tool_input: toolInput }); + await daemon.routeEvent({ hook_event_name: 'PermissionRequest', session_id: sid, tool_name: 'Bash', tool_input: toolInput }); + await daemon.drain('SIGTERM'); + await flushWeave(); + + const tool = exporter.getFinishedSpans().find(span => + span.attributes['gen_ai.tool.call.id'] === 'pending-tool'); + assert.ok(tool); + assert.equal(tool.attributes[ATTR.WEAVE_ORPHAN_REASON], 'daemon_shutdown'); + assert.equal(tool.status.code, 2); + assert.equal( + tool.events.some(event => event.name === ATTR.EVT_PERMISSION_RESOLVED), + false, + ); +}); diff --git a/tests/daemon-socket-ownership.test.ts b/tests/daemon-socket-ownership.test.ts new file mode 100644 index 0000000..5886dc5 --- /dev/null +++ b/tests/daemon-socket-ownership.test.ts @@ -0,0 +1,120 @@ +// SPDX-FileCopyrightText: 2026 CoreWeave, Inc. +// SPDX-License-Identifier: MIT +// SPDX-PackageName: weave-claude-code + +import { test } from 'node:test'; +import assert from 'node:assert/strict'; +import * as fs from 'node:fs'; +import * as net from 'node:net'; +import * as path from 'node:path'; +import { Daemon } from '../src/daemon.ts'; + +type DaemonInternals = { + bindSocketWithHerdProtection(): Promise; + releaseOwnedSocket(): void; + drain(reason: string): Promise; + hookHandler: { + waitForPendingEvents(): Promise; + }; + ownedSocketInode?: number; + server?: net.Server; +}; + +function makeDaemon(socketPath: string, dir: string): DaemonInternals { + return new Daemon(socketPath, path.join(dir, 'daemon.log'), { + weaveProject: null, + apiKey: null, + baseUrl: 'https://x', + agentName: 'claude-code', + debug: false, + }) as unknown as DaemonInternals; +} + +function listen(server: net.Server, socketPath: string): Promise { + return new Promise(resolve => server.listen(socketPath, resolve)); +} + +test('socket release never unlinks a successor inode', async (t) => { + const dir = fs.mkdtempSync('/tmp/wcp-socket-owner-'); + t.after(() => fs.rmSync(dir, { recursive: true, force: true })); + const socketPath = path.join(dir, 'daemon.sock'); + const daemon = makeDaemon(socketPath, dir); + const successor = net.createServer(); + t.after(() => { try { successor.close(); } catch { /* already closed */ } }); + t.after(() => { try { daemon.server?.close(); } catch { /* already closed */ } }); + + await daemon.bindSocketWithHerdProtection(); + fs.unlinkSync(socketPath); + await listen(successor, socketPath); + const successorInode = fs.statSync(socketPath).ino; + + daemon.releaseOwnedSocket(); + + assert.equal(fs.statSync(socketPath).ino, successorInode); +}); + +test('early release captures ownership from its listening server', async (t) => { + const dir = fs.mkdtempSync('/tmp/wcp-socket-early-'); + t.after(() => fs.rmSync(dir, { recursive: true, force: true })); + const socketPath = path.join(dir, 'daemon.sock'); + const daemon = makeDaemon(socketPath, dir); + t.after(() => { try { daemon.server?.close(); } catch { /* already closed */ } }); + + await daemon.bindSocketWithHerdProtection(); + daemon.ownedSocketInode = undefined; + daemon.releaseOwnedSocket(); + + assert.equal(fs.existsSync(socketPath), false); +}); + +test('drain releases socket ownership before waiting for queued hooks', async (t) => { + const dir = fs.mkdtempSync('/tmp/wcp-socket-drain-'); + t.after(() => fs.rmSync(dir, { recursive: true, force: true })); + const socketPath = path.join(dir, 'daemon.sock'); + const daemon = makeDaemon(socketPath, dir); + let releaseQueue!: () => void; + const queue = new Promise(resolve => { releaseQueue = resolve; }); + daemon.hookHandler.waitForPendingEvents = () => queue; + + await daemon.bindSocketWithHerdProtection(); + const draining = daemon.drain('test'); + await new Promise(resolve => setTimeout(resolve, 20)); + assert.equal(fs.existsSync(socketPath), false); + + const successor = net.createServer(); + t.after(() => { try { successor.close(); } catch { /* already closed */ } }); + await listen(successor, socketPath); + const successorInode = fs.statSync(socketPath).ino; + releaseQueue(); + await draining; + + assert.equal(fs.statSync(socketPath).ino, successorInode); +}); + +test('drain never removes a successor while old connections close', async (t) => { + const dir = fs.mkdtempSync('/tmp/wcp-socket-handoff-'); + t.after(() => fs.rmSync(dir, { recursive: true, force: true })); + const socketPath = path.join(dir, 'daemon.sock'); + const daemon = makeDaemon(socketPath, dir); + + await daemon.bindSocketWithHerdProtection(); + const oldClient = net.createConnection(socketPath); + t.after(() => oldClient.destroy()); + await new Promise((resolve, reject) => { + oldClient.once('connect', resolve); + oldClient.once('error', reject); + }); + + const draining = daemon.drain('test'); + await new Promise(resolve => setTimeout(resolve, 20)); + assert.equal(fs.existsSync(socketPath), false); + + const successor = net.createServer(); + t.after(() => { try { successor.close(); } catch { /* already closed */ } }); + await listen(successor, socketPath); + const successorInode = fs.statSync(socketPath).ino; + + oldClient.end(); + await draining; + assert.equal(fs.statSync(socketPath).ino, successorInode); +}); diff --git a/tests/permission-tracing.test.ts b/tests/permission-tracing.test.ts index 629c799..91153aa 100644 --- a/tests/permission-tracing.test.ts +++ b/tests/permission-tracing.test.ts @@ -167,6 +167,11 @@ test('Agent permission events stay on its invoke-agent span', async (t) => { const agent = exporter.getFinishedSpans().find(span => span.attributes[ATTR.AGENT_NAME] === 'Explore'); assert.ok(agent); + assert.equal(agent.attributes[ATTR.WEAVE_ORPHAN_REASON], undefined); + assert.equal( + agent.attributes[ATTR.OUTPUT_MESSAGES], + JSON.stringify([{ role: 'assistant', content: 'done' }]), + ); assert.ok(requestedPermission(agent)); assert.equal(resolvedPermission(agent)?.attributes[ATTR.EVT_PERMISSION_APPROVED], true); });