diff --git a/package.json b/package.json index c51a3cacb..fab54b631 100644 --- a/package.json +++ b/package.json @@ -1,6 +1,6 @@ { "name": "@bastani/atomic", - "version": "0.5.28-0", + "version": "0.5.28-1", "description": "Configuration management CLI and SDK for coding agents", "type": "module", "license": "MIT", diff --git a/src/commands/cli/claude-stop-hook.test.ts b/src/commands/cli/claude-stop-hook.test.ts index 85b13c1c4..4c40775c1 100644 --- a/src/commands/cli/claude-stop-hook.test.ts +++ b/src/commands/cli/claude-stop-hook.test.ts @@ -9,10 +9,11 @@ * and clean up in `afterEach` so test runs never collide with each other * or with real marker/queue/release files. * - * The hook's default wait for a queued follow-up prompt is 15 minutes. - * Every test here passes a short `waitTimeoutMs` so the hook exits quickly - * when no queue entry is present — we are testing the branching logic, - * not the real-world wait budget. + * The hook's default wait for a queued follow-up prompt is effectively + * unbounded (~24 days) so the workflow can take as long as it needs between + * turns. Every test here passes a short `waitTimeoutMs` so the hook exits + * quickly when no queue entry is present — we are testing the branching + * logic, not the real-world wait budget. */ import { describe, test, expect, afterEach, spyOn } from "bun:test"; @@ -20,7 +21,7 @@ import { access, rm, writeFile, mkdir } from "node:fs/promises"; import { join } from "node:path"; import { claudeStopHookCommand, claudeHookDirs } from "./claude-stop-hook.ts"; -const { marker: markerDir, queue: queueDir, release: releaseDir } = claudeHookDirs(); +const { marker: markerDir, queue: queueDir, release: releaseDir, pid: pidDir } = claudeHookDirs(); const SHORT_TIMEOUT_MS = 300; @@ -52,6 +53,7 @@ afterEach(async () => { rm(join(markerDir, id), { force: true }), rm(join(queueDir, id), { force: true }), rm(join(releaseDir, id), { force: true }), + rm(join(pidDir, id), { force: true }), ]); } sessionIdsToClean.length = 0; @@ -268,4 +270,48 @@ describe("claudeStopHookCommand", () => { // No block decision emitted. expect(stdoutChunks.join("")).toBe(""); }); + + // 9. Dead atomic PID → hook exits without waiting out the full timeout. + // + // Simulates the case where the atomic workflow was SIGKILL'd between + // turns: the pid file on disk points at a process that no longer exists, + // so the liveness check should fire and let the hook bail. We pick a + // deliberately-bogus PID (2^22 - 1) that is almost certainly unused. + test("dead atomic pid triggers liveness exit before the wait timeout", async () => { + const sessionId = crypto.randomUUID(); + sessionIdsToClean.push(sessionId); + + // Find a PID that doesn't currently exist. `process.kill(pid, 0)` throws + // ESRCH for free PIDs; we scan from a high number downward to dodge + // system-reserved low PIDs. + let deadPid = 4_194_303; + while (deadPid > 1) { + try { + process.kill(deadPid, 0); + deadPid -= 1; + } catch (e: unknown) { + if (e instanceof Error && "code" in e && (e as NodeJS.ErrnoException).code === "ESRCH") break; + deadPid -= 1; + } + } + + await mkdir(pidDir, { recursive: true }); + await writeFile(join(pidDir, sessionId), String(deadPid), "utf-8"); + + mockStdin(JSON.stringify({ session_id: sessionId })); + + // Use a long wait timeout so the test only passes if the liveness check + // short-circuits the wait. livenessIntervalMs is short so the test runs fast. + const started = Date.now(); + const code = await claudeStopHookCommand({ + waitTimeoutMs: 30_000, + pollIntervalMs: 10_000, + livenessIntervalMs: 50, + }); + const elapsed = Date.now() - started; + + expect(code).toBe(0); + expect(elapsed).toBeLessThan(5_000); + expect(await fileExists(join(markerDir, sessionId))).toBe(true); + }); }); diff --git a/src/commands/cli/claude-stop-hook.ts b/src/commands/cli/claude-stop-hook.ts index 0d24af8cd..8fb2751ba 100644 --- a/src/commands/cli/claude-stop-hook.ts +++ b/src/commands/cli/claude-stop-hook.ts @@ -29,6 +29,7 @@ */ import fs from "node:fs/promises"; +import { watch as watchDir } from "node:fs/promises"; import { existsSync } from "node:fs"; import path from "node:path"; import os from "node:os"; @@ -60,13 +61,25 @@ function isClaudeStopHookPayload(value: unknown): value is ClaudeStopHookPayload * * Exported so tests and `src/sdk/providers/claude.ts` share one source of truth. */ -export function claudeHookDirs(): { marker: string; queue: string; release: string; hil: string } { +export function claudeHookDirs(): { + marker: string; + queue: string; + release: string; + hil: string; + pid: string; +} { const base = path.join(os.homedir(), ".atomic"); return { marker: path.join(base, "claude-stop"), queue: path.join(base, "claude-queue"), release: path.join(base, "claude-release"), hil: path.join(base, "claude-hil"), + // Holds the PID of the atomic workflow process that owns each session. + // The Stop hook polls `process.kill(pid, 0)` against this value so that + // if atomic is SIGKILL'd (no chance to write a release marker), the hook + // can detect the orphaned session and self-exit instead of sitting in + // its wait loop for ~24 days. + pid: path.join(base, "claude-pid"), }; } @@ -74,12 +87,103 @@ export function claudeHookDirs(): { marker: string; queue: string; release: stri export interface ClaudeStopHookOptions { /** Maximum time the hook waits for a queued follow-up prompt before letting Claude stop. */ waitTimeoutMs?: number; - /** Polling interval for queue/release detection. */ + /** + * Interval for the polling fallback that runs alongside the `fs.watch` + * watchers in case an inotify/FSEvent notification gets dropped. In the + * happy path, watcher events fire on create and the poll never matches. + */ pollIntervalMs?: number; + /** + * Interval at which the hook checks whether the atomic workflow process + * that owns this session is still alive. Coarser than `pollIntervalMs` + * because atomic crashing is rare and `process.kill(pid, 0)` is a syscall. + */ + livenessIntervalMs?: number; } -const DEFAULT_WAIT_TIMEOUT_MS = 15 * 60 * 1000; +/** + * Effectively-unbounded default wait budget for the queue/release poll loop. + * + * The hook holds Claude Code in the Stop phase while the workflow runtime + * decides what to do next — either enqueueing a follow-up prompt (delivered + * back to Claude as `{decision:"block", reason:...}`) or writing a release + * marker on teardown. Any finite default here caps the time the workflow has + * between turns: when it expires, the hook exits 0, Claude stops, and the + * next `enqueuePrompt` writes to a file nobody's reading — the workflow + * hangs on `waitForIdle` for a turn that will never come. + * + * The Claude-side hook timeout (see `STOP_HOOK_TIMEOUT_SECONDS` in + * `src/sdk/providers/claude.ts`) is already set to ~24 days, so matching it + * here keeps the two bounds aligned — the hook either runs until the + * workflow releases it or until Claude Code itself gives up. Tests override + * `waitTimeoutMs` via options to keep runs fast. + * + * Expressed in ms: 2_147_483 s × 1000 = 2_147_483_000 ms, just under the + * max safe `setTimeout` value (2^31 - 1). + */ +const DEFAULT_WAIT_TIMEOUT_MS = 2_147_483_000; const DEFAULT_POLL_INTERVAL_MS = 100; +const DEFAULT_LIVENESS_INTERVAL_MS = 5_000; + +/** + * Read the atomic PID that owns this session from `~/.atomic/claude-pid/`, + * or return null if the file is missing / malformed. Missing is fine: older + * runtimes didn't write one, and we just skip the liveness check in that case. + */ +async function readAtomicPid(pidFilePath: string): Promise { + let raw: string; + try { + raw = await fs.readFile(pidFilePath, "utf-8"); + } catch { + return null; + } + const parsed = Number.parseInt(raw.trim(), 10); + return Number.isInteger(parsed) && parsed > 0 ? parsed : null; +} + +/** + * Sleep that resolves early when `signal` is aborted. Used by the hook's + * wait loops so `ac.abort()` unblocks everything immediately instead of + * waiting for the next wake-up tick — otherwise a task that detects a hit + * (e.g. liveness check) can't meaningfully cancel its siblings. + */ +function abortableSleep(ms: number, signal: AbortSignal): Promise { + return new Promise((resolve) => { + if (signal.aborted) { + resolve(); + return; + } + const timer = setTimeout(() => { + signal.removeEventListener("abort", onAbort); + resolve(); + }, ms); + const onAbort = (): void => { + clearTimeout(timer); + resolve(); + }; + signal.addEventListener("abort", onAbort, { once: true }); + }); +} + +/** + * True when a process with `pid` exists. Uses signal `0`, which performs the + * permission/existence check without delivering a signal. ESRCH means gone, + * EPERM means alive-but-not-ours (still alive for our purposes). + */ +function isProcessAlive(pid: number): boolean { + try { + process.kill(pid, 0); + return true; + } catch (e: unknown) { + if (e instanceof Error && "code" in e) { + const code = (e as NodeJS.ErrnoException).code; + if (code === "EPERM") return true; + if (code === "ESRCH") return false; + } + // Unknown error — assume alive to avoid false-positive teardown. + return true; + } +} /** * Handler for the hidden `_claude-stop-hook` subcommand. @@ -95,6 +199,8 @@ export async function claudeStopHookCommand( ): Promise { const waitTimeoutMs = options.waitTimeoutMs ?? DEFAULT_WAIT_TIMEOUT_MS; const pollIntervalMs = options.pollIntervalMs ?? DEFAULT_POLL_INTERVAL_MS; + const livenessIntervalMs = + options.livenessIntervalMs ?? DEFAULT_LIVENESS_INTERVAL_MS; // 1. Read stdin const raw = await Bun.stdin.text(); @@ -135,6 +241,7 @@ export async function claudeStopHookCommand( fs.mkdir(dirs.marker, { recursive: true }), fs.mkdir(dirs.queue, { recursive: true }), fs.mkdir(dirs.release, { recursive: true }), + fs.mkdir(dirs.pid, { recursive: true }), ]); // 4. Write the marker file directly. @@ -148,7 +255,7 @@ export async function claudeStopHookCommand( const markerPath = path.join(dirs.marker, payload.session_id); await Bun.write(markerPath, raw); - // 5. Block-poll for either a queued follow-up prompt or a release signal. + // 5. Wait for either a queued follow-up prompt or a release signal. // // The workflow's `waitForIdle` has already been unblocked by the marker // write above and is now returning control to the user's stage callback. @@ -164,34 +271,130 @@ export async function claudeStopHookCommand( // `~/.atomic/claude-release/`. We exit 0 with no stdout // payload and Claude stops as usual. // - // c. Neither happens within `waitTimeoutMs`. We exit 0 on timeout as a - // safety net — Claude stops rather than hanging its Stop hook forever. + // c. Neither happens within `waitTimeoutMs`. We exit 0 so Claude Code + // doesn't hang past its own per-hook timeout. The production default + // for `waitTimeoutMs` is aligned with the Claude-side hook timeout + // (~24 days), so this path is effectively unreachable in real runs — + // it only fires in tests that pass a short override. + // + // Delivery uses `fs.watch` on the queue and release dirs for ~0-latency + // wake-up on create events, with a slower `existsSync` polling fallback + // in case a watcher notification gets dropped under fs load (same pattern + // as `watchHILMarker` in `src/sdk/providers/claude.ts`). const queuePath = path.join(dirs.queue, payload.session_id); const releasePath = path.join(dirs.release, payload.session_id); - const deadline = Date.now() + waitTimeoutMs; - while (Date.now() <= deadline) { + type Hit = { kind: "release" } | { kind: "queue"; prompt: string }; + + const check = async (): Promise => { if (existsSync(releasePath)) { try { await fs.unlink(releasePath); } catch { /* ENOENT is fine */ } - return 0; + return { kind: "release" }; } if (existsSync(queuePath)) { let prompt: string; try { prompt = await fs.readFile(queuePath, "utf-8"); } catch { - return 0; + // Treat a failed read as a graceful release so the hook still exits. + return { kind: "release" }; } try { await fs.unlink(queuePath); } catch { /* ENOENT is fine */ } + return { kind: "queue", prompt }; + } + return null; + }; + + const emit = (hit: Hit): number => { + if (hit.kind === "queue") { process.stdout.write(JSON.stringify({ decision: "block", - reason: prompt, + reason: hit.prompt, })); - return 0; } - await Bun.sleep(pollIntervalMs); + return 0; + }; + + // Initial synchronous check — the runtime may have enqueued/released before + // we attached watchers, and without this the hook could hang until the + // polling fallback fires. + const early = await check(); + if (early) return emit(early); + + const ac = new AbortController(); + const overallTimer = setTimeout(() => ac.abort(), waitTimeoutMs); + let hit: Hit | null = null; + + // Read the atomic workflow's PID (if the runtime wrote one for this + // session). Used by the liveness task below to detect an atomic crash. + const atomicPid = await readAtomicPid( + path.join(dirs.pid, payload.session_id), + ); + + // Watch a single directory for change events and resolve `hit` on the + // first one that matches. `event.filename` is unreliable across OSes + // (see the comment in `watchHILMarker`), so disk state is authoritative. + const runWatcher = async (dir: string): Promise => { + try { + for await (const _event of watchDir(dir, { signal: ac.signal })) { + const result = await check(); + if (result) { + hit = result; + ac.abort(); + return; + } + } + } catch (e: unknown) { + if (!(e instanceof Error && e.name === "AbortError")) throw e; + } + }; + + // Polling fallback — catches the rare dropped inotify/FSEvent event. + // Only runs while the watchers are live; `ac.abort()` shuts it down. + const runPollFallback = async (): Promise => { + while (!ac.signal.aborted) { + await abortableSleep(pollIntervalMs, ac.signal); + if (ac.signal.aborted) return; + const result = await check(); + if (result) { + hit = result; + ac.abort(); + return; + } + } + }; + + // Liveness check — if the atomic workflow process died without writing a + // release marker (e.g. SIGKILL), this task abandons the wait and lets + // Claude stop. No-op when there's no pid file (older sessions or non- + // runtime spawns) so the hook still functions standalone. + const runLivenessCheck = async (): Promise => { + if (atomicPid === null) return; + while (!ac.signal.aborted) { + await abortableSleep(livenessIntervalMs, ac.signal); + if (ac.signal.aborted) return; + if (!isProcessAlive(atomicPid)) { + // hit stays null → the hook exits 0 without emitting a block decision. + ac.abort(); + return; + } + } + }; + + try { + await Promise.all([ + runWatcher(dirs.queue), + runWatcher(dirs.release), + runPollFallback(), + runLivenessCheck(), + ]); + } finally { + clearTimeout(overallTimer); + ac.abort(); } + if (hit) return emit(hit); + // Timeout — no queued prompt arrived. Let Claude stop normally. return 0; } diff --git a/src/sdk/providers/claude.ts b/src/sdk/providers/claude.ts index 6cd649c7c..15108fca4 100644 --- a/src/sdk/providers/claude.ts +++ b/src/sdk/providers/claude.ts @@ -71,6 +71,12 @@ export async function clearClaudeSession(paneId: string): Promise { // Best-effort — if release fails the hook will still exit on its // own safety timeout. } + try { + await unlinkAtomicPidFile(state.claudeSessionId); + } catch { + // Best-effort — stale pid file is inert; the next session writes a + // fresh one under its own UUID. + } } initializedPanes.delete(paneId); } @@ -258,6 +264,12 @@ export async function createClaudeSession(options: ClaudeSessionOptions): Promis chatFlags, readyTimeoutMs, }); + + // Write our PID so the Stop hook can detect an orphaned session if we + // crash/get SIGKILL'd without running teardown. Best-effort; failures just + // mean the hook falls back to waiting out Claude's own hook timeout. + await writeAtomicPidFile(claudeSessionId); + return claudeSessionId; } @@ -609,6 +621,38 @@ export async function releaseClaudeSession(claudeSessionId: string): Promise` so the Stop hook + * can use it as a liveness signal. If atomic is SIGKILL'd (no chance to run + * `clearClaudeSession`), the hook detects the dead PID via `process.kill(..,0)` + * and self-exits instead of parking Claude for the full 24-day timeout. + */ +async function writeAtomicPidFile(claudeSessionId: string): Promise { + await mkdir(pidDir(), { recursive: true }); + await writeFile(pidFilePath(claudeSessionId), String(process.pid), "utf-8"); +} + +/** Remove the pid file for a session. Idempotent — ENOENT is swallowed. */ +async function unlinkAtomicPidFile(claudeSessionId: string): Promise { + try { + await unlink(pidFilePath(claudeSessionId)); + } catch (e: unknown) { + if (!(e instanceof Error && "code" in e && (e as NodeJS.ErrnoException).code === "ENOENT")) { + throw e; + } + } +} + // --------------------------------------------------------------------------- // Idle detection via marker file watch // ---------------------------------------------------------------------------