diff --git a/docs/WORKFLOWS.md b/docs/WORKFLOWS.md index b0837827..e65b5a97 100644 --- a/docs/WORKFLOWS.md +++ b/docs/WORKFLOWS.md @@ -1010,6 +1010,36 @@ resume` takes a new generation and reruns only work after the last durable - Server status reports safe counts and timestamps. It does not report session IDs, project paths, prompts, payloads, tokens, process IDs, or credentials. +## Wake recovery + +The server survives laptop suspend and resume instead of restarting on every +wake: + +- The authenticated owner re-arms its own expired lease. Lease renewal matches + the recorded epoch, server id, token hash, pid, and process start identity, + so a lease that expired while the machine slept re-arms on the first + post-wake heartbeat and the server keeps serving. The server stops only when + its claim was genuinely superseded, meaning another server took the epoch. +- Takeover requires the recorded holder to be provably not serving. A starting + server probes the holder's socket with a bounded hello check before it may + remove the server lock file. A holder that answers keeps the lock, and the + new server exits with its `already running` message. A dead holder, or a live + holder that is mid-shutdown or frozen and does not answer, loses the lock and + the new server takes over. The epoch claim row stays the fencing authority: + a probe misfire costs one wasted start, never two serving servers. When the + superseded holder's process later exits, the operating system unlinks the + socket path its own listener bound, which is the path the replacement now + serves; the replacement re-creates the socket file within one poll tick + while its claim still names it, and clients retry through the brief window. +- Clients re-spawn failed replacements. When a spawned server exits without + becoming ready, the client starts the next one within the same start window, + capped at three attempts, so one lost race costs milliseconds instead of the + whole window. + +`expires_at` in the server claim answers whether the lease is currently valid +for takeover arbitration. It does not answer whether the owner process is +alive; ownership is proven by identity and token. + ## Run history retention Pi Workflows keeps a terminal root run and all its explicit restart descendants for 30 days from `finished_at`. The server can remove the tree after that point only when every descendant is terminal and no protected work or outside reference remains. diff --git a/docs/plans/2026-09-27-wake-recovery-plan.md b/docs/plans/2026-09-27-wake-recovery-plan.md new file mode 100644 index 00000000..6be0800f --- /dev/null +++ b/docs/plans/2026-09-27-wake-recovery-plan.md @@ -0,0 +1,232 @@ +--- +title: wake recovery plan +author: Onur Solmaz <2453968+osolmaz@users.noreply.github.com> +date: 2026-09-27 +--- + +# wake recovery plan + +## Purpose + +After macOS suspend and resume, every open Pi session logs +`Warning: Workflow server is unavailable: Workflow server did not become ready: +A workflow server is already running with PID X`, with a different PID on each +line across wake cycles. The goal is that after any sleep duration and any wake +pattern, the system recovers to one healthy server quickly, with at most one +transient reconnect per session and no `already running` warning bursts. A live +healthy server must never be killed or displaced, and a dead or shutdown-bound +holder must never block takeover for more than about a second. + +## Observed failure + +The failure is a chain, verified against the 0.17.4 code and reproduced in an +isolated sandbox. + +1. The server holds a 30 second lease in the `workflow_server_state` table, + renewed every 10 seconds. `renewServer()` in `src/server/state.ts` requires + `expires_at > now`, so after any sleep longer than the lease the first + heartbeat after thaw throws `Pi Workflows server claim lost`, and the + heartbeat catch in `src/server/server.ts` calls `stop()`. The server kills + itself on every wake. +2. During `stop()` the listener closes first, while the server lock file + `server.lock.json` is released only at the very end. `stop()` awaits runner + and channel supervisors (SIGTERM plus a 2 second kill grace each). If the + machine re-sleeps mid-shutdown, the process stays alive but deaf while + holding the lock for a long time. +3. On connect failure the client's `ensureAvailable()` in + `src/client/client.ts` spawns exactly one detached replacement and then only + retries connecting for 10 seconds. The replacement's `acquireServerLock()` in + `src/server/server.ts` sees a live process whose start identity matches the + lock record, writes `A workflow server is already running with PID X` to its + startup pipe, and exits. The client never spawns a second replacement within + that window, so every session that raced warns after its full 10 second + timeout. Freezing a live lock holder reproduced the exact warning after + 12.3 seconds. +4. Each wake or DarkWake cycle mints a new server generation, so warnings + accumulate as a chain of distinct PIDs, and the extension shows one warning + per session per outage. + +## Approach + +Implement layered wake recovery as one hard-cutover change with three +independent parts: the authenticated lease owner re-arms its own expired lease, +so a server that slept keeps serving; `acquireServerLock` probes the recorded +holder's socket and only reports `already running` for a provably serving +holder, with the epoch claim row as the fencing authority; and the client +re-spawns a replacement when a spawned child exits without becoming ready, +bounded by the existing 10 second deadline and a three-spawn cap. No schema or +protocol changes; the lock and lease schemas keep their versioned identifiers +and shapes. + +## Implementation steps + +1. **Suspend-aware lease self-renewal.** In `src/server/state.ts`, + `renewServer()` currently requires `expires_at > now`. Drop the + `AND expires_at > ?` predicate from the renewal UPDATE so it matches only on + `id=1`, `epoch`, `server_id`, `token_hash`, `pid`, and + `process_start_identity`. The token hash and process start identity remain + the fence, and a superseded epoch (another server claimed) or a wrong token + still fails the update and triggers the existing stop path. Update the + function's doc comment: expiry gates takeover, never kills the authenticated + owner. + Verify with unit tests: (a) rewind `expires_at` into the past on a row that + still matches identity, then renew succeeds and extends `expires_at`; (b) + after another server claims the epoch, renew throws and the caller stops; + (c) renew with a wrong token throws; (d) renew after `releaseServer` + (`server_id` null) throws. +2. **Heartbeat claim-lost semantics.** In `src/server/server.ts`, inside + `startTimers()`, keep the heartbeat catch semantics but reword the log and + comment from "lease expired" to "claim superseded": the catch now fires only + when the epoch row no longer carries this server's exact identity and token. + No behavioral change beyond step 1. + Verify with a unit test on the heartbeat path: an expired-but-unsuperseded + claim renews (no stop); a superseded claim stops the server and logs the + supersession reason. +3. **Serving-probe lock takeover.** In `src/server/lock.ts`, add a bounded + serving probe and use it in lock acquisition. New exported helper + `probeServerServing(socketPath, timeoutMs)`: `net.connect` to the socket + path; serving means hello data arrives within the timeout (the server writes + its hello on connection); not serving means a connect error (ENOENT, + ECONNREFUSED) or timeout; always destroy the probe socket. Extend + `acquireServerLock(lockPath, record, options)` to take + `{ socketPath, probeTimeoutMs?, probe? }` where the probe defaults to + `probeServerServing` and stays injectable for unit tests. New acquisition + logic: lock file missing or holder identity mismatched, remove and take over + (unchanged); holder identity matches and the probe says serving, throw the + existing `A workflow server is already running with PID X`; holder identity + matches and the probe says not serving, remove the lock and take over. Pass + `{ socketPath: this.socketPath }` at the `acquireServerLock` call site in + the server start path. A probe misfire during a holder's startup gap is + fenced by `acquireServer`'s existing + `A live Pi Workflows server already owns epoch N` check, so the worst case + is one wasted spawn, never two serving servers. + Verify with unit tests using an injected probe: a serving holder keeps the + `already running` error; a deaf holder is taken over; a dead holder is taken + over without probing. Verify with an integration test on a real temp unix + socket: a server that writes hello yields verdict serving; a bound socket + that accepts but stays silent yields takeover; a missing socket path yields + immediate takeover with no probe. +4. **Client re-spawn on dead replacements.** In `src/client/client.ts`, rework + `ensureAvailable()`: keep the overall `START_TIMEOUT_MS` deadline; wrap the + spawn in an outer loop capped at three spawn attempts; inside each attempt, + poll `connect()` every 50 ms as today but stop waiting on the current child + once it has exited (the startup failure resolved) plus a short 250 ms grace + for a competing holder to become connectable, then loop to spawn again; + after the cap, keep polling connect until the deadline because another + session's spawn may win. Preserve the final error contract: + `Workflow server did not become ready: ` followed by the last spawn's + diagnostic when any spawn ran, otherwise the last connect error. + Verify with integration tests using real spawned servers in temp dirs: (a) a + first spawn that exits with a diagnostic (server entry path pointed at a + failing stub) and a second real entry that becomes ready, all inside one + `ensureAvailable` call well under the deadline; (b) every spawn fails and + the thrown error carries the last diagnostic; (c) a serving holder reachable + by connect short-circuits without extra spawns. +5. **Extension notification path.** In `src/extension/index.ts` (the outage + notification path), no code change is expected: the warning text stays + accurate because `ensureAvailable`'s thrown message contract is preserved. + Only touch this file if the final error wording changes during + implementation; otherwise leave it alone. + Verify by grep: the notify message still composes from `errorMessage(error)` + with no format change, and the final change set has no diff in this file. +6. **Documentation.** In `docs/WORKFLOWS.md`, add a `Wake recovery` section + near the server lifecycle documentation covering the three behaviors: the + server re-arms its own expired lease after suspend and only stops when its + claim is genuinely superseded; takeover requires the recorded holder to be + provably not serving (bounded socket probe) with the epoch claim as the + fence; clients re-spawn failed replacements within a bounded budget. State + that `expires_at` answers "is the lease currently valid", not "is the owner + alive". + Verify that `npm run check` passes and the section reads as plain prose + consistent with the document's existing structure. +7. **Full repository validation.** From the repository root, `npm run check` + passes; `npx slophammer-ts@latest dry .` reports zero candidates; + `npx slophammer-ts@latest check . --only +ts.dependency-boundaries-required` passes; `TMPDIR=/tmp/e2e npm run +test:e2e` passes all files. The default macOS `TMPDIR` is too long for e2e + unix sockets and is a known environmental failure on this machine. + +## Contracts + +- `pi-workflows.server-lock.v1` keeps its exact record shape (`schema`, `pid`, + `startIdentity`, `serverId`); only acquisition semantics change (probe before + the `already running` verdict), so existing readers stay compatible. +- The `workflow_server_state` table keeps its schema; no migration. The refined + meaning: `expires_at` answers "is the lease currently valid for takeover + arbitration", not "is the owner process alive"; `serverStatus().live` keeps + its `expires_at`-based definition unchanged. +- `renewServer` keeps its signature and throw behavior + (`Pi Workflows server claim lost`) so the heartbeat catch and all callers + behave identically on genuine supersession. +- `acquireServerLock` gains an options parameter (`socketPath` required at the + server call site); it is internal to the server layer, not a published + interface, and the thrown `already running` message is unchanged for serving + holders. +- `WorkflowClient.ensureAvailable` keeps its promise contract (resolves with + hello or throws with the `Workflow server did not become ready: ` prefix); + internals may spawn up to three servers instead of one. +- The NDJSON client protocol and all message schemas are unchanged; no version + bumps are needed. +- Compatibility boundary: a 0.17.4 client against the new server, and the new + client against a 0.17.4 server, continue to work. Old clients may still spawn + losing replacements, which the new lock probe converts into fast takeovers, + and old servers keep their self-destruct behavior until adopted. + +## Tests + +- `renewServer`: an expired-but-identity-matching re-arm succeeds and extends + the lease; a superseded epoch fails; a wrong token fails; a released row + (null `server_id`) fails. +- Heartbeat: a superseded claim stops the server; an expired-but-ours claim + does not stop it. +- `acquireServerLock`: a missing lock takes over without probing; a dead holder + (identity mismatch) takes over without probing; a live serving holder keeps + the `already running` error; a live deaf holder is taken over. +- `probeServerServing` integration: a hello-writing server is serving; a silent + bound socket is not serving; a missing socket path is not serving without + connecting. +- `ensureAvailable`: a failing first spawn followed by a healthy second spawn + succeeds inside one call under the deadline; every spawn failing throws with + the last diagnostic; a serving holder short-circuits without extra spawns. +- Existing fencing and takeover tests pass unchanged in intent, and the e2e + suite passes with `TMPDIR=/tmp/e2e`. + +## Risks + +- A serving holder with a temporarily blocked event loop fails the 500 ms hello + probe and is robbed of the lock. The epoch claim table fences the outcome: + the robber fails `acquireServer` while the holder's lease is live, and if the + robber claims first, the holder's next renewal fails and it exits. The worst + case is one wasted spawn, never two serving servers. +- Dropping `expires_at > now` from renewal weakens assumptions elsewhere that + an expired lease means a dead owner. `acquireServer`'s takeover predicate is + unchanged; audit every `expiresAt` reader (`serverStatus`, recovery, views) + during implementation and keep their `expires_at`-based meanings, and add the + documentation note that `expires_at` is lease validity, not liveness. +- A handover window can briefly have two listening sockets (the old server + until its next heartbeat, the new server after takeover). This is a + pre-existing characteristic of lease-granularity fencing, unchanged by this + plan; the old server stops at its first failed renewal, and clients reaching + either server get valid responses until then. +- Incident paths spawn up to three node processes per `ensureAvailable` call + across several sessions. The three-spawn cap and the shared 10 second deadline + bound this; spawns are cheap losers that exit in milliseconds under the lock + and claim fencing. +- The probe adds latency to contended spawns. Only the contended live-holder + case probes; the common uncontended path (missing lock or dead holder) never + probes, keeping startup latency unchanged. + +## Boundaries + +- No changes to the OnurPi adoption or wrapper (`~/repos/onurpi`), including its + pinned dependency and sync flow; adoption is a separate follow-up. +- No npm publishing, version bumping, or release mechanics. +- No changes to Pi core (`@earendil-works/pi-coding-agent`), including its + tool-parameter transport. +- No changes to the Herdr plugin or Herdr-specific surfaces. +- No changes to resource-manager or decision-channel subsystems beyond what the + lease or lock fix directly touches (their leases live in the separate + `leases` table and are untouched). +- No new runtime dependencies, no new persisted state, no schema or protocol + version changes, and no work in any repository other than + `/Users/onur/repos/pi-workflows`. diff --git a/src/client/client.ts b/src/client/client.ts index d3f6c33b..4af22c82 100644 --- a/src/client/client.ts +++ b/src/client/client.ts @@ -35,6 +35,10 @@ const RECONNECT_BASE_DELAY_MS = 250; const RECONNECT_MAX_DELAY_MS = 10_000; /** After this many failed attempts the client reports a blocker and stops looping. */ const RECONNECT_MAX_ATTEMPTS = 12; +/** Spawn attempts within one start window; each loser exits under the lock fence. */ +const MAX_START_SPAWN_ATTEMPTS = 3; +/** After a spawned child exits, keep polling this long for a competing holder. */ +const SPAWN_EXIT_GRACE_MS = 250; const RESOLVER_TIMEOUT_MS = 30_000; const CLIENT_PACKAGE_VERSION = runtimePackageVersion(); @@ -239,6 +243,14 @@ export class WorkflowClient { return true; } + /** + * Start the server if needed and wait for a working connection. One spawn + * that exits without becoming ready, for example a replacement that lost the + * lock race to a holder that then died, must not consume the whole start + * window, so the loop spawns again within the deadline and a small attempt + * cap. The server lock and the epoch claim fence concurrent spawns, so a + * losing spawn exits in milliseconds. + */ async ensureAvailable(): Promise { try { return await this.connect(); @@ -248,20 +260,35 @@ export class WorkflowClient { // child that can never serve it. assertSocketPathSupported(this.endpoint); this.resetConnection(); - const startupFailure = this.startDetached(); - let startupError: Error | undefined; - void startupFailure.then((error) => { - startupError = error; - }); const deadline = Date.now() + START_TIMEOUT_MS; let lastError: unknown; - while (Date.now() < deadline) { - await delay(50); - try { - return await this.connect(); - } catch (error) { - lastError = error; - this.resetConnection(); + let startupError: Error | undefined; + let spawnAttempts = 0; + spawn: while (Date.now() < deadline) { + let exitedAt: number | undefined; + if (spawnAttempts < MAX_START_SPAWN_ATTEMPTS) { + spawnAttempts += 1; + const startupFailure = this.startDetached(); + void startupFailure.then((error) => { + startupError = error; + exitedAt = Date.now(); + }); + } + while (Date.now() < deadline) { + await delay(50); + try { + return await this.connect(); + } catch (error) { + lastError = error; + this.resetConnection(); + } + if ( + exitedAt !== undefined && + Date.now() >= exitedAt + SPAWN_EXIT_GRACE_MS && + spawnAttempts < MAX_START_SPAWN_ATTEMPTS + ) { + continue spawn; + } } } throw new Error( diff --git a/src/server/lock.ts b/src/server/lock.ts index 12730dcc..f73f5427 100644 --- a/src/server/lock.ts +++ b/src/server/lock.ts @@ -1,4 +1,5 @@ import fs from "node:fs"; +import net from "node:net"; import path from "node:path"; import { killProcessGroup, matchesProcessIdentity, type ProcessIdentity } from "./processes.js"; @@ -68,6 +69,94 @@ export function writeServerLock(lockPath: string, record: ServerLockRecord): voi }); } +/** One bounded socket probe: does a server accept connections here and answer? */ +export type ServerLockProbe = (socketPath: string, timeoutMs: number) => Promise; + +/** + * Ask the socket whether a server is really serving. Serving means the connect + * succeeds and the server's hello arrives within the timeout; a connect error or + * a silent window means the holder is alive but not serving, as during shutdown + * or while frozen. The probe socket is always destroyed. + */ +export async function probeServerServing(socketPath: string, timeoutMs: number): Promise { + return await new Promise((resolve) => { + const socket = net.connect(socketPath); + let settled = false; + const timer = setTimeout(() => finish(false), timeoutMs); + timer.unref?.(); + const finish = (serving: boolean): void => { + if (settled) return; + settled = true; + clearTimeout(timer); + socket.destroy(); + resolve(serving); + }; + socket.once("data", () => finish(true)); + socket.once("error", () => finish(false)); + socket.once("close", () => finish(false)); + }); +} + +/** + * Take the exclusive server lock for a starting server. A lock whose recorded + * process is gone, or whose process is alive but provably not serving, is stale + * and is removed. A holder that answers on the socket keeps the lock, and the + * caller must not start. Fencing against a second live server stays with the + * epoch claim, not with this file: a probe misfire costs one wasted start, and + * the loser exits when its claim fails. + * + * The probe awaits, so the lock record can change while it runs. After a + * not-serving verdict the record is re-read, and the lock is only removed when + * it still names the same silent holder; a record another starter already + * replaced makes this starter re-read and converge instead of clobbering the + * newer takeover. Creation stays exclusive (`wx`), so a racing creator loses + * and re-reads. Concurrent starters can therefore only displace a holder that + * is provably not serving, and the epoch claim fences whichever of them wins. + * + * The returned `displaced` record is the live holder this start replaced, so a + * starter whose claim is then fenced can restore the holder's lock file + * instead of leaving the serving holder without one. + */ +export async function acquireServerLock( + lockPath: string, + record: { pid: number; startIdentity: string; serverId: string }, + options: { socketPath: string; probeTimeoutMs?: number; probe?: ServerLockProbe }, +): Promise<{ displaced: ServerLockRecord | undefined }> { + const probe = options.probe ?? probeServerServing; + const timeoutMs = options.probeTimeoutMs ?? 500; + let displaced: ServerLockRecord | undefined; + for (let attempt = 0; attempt < 4; attempt += 1) { + const existing = readServerLock(lockPath); + if (existing !== undefined && matchesProcessIdentity(existing)) { + if (await probe(options.socketPath, timeoutMs)) { + throw new Error(`A workflow server is already running with PID ${existing.pid}`); + } + const current = readServerLock(lockPath); + if ( + current !== undefined && + matchesProcessIdentity(current) && + (current.pid !== existing.pid || + current.startIdentity !== existing.startIdentity || + current.serverId !== existing.serverId) + ) { + continue; + } + displaced = existing; + } else if (existing === undefined) { + displaced = undefined; + } + fs.rmSync(lockPath, { force: true }); + try { + writeServerLock(lockPath, record); + return { displaced }; + } catch (error) { + if ((error as NodeJS.ErrnoException)?.code !== "EEXIST") throw error; + // Another starter created the lock first; re-read and converge. + } + } + throw new Error(`Could not acquire the workflow server lock at ${lockPath}`); +} + /** * Stop the server named by the lock file when the recorded process still has its start * identity. A client whose package version differs from the running server cannot send a diff --git a/src/server/server.ts b/src/server/server.ts index e4a9ec5e..06bea17b 100644 --- a/src/server/server.ts +++ b/src/server/server.ts @@ -109,13 +109,14 @@ import { type ChannelEffectRecord, } from "./channel-effects.js"; import { ChannelAdapterSupervisor } from "./channel-supervisor.js"; -import { isServerLockRecord, writeServerLock } from "./lock.js"; import { - ServerProcessRegistry, - matchesProcessIdentity, - processParentPid, - processStartIdentity, -} from "./processes.js"; + acquireServerLock, + isServerLockRecord, + readServerLock, + writeServerLock, + type ServerLockRecord, +} from "./lock.js"; +import { ServerProcessRegistry, processParentPid, processStartIdentity } from "./processes.js"; import type { ResourceRunnerLaunchEnvelope, ResourceRunnerMessage, @@ -275,6 +276,10 @@ export class WorkflowServer { private readonly sockets = new Set(); private readonly connections = new Map(); private server: net.Server | null = null; + /** The socket path this server actually bound, or null while it has not. */ + private boundSocketPath: string | null = null; + /** The live holder's lock record this start displaced, if any. */ + private displacedLock: ServerLockRecord | undefined; private claim: ServerClaim | null = null; private heartbeatTimer: ReturnType | null = null; private pollTimer: ReturnType | null = null; @@ -292,6 +297,7 @@ export class WorkflowServer { private automaticStatePruneTimer: ReturnType | null = null; private automaticStatePruneScheduled = false; private automaticStatePruneDue = true; + private rebindingSocket = false; private lastAutomaticStatePruneAt: number | null = null; private nextAutomaticStatePruneAttemptAt = 0; private nextTerminalMessageReconciliationAt = 0; @@ -361,7 +367,12 @@ export class WorkflowServer { if (startIdentity === undefined) { throw new Error("Cannot attest the workflow server process start identity"); } - acquireServerLock(this.lockPath, { pid: process.pid, startIdentity, serverId: this.serverId }); + const { displaced } = await acquireServerLock( + this.lockPath, + { pid: process.pid, startIdentity, serverId: this.serverId }, + { socketPath: this.socketPath }, + ); + this.displacedLock = displaced; try { const reaped = this.registry.reapOrphans(); if (reaped.length > 0) this.log(`reaped ${reaped.length} exact orphan process(es)`); @@ -379,6 +390,7 @@ export class WorkflowServer { await this.listen(); this.startTimers(); this.started = true; + this.displacedLock = undefined; this.log(`ready on ${this.socketPath} at epoch ${this.claim.epoch}`); this.requestAutomaticStatePrune(); void this.expireTimedOutDecision().finally(() => this.resumeAutomaticStatePruneIfDue()); @@ -392,15 +404,17 @@ export class WorkflowServer { const server = this.server; this.server = null; await this.closeServer(server); + let released = false; if (this.claim !== null) { try { - this.serverState.releaseServer(this.claim); + released = this.serverState.releaseServer(this.claim); } catch { // Preserve the startup error when cleanup cannot release an already-lost claim. } } this.claim = null; - if (process.platform !== "win32") fs.rmSync(this.socketPath, { force: true }); + if (released) this.removeBoundSocket(); + else this.restoreDisplacedLock(); this.releaseLock(); throw error; } @@ -459,9 +473,19 @@ export class WorkflowServer { await Promise.allSettled([this.automaticStatePruneTask]); } this.registry.killAll(); - if (this.claim !== null) this.serverState.releaseServer(this.claim); - this.claim = null; - if (process.platform !== "win32") fs.rmSync(this.socketPath, { force: true }); + // The explicit socket-file removal is gated on releasing this server's own + // claim, and closing the listener is gated the same way inside + // closeServer, because closing unlinks the bound path. + let released = false; + if (this.claim !== null) { + try { + released = this.serverState.releaseServer(this.claim); + } catch (error) { + this.log(`claim release failed during stop: ${errorMessage(error)}`); + } + this.claim = null; + } + if (released) this.removeBoundSocket(); this.releaseLock(); this.runStore.close(); this.queue.close(); @@ -470,12 +494,49 @@ export class WorkflowServer { this.started = false; } + /** + * Re-create the socket file when a superseded predecessor's exit removed it: + * a process that exits unlinks the socket path its own listener bound, and + * after a takeover that path belongs to this server. The claim row must + * still name this server, or a newer server owns the path now. + */ + private ensureSocketFile(): void { + if (process.platform === "win32" || this.rebindingSocket) return; + if (fs.existsSync(this.socketPath)) return; + const row = this.serverState.serverStatus(); + if (row.serverId !== this.serverId || (this.claim !== null && row.epoch !== this.claim.epoch)) { + return; + } + this.rebindingSocket = true; + this.log("the socket file disappeared; rebinding"); + const previous = this.server; + this.server = null; + if (previous !== null) { + try { + previous.close(() => undefined); + } catch { + // The handle is already closed. + } + } + void this.listen() + .then(() => this.log(`socket rebound on ${this.socketPath}`)) + .catch((error) => { + this.rebindingSocket = false; + this.log(`socket rebind failed: ${errorMessage(error)}`); + void this.stop(); + }); + } + private async closeServer(server: net.Server | null): Promise { const closed = new Promise((resolve) => { if (server === null || !server.listening) { resolve(); return; } + // Closing the listener unlinks the bound path. After a takeover that + // path belongs to the replacement, so the unlink is transient: this + // server's poll re-creates the socket file while its claim still names + // it (ensureSocketFile), and clients retry through the brief window. try { server.close(() => resolve()); } catch { @@ -528,6 +589,9 @@ export class WorkflowServer { } } } catch (error) { + // Renewal fails only when the epoch row no longer carries this server's + // exact identity and token. A lease that expired while the process was + // frozen, as across a laptop suspend, re-arms instead of stopping. this.log(`server claim lost: ${errorMessage(error)}`); void this.stop(); } @@ -535,6 +599,7 @@ export class WorkflowServer { this.heartbeatTimer.unref?.(); this.pollTimer = setInterval(() => { try { + this.ensureSocketFile(); this.serverState.syncActiveTime(); this.recovery.sample(); this.reconcileSubmissionReminders(); @@ -555,6 +620,7 @@ export class WorkflowServer { } private async listen(): Promise { + this.rebindingSocket = false; if (process.platform !== "win32") fs.rmSync(this.socketPath, { force: true }); const server = net.createServer((socket) => this.handleConnection(socket)); this.server = server; @@ -583,6 +649,39 @@ export class WorkflowServer { }); server.on("error", (error) => this.log(`socket error: ${errorMessage(error)}`)); if (process.platform !== "win32") fs.chmodSync(this.socketPath, 0o600); + // Only a server that bound the socket may ever remove its file: a claim + // loser that cleans up after a fenced start must leave the winner's socket + // alone. + this.boundSocketPath = this.socketPath; + } + + /** + * Put back the lock record a fenced start displaced: the holder that kept + * serving must keep its lock file, or the version-mismatch stop path and + * later recovery lose track of it. A record another starter already wrote + * stays in place. + */ + private restoreDisplacedLock(): void { + const displaced = this.displacedLock; + if (displaced === undefined) return; + this.displacedLock = undefined; + const current = readServerLock(this.lockPath); + if (current !== undefined && current.serverId !== this.serverId) return; + fs.rmSync(this.lockPath, { force: true }); + try { + writeServerLock(this.lockPath, displaced); + } catch { + // A racing writer recreated the lock; its record is authoritative. + } + } + + /** Remove the socket file only when this server bound it. */ + private removeBoundSocket(): void { + if (process.platform === "win32") return; + if (this.boundSocketPath !== this.socketPath) return; + // DIAGNOSTIC + fs.rmSync(this.socketPath, { force: true }); + this.boundSocketPath = null; } private handleConnection(socket: Socket): void { @@ -5122,23 +5221,6 @@ function channelEventPayload(message: ChannelAdapterMessage): JsonValue { return payload as unknown as JsonValue; } -function acquireServerLock( - lockPath: string, - record: { pid: number; startIdentity: string; serverId: string }, -): void { - try { - const existing = JSON.parse(fs.readFileSync(lockPath, "utf8")) as unknown; - if (isServerLockRecord(existing) && matchesProcessIdentity(existing)) { - throw new Error(`A workflow server is already running with PID ${existing.pid}`); - } - fs.rmSync(lockPath, { force: true }); - } catch (error) { - if (error instanceof Error && error.message.startsWith("A workflow server is already")) - throw error; - } - writeServerLock(lockPath, record); -} - function runnerRunCommand( record: WorkflowRunQueueRecord, resumeInteractionAttemptId: string | undefined, diff --git a/src/server/state.ts b/src/server/state.ts index 05c80e63..39e50e94 100644 --- a/src/server/state.ts +++ b/src/server/state.ts @@ -151,6 +151,15 @@ export class ServerStateStore { }); } + /** + * Re-arm the lease of the authenticated owner. The update matches the full + * recorded identity (epoch, server id, token hash, pid, and process start + * identity) instead of requiring an unexpired lease: expiry gates takeover, + * it never kills the owner, so a lease that lapsed while the process was + * frozen, as across a laptop suspend, re-arms instead of stopping the + * server. A row another server claimed no longer carries this identity, so + * the update fails and the caller stops exactly as before. + */ renewServer(claim: ServerClaim, leaseMs: number, now: number = Date.now()): ServerClaim { requireLeaseMs(leaseMs); return this.state.transaction(() => { @@ -159,7 +168,7 @@ export class ServerStateStore { .prepare( `UPDATE workflow_server_state SET heartbeat_at = ?, expires_at = ? WHERE id = 1 AND epoch = ? AND server_id = ? AND token_hash = ? - AND pid = ? AND process_start_identity = ? AND expires_at > ?`, + AND pid = ? AND process_start_identity = ?`, ) .run( now, @@ -169,7 +178,6 @@ export class ServerStateStore { tokenHash(claim.token), claim.pid, claim.processStartIdentity, - now, ); if (changed.changes !== 1) throw new Error("Pi Workflows server claim lost"); return { ...claim, expiresAt }; diff --git a/test/client.test.ts b/test/client.test.ts index 80b7f27b..9f43d1a7 100644 --- a/test/client.test.ts +++ b/test/client.test.ts @@ -15,6 +15,7 @@ import { parseClientMessage, type ClientRequest, } from "../src/client/protocol.js"; +import { serverLockPath, stopRecordedServer } from "../src/server/lock.js"; import { WorkflowServer } from "../src/server/server.js"; import type { JsonValue } from "../src/state/json.js"; import { makeTempDir, waitUntil } from "./helpers.js"; @@ -147,6 +148,46 @@ describe("WorkflowClient", () => { } }); + it("re-spawns when the first replacement exits without becoming ready", async () => { + const databasePath = path.join(await makeTempDir("client-respawn"), "state.sqlite"); + const client = new WorkflowClient({ databasePath }); + const start = vi + .spyOn(client as unknown as { startDetached: () => Promise }, "startDetached") + // The first replacement loses the lock race and exits with a diagnostic. + .mockImplementationOnce(() => + Promise.resolve(new Error("A workflow server is already running with PID 4242")), + ); + try { + const hello = await client.ensureAvailable(); + expect(hello.type).toBe("hello"); + expect(start).toHaveBeenCalledTimes(2); + } finally { + await client.close(); + await stopRecordedServer(serverLockPath(databasePath)); + } + }, 30_000); + + it("caps re-spawn attempts and reports the last spawn diagnostic", async () => { + const databasePath = path.join(await makeTempDir("client-respawn-cap"), "state.sqlite"); + const client = new WorkflowClient({ databasePath }); + const start = vi + .spyOn(client as unknown as { startDetached: () => Promise }, "startDetached") + .mockResolvedValue(new Error("spawn lost the lock race again")); + vi.spyOn(client, "connect").mockRejectedValue(new Error("connect ENOENT server.sock")); + vi.useFakeTimers(); + try { + const unavailable = expect(client.ensureAvailable()).rejects.toThrow( + "Workflow server did not become ready: spawn lost the lock race again", + ); + await vi.advanceTimersByTimeAsync(10_050); + await unavailable; + expect(start).toHaveBeenCalledTimes(3); + } finally { + vi.useRealTimers(); + await client.close(); + } + }); + it("retries one durable invocation with the same idempotency key", async () => { const databasePath = path.join(await makeTempDir("client-durable-retry"), "state.sqlite"); const client = new WorkflowClient({ databasePath }); diff --git a/test/server-lock.test.ts b/test/server-lock.test.ts index dfec138f..4e64d552 100644 --- a/test/server-lock.test.ts +++ b/test/server-lock.test.ts @@ -1,8 +1,12 @@ import { spawn, type ChildProcess } from "node:child_process"; +import { once } from "node:events"; import fs from "node:fs"; +import net from "node:net"; import path from "node:path"; -import { describe, expect, it } from "vitest"; +import { describe, expect, it, vi } from "vitest"; import { + acquireServerLock, + probeServerServing, readServerLock, serverLockPath, stopRecordedServer, @@ -51,6 +55,18 @@ function recordFor( return { pid, startIdentity, serverId }; } +/** The test process itself: always alive, so its lock record always matches. */ +function selfRecord(serverId: string): { + pid: number; + startIdentity: string; + serverId: string; +} { + const startIdentity = processStartIdentity(process.pid); + if (startIdentity === undefined) + throw new Error(`The test process has no start identity: ${process.pid}`); + return { pid: process.pid, startIdentity, serverId }; +} + function killIfAlive(child: ChildProcess): void { const pid = child.pid; if (pid === undefined) return; @@ -163,3 +179,207 @@ describe("workflow server lock file", () => { } }); }); + +describe("server lock acquisition", () => { + it("takes over a missing lock or a dead holder without probing", async () => { + const { lockPath } = await lockDirectory("pw-lock-acquire-missing"); + const probe = vi.fn(); + const record = { pid: 1, startIdentity: "platform-start:0", serverId: "server-new" }; + const missing = await acquireServerLock(lockPath, record, { + socketPath: path.join("/tmp", "unused.sock"), + probe, + }); + expect(missing.displaced).toBeUndefined(); + expect(probe).not.toHaveBeenCalled(); + expect(readServerLock(lockPath)).toEqual(record); + + const exited = spawn(process.execPath, ["-e", ""], { stdio: "ignore" }); + await new Promise((resolve) => exited.once("exit", resolve)); + fs.rmSync(lockPath, { force: true }); + writeServerLock(lockPath, { + pid: exited.pid as number, + startIdentity: "platform-start:0", + serverId: "server-gone", + }); + const dead = await acquireServerLock(lockPath, record, { + socketPath: path.join("/tmp", "unused.sock"), + probe, + }); + expect(dead.displaced).toBeUndefined(); + expect(probe).not.toHaveBeenCalled(); + expect(readServerLock(lockPath)).toEqual(record); + }); + + it("keeps the already-running error for a holder that answers on the socket", async () => { + const { lockPath } = await lockDirectory("pw-lock-acquire-serving"); + const child = await startIdleProcess(); + try { + writeServerLock(lockPath, recordFor(child, "server-serving")); + await expect( + acquireServerLock(lockPath, recordFor(child, "server-next"), { + socketPath: path.join("/tmp", "unused.sock"), + probe: async () => true, + }), + ).rejects.toThrow(/already running with PID/); + expect(readServerLock(lockPath)?.serverId).toBe("server-serving"); + } finally { + killIfAlive(child); + } + }); + + it("takes over a live holder that is provably not serving", async () => { + const { lockPath } = await lockDirectory("pw-lock-acquire-deaf"); + const child = await startIdleProcess(); + try { + const holderRecord = recordFor(child, "server-shutdown-bound"); + writeServerLock(lockPath, holderRecord); + const probe = vi.fn(async () => false); + const record = recordFor(child, "server-next"); + const taken = await acquireServerLock(lockPath, record, { + socketPath: path.join("/tmp", "unused.sock"), + probe, + }); + expect(taken.displaced).toEqual(holderRecord); + expect(probe).toHaveBeenCalledOnce(); + expect(readServerLock(lockPath)).toEqual(record); + } finally { + killIfAlive(child); + } + }); + + it("revalidates the holder after the probe and keeps a newer serving starter", async () => { + const { lockPath } = await lockDirectory("pw-lock-revalidate"); + const holder = await startIdleProcess(); + try { + const holderRecord = recordFor(holder, "server-holder"); + const starterRecord = selfRecord("server-starter"); + writeServerLock(lockPath, holderRecord); + // A competing starter takes over the silent holder while this probe + // awaits, and answers on the socket before the caller re-reads. + const probe = vi.fn(async () => { + fs.rmSync(lockPath, { force: true }); + writeServerLock(lockPath, starterRecord); + return probe.mock.calls.length > 1; + }); + await expect( + acquireServerLock(lockPath, selfRecord("server-next"), { + socketPath: path.join("/tmp", "unused.sock"), + probe, + }), + ).rejects.toThrow(new RegExp(`already running with PID ${process.pid}`)); + expect(probe).toHaveBeenCalledTimes(2); + expect(readServerLock(lockPath)?.serverId).toBe("server-starter"); + } finally { + killIfAlive(holder); + } + }); + + it("converges when another starter replaces the holder during the probe", async () => { + const { lockPath } = await lockDirectory("pw-lock-converge"); + const holder = await startIdleProcess(); + try { + const holderRecord = recordFor(holder, "server-holder"); + const starterRecord = selfRecord("server-starter"); + writeServerLock(lockPath, holderRecord); + // The first probe sees the silent holder while a competing starter takes + // over; the second probe sees the new holder still silent, so the caller + // converges on the newest silent holder instead of clobbering blindly. + const probe = vi.fn(async () => { + if (readServerLock(lockPath)?.serverId === "server-holder") { + fs.rmSync(lockPath, { force: true }); + writeServerLock(lockPath, starterRecord); + } + return false; + }); + await acquireServerLock(lockPath, selfRecord("server-next"), { + socketPath: path.join("/tmp", "unused.sock"), + probe, + }); + expect(probe).toHaveBeenCalledTimes(2); + expect(readServerLock(lockPath)?.serverId).toBe("server-next"); + } finally { + killIfAlive(holder); + } + }); +}); + +describe("server serving probe", () => { + it("reports a hello-writing server as serving", async () => { + const socketPath = path.join(await makeTempDir("pw-probe-hello"), "s.sock"); + const server = net.createServer((socket) => { + socket.end('{"type":"hello"}\n'); + }); + server.listen(socketPath); + await once(server, "listening"); + try { + expect(await probeServerServing(socketPath, 1_000)).toBe(true); + } finally { + server.close(); + } + }); + + it("reports a silent bound socket as not serving", async () => { + const socketPath = path.join(await makeTempDir("pw-probe-silent"), "s.sock"); + const server = net.createServer(() => undefined); + server.listen(socketPath); + await once(server, "listening"); + try { + expect(await probeServerServing(socketPath, 300)).toBe(false); + } finally { + server.close(); + } + }); + + it("reports a missing socket path as not serving without connecting", async () => { + const socketPath = path.join(await makeTempDir("pw-probe-missing"), "gone.sock"); + expect(await probeServerServing(socketPath, 300)).toBe(false); + }); + + it("takes over a real deaf holder that still owns the lock file", async () => { + const directory = await makeTempDir("pw-probe-takeover"); + const lockPath = path.join(directory, "server.lock.json"); + const socketPath = path.join(directory, "s.sock"); + // The test process itself is the live holder: its start identity matches, + // and its bound socket accepts connections but never answers. + const holder = selfRecord("server-deaf"); + writeServerLock(lockPath, holder); + const server = net.createServer(() => undefined); + server.listen(socketPath); + await once(server, "listening"); + try { + const record = { ...holder, serverId: "server-next" }; + await acquireServerLock(lockPath, record, { socketPath, probeTimeoutMs: 300 }); + expect(readServerLock(lockPath)).toEqual(record); + } finally { + server.close(); + } + }); + + it("keeps the lock error against a real hello-writing holder", async () => { + const directory = await makeTempDir("pw-probe-serving"); + const lockPath = path.join(directory, "server.lock.json"); + const socketPath = path.join(directory, "s.sock"); + const holder = selfRecord("server-real"); + writeServerLock(lockPath, holder); + const server = net.createServer((socket) => { + socket.end('{"type":"hello"}\n'); + }); + server.listen(socketPath); + await once(server, "listening"); + try { + await expect( + acquireServerLock( + lockPath, + { ...holder, serverId: "server-next" }, + { + socketPath, + probeTimeoutMs: 1_000, + }, + ), + ).rejects.toThrow(/already running with PID/); + expect(readServerLock(lockPath)?.serverId).toBe("server-real"); + } finally { + server.close(); + } + }); +}); diff --git a/test/server-protocol-state.test.ts b/test/server-protocol-state.test.ts index 7af3b647..cbe5126e 100644 --- a/test/server-protocol-state.test.ts +++ b/test/server-protocol-state.test.ts @@ -82,7 +82,7 @@ describe("server protocol", () => { }); describe("server durable state", () => { - it("fences server epochs and never revives an expired server", async () => { + it("fences server epochs and never revives a superseded claim", async () => { const { server, queue } = await fixture(); const first = server.acquireServer({ serverId: "server-1", @@ -117,6 +117,55 @@ describe("server durable state", () => { queue.close(); }); + it("re-arms an expired lease while the row still carries the owner's identity", async () => { + const { server, queue } = await fixture(); + const claim = server.acquireServer({ + serverId: "server-1", + pid: 100, + processStartIdentity: "start-1", + leaseMs: 1_000, + now: 1_000, + }); + // A suspend let wall-clock time pass without a heartbeat: the lease + // expired, but no other server claimed the epoch. + server.state.connection + .prepare("UPDATE workflow_server_state SET expires_at = 500 WHERE id = 1") + .run(); + expect(server.renewServer(claim, 1_000, 10_000).expiresAt).toBe(11_000); + expect( + server.state.connection + .prepare("SELECT heartbeat_at, expires_at FROM workflow_server_state WHERE id = 1") + .get(), + ).toEqual({ heartbeat_at: 10_000, expires_at: 11_000 }); + server.close(); + queue.close(); + }); + + it("refuses renewal for a released or wrongly authenticated claim", async () => { + const { server, queue } = await fixture(); + const claim = server.acquireServer({ + serverId: "server-1", + pid: 100, + processStartIdentity: "start-1", + leaseMs: 1_000, + now: 1_000, + }); + expect(server.releaseServer(claim, 1_100)).toBe(true); + expect(() => server.renewServer(claim, 1_000, 1_200)).toThrow(/claim lost/); + const other = server.acquireServer({ + serverId: "server-2", + pid: 200, + processStartIdentity: "start-2", + leaseMs: 1_000, + now: 1_300, + }); + expect(() => + server.renewServer({ ...other, token: "not-the-real-token" }, 1_000, 1_400), + ).toThrow(/claim lost/); + server.close(); + queue.close(); + }); + it("stores command receipts and rejects request identity reuse", async () => { const { server, queue } = await fixture(); let calls = 0; diff --git a/test/server.test.ts b/test/server.test.ts index b5c7aef5..cf229bce 100644 --- a/test/server.test.ts +++ b/test/server.test.ts @@ -16,6 +16,7 @@ import { type ClientRequest, type ClientResponse, } from "../src/client/protocol.js"; +import { readServerLock, serverLockPath } from "../src/server/lock.js"; import { ServerProcessRegistry } from "../src/server/processes.js"; import { WorkflowServer } from "../src/server/server.js"; import { ServerStateStore } from "../src/server/state.js"; @@ -3735,6 +3736,215 @@ export { default } from ${JSON.stringify(path.resolve("examples/workflows/echo.w } }, 45_000); + /** Rewind or extend the server claim lease, retrying through SQLite busy windows. */ + async function updateServerLease(store: ServerStateStore, expiresAt: number): Promise { + const statement = store.state.connection.prepare( + "UPDATE workflow_server_state SET expires_at = ? WHERE id = 1", + ); + for (let attempt = 0; attempt < 5; attempt += 1) { + try { + statement.run(expiresAt); + return; + } catch (error) { + if ((error as NodeJS.ErrnoException)?.code !== "SQLITE_BUSY") throw error; + await new Promise((resolve) => setTimeout(resolve, 200)); + } + } + throw new Error("the server claim row stayed busy"); + } + + it("keeps serving when the lease expired while the process was frozen", async () => { + const databasePath = path.join(await makeTempDir("server-suspend-renew"), "state.sqlite"); + const server = new WorkflowServer({ databasePath, serverRenewMs: 100, serverLeaseMs: 500 }); + const client = new WorkflowClient({ databasePath }); + await server.start(); + const store = new ServerStateStore(databasePath); + try { + await waitUntil(() => store.serverStatus().live, 10_000); + // A suspend let wall-clock time pass without a heartbeat: rewind the lease + // into the past without changing ownership. + await updateServerLease(store, Date.now() - 10_000); + // The next heartbeat re-arms the expired-but-ours lease instead of + // stopping the server. + await waitUntil(() => store.serverStatus().live, 10_000); + const status = await client.request({ operation: "server.status" }); + expect(status.outcome).toBe("accepted"); + } finally { + store.close(); + await client.close(); + await server.stop(); + } + }, 30_000); + + it("keeps the serving socket when a fenced starter loses the claim", async () => { + const databasePath = path.join(await makeTempDir("server-fenced-starter"), "state.sqlite"); + const socketPath = path.join(path.dirname(serverLockPath(databasePath)), "server.sock"); + const client = new WorkflowClient({ databasePath }); + let holderPid: number | undefined; + try { + // Spawn the real holder and connect to it. + await client.ensureAvailable(); + const holder = readServerLock(serverLockPath(databasePath)); + holderPid = holder?.pid; + expect(typeof holderPid).toBe("number"); + // Extend the holder's lease while it is still running, then freeze it: + // a second starter passes the lock probe (the frozen holder does not + // answer) but must lose the claim. + const store = new ServerStateStore(databasePath); + await updateServerLease(store, Date.now() + 60_000); + process.kill(holderPid as number, "SIGSTOP"); + const second = new WorkflowServer({ databasePath }); + await expect(second.start()).rejects.toThrow(/live Pi Workflows server/); + // The fenced starter restores the holder's lock record. + expect(readServerLock(serverLockPath(databasePath))?.pid).toBe(holderPid); + // The fenced starter must leave the winner's bound socket alone. + expect(existsSync(socketPath)).toBe(true); + // The frozen holder wakes and keeps serving on its original socket. + process.kill(holderPid as number, "SIGCONT"); + const deadline = Date.now() + 15_000; + for (;;) { + let serving = false; + try { + await client.request({ operation: "server.status" }); + serving = true; + } catch { + serving = false; + } + if (serving) break; + if (Date.now() > deadline) throw new Error("the holder did not resume serving"); + await new Promise((resolve) => setTimeout(resolve, 50)); + } + } finally { + if (holderPid !== undefined) { + try { + process.kill(holderPid, "SIGCONT"); + } catch { + // The holder already exited. + } + try { + process.kill(holderPid, "SIGTERM"); + } catch { + // The holder already exited. + } + } + await client.close(); + } + }, 45_000); + + it("keeps the replacement serving when the superseded holder stops", async () => { + const databasePath = path.join(await makeTempDir("server-handover-socket"), "state.sqlite"); + const socketPath = path.join(path.dirname(serverLockPath(databasePath)), "server.sock"); + const client = new WorkflowClient({ databasePath }); + let holderPid: number | undefined; + let replacement: WorkflowServer | null = null; + try { + // Spawn the holder and connect to it. + await client.ensureAvailable(); + holderPid = readServerLock(serverLockPath(databasePath))?.pid; + expect(typeof holderPid).toBe("number"); + // Rewind the holder's lease while it is still running: a frozen process + // can hold the SQLite write lock, and the lease self-renewal re-arms an + // expired lease within one heartbeat, so freeze immediately after. + const store = new ServerStateStore(databasePath); + await updateServerLease(store, Date.now() - 10_000); + // Freeze the holder with an expired lease so the replacement takes over + // the claim and rebinds the socket path. + process.kill(holderPid as number, "SIGSTOP"); + replacement = new WorkflowServer({ databasePath, serverRenewMs: 100 }); + await replacement.start(); + // Wake the holder: its renewal fails against the new epoch and it stops. + process.kill(holderPid as number, "SIGCONT"); + const stopDeadline = Date.now() + 25_000; + for (;;) { + let alive = true; + try { + process.kill(holderPid as number, 0); + } catch { + alive = false; + } + if (!alive) break; + if (Date.now() > stopDeadline) throw new Error("the superseded holder did not stop"); + await new Promise((resolve) => setTimeout(resolve, 100)); + } + // The superseded holder's exit unlinks the socket path its own listener + // bound; the replacement re-creates the file within one poll tick and + // keeps serving on it. + const restoreDeadline = Date.now() + 10_000; + for (;;) { + if (existsSync(socketPath)) break; + if (Date.now() > restoreDeadline) { + throw new Error("the replacement did not restore the socket file"); + } + await new Promise((resolve) => setTimeout(resolve, 50)); + } + const deadline = Date.now() + 15_000; + for (;;) { + let serving = false; + try { + const status = await client.request({ operation: "server.status" }); + serving = status.outcome === "accepted"; + } catch { + serving = false; + } + if (serving) break; + if (Date.now() > deadline) throw new Error("the replacement did not keep serving"); + await new Promise((resolve) => setTimeout(resolve, 50)); + } + } finally { + if (holderPid !== undefined) { + try { + process.kill(holderPid, "SIGCONT"); + } catch { + // The holder already exited. + } + try { + process.kill(holderPid, "SIGTERM"); + } catch { + // The holder already exited. + } + } + if (replacement !== null) await replacement.stop(); + await client.close(); + } + }, 60_000); + + it("stops when another server claims the superseded epoch", async () => { + const databasePath = path.join(await makeTempDir("server-superseded-stop"), "state.sqlite"); + const server = new WorkflowServer({ databasePath, serverRenewMs: 100, serverLeaseMs: 500 }); + const client = new WorkflowClient({ databasePath }); + await server.start(); + const store = new ServerStateStore(databasePath); + try { + await waitUntil(() => store.serverStatus().live, 10_000); + await updateServerLease(store, Date.now() - 10_000); + const takeover = store.acquireServer({ + serverId: "server-takeover", + pid: 1, + processStartIdentity: "start-takeover", + leaseMs: 30_000, + }); + expect(takeover.epoch).toBeGreaterThan(1); + // The frozen server's next renewal fails against the new epoch, and the + // claim-lost path stops it. + const deadline = Date.now() + 10_000; + for (;;) { + let stopped = false; + try { + await client.request({ operation: "server.status" }); + } catch { + stopped = true; + } + if (stopped) break; + if (Date.now() > deadline) throw new Error("waitUntil timed out"); + await new Promise((resolve) => setTimeout(resolve, 50)); + } + } finally { + store.close(); + await client.close(); + await server.stop(); + } + }, 30_000); + it("retries an interrupted idempotent effect and adopts its durable reservation", async () => { const cwd = await makeTempDir("server-idempotent-effect-project"); const databasePath = path.join(