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 86ea0b7b6..242db0269 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,44 @@ export function createSidecarDeployRouter(deps: { // with the deployment. const activeSupervisors = new Map(); + // CL-7215: per-deployment FIFO serialization for the boot-restore + // 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`), 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 (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, + 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 +1670,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 +1708,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 +1837,91 @@ 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). 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, + async () => { + const onDisk = await readWorkflowDeploymentRecord( + dataDir, + deploymentId, + ); + if (onDisk === undefined) { + reclaimedDuringSpawn = true; + return false; + } + 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, + }); + } + 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"; } catch (cause) { @@ -1839,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 @@ -1849,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); @@ -1942,8 +2131,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 +2159,39 @@ 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 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`; + return; + } + const updated = await recordWorkflowDeploymentRestoreFailure( + dataDir, + deploymentId, + { 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}`; + } else { + logger.warn`Failed to restore workflow deployment ${deploymentId}: ${reason}`; + } + }); } }, ); 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..d9c67a97b --- /dev/null +++ b/apps/sidecar/test/workflow-restore-timeout-cancel.test.ts @@ -0,0 +1,247 @@ +// 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]); +}); + +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);