diff --git a/src/cron/service/active-run-cancellation.test.ts b/src/cron/service/active-run-cancellation.test.ts index 1d35d6eaead4..9c194849fe7a 100644 --- a/src/cron/service/active-run-cancellation.test.ts +++ b/src/cron/service/active-run-cancellation.test.ts @@ -1,8 +1,10 @@ import { describe, expect, it, vi } from "vitest"; import { + abortActiveCronTaskRuns, + cancelActiveCronTaskRun, getSuspensionVisibleCronTaskRunCount, + registerActiveCronTaskRun, retireActiveCronTaskRunTracking, - startActiveCronTaskRunSettlementGrace, trackActiveCronTaskRunSettlement, waitForActiveCronTaskRuns, } from "./active-run-cancellation.js"; @@ -38,16 +40,77 @@ describe("cron task cancellation tracking", () => { await vi.waitFor(() => expect(getSuspensionVisibleCronTaskRunCount()).toBe(0)); }); - it("drops never-settling cron promises after a bounded grace period", async () => { + it.each([ + { timing: "after tracking", abortBeforeTracking: false }, + { timing: "before tracking", abortBeforeTracking: true }, + ])( + "drops never-settling cron promises after an abort $timing", + async ({ abortBeforeTracking }) => { + vi.useFakeTimers(); + try { + resetActiveCronTaskRunsForTests(); + const controller = new AbortController(); + if (abortBeforeTracking) { + controller.abort(); + } + trackActiveCronTaskRunSettlement(new Promise(() => {}), controller.signal); + + await expect(waitForActiveCronTaskRuns(0)).resolves.toEqual({ + drained: false, + active: 1, + }); + + if (!abortBeforeTracking) { + await vi.advanceTimersByTimeAsync(CRON_TASK_RUN_SETTLEMENT_TRACKING_MAX_MS + 1); + + await expect(waitForActiveCronTaskRuns(0)).resolves.toEqual({ + drained: false, + active: 1, + }); + + controller.abort(); + } + await vi.advanceTimersByTimeAsync(CRON_TASK_RUN_SETTLEMENT_TRACKING_MAX_MS + 1); + + await expect(waitForActiveCronTaskRuns(0)).resolves.toEqual({ + drained: true, + active: 0, + }); + expect(getSuspensionVisibleCronTaskRunCount()).toBe(1); + } finally { + vi.useRealTimers(); + resetActiveCronTaskRunsForTests(); + } + }, + ); + + it("does not retire an unrelated active run when one registered run is cancelled", async () => { vi.useFakeTimers(); + const cancelledController = new AbortController(); + const unaffectedController = new AbortController(); + let settleCancelled = () => {}; + let settleUnaffected = () => {}; + const cancelledRun = new Promise((resolve) => { + settleCancelled = resolve; + }); + const unaffectedRun = new Promise((resolve) => { + settleUnaffected = resolve; + }); + try { resetActiveCronTaskRunsForTests(); - trackActiveCronTaskRunSettlement(new Promise(() => {})); - - await expect(waitForActiveCronTaskRuns(0)).resolves.toEqual({ - drained: false, - active: 1, + trackActiveCronTaskRunSettlement(cancelledRun, cancelledController.signal); + trackActiveCronTaskRunSettlement(unaffectedRun, unaffectedController.signal); + const release = registerActiveCronTaskRun({ + runId: "cancelled-run", + controller: cancelledController, }); + const existingTimerCount = vi.getTimerCount(); + + expect(cancelActiveCronTaskRun({ runId: "cancelled-run" })).toBe(true); + expect(vi.getTimerCount()).toBe(existingTimerCount + 1); + expect(cancelActiveCronTaskRun({ runId: "cancelled-run" })).toBe(false); + release?.(); await vi.advanceTimersByTimeAsync(CRON_TASK_RUN_SETTLEMENT_TRACKING_MAX_MS + 1); @@ -55,15 +118,51 @@ describe("cron task cancellation tracking", () => { drained: false, active: 1, }); + expect(getSuspensionVisibleCronTaskRunCount()).toBe(2); + } finally { + settleCancelled(); + settleUnaffected(); + await Promise.all([cancelledRun, unaffectedRun]); + vi.useRealTimers(); + resetActiveCronTaskRunsForTests(); + } + }); - startActiveCronTaskRunSettlementGrace(); + it.each([ + { scenario: "only handleless runs", cancellableRuns: 0, handlelessRuns: 1 }, + { scenario: "mixed active runs", cancellableRuns: 1, handlelessRuns: 1 }, + { scenario: "no active runs", cancellableRuns: 0, handlelessRuns: 0 }, + ])("retires $scenario during a bulk shutdown", async ({ cancellableRuns, handlelessRuns }) => { + vi.useFakeTimers(); + const controller = new AbortController(); + try { + resetActiveCronTaskRunsForTests(); + if (cancellableRuns > 0) { + trackActiveCronTaskRunSettlement(new Promise(() => {}), controller.signal); + } + if (handlelessRuns > 0) { + trackActiveCronTaskRunSettlement(new Promise(() => {})); + } + const release = + cancellableRuns > 0 + ? registerActiveCronTaskRun({ runId: "shutdown-run", controller }) + : undefined; + + expect(abortActiveCronTaskRuns()).toBe(cancellableRuns); + const retirementTimerCount = vi.getTimerCount(); + expect(abortActiveCronTaskRuns()).toBe(0); + expect(vi.getTimerCount()).toBe(retirementTimerCount); + if (cancellableRuns + handlelessRuns === 0) { + expect(retirementTimerCount).toBe(0); + } + release?.(); await vi.advanceTimersByTimeAsync(CRON_TASK_RUN_SETTLEMENT_TRACKING_MAX_MS + 1); await expect(waitForActiveCronTaskRuns(0)).resolves.toEqual({ drained: true, active: 0, }); - expect(getSuspensionVisibleCronTaskRunCount()).toBe(1); + expect(getSuspensionVisibleCronTaskRunCount()).toBe(cancellableRuns + handlelessRuns); } finally { vi.useRealTimers(); resetActiveCronTaskRunsForTests(); @@ -72,12 +171,13 @@ describe("cron task cancellation tracking", () => { it("keeps suspension blocked until a timed-out core actually settles", async () => { resetActiveCronTaskRunsForTests(); + const controller = new AbortController(); let settle = () => {}; const core = new Promise((resolve) => { settle = resolve; }); - trackActiveCronTaskRunSettlement(core); - startActiveCronTaskRunSettlementGrace(); + trackActiveCronTaskRunSettlement(core, controller.signal); + controller.abort(); expect(getSuspensionVisibleCronTaskRunCount()).toBe(1); settle(); diff --git a/src/cron/service/active-run-cancellation.ts b/src/cron/service/active-run-cancellation.ts index 2646778ee138..ee65bb504141 100644 --- a/src/cron/service/active-run-cancellation.ts +++ b/src/cron/service/active-run-cancellation.ts @@ -1,33 +1,25 @@ // Process-local cancellation handles for live cron task runs. -type CronTaskCancelHandle = { - controller: AbortController; - onCancel?: (reason: string) => void; -}; - -type SettlingCronTaskRun = { - retirementTimer?: NodeJS.Timeout; -}; - -const activeCronTaskRunsByRunId = new Map(); -const settlingCronTaskRuns = new Map, SettlingCronTaskRun>(); +const activeCronTaskRunsByRunId = new Map< + string, + { controller: AbortController; onCancel?: (reason: string) => void } +>(); +const settlingCronTaskRuns = new Map, { retirementTimer?: NodeJS.Timeout }>(); // Restart drain may retire an abort-ignoring core after a bounded grace, but a // host snapshot must keep refusing readiness until that core actually settles. const suspensionVisibleCronTaskRuns = new Set>(); const DEFAULT_CRON_TASK_RUN_DRAIN_POLL_MS = 25; const CRON_TASK_RUN_SETTLEMENT_TRACKING_MAX_MS = 60_000; -export function startActiveCronTaskRunSettlementGrace(): void { - for (const [promise, entry] of settlingCronTaskRuns) { - if (entry.retirementTimer) { - continue; - } - const retirementTimer = setTimeout(() => { - settlingCronTaskRuns.delete(promise); - }, CRON_TASK_RUN_SETTLEMENT_TRACKING_MAX_MS); - retirementTimer.unref?.(); - entry.retirementTimer = retirementTimer; +function startActiveCronTaskRunSettlementGrace(promise: Promise): void { + const entry = settlingCronTaskRuns.get(promise); + if (!entry || entry.retirementTimer) { + return; } + entry.retirementTimer = setTimeout(() => { + settlingCronTaskRuns.delete(promise); + }, CRON_TASK_RUN_SETTLEMENT_TRACKING_MAX_MS); + entry.retirementTimer.unref?.(); } export function registerActiveCronTaskRun(params: { @@ -60,18 +52,29 @@ export function abortActiveCronTaskRuns(reason = "Gateway restarting."): number handle.onCancel?.(reason); aborted += 1; } - if (aborted > 0) { - startActiveCronTaskRunSettlementGrace(); + // Shutdown also retires main-session runs without cancellation handles. + for (const promise of settlingCronTaskRuns.keys()) { + startActiveCronTaskRunSettlementGrace(promise); } return aborted; } -export function trackActiveCronTaskRunSettlement(promise: Promise): void { +export function trackActiveCronTaskRunSettlement( + promise: Promise, + abortSignal?: AbortSignal, +): void { settlingCronTaskRuns.set(promise, {}); suspensionVisibleCronTaskRuns.add(promise); + // Cancellation belongs to this core only; sibling jobs must remain drain-visible. + const startSettlementGrace = () => startActiveCronTaskRunSettlementGrace(promise); + abortSignal?.addEventListener("abort", startSettlementGrace, { once: true }); + if (abortSignal?.aborted) { + startSettlementGrace(); + } void promise .catch(() => undefined) .finally(() => { + abortSignal?.removeEventListener("abort", startSettlementGrace); const entry = settlingCronTaskRuns.get(promise); if (entry?.retirementTimer) { clearTimeout(entry.retirementTimer); @@ -131,7 +134,6 @@ export function cancelActiveCronTaskRun(params: { const reason = params.reason?.trim() || "Cancelled by operator."; handle.controller.abort(reason); handle.onCancel?.(reason); - startActiveCronTaskRunSettlementGrace(); return true; } diff --git a/src/cron/service/timer-job-runner.ts b/src/cron/service/timer-job-runner.ts index cc1eddd943e7..8740c63c0fe9 100644 --- a/src/cron/service/timer-job-runner.ts +++ b/src/cron/service/timer-job-runner.ts @@ -6,7 +6,6 @@ import { createCronRunDiagnosticsFromError } from "../run-diagnostics.js"; import type { CronAgentExecutionStarted, CronJob } from "../types.js"; import { registerActiveCronTaskRun, - startActiveCronTaskRunSettlementGrace, trackActiveCronTaskRunSettlement, } from "./active-run-cancellation.js"; import { @@ -250,7 +249,7 @@ export async function executeJobCoreWithTimeout( progress.completedCoreResult = result; return await deliverPrimaryWebhook(state, job, result, runAbortController.signal, progress); }); - trackActiveCronTaskRunSettlement(runPromise); + trackActiveCronTaskRunSettlement(runPromise, runAbortController.signal); void runPromise.catch((err: unknown) => { if (runAbortController.signal.aborted) { state.deps.log.warn( @@ -263,7 +262,6 @@ export async function executeJobCoreWithTimeout( if (first !== operatorCancellationMarker) { return first; } - startActiveCronTaskRunSettlementGrace(); const settled = resolveInterruptedRunProgress({ progress, job, @@ -334,7 +332,7 @@ export async function executeJobCoreWithTimeout( watchdog.deadlineAtMs(), ); }); - trackActiveCronTaskRunSettlement(runPromise); + trackActiveCronTaskRunSettlement(runPromise, runAbortController.signal); void runPromise.catch((err: unknown) => { if (runAbortController.signal.aborted) { state.deps.log.warn( @@ -346,7 +344,6 @@ export async function executeJobCoreWithTimeout( try { const first = await Promise.race([runPromise, timeoutPromise, operatorCancellationPromise]); if (first === operatorCancellationMarker) { - startActiveCronTaskRunSettlementGrace(); const settled = resolveInterruptedRunProgress({ progress, job, @@ -360,7 +357,6 @@ export async function executeJobCoreWithTimeout( if (first !== timeoutMarker) { return first; } - startActiveCronTaskRunSettlementGrace(); const activeExecution = watchdog.activeExecution(); const settled = resolveInterruptedRunProgress({ progress,