{"version":3,"file":"control-channel.d.ts","sourceRoot":"","sources":["../../../../src/runs/background/control-channel.ts"],"names":[],"mappings":"AAAA;;;;;;;;;;;;;;GAcG;AAGH,OAAO,KAAK,EAAE,MAAM,SAAS,CAAC;AAM9B;;;;GAIG;AACH,eAAO,MAAM,gBAAgB,EAAE,MAAM,CAAC,OAA+D,CAAC;AAEtG,MAAM,MAAM,gBAAgB,GAAG,IAAI,CAClC,OAAO,EAAE,EACT,WAAW,GAAG,YAAY,GAAG,QAAQ,GAAG,OAAO,GAAG,aAAa,GAAG,cAAc,GAAG,cAAc,CACjG,CAAC;AACF,MAAM,MAAM,oBAAoB,GAAG;IAAE,WAAW,EAAE,OAAO,WAAW,CAAC;IAAC,aAAa,EAAE,OAAO,aAAa,CAAA;CAAE,CAAC;AAC5G,KAAK,MAAM,GAAG,CAAC,GAAG,EAAE,MAAM,EAAE,MAAM,CAAC,EAAE,MAAM,CAAC,OAAO,GAAG,CAAC,KAAK,OAAO,CAAC;AAEpE,MAAM,WAAW,gBAAgB;IAChC,IAAI,EAAE,WAAW,CAAC;IAClB,EAAE,CAAC,EAAE,MAAM,CAAC;IACZ,MAAM,CAAC,EAAE,MAAM,CAAC;IAChB,MAAM,CAAC,EAAE,MAAM,CAAC;CAChB;AAED,MAAM,WAAW,cAAc;IAC9B,IAAI,EAAE,SAAS,CAAC;IAChB,EAAE,CAAC,EAAE,MAAM,CAAC;IACZ,MAAM,CAAC,EAAE,MAAM,CAAC;IAChB,MAAM,CAAC,EAAE,MAAM,CAAC;CAChB;AAED,MAAM,WAAW,WAAW;IAC3B,IAAI,EAAE,MAAM,CAAC;IACb,EAAE,CAAC,EAAE,MAAM,CAAC;IACZ,MAAM,CAAC,EAAE,MAAM,CAAC;IAChB,MAAM,CAAC,EAAE,MAAM,CAAC;CAChB;AAED,MAAM,WAAW,yBAAyB;IACzC,IAAI,EAAE,oBAAoB,GAAG,mBAAmB,CAAC;IACjD,EAAE,CAAC,EAAE,MAAM,CAAC;IACZ,MAAM,CAAC,EAAE,MAAM,CAAC;IAChB,MAAM,CAAC,EAAE,MAAM,CAAC;CAChB;AAED,MAAM,MAAM,iBAAiB,GAAG,OAAO,GAAG,WAAW,GAAG,MAAM,CAAC;AAC/D,MAAM,MAAM,mBAAmB,GAAG,WAAW,GAAG,QAAQ,CAAC;AAEzD,MAAM,WAAW,YAAY;IAC5B,IAAI,EAAE,OAAO,CAAC;IACd,EAAE,EAAE,MAAM,CAAC;IACX,EAAE,EAAE,MAAM,CAAC;IACX,OAAO,EAAE,MAAM,CAAC;IAChB,IAAI,CAAC,EAAE,iBAAiB,CAAC;IACzB,WAAW,CAAC,EAAE,MAAM,CAAC;IACrB,aAAa,CAAC,EAAE,MAAM,EAAE,CAAC;IACzB,MAAM,CAAC,EAAE,MAAM,CAAC;CAChB;AAED,MAAM,WAAW,eAAe;IAC/B,IAAI,EAAE,kBAAkB,CAAC;IACzB,eAAe,EAAE,CAAC,CAAC;IACnB,KAAK,EAAE,MAAM,CAAC;IACd,GAAG,EAAE,MAAM,CAAC;IACZ,OAAO,EAAE,MAAM,CAAC;IAChB,SAAS,EAAE,OAAO,CAAC;CACnB;AAED,MAAM,WAAW,QAAQ;IACxB,IAAI,EAAE,WAAW,CAAC;IAClB,eAAe,EAAE,CAAC,CAAC;IACnB,SAAS,EAAE,MAAM,CAAC;IAClB,KAAK,EAAE,MAAM,CAAC;IACd,EAAE,EAAE,MAAM,CAAC;IACX,KAAK,EAAE,WAAW,GAAG,QAAQ,GAAG,QAAQ,CAAC;IACzC,cAAc,CAAC,EAAE,mBAAmB,CAAC;IACrC,OAAO,EAAE,MAAM,CAAC;CAChB;AAID,eAAO,MAAM,oBAAoB,KAAK,CAAC;AAQvC,uDAAuD;AACvD,wBAAgB,eAAe,CAAC,QAAQ,EAAE,MAAM,GAAG,MAAM,CAExD;AAED,mDAAmD;AACnD,wBAAgB,oBAAoB,CAAC,QAAQ,EAAE,MAAM,GAAG,MAAM,CAE7D;AAED,iDAAiD;AACjD,wBAAgB,kBAAkB,CAAC,QAAQ,EAAE,MAAM,GAAG,MAAM,CAE3D;AAED,qDAAqD;AACrD,wBAAgB,eAAe,CAAC,QAAQ,EAAE,MAAM,GAAG,MAAM,CAExD;AAED,wBAAgB,4BAA4B,CAAC,QAAQ,EAAE,MAAM,GAAG,MAAM,CAErE;AAED,wBAAgB,2BAA2B,CAAC,QAAQ,EAAE,MAAM,GAAG,MAAM,CAEpE;AAED,uDAAuD;AACvD,wBAAgB,gBAAgB,CAAC,QAAQ,EAAE,MAAM,GAAG,MAAM,CAEzD;AAED,wBAAgB,oBAAoB,CAAC,QAAQ,EAAE,MAAM,GAAG,MAAM,CAE7D;AAED,wBAAgB,eAAe,CAAC,QAAQ,EAAE,MAAM,EAAE,KAAK,EAAE,MAAM,GAAG,IAAI,CAErE;AAED,kFAAkF;AAClF,wBAAgB,iBAAiB,CAAC,QAAQ,EAAE,MAAM,EAAE,KAAK,EAAE,MAAM,GAAG,MAAM,CAGzE;AAED,wBAAgB,oBAAoB,CAAC,QAAQ,EAAE,MAAM,GAAG,MAAM,CAE7D;AAED,wBAAgB,mBAAmB,CAAC,QAAQ,EAAE,MAAM,EAAE,KAAK,EAAE,MAAM,GAAG,MAAM,CAG3E;AAED,wBAAgB,YAAY,CAAC,QAAQ,EAAE,MAAM,EAAE,KAAK,EAAE,MAAM,GAAG,MAAM,CAGpE;AAMD,wBAAgB,mBAAmB,CAAC,GAAG,EAAE,MAAM,EAAE,SAAS,EAAE,MAAM,GAAG,MAAM,CAI1E;AAyCD,wBAAgB,sBAAsB,CAAC,GAAG,EAAE,MAAM,EAAE,OAAO,EAAE,YAAY,GAAG,MAAM,CAKjF;AAED,wBAAgB,sBAAsB,CACrC,QAAQ,EAAE,MAAM,EAChB,UAAU,EAAE,IAAI,CAAC,eAAe,EAAE,MAAM,GAAG,iBAAiB,CAAC,GAC3D,MAAM,CASR;AAED,wBAAgB,oBAAoB,CACnC,QAAQ,EAAE,MAAM,EAChB,UAAU,EAAE,IAAI,CAAC,eAAe,EAAE,MAAM,GAAG,iBAAiB,CAAC,GAC3D,MAAM,CAER;AAgBD,wBAAgB,eAAe,CAAC,QAAQ,EAAE,MAAM,EAAE,GAAG,EAAE,IAAI,CAAC,QAAQ,EAAE,MAAM,GAAG,iBAAiB,CAAC,GAAG,MAAM,CAUzG;AAED,wBAAgB,aAAa,CAAC,QAAQ,EAAE,MAAM,EAAE,GAAG,EAAE,IAAI,CAAC,QAAQ,EAAE,MAAM,GAAG,iBAAiB,CAAC,GAAG,MAAM,CAEvG;AAED;;;GAGG;AACH,wBAAgB,qBAAqB,CACpC,QAAQ,EAAE,MAAM,EAChB,OAAO,GAAE,IAAI,CAAC,gBAAgB,EAAE,MAAM,CAAM,EAC5C,IAAI,GAAE;IAAE,GAAG,CAAC,EAAE,MAAM,MAAM,CAAA;CAAO,GAC/B,MAAM,CAKR;AAED,wBAAgB,mBAAmB,CAClC,QAAQ,EAAE,MAAM,EAChB,OAAO,GAAE,IAAI,CAAC,cAAc,EAAE,MAAM,CAAM,EAC1C,IAAI,GAAE;IAAE,GAAG,CAAC,EAAE,MAAM,MAAM,CAAA;CAAO,GAC/B,MAAM,CAKR;AAED,wBAAgB,gBAAgB,CAC/B,QAAQ,EAAE,MAAM,EAChB,OAAO,GAAE,IAAI,CAAC,WAAW,EAAE,MAAM,CAAM,EACvC,IAAI,GAAE;IAAE,GAAG,CAAC,EAAE,MAAM,MAAM,CAAA;CAAO,GAC/B,MAAM,CAKR;AAED,wBAAgB,8BAA8B,CAC7C,QAAQ,EAAE,MAAM,EAChB,IAAI,EAAE,yBAAyB,CAAC,MAAM,CAAC,EACvC,OAAO,GAAE,IAAI,CAAC,yBAAyB,EAAE,MAAM,CAAM,EACrD,IAAI,GAAE;IAAE,GAAG,CAAC,EAAE,MAAM,MAAM,CAAA;CAAO,GAC/B,MAAM,CAMR;AAED,wBAAgB,iBAAiB,CAChC,QAAQ,EAAE,MAAM,EAChB,OAAO,EAAE;IACR,OAAO,EAAE,MAAM,CAAC;IAChB,IAAI,CAAC,EAAE,iBAAiB,CAAC;IACzB,WAAW,CAAC,EAAE,MAAM,CAAC;IACrB,aAAa,CAAC,EAAE,MAAM,EAAE,CAAC;IACzB,MAAM,CAAC,EAAE,MAAM,CAAC;IAChB,EAAE,CAAC,EAAE,MAAM,CAAC;IACZ,EAAE,CAAC,EAAE,MAAM,CAAC;CACZ,EACD,IAAI,GAAE;IAAE,GAAG,CAAC,EAAE,MAAM,MAAM,CAAC;IAAC,QAAQ,CAAC,EAAE,MAAM,MAAM,CAAA;CAAO,GACxD,MAAM,CA0CR;AAED,wBAAgB,gBAAgB,CAAC,QAAQ,EAAE,MAAM,EAAE,KAAK,EAAE,MAAM,EAAE,OAAO,EAAE,YAAY,GAAG,MAAM,CAQ/F;AA2DD,wBAAgB,mBAAmB,CAAC,QAAQ,EAAE,MAAM,EAAE,KAAK,EAAE,MAAM,GAAG,eAAe,GAAG,SAAS,CAMhG;AAED,wBAAgB,wBAAwB,CACvC,QAAQ,EAAE,MAAM,EAChB,MAAM,GAAE,IAAI,CAAC,OAAO,EAAE,EAAE,YAAY,GAAG,aAAa,GAAG,cAAc,CAAM,GACzE,eAAe,EAAE,CAgBnB;AAED,wBAAgB,gBAAgB,CAC/B,QAAQ,EAAE,MAAM,EAChB,MAAM,GAAE,IAAI,CAAC,OAAO,EAAE,EAAE,YAAY,GAAG,aAAa,GAAG,cAAc,GAAG,QAAQ,CAAM,GACpF,QAAQ,EAAE,CAsCZ;AAkBD,wBAAgB,2BAA2B,CAC1C,GAAG,EAAE,MAAM,EACX,MAAM,GAAE,IAAI,CAAC,OAAO,EAAE,EAAE,YAAY,GAAG,QAAQ,GAAG,aAAa,GAAG,cAAc,CAAM,GACpF,YAAY,EAAE,CA8BhB;AAED,wBAAgB,oBAAoB,CACnC,QAAQ,EAAE,MAAM,EAChB,MAAM,GAAE,IAAI,CAAC,OAAO,EAAE,EAAE,YAAY,GAAG,QAAQ,GAAG,aAAa,GAAG,cAAc,CAAM,GACpF,YAAY,EAAE,CAEhB;AAED,wBAAgB,iBAAiB,CAAC,QAAQ,EAAE,MAAM,EAAE,OAAO,EAAE,YAAY,GAAG,MAAM,CAKjF;AAED,wBAAgB,iBAAiB,CAAC,QAAQ,EAAE,MAAM,GAAG,KAAK,CAAC;IAAE,OAAO,EAAE,YAAY,CAAC;IAAC,IAAI,EAAE,MAAM,CAAA;CAAE,CAAC,CAgBlG;AAED;;;GAGG;AACH,wBAAgB,uBAAuB,CACtC,QAAQ,EAAE,MAAM,EAChB,MAAM,GAAE,IAAI,CAAC,OAAO,EAAE,EAAE,YAAY,GAAG,QAAQ,CAAM,GACnD,OAAO,CAST;AAED,wBAAgB,qBAAqB,CACpC,QAAQ,EAAE,MAAM,EAChB,MAAM,GAAE,IAAI,CAAC,OAAO,EAAE,EAAE,YAAY,GAAG,QAAQ,CAAM,GACnD,OAAO,CAST;AAED,wBAAgB,kBAAkB,CAAC,QAAQ,EAAE,MAAM,EAAE,MAAM,GAAE,IAAI,CAAC,OAAO,EAAE,EAAE,YAAY,GAAG,QAAQ,CAAM,GAAG,OAAO,CASnH;AAED,wBAAgB,gCAAgC,CAC/C,QAAQ,EAAE,MAAM,EAChB,MAAM,GAAE,IAAI,CAAC,OAAO,EAAE,EAAE,YAAY,GAAG,QAAQ,CAAM,GACnD,UAAU,GAAG,UAAU,GAAG,SAAS,CAcrC;AAED;;;;;;GAMG;AACH,wBAAgB,uBAAuB,CAAC,KAAK,EAAE;IAC9C,QAAQ,EAAE,MAAM,CAAC;IACjB,GAAG,CAAC,EAAE,MAAM,CAAC;IACb,IAAI,CAAC,EAAE,MAAM,CAAC;IACd,MAAM,CAAC,EAAE,MAAM,CAAC,OAAO,CAAC;IACxB,GAAG,CAAC,EAAE,MAAM,MAAM,CAAC;IACnB,MAAM,CAAC,EAAE,MAAM,CAAC;CAChB,GAAG,IAAI,CAoBP;AAED,wBAAgB,qBAAqB,CAAC,KAAK,EAAE;IAC5C,QAAQ,EAAE,MAAM,CAAC;IACjB,GAAG,CAAC,EAAE,MAAM,CAAC;IACb,IAAI,CAAC,EAAE,MAAM,CAAC;IACd,MAAM,CAAC,EAAE,MAAM,CAAC,OAAO,CAAC;IACxB,GAAG,CAAC,EAAE,MAAM,MAAM,CAAC;IACnB,MAAM,CAAC,EAAE,MAAM,CAAC;CAChB,GAAG,IAAI,CAEP;AAED,wBAAgB,kBAAkB,CAAC,KAAK,EAAE;IACzC,QAAQ,EAAE,MAAM,CAAC;IACjB,GAAG,CAAC,EAAE,MAAM,CAAC;IACb,IAAI,CAAC,EAAE,MAAM,CAAC;IACd,MAAM,CAAC,EAAE,MAAM,CAAC,OAAO,CAAC;IACxB,GAAG,CAAC,EAAE,MAAM,MAAM,CAAC;IACnB,MAAM,CAAC,EAAE,MAAM,CAAC;CAChB,GAAG,IAAI,CAEP;AAED,wBAAgB,gCAAgC,CAAC,KAAK,EAAE;IACvD,QAAQ,EAAE,MAAM,CAAC;IACjB,QAAQ,EAAE,UAAU,GAAG,UAAU,CAAC;IAClC,GAAG,CAAC,EAAE,MAAM,MAAM,CAAC;IACnB,MAAM,CAAC,EAAE,MAAM,CAAC;IAChB,MAAM,CAAC,EAAE,MAAM,CAAC;CAChB,GAAG,IAAI,CAOP;AAED;;;;;GAKG;AACH,wBAAgB,sBAAsB,CACrC,QAAQ,EAAE,MAAM,EAChB,IAAI,EAAE;IACL,WAAW,EAAE,MAAM,IAAI,CAAC;IACxB,SAAS,CAAC,EAAE,MAAM,IAAI,CAAC;IACvB,MAAM,CAAC,EAAE,MAAM,IAAI,CAAC;IACpB,OAAO,CAAC,EAAE,CAAC,OAAO,EAAE,YAAY,KAAK,IAAI,CAAC;IAC1C,oBAAoB,CAAC,EAAE,CAAC,QAAQ,EAAE,UAAU,GAAG,UAAU,KAAK,IAAI,CAAC;IACnE,iBAAiB,CAAC,EAAE,CAAC,UAAU,EAAE,eAAe,KAAK,IAAI,CAAC;IAC1D,UAAU,CAAC,EAAE,CAAC,GAAG,EAAE,QAAQ,KAAK,IAAI,CAAC;IACrC,cAAc,CAAC,EAAE,MAAM,CAAC;IACxB,EAAE,CAAC,EAAE,gBAAgB,CAAC;IACtB,MAAM,CAAC,EAAE,oBAAoB,CAAC;CAC9B,GACC,MAAM,IAAI,CAqDZ","sourcesContent":["/**\n * Cross-OS control channel for async subagent runs.\n *\n * Background runs are detached OS processes. The original control path delivered\n * an interrupt with `process.kill(pid, SIGUSR2|SIGBREAK)`, but Windows cannot\n * deliver those signals cross-process via `process.kill` and throws `ENOSYS`,\n * which left async runs uninterruptible (no stop, no live steer) on Windows.\n *\n * This module adds a portable, file-based control inbox inside the run directory.\n * The parent drops an interrupt request file; the runner watches the inbox and\n * routes the request into its existing graceful `interruptRunner()` (pause +\n * resumable), identically on every platform. The OS signal is kept only as an\n * opportunistic fast-path; its failure is non-fatal because the file inbox is\n * authoritative.\n */\n\nimport { randomUUID } from \"node:crypto\";\nimport * as fs from \"node:fs\";\nimport * as path from \"node:path\";\nimport { writeAtomicJson } from \"../../shared/atomic-json.ts\";\nimport { POLL_INTERVAL_MS } from \"../../shared/types.ts\";\nimport { resolveWatchPath } from \"../../shared/utils.ts\";\n\n/**\n * Opportunistic fast-path interrupt signal. On Unix `SIGUSR2` is trapped by the\n * runner; on Windows `process.kill(pid, \"SIGBREAK\")` is not deliverable\n * cross-process and throws `ENOSYS`, so the file inbox below is the real channel.\n */\nexport const INTERRUPT_SIGNAL: NodeJS.Signals = process.platform === \"win32\" ? \"SIGBREAK\" : \"SIGUSR2\";\n\nexport type ControlChannelFs = Pick<\n\ttypeof fs,\n\t\"mkdirSync\" | \"existsSync\" | \"rmSync\" | \"watch\" | \"readdirSync\" | \"readFileSync\" | \"realpathSync\"\n>;\nexport type ControlChannelTimers = { setInterval: typeof setInterval; clearInterval: typeof clearInterval };\ntype KillFn = (pid: number, signal?: NodeJS.Signals | 0) => unknown;\n\nexport interface InterruptRequest {\n\ttype: \"interrupt\";\n\tts?: number;\n\tsource?: string;\n\treason?: string;\n}\n\nexport interface TimeoutRequest {\n\ttype: \"timeout\";\n\tts?: number;\n\tsource?: string;\n\treason?: string;\n}\n\nexport interface StopRequest {\n\ttype: \"stop\";\n\tts?: number;\n\tsource?: string;\n\treason?: string;\n}\n\nexport interface CheckpointDecisionRequest {\n\ttype: \"approve-checkpoint\" | \"reject-checkpoint\";\n\tts?: number;\n\tsource?: string;\n\treason?: string;\n}\n\nexport type SteerDeliveryMode = \"steer\" | \"follow_up\" | \"auto\";\nexport type SteerDeliveryStatus = \"delivered\" | \"queued\";\n\nexport interface SteerRequest {\n\ttype: \"steer\";\n\tid: string;\n\tts: number;\n\tmessage: string;\n\tmode?: SteerDeliveryMode;\n\ttargetIndex?: number;\n\ttargetIndexes?: number[];\n\tsource?: string;\n}\n\nexport interface SteerCapability {\n\ttype: \"steer-capability\";\n\tprotocolVersion: 1;\n\tindex: number;\n\tpid: number;\n\treadyAt: number;\n\tsupported: boolean;\n}\n\nexport interface SteerAck {\n\ttype: \"steer-ack\";\n\tprotocolVersion: 1;\n\trequestId: string;\n\tindex: number;\n\tts: number;\n\tstate: \"delivered\" | \"queued\" | \"failed\";\n\tdeliveryStatus?: SteerDeliveryStatus;\n\tmessage: string;\n}\n\nconst STEER_REQUESTS_DIR = \"steer-requests\";\nconst REVIVAL_BRIEFS_DIR = \"revival-briefs\";\nexport const MAX_STEER_QUEUE_SIZE = 20;\nconst STEER_TARGETS_DIR = \"steer-targets\";\nconst STEER_CAPABILITIES_DIR = \"steer-capabilities\";\nconst STEER_ACKS_DIR = \"steer-acks\";\nconst STEER_INBOX_CLOSED_FILE = \"steer-inbox-closed.json\";\nconst MAX_STEER_MESSAGE_BYTES = 128 * 1024;\nconst MAX_STEER_REQUEST_ID_LENGTH = 256;\n\n/** Control inbox directory inside an async run dir. */\nexport function controlInboxDir(asyncDir: string): string {\n\treturn path.join(asyncDir, \"control\");\n}\n\n/** Path of the portable interrupt request file. */\nexport function interruptRequestPath(asyncDir: string): string {\n\treturn path.join(controlInboxDir(asyncDir), \"interrupt.json\");\n}\n\n/** Path of the portable timeout request file. */\nexport function timeoutRequestPath(asyncDir: string): string {\n\treturn path.join(controlInboxDir(asyncDir), \"timeout.json\");\n}\n\n/** Path of the portable manual stop request file. */\nexport function stopRequestPath(asyncDir: string): string {\n\treturn path.join(controlInboxDir(asyncDir), \"stop.json\");\n}\n\nexport function approveCheckpointRequestPath(asyncDir: string): string {\n\treturn path.join(controlInboxDir(asyncDir), \"approve-checkpoint.json\");\n}\n\nexport function rejectCheckpointRequestPath(asyncDir: string): string {\n\treturn path.join(controlInboxDir(asyncDir), \"reject-checkpoint.json\");\n}\n\n/** Directory of parent-to-runner steering requests. */\nexport function steerRequestsDir(asyncDir: string): string {\n\treturn path.join(controlInboxDir(asyncDir), STEER_REQUESTS_DIR);\n}\n\nexport function steerInboxClosedPath(asyncDir: string): string {\n\treturn path.join(controlInboxDir(asyncDir), STEER_INBOX_CLOSED_FILE);\n}\n\nexport function closeSteerInbox(asyncDir: string, state: string): void {\n\twriteAtomicJson(steerInboxClosedPath(asyncDir), { version: 1, closedAt: Date.now(), state });\n}\n\n/** Per-child inbox consumed by the child prompt runtime inside the Pi process. */\nexport function stepSteerInboxDir(asyncDir: string, index: number): string {\n\tassertChildIndex(index);\n\treturn path.join(controlInboxDir(asyncDir), STEER_TARGETS_DIR, String(index));\n}\n\nexport function steerCapabilitiesDir(asyncDir: string): string {\n\treturn path.join(controlInboxDir(asyncDir), STEER_CAPABILITIES_DIR);\n}\n\nexport function steerCapabilityPath(asyncDir: string, index: number): string {\n\tassertChildIndex(index);\n\treturn path.join(steerCapabilitiesDir(asyncDir), `${index}.json`);\n}\n\nexport function steerAcksDir(asyncDir: string, index: number): string {\n\tassertChildIndex(index);\n\treturn path.join(controlInboxDir(asyncDir), STEER_ACKS_DIR, String(index));\n}\n\nfunction steerAckFileName(requestId: string): string {\n\treturn `${Buffer.from(requestId).toString(\"base64url\")}.json`;\n}\n\nexport function steerAckPathFromDir(dir: string, requestId: string): string {\n\tif (!/^[^\\s]+$/.test(requestId) || requestId.length > 256)\n\t\tthrow new Error(\"steer acknowledgment requestId is invalid.\");\n\treturn path.join(dir, steerAckFileName(requestId));\n}\n\nfunction assertChildIndex(index: number): void {\n\tif (!Number.isInteger(index) || index < 0 || index > 1_000_000)\n\t\tthrow new Error(\"steer child index must be a non-negative integer.\");\n}\n\nfunction steerRequestFileName(request: SteerRequest): string {\n\treturn `${String(request.ts).padStart(13, \"0\")}-${Buffer.from(request.id).toString(\"base64url\")}.json`;\n}\n\nfunction validSteerRequest(request: Partial<SteerRequest>): request is SteerRequest {\n\treturn (\n\t\trequest.type === \"steer\" &&\n\t\ttypeof request.id === \"string\" &&\n\t\t/^[^\\s]+$/.test(request.id) &&\n\t\trequest.id.length <= MAX_STEER_REQUEST_ID_LENGTH &&\n\t\ttypeof request.ts === \"number\" &&\n\t\tNumber.isFinite(request.ts) &&\n\t\trequest.ts > 0 &&\n\t\ttypeof request.message === \"string\" &&\n\t\tBoolean(request.message.trim()) &&\n\t\tBuffer.byteLength(request.message, \"utf8\") <= MAX_STEER_MESSAGE_BYTES &&\n\t\t(request.mode === undefined ||\n\t\t\trequest.mode === \"steer\" ||\n\t\t\trequest.mode === \"follow_up\" ||\n\t\t\trequest.mode === \"auto\") &&\n\t\t(request.targetIndex === undefined ||\n\t\t\t(Number.isInteger(request.targetIndex) && request.targetIndex >= 0 && request.targetIndex <= 1_000_000)) &&\n\t\t(request.targetIndexes === undefined ||\n\t\t\t(request.targetIndex === undefined &&\n\t\t\t\tArray.isArray(request.targetIndexes) &&\n\t\t\t\trequest.targetIndexes.length > 0 &&\n\t\t\t\trequest.targetIndexes.length <= 1_000 &&\n\t\t\t\trequest.targetIndexes.every((index) => Number.isInteger(index) && index >= 0 && index <= 1_000_000) &&\n\t\t\t\tnew Set(request.targetIndexes).size === request.targetIndexes.length)) &&\n\t\t(request.source === undefined ||\n\t\t\t(typeof request.source === \"string\" && Boolean(request.source.trim()) && request.source.length <= 256))\n\t);\n}\n\nexport function writeSteerRequestToDir(dir: string, request: SteerRequest): string {\n\tif (!validSteerRequest(request)) throw new Error(\"steer request is malformed or exceeds transport limits.\");\n\tconst requestPath = path.join(dir, steerRequestFileName(request));\n\twriteAtomicJson(requestPath, request);\n\treturn requestPath;\n}\n\nexport function writeSteerCapabilityAt(\n\tfilePath: string,\n\tcapability: Omit<SteerCapability, \"type\" | \"protocolVersion\">,\n): string {\n\tassertChildIndex(capability.index);\n\tif (!Number.isInteger(capability.pid) || capability.pid <= 0)\n\t\tthrow new Error(\"steer capability pid must be a positive integer.\");\n\tif (!Number.isFinite(capability.readyAt) || capability.readyAt <= 0)\n\t\tthrow new Error(\"steer capability readyAt must be a finite timestamp.\");\n\tconst record: SteerCapability = { type: \"steer-capability\", protocolVersion: 1, ...capability };\n\twriteAtomicJson(filePath, record);\n\treturn filePath;\n}\n\nexport function writeSteerCapability(\n\tasyncDir: string,\n\tcapability: Omit<SteerCapability, \"type\" | \"protocolVersion\">,\n): string {\n\treturn writeSteerCapabilityAt(steerCapabilityPath(asyncDir, capability.index), capability);\n}\n\nfunction steerAckWritePath(filePath: string, ack: Omit<SteerAck, \"type\" | \"protocolVersion\">): string {\n\tconst parsed = path.parse(filePath);\n\tconst stateOrder = ack.state === \"queued\" ? \"0\" : ack.state === \"delivered\" ? \"1\" : \"2\";\n\tconst timestamp = String(Math.trunc(ack.ts)).padStart(13, \"0\");\n\tfor (let suffix = 0; suffix < 1_000; suffix += 1) {\n\t\tconst candidate = path.join(\n\t\t\tparsed.dir,\n\t\t\t`${parsed.name}-${timestamp}-${stateOrder}-${ack.state}${suffix === 0 ? \"\" : `-${suffix}`}${parsed.ext}`,\n\t\t);\n\t\tif (!fs.existsSync(candidate)) return candidate;\n\t}\n\tthrow new Error(\"steer acknowledgment queue is full.\");\n}\n\nexport function writeSteerAckAt(filePath: string, ack: Omit<SteerAck, \"type\" | \"protocolVersion\">): string {\n\tassertChildIndex(ack.index);\n\tif (!/^[^\\s]+$/.test(ack.requestId) || ack.requestId.length > 256)\n\t\tthrow new Error(\"steer acknowledgment requestId is invalid.\");\n\tif (!Number.isFinite(ack.ts) || ack.ts <= 0) throw new Error(\"steer acknowledgment ts must be a finite timestamp.\");\n\tif (!ack.message.trim() || ack.message.length > 1000) throw new Error(\"steer acknowledgment message is invalid.\");\n\tconst record: SteerAck = { type: \"steer-ack\", protocolVersion: 1, ...ack, message: ack.message.trim() };\n\tconst ackPath = steerAckWritePath(filePath, ack);\n\twriteAtomicJson(ackPath, record);\n\treturn ackPath;\n}\n\nexport function writeSteerAck(asyncDir: string, ack: Omit<SteerAck, \"type\" | \"protocolVersion\">): string {\n\treturn writeSteerAckAt(path.join(steerAcksDir(asyncDir, ack.index), steerAckFileName(ack.requestId)), ack);\n}\n\n/**\n * Parent side: drop a portable interrupt request the runner's inbox watcher will\n * pick up regardless of OS. Written atomically (temp + rename), dir auto-created.\n */\nexport function requestAsyncInterrupt(\n\tasyncDir: string,\n\tpayload: Omit<InterruptRequest, \"type\"> = {},\n\tdeps: { now?: () => number } = {},\n): string {\n\tconst requestPath = interruptRequestPath(asyncDir);\n\tconst request: InterruptRequest = { ...payload, ts: payload.ts ?? deps.now?.() ?? Date.now(), type: \"interrupt\" };\n\twriteAtomicJson(requestPath, request);\n\treturn requestPath;\n}\n\nexport function requestAsyncTimeout(\n\tasyncDir: string,\n\tpayload: Omit<TimeoutRequest, \"type\"> = {},\n\tdeps: { now?: () => number } = {},\n): string {\n\tconst requestPath = timeoutRequestPath(asyncDir);\n\tconst request: TimeoutRequest = { ...payload, ts: payload.ts ?? deps.now?.() ?? Date.now(), type: \"timeout\" };\n\twriteAtomicJson(requestPath, request);\n\treturn requestPath;\n}\n\nexport function requestAsyncStop(\n\tasyncDir: string,\n\tpayload: Omit<StopRequest, \"type\"> = {},\n\tdeps: { now?: () => number } = {},\n): string {\n\tconst requestPath = stopRequestPath(asyncDir);\n\tconst request: StopRequest = { ...payload, ts: payload.ts ?? deps.now?.() ?? Date.now(), type: \"stop\" };\n\twriteAtomicJson(requestPath, request);\n\treturn requestPath;\n}\n\nexport function requestAsyncCheckpointDecision(\n\tasyncDir: string,\n\ttype: CheckpointDecisionRequest[\"type\"],\n\tpayload: Omit<CheckpointDecisionRequest, \"type\"> = {},\n\tdeps: { now?: () => number } = {},\n): string {\n\tconst requestPath =\n\t\ttype === \"approve-checkpoint\" ? approveCheckpointRequestPath(asyncDir) : rejectCheckpointRequestPath(asyncDir);\n\tconst request: CheckpointDecisionRequest = { ...payload, ts: payload.ts ?? deps.now?.() ?? Date.now(), type };\n\twriteAtomicJson(requestPath, request);\n\treturn requestPath;\n}\n\nexport function requestAsyncSteer(\n\tasyncDir: string,\n\tpayload: {\n\t\tmessage: string;\n\t\tmode?: SteerDeliveryMode;\n\t\ttargetIndex?: number;\n\t\ttargetIndexes?: number[];\n\t\tsource?: string;\n\t\tid?: string;\n\t\tts?: number;\n\t},\n\tdeps: { now?: () => number; randomId?: () => string } = {},\n): string {\n\tconst message = payload.message.trim();\n\tif (!message) throw new Error(\"steer message must not be empty.\");\n\tif (Buffer.byteLength(message, \"utf8\") > MAX_STEER_MESSAGE_BYTES)\n\t\tthrow new Error(`steer message exceeds ${MAX_STEER_MESSAGE_BYTES} UTF-8 bytes.`);\n\tif (\n\t\tpayload.targetIndex !== undefined &&\n\t\t(!Number.isInteger(payload.targetIndex) || payload.targetIndex < 0 || payload.targetIndex > 1_000_000)\n\t) {\n\t\tthrow new Error(\"steer targetIndex must be an integer between 0 and 1000000.\");\n\t}\n\tif (\n\t\tpayload.targetIndexes !== undefined &&\n\t\t(!Array.isArray(payload.targetIndexes) ||\n\t\t\tpayload.targetIndex !== undefined ||\n\t\t\tpayload.targetIndexes.length === 0 ||\n\t\t\tpayload.targetIndexes.length > 1_000 ||\n\t\t\tpayload.targetIndexes.some((index) => !Number.isInteger(index) || index < 0 || index > 1_000_000) ||\n\t\t\tnew Set(payload.targetIndexes).size !== payload.targetIndexes.length)\n\t) {\n\t\tthrow new Error(\n\t\t\t\"steer targetIndexes must contain 1-1000 unique non-negative integers and cannot be combined with targetIndex.\",\n\t\t);\n\t}\n\tconst closedPath = steerInboxClosedPath(asyncDir);\n\tif (fs.existsSync(closedPath)) throw new Error(\"Async run no longer accepts steering requests.\");\n\tconst request: SteerRequest = {\n\t\ttype: \"steer\",\n\t\tid: payload.id ?? deps.randomId?.() ?? randomUUID(),\n\t\tts: payload.ts ?? deps.now?.() ?? Date.now(),\n\t\tmessage,\n\t\t...(payload.mode && payload.mode !== \"steer\" ? { mode: payload.mode } : {}),\n\t\t...(payload.targetIndex !== undefined ? { targetIndex: payload.targetIndex } : {}),\n\t\t...(payload.targetIndexes !== undefined ? { targetIndexes: [...payload.targetIndexes] } : {}),\n\t\t...(payload.source ? { source: payload.source } : {}),\n\t};\n\tconst requestPath = writeSteerRequestToDir(steerRequestsDir(asyncDir), request);\n\tif (fs.existsSync(closedPath)) {\n\t\tfs.rmSync(requestPath, { force: true });\n\t\tthrow new Error(\"Async run stopped accepting steering before the request was committed.\");\n\t}\n\treturn requestPath;\n}\n\nexport function enqueueStepSteer(asyncDir: string, index: number, request: SteerRequest): string {\n\tassertChildIndex(index);\n\tconst { targetIndexes: _targetIndexes, ...singleTargetRequest } = request;\n\treturn writeSteerRequestToDir(stepSteerInboxDir(asyncDir, index), {\n\t\t...singleTargetRequest,\n\t\ttargetIndex: index,\n\t\ttype: \"steer\",\n\t});\n}\n\nfunction parseSteerCapability(raw: unknown): SteerCapability | undefined {\n\tif (!raw || typeof raw !== \"object\" || Array.isArray(raw)) return undefined;\n\tconst input = raw as Partial<SteerCapability>;\n\tif (input.type !== \"steer-capability\" || input.protocolVersion !== 1) return undefined;\n\tconst { index, pid, readyAt, supported } = input;\n\tif (typeof index !== \"number\" || !Number.isInteger(index) || index < 0 || index > 1_000_000) return undefined;\n\tif (\n\t\ttypeof pid !== \"number\" ||\n\t\ttypeof readyAt !== \"number\" ||\n\t\t!Number.isInteger(pid) ||\n\t\tpid <= 0 ||\n\t\t!Number.isFinite(readyAt) ||\n\t\treadyAt <= 0 ||\n\t\ttypeof supported !== \"boolean\"\n\t)\n\t\treturn undefined;\n\treturn { type: \"steer-capability\", protocolVersion: 1, index, pid, readyAt, supported };\n}\n\nfunction parseSteerAck(raw: unknown): SteerAck | undefined {\n\tif (!raw || typeof raw !== \"object\" || Array.isArray(raw)) return undefined;\n\tconst input = raw as Partial<SteerAck>;\n\tif (\n\t\tinput.type !== \"steer-ack\" ||\n\t\tinput.protocolVersion !== 1 ||\n\t\ttypeof input.requestId !== \"string\" ||\n\t\t!/^[^\\s]+$/.test(input.requestId) ||\n\t\tinput.requestId.length > 256\n\t)\n\t\treturn undefined;\n\tconst { index, ts, state, message } = input;\n\tif (\n\t\ttypeof index !== \"number\" ||\n\t\ttypeof ts !== \"number\" ||\n\t\t!Number.isInteger(index) ||\n\t\tindex < 0 ||\n\t\tindex > 1_000_000 ||\n\t\t!Number.isFinite(ts) ||\n\t\tts <= 0\n\t)\n\t\treturn undefined;\n\tif (state !== \"delivered\" && state !== \"queued\" && state !== \"failed\") return undefined;\n\tif (input.deliveryStatus !== undefined && input.deliveryStatus !== \"delivered\" && input.deliveryStatus !== \"queued\")\n\t\treturn undefined;\n\tif (typeof message !== \"string\" || !message.trim() || message.length > 1000) return undefined;\n\treturn {\n\t\ttype: \"steer-ack\",\n\t\tprotocolVersion: 1,\n\t\trequestId: input.requestId,\n\t\tindex,\n\t\tts,\n\t\tstate,\n\t\t...(input.deliveryStatus ? { deliveryStatus: input.deliveryStatus } : {}),\n\t\tmessage: message.trim(),\n\t};\n}\n\nexport function readSteerCapability(asyncDir: string, index: number): SteerCapability | undefined {\n\ttry {\n\t\treturn parseSteerCapability(JSON.parse(fs.readFileSync(steerCapabilityPath(asyncDir, index), \"utf-8\")));\n\t} catch {\n\t\treturn undefined;\n\t}\n}\n\nexport function consumeSteerCapabilities(\n\tasyncDir: string,\n\tfsImpl: Pick<typeof fs, \"existsSync\" | \"readdirSync\" | \"readFileSync\"> = fs,\n): SteerCapability[] {\n\tconst dir = steerCapabilitiesDir(asyncDir);\n\tif (!fsImpl.existsSync(dir)) return [];\n\tconst capabilities: SteerCapability[] = [];\n\tfor (const entry of fsImpl\n\t\t.readdirSync(dir)\n\t\t.filter((name) => /^\\d+\\.json$/.test(name))\n\t\t.sort()) {\n\t\ttry {\n\t\t\tconst capability = parseSteerCapability(JSON.parse(fsImpl.readFileSync(path.join(dir, entry), \"utf-8\")));\n\t\t\tif (capability) capabilities.push(capability);\n\t\t} catch {\n\t\t\t// A partially written or malformed capability is ignored until a valid one arrives.\n\t\t}\n\t}\n\treturn capabilities;\n}\n\nexport function consumeSteerAcks(\n\tasyncDir: string,\n\tfsImpl: Pick<typeof fs, \"existsSync\" | \"readdirSync\" | \"readFileSync\" | \"rmSync\"> = fs,\n): SteerAck[] {\n\tconst root = path.join(controlInboxDir(asyncDir), STEER_ACKS_DIR);\n\tif (!fsImpl.existsSync(root)) return [];\n\tconst acks: SteerAck[] = [];\n\tlet indexNames: string[];\n\ttry {\n\t\tindexNames = fsImpl.readdirSync(root).filter((name) => /^\\d+$/.test(name));\n\t} catch {\n\t\treturn [];\n\t}\n\tfor (const indexName of indexNames) {\n\t\tconst dir = path.join(root, indexName);\n\t\tlet entries: string[];\n\t\ttry {\n\t\t\tentries = fsImpl\n\t\t\t\t.readdirSync(dir)\n\t\t\t\t.filter((name) => name.endsWith(\".json\"))\n\t\t\t\t.sort();\n\t\t} catch {\n\t\t\tcontinue;\n\t\t}\n\t\tfor (const entry of entries) {\n\t\t\tconst target = path.join(dir, entry);\n\t\t\tlet ack: SteerAck | undefined;\n\t\t\ttry {\n\t\t\t\tack = parseSteerAck(JSON.parse(fsImpl.readFileSync(target, \"utf-8\")));\n\t\t\t} catch {\n\t\t\t\tack = undefined;\n\t\t\t}\n\t\t\ttry {\n\t\t\t\tfsImpl.rmSync(target, { force: true });\n\t\t\t} catch {\n\t\t\t\tcontinue;\n\t\t\t}\n\t\t\tif (ack) acks.push(ack);\n\t\t}\n\t}\n\treturn acks;\n}\n\nfunction parseSteerRequest(raw: unknown): SteerRequest | undefined {\n\tif (!raw || typeof raw !== \"object\" || Array.isArray(raw)) return undefined;\n\tconst input = raw as Partial<SteerRequest>;\n\tif (!validSteerRequest(input)) return undefined;\n\treturn {\n\t\ttype: \"steer\",\n\t\tid: input.id.trim(),\n\t\tts: input.ts,\n\t\tmessage: input.message.trim(),\n\t\t...(input.mode ? { mode: input.mode } : {}),\n\t\t...(input.targetIndex !== undefined ? { targetIndex: input.targetIndex } : {}),\n\t\t...(input.targetIndexes !== undefined ? { targetIndexes: [...input.targetIndexes] } : {}),\n\t\t...(typeof input.source === \"string\" && input.source.trim() ? { source: input.source } : {}),\n\t};\n}\n\nexport function consumeSteerRequestsFromDir(\n\tdir: string,\n\tfsImpl: Pick<typeof fs, \"existsSync\" | \"rmSync\" | \"readdirSync\" | \"readFileSync\"> = fs,\n): SteerRequest[] {\n\tif (!fsImpl.existsSync(dir)) return [];\n\tlet entries: string[];\n\ttry {\n\t\tentries = fsImpl\n\t\t\t.readdirSync(dir)\n\t\t\t.filter((name) => name.endsWith(\".json\"))\n\t\t\t.sort();\n\t} catch {\n\t\t// Leave requests in place so the periodic poll can retry the scan.\n\t\treturn [];\n\t}\n\tconst requests: SteerRequest[] = [];\n\tfor (const entry of entries) {\n\t\tconst requestPath = path.join(dir, entry);\n\t\tlet parsed: SteerRequest | undefined;\n\t\ttry {\n\t\t\tparsed = parseSteerRequest(JSON.parse(fsImpl.readFileSync(requestPath, \"utf-8\")));\n\t\t} catch {\n\t\t\tparsed = undefined;\n\t\t}\n\t\ttry {\n\t\t\tfsImpl.rmSync(requestPath, { recursive: true });\n\t\t} catch {\n\t\t\t// Already removed by a concurrent check — do not execute it twice.\n\t\t\tcontinue;\n\t\t}\n\t\tif (parsed) requests.push(parsed);\n\t}\n\treturn requests.sort((left, right) => left.ts - right.ts || left.id.localeCompare(right.id));\n}\n\nexport function consumeSteerRequests(\n\tasyncDir: string,\n\tfsImpl: Pick<typeof fs, \"existsSync\" | \"rmSync\" | \"readdirSync\" | \"readFileSync\"> = fs,\n): SteerRequest[] {\n\treturn consumeSteerRequestsFromDir(steerRequestsDir(asyncDir), fsImpl);\n}\n\nexport function queueRevivalBrief(asyncDir: string, request: SteerRequest): string {\n\tconst dir = path.join(controlInboxDir(asyncDir), REVIVAL_BRIEFS_DIR);\n\tconst queued = fs.existsSync(dir) ? fs.readdirSync(dir).filter((entry) => entry.endsWith(\".json\")).length : 0;\n\tif (queued >= MAX_STEER_QUEUE_SIZE) throw new Error(`Follow-up queue is full (${MAX_STEER_QUEUE_SIZE} messages).`);\n\treturn writeSteerRequestToDir(dir, { ...request, mode: \"follow_up\" });\n}\n\nexport function readRevivalBriefs(asyncDir: string): Array<{ request: SteerRequest; path: string }> {\n\tconst dir = path.join(controlInboxDir(asyncDir), REVIVAL_BRIEFS_DIR);\n\tif (!fs.existsSync(dir)) return [];\n\treturn fs\n\t\t.readdirSync(dir)\n\t\t.filter((entry) => entry.endsWith(\".json\"))\n\t\t.sort()\n\t\t.flatMap((entry) => {\n\t\t\tconst filePath = path.join(dir, entry);\n\t\t\ttry {\n\t\t\t\tconst request = parseSteerRequest(JSON.parse(fs.readFileSync(filePath, \"utf-8\")));\n\t\t\t\treturn request ? [{ request, path: filePath }] : [];\n\t\t\t} catch {\n\t\t\t\treturn [];\n\t\t\t}\n\t\t});\n}\n\n/**\n * Runner side: consume a pending interrupt request. Idempotent — removes the file\n * so each distinct request fires exactly once. Returns whether one was pending.\n */\nexport function consumeInterruptRequest(\n\tasyncDir: string,\n\tfsImpl: Pick<typeof fs, \"existsSync\" | \"rmSync\"> = fs,\n): boolean {\n\tconst requestPath = interruptRequestPath(asyncDir);\n\tif (!fsImpl.existsSync(requestPath)) return false;\n\ttry {\n\t\tfsImpl.rmSync(requestPath, { force: true, recursive: true });\n\t} catch {\n\t\t// Already removed by a concurrent check — still counts as consumed.\n\t}\n\treturn true;\n}\n\nexport function consumeTimeoutRequest(\n\tasyncDir: string,\n\tfsImpl: Pick<typeof fs, \"existsSync\" | \"rmSync\"> = fs,\n): boolean {\n\tconst requestPath = timeoutRequestPath(asyncDir);\n\tif (!fsImpl.existsSync(requestPath)) return false;\n\ttry {\n\t\tfsImpl.rmSync(requestPath, { force: true, recursive: true });\n\t} catch {\n\t\t// Already removed by a concurrent check — still counts as consumed.\n\t}\n\treturn true;\n}\n\nexport function consumeStopRequest(asyncDir: string, fsImpl: Pick<typeof fs, \"existsSync\" | \"rmSync\"> = fs): boolean {\n\tconst requestPath = stopRequestPath(asyncDir);\n\tif (!fsImpl.existsSync(requestPath)) return false;\n\ttry {\n\t\tfsImpl.rmSync(requestPath, { force: true, recursive: true });\n\t} catch {\n\t\t// Already removed by a concurrent check — still counts as consumed.\n\t}\n\treturn true;\n}\n\nexport function consumeCheckpointDecisionRequest(\n\tasyncDir: string,\n\tfsImpl: Pick<typeof fs, \"existsSync\" | \"rmSync\"> = fs,\n): \"approved\" | \"rejected\" | undefined {\n\tif (fsImpl.existsSync(rejectCheckpointRequestPath(asyncDir))) {\n\t\ttry {\n\t\t\tfsImpl.rmSync(rejectCheckpointRequestPath(asyncDir), { force: true, recursive: true });\n\t\t} catch {}\n\t\treturn \"rejected\";\n\t}\n\tif (fsImpl.existsSync(approveCheckpointRequestPath(asyncDir))) {\n\t\ttry {\n\t\t\tfsImpl.rmSync(approveCheckpointRequestPath(asyncDir), { force: true, recursive: true });\n\t\t} catch {}\n\t\treturn \"approved\";\n\t}\n\treturn undefined;\n}\n\n/**\n * Parent side: portable interrupt = authoritative file request + best-effort OS\n * signal. The signal is only a latency optimization on Unix; ENOSYS on Windows\n * is swallowed because the file inbox is authoritative there. Other signal\n * failures are surfaced because they usually mean the runner is not alive to\n * consume the request.\n */\nexport function deliverInterruptRequest(input: {\n\tasyncDir: string;\n\tpid?: number;\n\tkill?: KillFn;\n\tsignal?: NodeJS.Signals;\n\tnow?: () => number;\n\tsource?: string;\n}): void {\n\tconst requestPath = requestAsyncInterrupt(input.asyncDir, input.source ? { source: input.source } : {}, {\n\t\tnow: input.now,\n\t});\n\tif (typeof input.pid === \"number\" && input.pid > 0) {\n\t\ttry {\n\t\t\t(input.kill ?? process.kill)(input.pid, input.signal ?? INTERRUPT_SIGNAL);\n\t\t} catch (error) {\n\t\t\tif ((error as NodeJS.ErrnoException | undefined)?.code === \"ENOSYS\") {\n\t\t\t\t// File inbox is authoritative when custom cross-process signals are unavailable.\n\t\t\t\treturn;\n\t\t\t}\n\t\t\ttry {\n\t\t\t\tfs.rmSync(requestPath, { force: true });\n\t\t\t} catch {\n\t\t\t\t// Best effort cleanup; the caller still gets the signal failure.\n\t\t\t}\n\t\t\tthrow error;\n\t\t}\n\t}\n}\n\nexport function deliverTimeoutRequest(input: {\n\tasyncDir: string;\n\tpid?: number;\n\tkill?: KillFn;\n\tsignal?: NodeJS.Signals;\n\tnow?: () => number;\n\tsource?: string;\n}): void {\n\trequestAsyncTimeout(input.asyncDir, input.source ? { source: input.source } : {}, { now: input.now });\n}\n\nexport function deliverStopRequest(input: {\n\tasyncDir: string;\n\tpid?: number;\n\tkill?: KillFn;\n\tsignal?: NodeJS.Signals;\n\tnow?: () => number;\n\tsource?: string;\n}): void {\n\trequestAsyncStop(input.asyncDir, input.source ? { source: input.source } : {}, { now: input.now });\n}\n\nexport function deliverCheckpointDecisionRequest(input: {\n\tasyncDir: string;\n\tdecision: \"approved\" | \"rejected\";\n\tnow?: () => number;\n\tsource?: string;\n\treason?: string;\n}): void {\n\trequestAsyncCheckpointDecision(\n\t\tinput.asyncDir,\n\t\tinput.decision === \"approved\" ? \"approve-checkpoint\" : \"reject-checkpoint\",\n\t\t{ ...(input.source ? { source: input.source } : {}), ...(input.reason ? { reason: input.reason } : {}) },\n\t\t{ now: input.now },\n\t);\n}\n\n/**\n * Runner side: watch the control inbox and route interrupt requests into\n * `onInterrupt`. Uses `fs.watch` when available plus an interval poll as a\n * portable safety net (covers filesystems/platforms where `fs.watch` is\n * unreliable). Fires once per distinct request. Returns a disposer.\n */\nexport function watchAsyncControlInbox(\n\tasyncDir: string,\n\topts: {\n\t\tonInterrupt: () => void;\n\t\tonTimeout?: () => void;\n\t\tonStop?: () => void;\n\t\tonSteer?: (request: SteerRequest) => void;\n\t\tonCheckpointDecision?: (decision: \"approved\" | \"rejected\") => void;\n\t\tonSteerCapability?: (capability: SteerCapability) => void;\n\t\tonSteerAck?: (ack: SteerAck) => void;\n\t\tpollIntervalMs?: number;\n\t\tfs?: ControlChannelFs;\n\t\ttimers?: ControlChannelTimers;\n\t},\n): () => void {\n\tconst fsImpl = opts.fs ?? fs;\n\tconst timers = opts.timers ?? { setInterval, clearInterval };\n\tconst dir = controlInboxDir(asyncDir);\n\ttry {\n\t\tfsImpl.mkdirSync(dir, { recursive: true });\n\t} catch {\n\t\t// Best effort — the poll/watch below tolerates a missing dir.\n\t}\n\n\tlet disposed = false;\n\tconst check = (): void => {\n\t\tif (disposed) return;\n\t\ttry {\n\t\t\tif (consumeStopRequest(asyncDir, fsImpl)) opts.onStop?.();\n\t\t\tif (consumeTimeoutRequest(asyncDir, fsImpl)) opts.onTimeout?.();\n\t\t\tif (consumeInterruptRequest(asyncDir, fsImpl)) opts.onInterrupt();\n\t\t\tconst checkpointDecision = consumeCheckpointDecisionRequest(asyncDir, fsImpl);\n\t\t\tif (checkpointDecision) opts.onCheckpointDecision?.(checkpointDecision);\n\t\t\tfor (const request of consumeSteerRequests(asyncDir, fsImpl)) opts.onSteer?.(request);\n\t\t\tfor (const capability of consumeSteerCapabilities(asyncDir, fsImpl)) opts.onSteerCapability?.(capability);\n\t\t\tfor (const ack of consumeSteerAcks(asyncDir, fsImpl)) opts.onSteerAck?.(ack);\n\t\t} catch {\n\t\t\t// Never let inbox errors crash the runner.\n\t\t}\n\t};\n\n\t// Handle a request that may have arrived before the watcher started.\n\tcheck();\n\n\tlet watcher: fs.FSWatcher | undefined;\n\ttry {\n\t\twatcher = fsImpl.watch(resolveWatchPath(dir, fsImpl.realpathSync.native), () => check());\n\t\twatcher.on?.(\"error\", () => {\n\t\t\t// fs.watch can emit on transient FS errors; the interval poll keeps us live.\n\t\t});\n\t} catch {\n\t\twatcher = undefined;\n\t}\n\n\tconst interval = timers.setInterval(check, opts.pollIntervalMs ?? POLL_INTERVAL_MS);\n\tinterval.unref?.();\n\n\treturn () => {\n\t\tif (disposed) return;\n\t\tdisposed = true;\n\t\ttry {\n\t\t\twatcher?.close();\n\t\t} catch {\n\t\t\t// ignore\n\t\t}\n\t\ttimers.clearInterval(interval);\n\t};\n}\n"]}