From e90d9da4928baebb14c4279ad909047e15d7bb1f Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Sun, 30 Aug 2026 06:06:53 -0700 Subject: [PATCH 1/3] Add tests for restore-timeout cancellation withRestoreTimeout races a boot-restore attempt against a timer and, on timeout, abandons the still-running attempt outright. If restoreDeploymentFromRecord later succeeds anyway, spawnWorkflowDeployment registers a live supervisor for an address the boot loop already recorded as a restore failure -- a genuinely live deployment whose durable record says it failed to restore. This asserts the fix ahead of the implementation: a restore that outlasts its timeout must still finish (never abandoned), and once it does, its record must be corrected rather than left claiming a live deployment failed. createHandshakeSpawner is extracted out of makeLifecycleFixture so this test can drive a real spawn/ready handshake on its own timing instead of one that throws or resolves synchronously; the fixture also grows optional materializeDeploymentClosure and restoreAttemptTimeoutMs overrides so a boot-restore test can use a multi-step closure and a millisecond-scale timeout instead of the real 30s default. --- .../support/workflow-lifecycle-fixture.ts | 78 ++++++--- .../workflow-restore-timeout-cancel.test.ts | 155 ++++++++++++++++++ 2 files changed, 210 insertions(+), 23 deletions(-) create mode 100644 apps/sidecar/test/workflow-restore-timeout-cancel.test.ts diff --git a/apps/sidecar/test/support/workflow-lifecycle-fixture.ts b/apps/sidecar/test/support/workflow-lifecycle-fixture.ts index 976915ccb..fafd4c8d2 100644 --- a/apps/sidecar/test/support/workflow-lifecycle-fixture.ts +++ b/apps/sidecar/test/support/workflow-lifecycle-fixture.ts @@ -195,23 +195,17 @@ export type Fixture = { multistepCredentialsRouter: MultistepCredentialsRouter; }; -export async function makeLifecycleFixture(opts?: { - /** - * Reuse a prior fixture's data dir, so a test can construct a SECOND - * router against the SAME on-disk state to model a sidecar process - * restart (boot-time restore). - */ - dataDir?: string; - /** - * Override the default in-memory fake key store with a real one (e.g. - * `@intx/hub-agent`'s `createAgentKeyStore` bound to the same - * `dataDir`), for a test that needs the actual on-disk key-persistence - * behavior the fake bypasses entirely. - */ - keyStore?: Parameters[0]["keyStore"]; -}): Promise { - const spawns: Spawn[] = []; - const spawner: SubprocessSpawner = ({ env }) => { +/** + * A `SubprocessSpawner` whose child never signals `ready` on its own -- + * `answerReadyHandshake` drives that half explicitly -- so a caller + * controls exactly when a mocked `wired.supervisor.spawn()` call resolves. + * Exported (not just used internally by `makeLifecycleFixture`) so a test + * that builds its own router directly (rather than through the fixture) + * can still exercise a real spawn/ready handshake instead of a spawner + * that throws or resolves synchronously. + */ +export function createHandshakeSpawner(spawns: Spawn[]): SubprocessSpawner { + return ({ env }) => { const supervisorToChild = createMemoryNdjsonStream(); const childToSupervisor = createMemoryNdjsonStream(); const eventChildToSupervisor = createMemoryFrameStream(); @@ -256,6 +250,39 @@ export async function makeLifecycleFixture(opts?: { spawns.push(entry); return handle; }; +} + +export async function makeLifecycleFixture(opts?: { + /** + * Reuse a prior fixture's data dir, so a test can construct a SECOND + * router against the SAME on-disk state to model a sidecar process + * restart (boot-time restore). + */ + dataDir?: string; + /** + * Override the default in-memory fake key store with a real one (e.g. + * `@intx/hub-agent`'s `createAgentKeyStore` bound to the same + * `dataDir`), for a test that needs the actual on-disk key-persistence + * behavior the fake bypasses entirely. + */ + keyStore?: Parameters[0]["keyStore"]; + /** + * Override the default single-step `LIFECYCLE_CLOSURE_DEFINITION` closure + * materializer -- e.g. with a multi-step definition, for a test that needs + * the eager (non-deferred-to-wake) boot-restore path. + */ + materializeDeploymentClosure?: Parameters< + typeof createSidecarDeployRouter + >[0]["materializeDeploymentClosure"]; + /** + * Override the boot-restore attempt's timeout ceiling, so a test can + * exercise the timeout path in milliseconds instead of the real 30s + * default. + */ + restoreAttemptTimeoutMs?: number; +}): Promise { + const spawns: Spawn[] = []; + const spawner = createHandshakeSpawner(spawns); const transport = createInMemoryTransport(); const keyPair = await generateKeyPair(); @@ -321,17 +348,22 @@ export async function makeLifecycleFixture(opts?: { // Stand in for the real fetch + SRI-verify + layout + evaluate pass: a // fixture cannot publish a package, so every deploy evaluates to the one // lifecycle definition below. - materializeDeploymentClosure: ({ deploymentId }) => - Promise.resolve({ - definition: LIFECYCLE_CLOSURE_DEFINITION, - packageDir: path.join(dataDir, "closure-package", deploymentId), - deployDir: path.join(dataDir, "closure-deploy", deploymentId), - }), + materializeDeploymentClosure: + opts?.materializeDeploymentClosure ?? + (({ deploymentId }) => + Promise.resolve({ + definition: LIFECYCLE_CLOSURE_DEFINITION, + packageDir: path.join(dataDir, "closure-package", deploymentId), + deployDir: path.join(dataDir, "closure-deploy", deploymentId), + })), multistepMailRouter, multistepSignalRouter, multistepDrainRouter, multistepSourcesRouter, multistepCredentialsRouter, + ...(opts?.restoreAttemptTimeoutMs !== undefined + ? { restoreAttemptTimeoutMs: opts.restoreAttemptTimeoutMs } + : {}), }); return { router, diff --git a/apps/sidecar/test/workflow-restore-timeout-cancel.test.ts b/apps/sidecar/test/workflow-restore-timeout-cancel.test.ts new file mode 100644 index 000000000..baaa310a3 --- /dev/null +++ b/apps/sidecar/test/workflow-restore-timeout-cancel.test.ts @@ -0,0 +1,155 @@ +// CL-7215: `withRestoreTimeout` used to race a boot-restore attempt +// against a timer and, on timeout, simply abandon the still-running +// attempt -- if it later succeeded, `spawnWorkflowDeployment` registered +// a live supervisor for an address the boot loop had already recorded as +// a restore FAILURE. This exercises the fix: a restore that overshoots +// its timeout still gets an `AbortSignal`, and once it actually finishes +// (late), it corrects its own durable record instead of leaving a live +// deployment marked failed. + +import { mkdtemp, rm } from "node:fs/promises"; +import { tmpdir } from "node:os"; +import path from "node:path"; +import { afterEach, expect, test } from "bun:test"; +import { defineWorkflow, step, type WorkflowDefinition } from "@intx/workflow"; +import { buildSingleStepAgentDefinition } from "@intx/workflow-deploy"; +import type { InferenceSource } from "@intx/types/runtime"; + +import { deriveDeploymentId } from "../src/workflow-host-wiring"; +import { + readWorkflowDeploymentRecord, + writeWorkflowDeploymentRecord, + type WorkflowDeploymentRecord, +} from "../src/workflow-deployment-record"; +import { + answerReadyHandshake, + makeLifecycleFixture, +} from "./support/workflow-lifecycle-fixture"; + +const tempDirs: string[] = []; + +afterEach(async () => { + await Promise.all( + tempDirs.splice(0).map((dir) => rm(dir, { recursive: true, force: true })), + ); +}); + +async function makeDataDir(): Promise { + const dir = await mkdtemp(path.join(tmpdir(), "sidecar-restore-timeout-")); + tempDirs.push(dir); + return dir; +} + +// Two steps, deliberately: a single-step ("warm-keep") record defers to its +// wake path instead of eagerly restoring (CL-6648), so it never reaches +// `spawnWorkflowDeployment` and could never exercise this timeout path. +const TWO_STEP_DEFINITION: WorkflowDefinition = defineWorkflow({ + id: "wf-restore-timeout", + trigger: { type: "mail", to: "wf-restore-timeout@example.com" }, + steps: { + "step-1": step({ + agent: buildSingleStepAgentDefinition({ + id: "step-1", + systemPrompt: "", + inferencePreferences: [], + toolFactories: [], + }), + triggers: "unbounded", + }), + "step-2": step({ + agent: buildSingleStepAgentDefinition({ + id: "step-2", + systemPrompt: "", + inferencePreferences: [], + toolFactories: [], + }), + triggers: "unbounded", + }), + }, +}); + +function makeSource(): InferenceSource { + return { + id: "source-1", + provider: "buildable", + baseURL: "https://inference.example.com", + apiKey: "key", + model: "model-1", + }; +} + +const SOURCE_REF: WorkflowDeploymentRecord["sourceRef"] = { + source: { kind: "registry", registry: "npm" }, + closure: { schemaVersion: "1", topLevel: [], entries: [] }, +}; + +test("a restore that outlasts its timeout corrects its record once the spawn actually finishes, instead of leaving a live deployment marked failed", async () => { + const dataDir = await makeDataDir(); + const restoreAttemptTimeoutMs = 25; + + const { router, spawns } = await makeLifecycleFixture({ + dataDir, + restoreAttemptTimeoutMs, + materializeDeploymentClosure: ({ deploymentId }) => + Promise.resolve({ + definition: TWO_STEP_DEFINITION, + packageDir: path.join(dataDir, "closure-package", deploymentId), + deployDir: path.join(dataDir, "closure-deploy", deploymentId), + }), + }); + + const agentAddress = "run_late-restore@example.com"; + const deploymentId = deriveDeploymentId(agentAddress); + const record: WorkflowDeploymentRecord = { + version: 1, + agentAddress, + definitionId: "def_1", + sources: { + "step-1": [makeSource()], + "step-2": [makeSource()], + }, + approvedWireHash: "d".repeat(64), + sourceRef: SOURCE_REF, + }; + await writeWorkflowDeploymentRecord(dataDir, deploymentId, record); + + // Don't answer the ready handshake yet: `wired.supervisor.spawn()` stays + // pending well past `restoreAttemptTimeoutMs`, so the boot loop's own + // timeout fires and records this attempt as a boot failure first. + const restorePromise = router.restoreWorkflowDeployments(); + await restorePromise; + + const afterTimeout = await readWorkflowDeploymentRecord( + dataDir, + deploymentId, + ); + expect(afterTimeout?.restoreFailure).toBeDefined(); + expect(router.activeAddresses()).toEqual([]); + + // The underlying restore was never abandoned: answering the handshake now + // lets `spawnWorkflowDeployment` actually finish, well after the boot + // loop already gave up on it. + await answerReadyHandshake(spawns, 0); + + // Both the live registration AND the record correction are background + // continuations of the restore attempt, not something + // `restoreWorkflowDeployments()` awaits -- and the registration lands + // strictly before the correction (the correction only runs once + // `spawnWorkflowDeployment` has already resolved). Poll on the record + // correction specifically, since it's the last of the two to land. + let afterLateRestore: WorkflowDeploymentRecord | undefined; + const deadline = Date.now() + 5000; + while (Date.now() < deadline) { + afterLateRestore = await readWorkflowDeploymentRecord( + dataDir, + deploymentId, + ); + if (afterLateRestore?.restoreFailure === undefined) break; + await new Promise((r) => setTimeout(r, 5)); + } + + // The deployment is genuinely live now -- the record must no longer + // claim it failed to restore. + expect(afterLateRestore?.restoreFailure).toBeUndefined(); + expect(router.activeAddresses()).toEqual([agentAddress]); +}); From 5410ecd05eb3f248a7097207447b31cbe3425e0a Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Sun, 30 Aug 2026 06:07:09 -0700 Subject: [PATCH 2/3] Give boot-time restore a real cancellation path instead of an abandoning timeout withRestoreTimeout raced restoreDeploymentFromRecord against a timer and walked away on timeout, with no way to tell the loser it had lost. The underlying restore kept running regardless -- including spawnWorkflowDeployment, which registers a live supervisor and spawns a real workflow-process child -- so a restore that finished after the timeout already fired left a genuinely live deployment whose durable record claimed it had failed to restore (CL-7215). withRestoreTimeout now hands the restore an AbortSignal that fires the instant the deadline wins. restoreDeploymentFromRecord checks it at two points where real cancellation is still possible: before the closure fetch (applyAtomic has no signal of its own, so bailing here is the only way to skip it) and before the spawn (Supervisor.spawn has none either). Past that second checkpoint the restore cannot be cancelled -- so once spawnWorkflowDeployment resolves, it checks whether the caller's own timeout had already fired and, if so, corrects the durable record instead of leaving that lie in place. The correction and the boot loop's own failure-recording catch both touch the same on-disk record, so they're serialized through a per-deployment in-memory lock (withDeploymentRecordLock): without it, the correction's disk read could land before the catch's disk write, observe nothing to clear, and no-op -- leaving the catch's write to permanently mismark a live deployment. The catch also consults activeSupervisors (in-memory, always set synchronously the instant a spawn actually succeeds) while holding that same lock, so a restore that raced ahead of its own timeout never gets a boot failure recorded against it in the first place. Together the two make the final on-disk state correct regardless of which side finishes first. RESTORE_ATTEMPT_TIMEOUT_MS is now overridable per router (restoreAttemptTimeoutMs) so a test can exercise this in milliseconds rather than the real 30s ceiling. --- .../sidecar/src/workflow-host-wiring/index.ts | 198 ++++++++++++++++-- 1 file changed, 175 insertions(+), 23 deletions(-) diff --git a/apps/sidecar/src/workflow-host-wiring/index.ts b/apps/sidecar/src/workflow-host-wiring/index.ts index 86ea0b7b6..5f8fda3fe 100644 --- a/apps/sidecar/src/workflow-host-wiring/index.ts +++ b/apps/sidecar/src/workflow-host-wiring/index.ts @@ -58,6 +58,7 @@ import { isWorkflowDeploymentRestoreQuarantined, markWorkflowDeploymentRecordParked, partitionScannedDeployments, + readWorkflowDeploymentRecord, recordWorkflowDeploymentRestoreFailure, scanWorkflowDeploymentRecords, writeWorkflowDeploymentRecord, @@ -160,28 +161,46 @@ export const RESTORE_CONCURRENCY = 8; * transient failure: counted on the record's `restoreFailure` counter, * logged, and skipped, so the boot moves on and the record gets another * attempt next boot (or quarantines after `RESTORE_QUARANTINE_THRESHOLD` - * consecutive permanent failures). + * consecutive permanent failures). Overridable per router + * (`deps.restoreAttemptTimeoutMs`) so a test can exercise the timeout path + * in milliseconds instead of waiting out the real ceiling. */ export const RESTORE_ATTEMPT_TIMEOUT_MS = 30_000; +/** + * Bounds a boot-restore attempt to `timeoutMs` -- so `restoreWorkflowDeployments`'s + * bounded worker pool frees this slot and moves on to the next record rather + * than waiting on a wedged restore forever -- WITHOUT quietly abandoning the + * restore itself (CL-7215). `work` receives an `AbortSignal` that fires the + * instant the deadline wins, so a restore attempt that has not yet started + * its closure fetch or its spawn can cancel itself for real rather than + * running either to completion unobserved. Neither the closure fetch + * (`@intx/tool-packaging`'s `applyAtomic`) nor the workflow-process spawn + * (`@intx/workflow-host`'s `Supervisor.spawn`) accept a signal of their own, + * so a restore already past both of `restoreDeploymentFromRecord`'s + * checkpoints when the deadline fires keeps running -- but it corrects its + * own durable record once it learns the true outcome instead of leaving a + * live deployment recorded as a boot failure; see the `signal.aborted` + * check around `spawnWorkflowDeployment` in `restoreDeploymentFromRecord`. + */ function withRestoreTimeout( - promise: Promise, + work: (signal: AbortSignal) => Promise, deploymentId: string, + timeoutMs: number, ): Promise { + const controller = new AbortController(); return new Promise((resolve, reject) => { + const message = `restore of ${deploymentId} exceeded ${timeoutMs}ms`; const timer = setTimeout(() => { - reject( - new Error( - `restore of ${deploymentId} exceeded ${RESTORE_ATTEMPT_TIMEOUT_MS}ms`, - ), - ); - }, RESTORE_ATTEMPT_TIMEOUT_MS); - promise.then( + controller.abort(new Error(message)); + reject(new Error(message)); + }, timeoutMs); + work(controller.signal).then( (value) => { clearTimeout(timer); resolve(value); }, - (error) => { + (error: unknown) => { clearTimeout(timer); reject(error); }, @@ -189,6 +208,24 @@ function withRestoreTimeout( }); } +/** + * Cooperative cancellation checkpoint for `restoreDeploymentFromRecord` + * (CL-7215): called immediately before starting a phase of work that is + * still avoidable -- the closure fetch, or the spawn itself -- so a restore + * whose caller has already timed out stops advancing instead of spending + * that work on an attempt nobody is waiting on anymore. + */ +function throwIfRestoreAborted( + signal: AbortSignal, + deploymentId: string, +): void { + if (signal.aborted) { + throw new Error( + `restore of ${deploymentId} was cancelled after its boot-restore attempt timed out`, + ); + } +} + /** * Await a supervisor's graceful `shutdown()`, escalating to a direct * SIGKILL of its child if `shutdown()` hasn't settled within @@ -541,6 +578,14 @@ export function createSidecarDeployRouter(deps: { * the deploy or boot-time restore. */ readyTimeoutMs?: number; + /** + * Ceiling (ms) on a single boot-time `restoreDeploymentFromRecord` + * attempt; see `RESTORE_ATTEMPT_TIMEOUT_MS`'s doc comment for what the + * bound protects. Defaults to `RESTORE_ATTEMPT_TIMEOUT_MS`; overridable + * so a test can exercise the timeout-and-late-settle path in + * milliseconds instead of the real 30s ceiling. + */ + restoreAttemptTimeoutMs?: number; /** * Deployment-record writer, injectable so a test can block or fail the * persist at a controlled point -- the natural seam for exercising a @@ -589,6 +634,8 @@ export function createSidecarDeployRouter(deps: { deps.writeWorkflowDeploymentRecord ?? writeWorkflowDeploymentRecord; const applyClosure = deps.materializeDeploymentClosure ?? materializeDeploymentClosure; + const restoreAttemptTimeoutMs = + deps.restoreAttemptTimeoutMs ?? RESTORE_ATTEMPT_TIMEOUT_MS; const multistepSpawner = deps.multistepSubprocessSpawner ?? defaultSubprocessSpawner; const multistepDeriveStepAddress: DeriveStepAddress = @@ -603,6 +650,37 @@ export function createSidecarDeployRouter(deps: { // with the deployment. const activeSupervisors = new Map(); + // CL-7215: per-deployment FIFO serialization for the boot-restore + // failure record. Two independent writers touch the same on-disk + // `deployment.json` around a timed-out restore -- the boot loop's own + // catch (`recordWorkflowDeploymentRestoreFailure`) and, if the + // abandoned restore later finishes anyway, `restoreDeploymentFromRecord`'s + // late-settle correction (`clearWorkflowDeploymentRestoreFailure`) -- and + // without serialization a correction's disk READ can land before the + // catch's disk WRITE, observe nothing to clear, and no-op, leaving the + // catch's later write to permanently mismark a live deployment as + // failed. Queuing both through this lock, keyed by deployment id, makes + // whichever runs second observe the first's completed write; combined + // with the catch consulting `activeSupervisors` (see below) while + // holding the lock, the final on-disk state is correct regardless of + // which side actually finishes first. + const deploymentRecordLocks = new Map>(); + function withDeploymentRecordLock( + deploymentId: string, + fn: () => Promise, + ): Promise { + const prior = deploymentRecordLocks.get(deploymentId) ?? Promise.resolve(); + const settled = prior.then(fn, fn); + deploymentRecordLocks.set( + deploymentId, + settled.then( + () => undefined, + () => undefined, + ), + ); + return settled; + } + // Synchronous single-flight guard for the deploy path. The real supervisor // does not exist until inside `spawnWorkflowDeployment`, so `deployMultiStep` // cannot reserve its `activeSupervisors` slot up front; instead it records @@ -1585,6 +1663,7 @@ export function createSidecarDeployRouter(deps: { dataDir: string, deploymentId: string, record: WorkflowDeploymentRecord, + signal: AbortSignal, ): Promise<"restored" | "deferred-to-wake" | "pruned"> { // Integrity: the stored address must re-derive to its own directory // name. A mismatch means a corrupt or misplaced record -- permanent, @@ -1622,6 +1701,11 @@ export function createSidecarDeployRouter(deps: { return "deferred-to-wake"; } + // CL-7215: the closure fetch below is real, avoidable work -- bail + // before starting it rather than spending it on an attempt the caller + // has already timed out and recorded as failed. + throwIfRestoreAborted(signal, deploymentId); + // Reconstruct this deployment's runnable definition: re-materialize the // pinned closure and evaluate the pinned code, then project it to the // inert wire shape -- the SAME computation the deploy path applies. The @@ -1746,7 +1830,54 @@ export function createSidecarDeployRouter(deps: { slugClaims.get(deploymentId) !== record.agentAddress; claimSlug(deploymentId, record.agentAddress); try { + // CL-7215: checked INSIDE this try, after the slug is claimed, so a + // cancellation here unwinds through the same catch below that + // releases it -- an abort must never leak a claimed slug. + // `spawnWorkflowDeployment` itself has no cancellation point (no + // signal reaches `@intx/workflow-host`'s `Supervisor.spawn`), so this + // is the last point this restore can still avoid starting it. + throwIfRestoreAborted(signal, deploymentId); await spawnWorkflowDeployment(spec); + if (signal.aborted) { + // The caller's own timeout already fired -- and its catch may + // already have recorded this attempt as a boot failure -- but + // `spawnWorkflowDeployment` just registered a live supervisor for + // this address regardless (CL-7215). Correct the durable record so + // it stops claiming a genuinely live deployment failed to restore. + // Serialized through `withDeploymentRecordLock` against that same + // catch: without it, this read could land BEFORE the catch's write + // does, observe nothing to clear, and no-op -- leaving the catch's + // write to permanently mismark a live deployment once it lands + // afterward. The lock guarantees whichever of the two runs second + // observes the first's completed write. + try { + const corrected = await withDeploymentRecordLock( + deploymentId, + async () => { + const onDisk = await readWorkflowDeploymentRecord( + dataDir, + deploymentId, + ); + if (onDisk?.restoreFailure === undefined) return false; + await clearWorkflowDeploymentRestoreFailure( + dataDir, + deploymentId, + onDisk, + ); + return true; + }, + ); + if (corrected) { + logger.warn`Workflow deployment ${record.agentAddress} finished restoring after its boot-restore attempt had already timed out; corrected its record so it no longer claims the deployment failed to restore`; + } + } catch (correctionError) { + reportError(correctionError, { + operation: + "workflow-host-wiring.restoreDeploymentFromRecord.lateRestoreCorrection", + agentId: record.agentAddress, + }); + } + } logger.info`Restored workflow deployment for ${record.agentAddress}`; return "restored"; } catch (cause) { @@ -1942,8 +2073,15 @@ export function createSidecarDeployRouter(deps: { } try { const outcome = await withRestoreTimeout( - restoreDeploymentFromRecord(dataDir, deploymentId, record), + (signal) => + restoreDeploymentFromRecord( + dataDir, + deploymentId, + record, + signal, + ), deploymentId, + restoreAttemptTimeoutMs, ); if (outcome === "deferred-to-wake") { deferredToWakeCount += 1; @@ -1963,18 +2101,32 @@ export function createSidecarDeployRouter(deps: { cause instanceof WorkflowRestoreFailure ? cause.kind : "transient"; - const updated = await recordWorkflowDeploymentRestoreFailure( - dataDir, - deploymentId, - record, - { kind, reason }, - ); - if (isWorkflowDeploymentRestoreQuarantined(updated)) { - const attempts = updated.restoreFailure?.attempts ?? 0; - logger.warn`Workflow deployment ${deploymentId} failed to restore ${attempts} consecutive times and is now quarantined -- it will not be retried again until the address is undeployed. Last failure: ${reason}`; - } else { - logger.warn`Failed to restore workflow deployment ${deploymentId}: ${reason}`; - } + // CL-7215: serialized against `restoreDeploymentFromRecord`'s + // own late-settle correction via the same lock, and gated on + // `activeSupervisors` -- in-memory, always set synchronously the + // instant `spawnWorkflowDeployment` actually succeeds, unlike a + // disk snapshot -- so a restore that timed out here but has + // ALREADY gone live by the time this runs never gets a boot + // failure recorded against it in the first place. See the lock's + // own doc comment for why both writers must go through it. + await withDeploymentRecordLock(deploymentId, async () => { + if (activeSupervisors.has(record.agentAddress)) { + logger.warn`Workflow deployment ${deploymentId} timed out during boot restore but finished spawning before its failure could be recorded; leaving its record as a live, successful restore`; + return; + } + const updated = await recordWorkflowDeploymentRestoreFailure( + dataDir, + deploymentId, + record, + { kind, reason }, + ); + if (isWorkflowDeploymentRestoreQuarantined(updated)) { + const attempts = updated.restoreFailure?.attempts ?? 0; + logger.warn`Workflow deployment ${deploymentId} failed to restore ${attempts} consecutive times and is now quarantined -- it will not be retried again until the address is undeployed. Last failure: ${reason}`; + } else { + logger.warn`Failed to restore workflow deployment ${deploymentId}: ${reason}`; + } + }); } }, ); From 06ae45e0a5c4f94891c3f8ddfb413ceddc922c81 Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Sun, 30 Aug 2026 09:19:28 -0700 Subject: [PATCH 3/3] Serialize teardown's record writes against the boot-restore lock teardownDeployment's own two on-disk deployment-record writes (deleteWorkflowDeploymentRecord for a reclaiming undeploy, markWorkflowDeploymentRecordParked for a hibernate) were a third, unlocked writer touching deployment.json for a given deploymentId, alongside the two the CL-7215 boot-restore fix already serializes through withDeploymentRecordLock: the boot loop's restoreFailure write and a dangling (timed-out but still-running) restore's late-settle correction. Routing both teardown writes through the same lock closes the write-ordering half of that race. Closing the write ordering alone was not enough: a reclaiming teardown racing a dangling restore can still delete the record between when the boot loop scans it and when either writer runs, and two blind writers mishandled that: - recordWorkflowDeploymentRestoreFailure wrote the caller's stale in-memory record unconditionally, so a teardown that reclaimed the deployment first got its delete silently undone -- the boot loop's catch would resurrect deployment.json with a false restoreFailure for a deployment that no longer exists. It now re-reads from disk under the lock and no-ops when the record is already gone, matching markWorkflowDeploymentRecordParked's sibling shape. - restoreDeploymentFromRecord's late-settle correction only handled "record still exists but wrongly marked failed"; it had no branch for "record is gone entirely". If the reclaimed restore's own spawn later succeeded anyway, spawnWorkflowDeployment would register a live supervisor and re-register its routers with nothing durable backing them -- an orphaned live deployment invisible to any future boot scan, arguably worse than the original false-failure bug. It now distinguishes that outcome and unwinds the orphaned supervisor through the ordinary teardownDeployment(reclaimDirs: true) path. Two tests exercise the fixed race: a reclaiming teardown ahead of the boot loop's own failure write must not have its delete undone, and a teardown racing a restore whose spawn resolves afterward must not leave a live, record-less deployment behind. --- .../src/workflow-deployment-record.test.ts | 70 ++++++---- .../sidecar/src/workflow-deployment-record.ts | 27 ++-- .../sidecar/src/workflow-host-wiring/index.ts | 121 ++++++++++++++---- .../workflow-restore-timeout-cancel.test.ts | 92 +++++++++++++ 4 files changed, 246 insertions(+), 64 deletions(-) diff --git a/apps/sidecar/src/workflow-deployment-record.test.ts b/apps/sidecar/src/workflow-deployment-record.test.ts index 21ef2034d..4e0a769c1 100644 --- a/apps/sidecar/src/workflow-deployment-record.test.ts +++ b/apps/sidecar/src/workflow-deployment-record.test.ts @@ -184,15 +184,14 @@ describe("recordWorkflowDeploymentRestoreFailure", () => { const updated = await recordWorkflowDeploymentRestoreFailure( dataDir, "dep_1", - baseRecord, { kind: "permanent", reason: "address derives a different slug" }, ); - expect(updated.restoreFailure).toEqual({ + expect(updated?.restoreFailure).toEqual({ kind: "permanent", attempts: 1, reason: "address derives a different slug", - lastAttemptAt: updated.restoreFailure?.lastAttemptAt ?? "", + lastAttemptAt: updated?.restoreFailure?.lastAttemptAt ?? "", }); const onDisk = await readWorkflowDeploymentRecord(dataDir, "dep_1"); expect(onDisk?.restoreFailure?.attempts).toBe(1); @@ -202,47 +201,59 @@ describe("recordWorkflowDeploymentRestoreFailure", () => { const dataDir = await makeDataDir(); await writeWorkflowDeploymentRecord(dataDir, "dep_1", baseRecord); - let record = baseRecord; + let record: WorkflowDeploymentRecord | undefined; for (let i = 0; i < 3; i++) { - record = await recordWorkflowDeploymentRestoreFailure( - dataDir, - "dep_1", - record, - { kind: "permanent", reason: "still malformed" }, - ); + record = await recordWorkflowDeploymentRestoreFailure(dataDir, "dep_1", { + kind: "permanent", + reason: "still malformed", + }); } - expect(record.restoreFailure?.attempts).toBe(3); - expect(record.restoreFailure?.kind).toBe("permanent"); + expect(record?.restoreFailure?.attempts).toBe(3); + expect(record?.restoreFailure?.kind).toBe("permanent"); }); test("a kind change resets the counter rather than adding to the other kind's count", async () => { const dataDir = await makeDataDir(); await writeWorkflowDeploymentRecord(dataDir, "dep_1", baseRecord); + await recordWorkflowDeploymentRestoreFailure(dataDir, "dep_1", { + kind: "permanent", + reason: "malformed", + }); let record = await recordWorkflowDeploymentRestoreFailure( dataDir, "dep_1", - baseRecord, - { kind: "permanent", reason: "malformed" }, - ); - record = await recordWorkflowDeploymentRestoreFailure( - dataDir, - "dep_1", - record, - { kind: "permanent", reason: "still malformed" }, + { + kind: "permanent", + reason: "still malformed", + }, ); - expect(record.restoreFailure?.attempts).toBe(2); + expect(record?.restoreFailure?.attempts).toBe(2); + + record = await recordWorkflowDeploymentRestoreFailure(dataDir, "dep_1", { + kind: "transient", + reason: "provider not registered", + }); - record = await recordWorkflowDeploymentRestoreFailure( + expect(record?.restoreFailure?.kind).toBe("transient"); + expect(record?.restoreFailure?.attempts).toBe(1); + }); + + test("is a no-op that returns undefined when the record is already gone (CL-7215: a reclaiming teardown raced ahead of this write)", async () => { + const dataDir = await makeDataDir(); + // Deliberately never written: simulates a teardown deleting the record + // before the boot loop's failure-recording catch acquires the lock. + + const updated = await recordWorkflowDeploymentRestoreFailure( dataDir, - "dep_1", - record, + "dep_missing", { kind: "transient", reason: "provider not registered" }, ); - expect(record.restoreFailure?.kind).toBe("transient"); - expect(record.restoreFailure?.attempts).toBe(1); + expect(updated).toBeUndefined(); + const onDisk = await readWorkflowDeploymentRecord(dataDir, "dep_missing"); + expect(onDisk).toBeUndefined(); }); }); @@ -253,11 +264,14 @@ describe("clearWorkflowDeploymentRestoreFailure", () => { const failed = await recordWorkflowDeploymentRestoreFailure( dataDir, "dep_1", - baseRecord, { kind: "transient", reason: "provider not registered" }, ); - await clearWorkflowDeploymentRestoreFailure(dataDir, "dep_1", failed); + await clearWorkflowDeploymentRestoreFailure( + dataDir, + "dep_1", + failed as WorkflowDeploymentRecord, + ); const onDisk = await readWorkflowDeploymentRecord(dataDir, "dep_1"); expect(onDisk?.restoreFailure).toBeUndefined(); diff --git a/apps/sidecar/src/workflow-deployment-record.ts b/apps/sidecar/src/workflow-deployment-record.ts index 53e811572..3514d35ae 100644 --- a/apps/sidecar/src/workflow-deployment-record.ts +++ b/apps/sidecar/src/workflow-deployment-record.ts @@ -227,24 +227,35 @@ export async function markWorkflowDeploymentRecordParked( * `attempts` (reset to 1 when `failure.kind` differs from the previously * recorded kind, so a transient streak can never inflate the permanent * counter or vice versa), stamp `reason`/`lastAttemptAt`, and persist. - * Returns the updated record so the caller can check `isWorkflowDeploymentRestoreQuarantined` - * without a re-read. Takes the caller's in-memory `record` rather than - * re-reading it, matching `markWorkflowDeploymentRecordParked`'s sibling - * shape but avoiding a redundant read on the hot boot-restore path. + * Returns the updated record so the caller can check + * `isWorkflowDeploymentRestoreQuarantined` without a re-read. + * + * Re-reads the record from disk rather than trusting the caller's + * in-memory copy (CL-7215): the boot-restore path's per-deployment lock + * serializes this against a reclaiming teardown that can delete the + * record out from under a dangling (timed-out but still-running) restore + * between when the boot loop scanned it and when this runs. Blindly + * writing the stale in-memory `record` would resurrect `deployment.json` + * with a false failure stamp for a deployment that was fully reclaimed -- + * the exact "durable record disagrees with reality" failure this ticket + * exists to close, just from the opposite direction. Returns `undefined` + * (no write) when the record is already gone, matching + * `markWorkflowDeploymentRecordParked`'s no-op-on-missing shape. */ export async function recordWorkflowDeploymentRestoreFailure( dataDir: string, deploymentId: string, - record: WorkflowDeploymentRecord, failure: { kind: RestoreFailureKind; reason: string }, -): Promise { - const previous = record.restoreFailure; +): Promise { + const existing = await readWorkflowDeploymentRecord(dataDir, deploymentId); + if (existing === undefined) return undefined; + const previous = existing.restoreFailure; const attempts = previous !== undefined && previous.kind === failure.kind ? previous.attempts + 1 : 1; const updated: WorkflowDeploymentRecord = { - ...record, + ...existing, restoreFailure: { kind: failure.kind, attempts, diff --git a/apps/sidecar/src/workflow-host-wiring/index.ts b/apps/sidecar/src/workflow-host-wiring/index.ts index 5f8fda3fe..242db0269 100644 --- a/apps/sidecar/src/workflow-host-wiring/index.ts +++ b/apps/sidecar/src/workflow-host-wiring/index.ts @@ -651,19 +651,26 @@ export function createSidecarDeployRouter(deps: { const activeSupervisors = new Map(); // CL-7215: per-deployment FIFO serialization for the boot-restore - // failure record. Two independent writers touch the same on-disk + // failure record. Three independent writers can touch the same on-disk // `deployment.json` around a timed-out restore -- the boot loop's own - // catch (`recordWorkflowDeploymentRestoreFailure`) and, if the - // abandoned restore later finishes anyway, `restoreDeploymentFromRecord`'s - // late-settle correction (`clearWorkflowDeploymentRestoreFailure`) -- and + // catch (`recordWorkflowDeploymentRestoreFailure`), if the abandoned + // restore later finishes anyway `restoreDeploymentFromRecord`'s + // late-settle correction (`clearWorkflowDeploymentRestoreFailure`), and + // `teardownDeployment`'s reclaim/park writes if an operator undeploys + // or hibernates the address while a restore is still dangling -- and // without serialization a correction's disk READ can land before the // catch's disk WRITE, observe nothing to clear, and no-op, leaving the // catch's later write to permanently mismark a live deployment as - // failed. Queuing both through this lock, keyed by deployment id, makes - // whichever runs second observe the first's completed write; combined - // with the catch consulting `activeSupervisors` (see below) while - // holding the lock, the final on-disk state is correct regardless of - // which side actually finishes first. + // failed (or, symmetrically, the catch's write can land after a + // reclaiming teardown deleted the record, resurrecting it with a false + // failure for a deployment that no longer exists -- + // `recordWorkflowDeploymentRestoreFailure` re-reads from disk and + // no-ops on a missing record specifically to make that resurrection + // impossible). Queuing every writer through this lock, keyed by + // deployment id, makes whichever runs last observe every prior one's + // completed write; combined with the catch consulting + // `activeSupervisors` (see below) while holding the lock, the final + // on-disk state is correct regardless of ordering. const deploymentRecordLocks = new Map>(); function withDeploymentRecordLock( deploymentId: string, @@ -1842,14 +1849,29 @@ export function createSidecarDeployRouter(deps: { // The caller's own timeout already fired -- and its catch may // already have recorded this attempt as a boot failure -- but // `spawnWorkflowDeployment` just registered a live supervisor for - // this address regardless (CL-7215). Correct the durable record so - // it stops claiming a genuinely live deployment failed to restore. - // Serialized through `withDeploymentRecordLock` against that same - // catch: without it, this read could land BEFORE the catch's write - // does, observe nothing to clear, and no-op -- leaving the catch's - // write to permanently mismark a live deployment once it lands - // afterward. The lock guarantees whichever of the two runs second - // observes the first's completed write. + // this address regardless (CL-7215). Reconcile the durable record + // against whatever the current on-disk truth is, under the same + // lock the boot loop's catch and `teardownDeployment` use: without + // it, this read could land BEFORE the catch's write does, observe + // nothing to correct, and no-op -- leaving the catch's write to + // permanently mismark a live deployment once it lands afterward. + // The lock guarantees whichever writer runs last observes every + // prior one's completed write. + // + // Three outcomes are possible for what this read finds: + // - a `restoreFailure` stamp: the boot loop's catch ran first and + // falsely marked this now-live deployment failed -- clear it. + // - the record itself is GONE: a reclaiming `teardownDeployment` + // ran while this spawn was in flight. `activeSupervisors` had no + // entry for this address yet at that point (this spawn had not + // resolved), so the teardown's own `wired` check missed it and + // left the routers this spawn just (re-)registered live, with no + // durable record backing them -- an orphaned live deployment, + // arguably worse than the original false-failure bug this ticket + // started from. Unwind it (below, outside the lock). + // - neither: nothing raced this restore; leave it as a normal + // successful (if late) restore. + let reclaimedDuringSpawn = false; try { const corrected = await withDeploymentRecordLock( deploymentId, @@ -1858,7 +1880,11 @@ export function createSidecarDeployRouter(deps: { dataDir, deploymentId, ); - if (onDisk?.restoreFailure === undefined) return false; + if (onDisk === undefined) { + reclaimedDuringSpawn = true; + return false; + } + if (onDisk.restoreFailure === undefined) return false; await clearWorkflowDeploymentRestoreFailure( dataDir, deploymentId, @@ -1877,6 +1903,24 @@ export function createSidecarDeployRouter(deps: { agentId: record.agentAddress, }); } + if (reclaimedDuringSpawn) { + // Called OUTSIDE the lock above: `teardownDeployment` acquires + // the same per-deployment lock itself, and this restore's own + // correction closure has already returned by this point, so + // there is no reentrant deadlock. + logger.warn`Workflow deployment ${record.agentAddress} finished spawning after being torn down while its boot-restore attempt was still in flight; unwinding the orphaned live supervisor`; + try { + await teardownDeployment(record.agentAddress, { + reclaimDirs: true, + }); + } catch (unwindError) { + reportError(unwindError, { + operation: + "workflow-host-wiring.restoreDeploymentFromRecord.unwindReclaimedDuringSpawn", + agentId: record.agentAddress, + }); + } + } } logger.info`Restored workflow deployment for ${record.agentAddress}`; return "restored"; @@ -1970,7 +2014,16 @@ export function createSidecarDeployRouter(deps: { // crash-interrupted deploy is reclaimed too. Skipped for a hibernate: // the record IS the durable state a later relaunch reads to resume. if (opts.reclaimDirs && stepStateDataDir !== undefined) { - await deleteWorkflowDeploymentRecord(stepStateDataDir, deploymentId); + // CL-7215: serialized through the same per-deployment lock the + // boot-restore path uses. Without it, a teardown racing a dangling + // (timed-out but still-running) restore's correction could delete + // the record out from under it, or clobber a `restoreFailure` write + // it lands mid-teardown -- neither of the boot-restore path's own + // two writers accounts for a third, unlocked writer touching the + // same file. + await withDeploymentRecordLock(deploymentId, () => + deleteWorkflowDeploymentRecord(stepStateDataDir, deploymentId), + ); } // Stamp the kept record as parked. This is the ONLY durable signal that // distinguishes "the hub parked this deployment on purpose" from "this @@ -1980,7 +2033,12 @@ export function createSidecarDeployRouter(deps: { // counts today; it does not yet change what it spawns (that cutover is // CL-6282). if (!opts.reclaimDirs && stepStateDataDir !== undefined) { - await markWorkflowDeploymentRecordParked(stepStateDataDir, deploymentId); + // CL-7215: same lock as above, for the same reason -- a hibernate's + // parked-stamp read-modify-write must not interleave with a dangling + // restore's own record writes for this deploymentId. + await withDeploymentRecordLock(deploymentId, () => + markWorkflowDeploymentRecordParked(stepStateDataDir, deploymentId), + ); } if (opts.reclaimDirs) { releaseSlug(deploymentId, agentAddress); @@ -2102,13 +2160,14 @@ export function createSidecarDeployRouter(deps: { ? cause.kind : "transient"; // CL-7215: serialized against `restoreDeploymentFromRecord`'s - // own late-settle correction via the same lock, and gated on - // `activeSupervisors` -- in-memory, always set synchronously the - // instant `spawnWorkflowDeployment` actually succeeds, unlike a - // disk snapshot -- so a restore that timed out here but has - // ALREADY gone live by the time this runs never gets a boot - // failure recorded against it in the first place. See the lock's - // own doc comment for why both writers must go through it. + // own late-settle correction AND `teardownDeployment`'s record + // writes via the same lock, and gated on `activeSupervisors` -- + // in-memory, always set synchronously the instant + // `spawnWorkflowDeployment` actually succeeds, unlike a disk + // snapshot -- so a restore that timed out here but has ALREADY + // gone live by the time this runs never gets a boot failure + // recorded against it in the first place. See the lock's own + // doc comment for why every writer must go through it. await withDeploymentRecordLock(deploymentId, async () => { if (activeSupervisors.has(record.agentAddress)) { logger.warn`Workflow deployment ${deploymentId} timed out during boot restore but finished spawning before its failure could be recorded; leaving its record as a live, successful restore`; @@ -2117,9 +2176,15 @@ export function createSidecarDeployRouter(deps: { const updated = await recordWorkflowDeploymentRestoreFailure( dataDir, deploymentId, - record, { kind, reason }, ); + if (updated === undefined) { + // The record is gone -- a concurrent reclaiming teardown + // won the lock first and deleted it. Nothing to mark: the + // deployment was torn down on purpose, not left claiming a + // false restore failure. + return; + } if (isWorkflowDeploymentRestoreQuarantined(updated)) { const attempts = updated.restoreFailure?.attempts ?? 0; logger.warn`Workflow deployment ${deploymentId} failed to restore ${attempts} consecutive times and is now quarantined -- it will not be retried again until the address is undeployed. Last failure: ${reason}`; diff --git a/apps/sidecar/test/workflow-restore-timeout-cancel.test.ts b/apps/sidecar/test/workflow-restore-timeout-cancel.test.ts index baaa310a3..d9c67a97b 100644 --- a/apps/sidecar/test/workflow-restore-timeout-cancel.test.ts +++ b/apps/sidecar/test/workflow-restore-timeout-cancel.test.ts @@ -153,3 +153,95 @@ test("a restore that outlasts its timeout corrects its record once the spawn act expect(afterLateRestore?.restoreFailure).toBeUndefined(); expect(router.activeAddresses()).toEqual([agentAddress]); }); + +test("a reclaiming teardown racing a dangling restore's boot-failure write is not resurrected with a false restoreFailure, and the restore's late spawn does not orphan a live supervisor with no record", async () => { + // The unwind this test exercises routes through `teardownDeployment`'s + // real `drain`/kill-escalation timers (`TEARDOWN_DRAIN_DEADLINE_MS` = + // 5000ms, `CHILD_KILL_ESCALATION_MS` = 3000ms). Bun's 5000ms default test + // timeout leaves no margin against that combined worst case, so both this + // timeout and the poll deadline below give it headroom past it. + const dataDir = await makeDataDir(); + const restoreAttemptTimeoutMs = 25; + + const { router, spawns } = await makeLifecycleFixture({ + dataDir, + restoreAttemptTimeoutMs, + materializeDeploymentClosure: ({ deploymentId }) => + Promise.resolve({ + definition: TWO_STEP_DEFINITION, + packageDir: path.join(dataDir, "closure-package", deploymentId), + deployDir: path.join(dataDir, "closure-deploy", deploymentId), + }), + }); + + const agentAddress = "run_teardown-race@example.com"; + const deploymentId = deriveDeploymentId(agentAddress); + const record: WorkflowDeploymentRecord = { + version: 1, + agentAddress, + definitionId: "def_1", + sources: { + "step-1": [makeSource()], + "step-2": [makeSource()], + }, + approvedWireHash: "d".repeat(64), + sourceRef: SOURCE_REF, + }; + await writeWorkflowDeploymentRecord(dataDir, deploymentId, record); + + // Don't answer the ready handshake yet: the restore's spawn stays pending + // past `restoreAttemptTimeoutMs`, leaving a window for a concurrent + // teardown to race it. + const restorePromise = router.restoreWorkflowDeployments(); + + // Wait for the restore to actually reach the spawn (i.e. the mock + // spawner has been invoked) before racing a teardown against it. + // Racing it any earlier collides with a DIFFERENT, unrelated race: the + // boot scan's own directory read of this same record file, which the + // scan already handles by skipping a record it can no longer read -- + // that would starve this test of the spawn-in-flight race it exists to + // exercise. + const spawnDeadline = Date.now() + 2000; + while (Date.now() < spawnDeadline && spawns.length === 0) { + await new Promise((r) => setTimeout(r, 1)); + } + expect(spawns).toHaveLength(1); + + // Race a reclaiming teardown against the still-in-flight restore, before + // its boot-restore timeout fires. `activeSupervisors` has no entry for + // this address yet (the spawn handshake is unanswered), so this exercises + // the ordinary "operator undeploys mid-restore" path, not the + // already-live-supervisor guard the other test above covers. + await router.teardownDeployment(agentAddress, { reclaimDirs: true }); + + await restorePromise; + + // The teardown deleted the record before the boot loop's own timeout + // catch could write to it. `recordWorkflowDeploymentRestoreFailure` + // re-reads from disk and must no-op on a missing record rather than + // resurrecting `deployment.json` with a false restoreFailure for a + // deployment that was fully reclaimed (CL-7215). + const afterRace = await readWorkflowDeploymentRecord(dataDir, deploymentId); + expect(afterRace).toBeUndefined(); + + // The underlying restore was never abandoned: answering the handshake now + // lets `spawnWorkflowDeployment` finish, well after the teardown already + // reclaimed this address. Without reconciling against the now-missing + // record, this would register a live supervisor with no durable record + // behind it -- an orphaned deployment, invisible to any future boot scan. + await answerReadyHandshake(spawns, 0); + + // Poll for the unwind: it is a background continuation of the restore + // attempt, not something `restoreWorkflowDeployments()` awaits. + const deadline = Date.now() + 10_000; + while ( + Date.now() < deadline && + router.activeAddresses().includes(agentAddress) + ) { + await new Promise((r) => setTimeout(r, 5)); + } + + expect(router.activeAddresses()).toEqual([]); + const afterUnwind = await readWorkflowDeploymentRecord(dataDir, deploymentId); + expect(afterUnwind).toBeUndefined(); +}, 15_000);