Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
20 commits
Select commit Hold shift + click to select a range
9f35455
fix: recover stale run coordination locks
RobertTLange Aug 29, 2026
8132ddc
fix: retry Windows run state replacement
RobertTLange Aug 29, 2026
dacc0c9
ci: test run storage on Windows
RobertTLange Aug 29, 2026
d596cf3
test: cover run state replacement defaults
RobertTLange Aug 29, 2026
530a566
test: pass Windows lock path through environment
RobertTLange Aug 29, 2026
fc8ddb4
test: preserve Windows replacement failures
RobertTLange Aug 29, 2026
cbf13f9
fix: extend Windows state replacement retries
RobertTLange Aug 29, 2026
0cafee6
fix: preserve live node lock owners
RobertTLange Aug 29, 2026
f1da717
fix: retain async agent lock ownership
RobertTLange Aug 29, 2026
37a0d67
fix: roll back failed async message starts
RobertTLange Aug 29, 2026
8df0d75
fix: fail closed on process probe ambiguity
RobertTLange Aug 29, 2026
46cb02b
fix: preserve Windows descendants after PID reuse
RobertTLange Aug 29, 2026
70ce60e
ci: run process probes on Windows
RobertTLange Aug 29, 2026
ea5d952
fix: repair interrupted async message starts
RobertTLange Aug 29, 2026
73113db
fix: pass Windows probe PIDs reliably
RobertTLange Aug 29, 2026
a3e84f0
fix: ignore Windows probe process descendants
RobertTLange Aug 29, 2026
134eb8b
test: isolate Windows PID reuse fixture
RobertTLange Aug 29, 2026
a9ee571
docs: update changelog for run lock recovery
RobertTLange Aug 30, 2026
8fbde30
fix: tolerate slow Windows process identity probes
RobertTLange Aug 30, 2026
ca588ed
fix: allow slow Windows owner identity capture
RobertTLange Aug 30, 2026
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
29 changes: 29 additions & 0 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -44,6 +44,35 @@ jobs:
npm pack --pack-destination /tmp
npx -y --package /tmp/roberttlange-headless-*.tgz headless --help

windows-run-storage:
name: Windows run storage / Node ${{ matrix.node-version }}
runs-on: windows-latest
strategy:
fail-fast: false
matrix:
node-version:
- 22
- 24

steps:
- name: Checkout
uses: actions/checkout@v4

- name: Setup Node
uses: actions/setup-node@v4
with:
node-version: ${{ matrix.node-version }}
cache: npm

- name: Install dependencies
run: npm ci

- name: Build
run: npm run build

- name: Test run storage
run: node --import tsx --test tests/run-storage.test.ts tests/run-storage-process.test.ts

python:
name: Python ${{ matrix.python-version }}
runs-on: ubuntu-latest
Expand Down
2 changes: 2 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,8 @@

## TBD

- Fixed run coordination to recover stale run and node locks after owner crashes, retain async ownership for the full detached process tree, and handle Windows process probing and state replacement safely (#28).

## 0.6.1 - 2026-08-12

- Fixed Modal runs to use the immutable, verified v0.6.0 runtime image digest instead of a stale cached `latest` tag.
Expand Down
207 changes: 207 additions & 0 deletions src/async-run-message.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,207 @@
import { execFileSync, spawn, type ChildProcess } from "node:child_process";

import { windowsTaskkillPath } from "./process-tree.js";
import {
handoffNodeStoreLockOwner,
nodeStoreLockHasLiveSuccessor,
removeRunStoreLockSignalListeners,
} from "./run-storage.js";
import {
acquireNodeLock,
appendNodeLog,
nodeLockPath,
updateNodeStatus,
} from "./runs.js";
import type { Env } from "./types.js";

export interface AsyncRunMessageTask {
command: { command: string; args: string[] };
cwd?: string;
runId: string;
nodeId: string;
}

export type AsyncRunMessageRequest =
| { type: "task"; task: AsyncRunMessageTask }
| { type: "start" }
| { type: "cancel" };

export type AsyncRunMessageResponse =
| { type: "ready" }
| { type: "started" }
| { type: "error"; message: string };

let task: AsyncRunMessageTask | undefined;
let releaseLock: (() => void) | undefined;
let activeChild: ChildProcess | undefined;
let started = false;
let finished = false;

if (process.send) {
process.on("message", (message: AsyncRunMessageRequest) => {
void handleRequest(message);
});
process.once("disconnect", () => {
if (!started && !finished) failBeforeStart("async message parent disconnected before agent startup");
});
for (const signal of forwardedSignals()) {
process.on(signal, () => terminateActiveChild(signal));
}
}

async function handleRequest(request: AsyncRunMessageRequest): Promise<void> {
if (request.type === "task") {
await prepare(request.task);
return;
}
if (request.type === "cancel") {
if (started) {
terminateActiveChild("SIGTERM");
} else {
finish(1);
}
return;
}
if (request.type === "start" && task && releaseLock && !started) {
started = true;
await execute(task);
}
}

function terminateActiveChild(signal: NodeJS.Signals): void {
if (!activeChild) return;
if (process.platform === "win32" && activeChild.pid) {
try {
execFileSync(windowsTaskkillPath(), ["/pid", String(activeChild.pid), "/T", "/F"], {
stdio: "ignore",
timeout: 3_000,
windowsHide: true,
});
return;
} catch {
return;
}
}
activeChild.kill(signal);
}

function failBeforeStart(message: string): void {
if (finished) return;
if (task) {
try {
updateNodeStatus(process.env as Env, task.runId, task.nodeId, "failed", message);
} catch {
// Lock cleanup must still run when status persistence fails.
}
}
finish(1);
}

async function prepare(nextTask: AsyncRunMessageTask): Promise<void> {
if (task || finished) return;
try {
releaseLock = acquireNodeLock(process.env as Env, nextTask.runId, nextTask.nodeId, {
processTreeRootPid: process.pid,
});
removeRunStoreLockSignalListeners();
task = nextTask;
send({ type: "ready" });
} catch (error) {
send({ type: "error", message: errorMessage(error) });
finish(2);
}
}

async function execute(currentTask: AsyncRunMessageTask): Promise<void> {
let code = 1;
try {
code = await runChild(currentTask);
if (code !== 0) {
appendNodeLog(
process.env as Env,
currentTask.runId,
currentTask.nodeId,
"stderr",
`async child exited with code ${code}\n`,
);
}
updateNodeStatus(
process.env as Env,
currentTask.runId,
currentTask.nodeId,
code === 0 ? "idle" : "failed",
);
} catch (error) {
const message = errorMessage(error);
appendNodeLog(process.env as Env, currentTask.runId, currentTask.nodeId, "stderr", `${message}\n`);
updateNodeStatus(process.env as Env, currentTask.runId, currentTask.nodeId, "failed", message);
} finally {
finish(code);
}
}

function runChild(currentTask: AsyncRunMessageTask): Promise<number> {
return new Promise((resolve, reject) => {
const child = spawn(currentTask.command.command, currentTask.command.args, {
cwd: currentTask.cwd,
detached: process.platform !== "win32",
env: {
...process.env,
HEADLESS_ASYNC_MESSAGE_OWNER_PID: String(process.pid),
HEADLESS_ASYNC_MESSAGE_WORKER: "1",
},
stdio: ["ignore", "ignore", "inherit"],
});
activeChild = child;
child.once("error", reject);
child.once("spawn", () => {
try {
if (!child.pid) throw new Error("async message CLI started without a process ID");
handoffNodeStoreLockOwner(
nodeLockPath(process.env as Env, currentTask.runId, currentTask.nodeId),
process.pid,
{ processTreeRootPid: child.pid },
);
send({ type: "started" });
} catch (error) {
terminateActiveChild("SIGKILL");
reject(error);
}
});
child.once("close", (code, signal) => {
activeChild = undefined;
resolve(signal ? 1 : (code ?? 1));
});
});
}

function finish(code = 0): void {
if (finished) return;
finished = true;
try {
if (!task || !nodeStoreLockHasLiveSuccessor(
nodeLockPath(process.env as Env, task.runId, task.nodeId),
process.pid,
)) {
releaseLock?.();
}
} finally {
releaseLock = undefined;
if (process.connected) process.disconnect();
process.exitCode = code;
}
}

function send(response: AsyncRunMessageResponse): void {
if (process.send) process.send(response);
}

function errorMessage(error: unknown): string {
return error instanceof Error ? error.message : String(error);
}

function forwardedSignals(): NodeJS.Signals[] {
return process.platform === "win32"
? ["SIGINT", "SIGTERM", "SIGBREAK"]
: ["SIGHUP", "SIGINT", "SIGTERM", "SIGQUIT"];
}
46 changes: 40 additions & 6 deletions src/cli.ts
Original file line number Diff line number Diff line change
Expand Up @@ -102,9 +102,11 @@ import { handleCronCommand as handleCronCommandImpl, type CronCommand } from "./
import { runCronDaemon } from "./cron.js";
import { extractRunNodeMetrics } from "./run-metrics.js";
import { createRunStatusReporter, parseRunStatusIntervalMs } from "./run-status.js";
import { isRunStoreLockSignalListener, waitForNodeStoreLockOwner } from "./run-storage.js";
import {
appendNodeLog,
completeIdleRunNodes,
nodeLockPath,
readRun,
registerNode,
runDirectory,
Expand Down Expand Up @@ -1887,14 +1889,25 @@ async function executeCommand(
result.stdoutEndsWithNewline = stdoutEndsWithNewline;
resolve(result);
};
const ownsChildProcessGroup =
process.platform !== "win32" &&
(options.timeoutSeconds !== undefined || options.cleanupBeforeParentSignalExit !== undefined);
const handlesParentSignals = ownsChildProcessGroup || options.cleanupBeforeParentSignalExit !== undefined;
const asyncMessageWorker = env.HEADLESS_ASYNC_MESSAGE_WORKER === "1";
const ownsChildProcessGroup = !asyncMessageWorker && process.platform !== "win32" && (
options.timeoutSeconds !== undefined
|| options.cleanupBeforeParentSignalExit !== undefined
);
const handlesParentSignals = asyncMessageWorker
|| ownsChildProcessGroup
|| options.cleanupBeforeParentSignalExit !== undefined;
waitForAsyncMessageOwnership(env, command);
let childEnv = commandEnv(env, command);
if (asyncMessageWorker) {
childEnv = { ...childEnv };
delete childEnv.HEADLESS_ASYNC_MESSAGE_OWNER_PID;
delete childEnv.HEADLESS_ASYNC_MESSAGE_WORKER;
}
const child = spawn(command.command, command.args, {
cwd,
detached: ownsChildProcessGroup,
env: commandEnv(env, command) as NodeJS.ProcessEnv,
env: childEnv as NodeJS.ProcessEnv,
stdio,
});

Expand Down Expand Up @@ -3225,6 +3238,23 @@ function withRunEnvironment(command: BuiltCommand, runId: string | undefined, no
};
}

function waitForAsyncMessageOwnership(env: Env, command: BuiltCommand): void {
const ownerValue = env.HEADLESS_ASYNC_MESSAGE_OWNER_PID;
if (ownerValue === undefined) return;
const expectedProcessTreeRootPid = Number(ownerValue);
const runId = command.env?.HEADLESS_RUN_ID;
const nodeId = command.env?.HEADLESS_RUN_NODE;
if (
!Number.isSafeInteger(expectedProcessTreeRootPid)
|| expectedProcessTreeRootPid <= 0
|| !runId
|| !nodeId
) {
throw new Error("invalid async message ownership context");
}
waitForNodeStoreLockOwner(nodeLockPath(env, runId, nodeId), process.pid);
}

async function executeStoredNode(
node: {
agent: AgentName;
Expand Down Expand Up @@ -3370,7 +3400,11 @@ export async function runCli(argv: string[], deps: CliDeps = {}): Promise<number
const inheritedSignalListeners = new Map(
parentExitSignals().map((signal) => [
signal,
new Set(process.listeners(signal).filter((listener) => !isDockerSessionLockSignalListener(listener))),
new Set(
process.listeners(signal).filter(
(listener) => !isDockerSessionLockSignalListener(listener) && !isRunStoreLockSignalListener(listener),
),
),
] as const),
);
let registeredRunNode: { runId: string; nodeId: string } | undefined;
Expand Down
Loading
Loading