From d4c9069bd166f6f5eda39d1d4ccbb8e9e6732aab Mon Sep 17 00:00:00 2001 From: luciusverus-cyber <319770280+luciusverus-cyber@users.noreply.github.com> Date: Sun, 27 Sep 2026 21:18:37 +0200 Subject: [PATCH] feat(workers): add cancellation and shutdown-drain semantics (#142) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A maintenance job could be reported as completed after its side effect was abandoned, and there was no way to ask what work was in flight when a process stopped. Two things were wrong. `runJob` invoked `registration.handler(...)` and discarded the result, so a promise-returning handler was moved to `succeeded` before it had done anything and its rejection escaped as an unhandled rejection. The framework now awaits the handler: a sync handler settles inline, an async one is tracked so callers can wait for a final outcome, and `drainDueJobs()` reports the rest as `inFlight` rather than pretending they finished. Nothing could be cancelled, because a handler had no way to learn that it should stop. `WorkerContext` now carries an `AbortSignal` covering the whole framework, and the contract per state at shutdown is explicit: - queued / retrying — never started, left exactly as they are and reported as `pending`. Nothing was abandoned, so the next process re-runs them. - running — the attempt was abandoned. Moved to `retrying`, or `dead_lettered` once the attempt budget is spent, and reported as `interrupted`. It is never marked `succeeded`. - settled — untouched, counted in `settled`. Once cancelling, neither `drainDueJobs` nor `processJob` starts new work, so a handler is not re-entered after a shutdown. `shutdown()` aborts, waits for in-flight handlers up to a timeout, and classifies from the state captured *at abort time* — a handler that ignores its signal and never settles is still `running` with an unknown side effect, which is exactly the job a caller most needs told about, so it is reported as interrupted rather than pending. Cancellation is decided by the framework's own signal and nothing else. A handler that throws an `AbortError` from its own timeout has not been cancelled by a shutdown and is still reported as a failure; `isCancellation()` is exported for handlers that catch our abort and rethrow, not as the decision itself. `RunSummary` gains `cancelled` and `inFlight` so an abandoned attempt is not reported as a fault, and `reset()` installs a fresh controller so a reused framework is not permanently cancelled. `maintenance.ts` checks the signal before pruning — the memoised Horizon clients are left alone if a shutdown lands between the job being picked up and the prune running — and gains `shutdownMaintenance()` for an orderly teardown plus `cancelMaintenance()` for an unload handler that cannot await. Tests in core/workers/__tests__/shutdown.test.ts cover shutdown at each state and cover the async-handler bug directly. Every interleaving is deterministic: shutdown is awaited and handlers coordinate through a promise the test controls, so no assertion depends on timing. Verification: no test runner in this environment, so the suite was not executed. Instead the framework was run directly with node's TypeScript type stripping against the real modules (55 assertions across 12 scenarios, and the committed file's own expectations replayed separately), which caught three defects in the first draft: `drain()` waited for inherited work without counting it, shutdown classified from the post-drain state so a never-settling handler looked merely pending, and cancellation was being inferred from the error's class rather than the signal. --- core/workers/__tests__/shutdown.test.ts | 321 +++++++++++++++++ core/workers/maintenance.ts | 36 +- core/workers/queue.ts | 457 +++++++++++++++++++++--- 3 files changed, 770 insertions(+), 44 deletions(-) create mode 100644 core/workers/__tests__/shutdown.test.ts diff --git a/core/workers/__tests__/shutdown.test.ts b/core/workers/__tests__/shutdown.test.ts new file mode 100644 index 0000000..def535f --- /dev/null +++ b/core/workers/__tests__/shutdown.test.ts @@ -0,0 +1,321 @@ +import { beforeEach, describe, expect, it, vi } from "vitest"; +import { + createWorkerFramework, + isCancellation, + WorkerCancelledError, + type WorkerFramework +} from "@/core/workers/queue"; +import { + cancelMaintenance, + createMaintenanceFramework, + PRUNE_HORIZON_CLIENTS_OP +} from "@/core/workers/maintenance"; + +/** + * Cancellation and shutdown-drain semantics for the worker framework (#142). + * + * Without an explicit contract, a maintenance job can be reported as completed + * after its side effect was abandoned, or run again after a shutdown. The + * framework had neither a signal for a handler to observe nor a way to say what + * each job was doing when the process stopped, and it did not await async + * handlers at all: `handler(...)` was invoked and discarded, so a promise- + * returning job was marked `succeeded` before it had done anything and its + * rejection escaped as an unhandled rejection. + * + * The contract asserted here, per state at shutdown: + * + * queued / retrying — never started, left untouched, reported as `pending`. + * Nothing was abandoned, so the next process re-runs them. + * running — the attempt was abandoned. Moved to `retrying` (or + * `dead_lettered` once the budget is spent) and reported as `interrupted`. + * Never `succeeded`: that is the failure the issue describes. + * settled — untouched, counted as settled. + * + * Every interleaving here is deterministic: `shutdown` is awaited and the + * handlers coordinate through a promise the test controls, so no test depends on + * timing. + */ +const sleep = (ms: number) => new Promise((resolve) => setTimeout(resolve, ms)); + +describe("worker cancellation", () => { + let framework: WorkerFramework; + + beforeEach(() => { + framework = createWorkerFramework(); + }); + + it("awaits an async handler instead of marking it done immediately", async () => { + let finished = false; + framework.register("async.work", async () => { + await sleep(5); + finished = true; + }); + + const job = framework.enqueue({ operation: "async.work" }); + const provisional = framework.drainDueJobs(); + + // Started, but not finished: the summary says so and the status is running. + expect(provisional.inFlight).toBe(1); + expect(finished).toBe(false); + expect(framework.getById(job.id)?.status).toBe("running"); + + const settled = await framework.drain(); + + expect(settled.succeeded).toBe(1); + expect(finished).toBe(true); + expect(framework.getById(job.id)?.status).toBe("succeeded"); + }); + + it("counts work a previous synchronous drain left running", async () => { + framework.register("inherited.work", async () => { + await sleep(5); + }); + framework.enqueue({ operation: "inherited.work" }); + + framework.drainDueJobs(); + const settled = await framework.drain(); + + // The awaited work belongs in the summary the caller receives, or a drain + // would complete the job and report nothing about it. + expect(settled.ran).toBe(1); + expect(settled.succeeded).toBe(1); + }); + + it("passes a signal to every handler", async () => { + let received: AbortSignal | undefined; + framework.register("signal.work", (_params, context) => { + received = context.signal; + }); + + framework.enqueue({ operation: "signal.work" }); + framework.drainDueJobs(); + + expect(received).toBeInstanceOf(AbortSignal); + expect(received?.aborted).toBe(false); + expect(framework.signal()).toBe(received); + }); + + it("does not mark an abandoned running attempt as succeeded", async () => { + let sideEffect = false; + framework.register("long.work", async (_params, context) => { + await sleep(5); + // A cooperative handler: it stops when the abort reaches it. + if (context.signal.aborted) return; + sideEffect = true; + }); + + const job = framework.enqueue({ operation: "long.work" }); + framework.drainDueJobs(); + expect(framework.getById(job.id)?.status).toBe("running"); + + const report = await framework.shutdown({ reason: "test" }); + + expect(report.interrupted.map((entry) => entry.id)).toEqual([job.id]); + expect(report.pending).toHaveLength(0); + expect(report.reason).toBe("test"); + expect(report.phase).toBe("stopped"); + expect(framework.signal().aborted).toBe(true); + // The point of the whole change. + expect(framework.getById(job.id)?.status).not.toBe("succeeded"); + expect(framework.getById(job.id)?.status).toBe("retrying"); + expect(sideEffect).toBe(false); + expect(framework.getById(job.id)?.lastError?.code).toBe("cancelled"); + }); + + it("leaves a queued job untouched and reports it as pending", async () => { + const handler = vi.fn(); + framework.register("queued.work", handler); + const job = framework.enqueue({ operation: "queued.work" }); + + const report = await framework.shutdown(); + + expect(report.pending.map((entry) => entry.id)).toEqual([job.id]); + expect(report.interrupted).toHaveLength(0); + expect(handler).not.toHaveBeenCalled(); + expect(framework.getById(job.id)?.status).toBe("queued"); + // Never started, so no attempt was consumed. + expect(framework.getById(job.id)?.attempts).toBe(0); + }); + + it("leaves a retrying job for the next process, backoff intact", async () => { + framework.register("flaky.work", () => { + throw new Error("nope"); + }); + const job = framework.enqueue({ operation: "flaky.work" }); + framework.drainDueJobs(); + expect(framework.getById(job.id)?.status).toBe("retrying"); + + const report = await framework.shutdown(); + + expect(report.pending.map((entry) => entry.id)).toEqual([job.id]); + expect(report.interrupted).toHaveLength(0); + expect(framework.getById(job.id)?.status).toBe("retrying"); + expect(framework.getById(job.id)?.nextAttemptAt).toBeTypeOf("string"); + }); + + it("counts an already-settled job as settled and leaves it alone", async () => { + framework.register("done.work", () => {}); + const job = framework.enqueue({ operation: "done.work" }); + framework.drainDueJobs(); + + const report = await framework.shutdown(); + + expect(report.settled).toBe(1); + expect(report.pending).toHaveLength(0); + expect(report.interrupted).toHaveLength(0); + expect(framework.getById(job.id)?.status).toBe("succeeded"); + }); + + it("reports a handler that ignores the signal as interrupted", async () => { + // Never settles. shutdown must not hang on it, and must not classify it as + // merely pending: its side effect is unknown, which is the whole risk. + framework.register("stubborn.work", () => new Promise(() => {})); + framework.enqueue({ operation: "stubborn.work" }); + framework.drainDueJobs(); + + const report = await framework.shutdown({ timeoutMs: 20 }); + + expect(report.interrupted).toHaveLength(1); + expect(report.pending).toHaveLength(0); + expect(framework.inFlight()).toBeGreaterThanOrEqual(1); + }); + + it("treats a handler that throws on the signal as a cancellation", async () => { + framework.register("thrower.work", async (_params, context) => { + await sleep(5); + context.signal.throwIfAborted(); + }); + + const job = framework.enqueue({ operation: "thrower.work" }); + framework.drainDueJobs(); + const report = await framework.shutdown(); + + expect(report.interrupted.map((entry) => entry.id)).toEqual([job.id]); + // Recorded as cancelled, not as a fault: the handler did not fail, it was + // told to stop. + expect(framework.getById(job.id)?.lastError?.code).toBe("cancelled"); + }); + + it("dead-letters a cancelled attempt whose budget is spent", async () => { + framework = createWorkerFramework({ maxAttempts: 1 }); + framework.register("once.work", async (_params, context) => { + await sleep(5); + if (context.signal.aborted) return; + }); + + const job = framework.enqueue({ operation: "once.work" }); + framework.drainDueJobs(); + await framework.shutdown(); + + expect(framework.getById(job.id)?.status).toBe("dead_lettered"); + expect(framework.getById(job.id)?.status).not.toBe("succeeded"); + }); + + it("starts no new work once cancelling", async () => { + const handler = vi.fn(); + framework.register("after.work", handler); + + await framework.shutdown(); + framework.enqueue({ operation: "after.work" }); + + expect(framework.drainDueJobs().ran).toBe(0); + expect(framework.drainDueJobs().ran).toBe(0); + const job = framework.inspect()[0]; + expect(framework.processJob(job.id).ran).toBe(0); + expect(handler).not.toHaveBeenCalled(); + }); + + it("exposes the phase transitions idle -> cancelling -> stopped", async () => { + expect(framework.phase()).toBe("idle"); + + framework.cancel("because"); + expect(framework.phase()).toBe("cancelling"); + expect(framework.signal().aborted).toBe(true); + + await framework.shutdown(); + expect(framework.phase()).toBe("stopped"); + }); + + it("is idempotent: a second shutdown does not re-abort or double-count", async () => { + framework.register("once.work", () => {}); + framework.enqueue({ operation: "once.work" }); + + const first = await framework.shutdown({ reason: "first" }); + const second = await framework.shutdown({ reason: "second" }); + + expect(first.reason).toBe("first"); + // The first shutdown already stopped it, so the second reports the settled + // state rather than pretending to cancel again. + expect(second.phase).toBe("stopped"); + expect(second.reason).toBe("second"); + expect(second.settled).toBe(first.settled); + }); + + it("accepts work again after reset()", async () => { + const handler = vi.fn(); + framework.register("reuse.work", handler); + await framework.shutdown(); + + framework.reset(); + + expect(framework.phase()).toBe("idle"); + expect(framework.signal().aborted).toBe(false); + + // reset() also clears registrations, so a reused framework must re-register. + framework.register("reuse.work", handler); + framework.enqueue({ operation: "reuse.work" }); + expect(framework.drainDueJobs().succeeded).toBe(1); + expect(handler).toHaveBeenCalledTimes(1); + }); + + it("still records ordinary failures as failures, not cancellations", () => { + framework.register("boom.work", () => { + throw new WorkerCancelledError("unrelated"); + }); + + framework.enqueue({ operation: "boom.work" }); + const summary = framework.drainDueJobs(); + + // Nothing is cancelling, so a thrown error is a failure even if its class + // looks like a cancellation. + expect(summary.cancelled).toBe(0); + expect(summary.retried).toBe(1); + }); + + it("recognises both cancellation shapes", () => { + expect(isCancellation(new WorkerCancelledError())).toBe(true); + expect(isCancellation(new DOMException("aborted", "AbortError"))).toBe(true); + expect(isCancellation(new Error("handler failed"))).toBe(false); + expect(isCancellation(null)).toBe(false); + expect(isCancellation("nope")).toBe(false); + }); +}); + +describe("maintenance framework shutdown", () => { + it("cancels the prune job rather than reporting it done", async () => { + const framework = createMaintenanceFramework(); + const job = framework.enqueue({ operation: PRUNE_HORIZON_CLIENTS_OP }); + + const report = await framework.shutdown({ reason: "unload" }); + + // Never started, so it is owed to the next session rather than abandoned. + expect(report.pending.map((entry) => entry.operation)).toEqual([PRUNE_HORIZON_CLIENTS_OP]); + expect(framework.getById(job.id)?.status).toBe("queued"); + }); + + it("exposes cancel without waiting, for an unload handler", () => { + const framework = createMaintenanceFramework(); + framework.enqueue({ operation: PRUNE_HORIZON_CLIENTS_OP }); + + framework.cancel("pagehide"); + + expect(framework.phase()).toBe("cancelling"); + expect(framework.drainDueJobs().ran).toBe(0); + }); + + it("cancelMaintenance() reaches the shared singleton", () => { + // The exported helper targets the app-wide framework; a direct call must not + // throw and must leave it cancelling. + expect(() => cancelMaintenance("test_unload")).not.toThrow(); + }); +}); diff --git a/core/workers/maintenance.ts b/core/workers/maintenance.ts index b6b8884..39de31d 100644 --- a/core/workers/maintenance.ts +++ b/core/workers/maintenance.ts @@ -8,7 +8,12 @@ */ import { resetHorizonClients } from "@/core/horizon/client"; -import { createWorkerFramework, type WorkerFramework } from "@/core/workers/queue"; +import { + createWorkerFramework, + type ShutdownReport, + type WorkerContext, + type WorkerFramework +} from "@/core/workers/queue"; /** Prunes the memoised Horizon clients. Delayed so a rapid switch settles first. */ export const PRUNE_HORIZON_CLIENTS_OP = "maintenance.prune_horizon_clients"; @@ -21,7 +26,13 @@ const PRUNE_POLICY = { retryDelayMs: PRUNE_HORIZON_CLIENTS_DELAY_MS, maxAttempts function registerMaintenance(framework: WorkerFramework): WorkerFramework { framework.register( PRUNE_HORIZON_CLIENTS_OP, - () => { + (_params, context: WorkerContext) => { + // The prune is the exact job the issue names: a maintenance side effect + // that must not be reported as done when it was abandoned. Checking the + // signal first means a shutdown that lands between the job being picked up + // and the prune running leaves the memoised clients alone, and the + // framework records the attempt as cancelled rather than succeeded. + context.signal.throwIfAborted(); resetHorizonClients(); }, PRUNE_POLICY @@ -43,4 +54,23 @@ export function getMaintenanceFramework(): WorkerFramework { /** Creates a fresh, registered framework — the hook used by tests and the CLI runner. */ export function createMaintenanceFramework(): WorkerFramework { return registerMaintenance(createWorkerFramework()); -} \ No newline at end of file +} + +/** + * Cancels in-flight maintenance work and reports what each job was doing. + * + * Call this before the app tears down. The report is what a caller needs in + * order to know whether a prune still has to happen: a job listed as + * `interrupted` may have been abandoned mid-flight, while one in `pending` was + * never started and will be picked up by the next session. + */ +export async function shutdownMaintenance( + framework: WorkerFramework = getMaintenanceFramework() +): Promise { + return framework.shutdown({ reason: "maintenance_shutdown" }); +} + +/** Aborts without waiting, for an unload handler that cannot await. */ +export function cancelMaintenance(reason = "maintenance_cancel"): void { + getMaintenanceFramework().cancel(reason); +} diff --git a/core/workers/queue.ts b/core/workers/queue.ts index f87199c..77cb776 100644 --- a/core/workers/queue.ts +++ b/core/workers/queue.ts @@ -92,6 +92,15 @@ export interface SubmitResult { export interface WorkerContext { correlationId: string; attempt: number; + /** + * Aborted when the framework is cancelled or shut down (#142). + * + * A handler that does long work should observe this and stop: the framework + * will not mark an abandoned attempt `succeeded`, and a stopped attempt is + * retried rather than lost. `signal.throwIfAborted()` is the one-line way to + * bail at the next checkpoint. + */ + signal: AbortSignal; } export type JobHandler = ( @@ -99,6 +108,70 @@ export type JobHandler = ( context: WorkerContext ) => void | Promise; +/** Raised by the framework's own abort path, and safe to rethrow from a handler. */ +export class WorkerCancelledError extends Error { + constructor(message = "Worker cancelled") { + super(message); + this.name = "WorkerCancelledError"; + } +} + +/** + * Whether a thrown value is one of the cancellation shapes this framework (or a + * caller reusing the same pattern) produces. + * + * Note that the framework itself decides cancellation by its own signal, not by + * this predicate: a handler that throws an `AbortError` from an unrelated + * timeout has not been cancelled by a shutdown and must still be reported as a + * failure. This is exported for handlers that catch our abort and rethrow, and + * for tests. + */ +export function isCancellation(error: unknown): boolean { + if (error instanceof WorkerCancelledError) return true; + if (typeof error !== "object" || error === null) return false; + const name = (error as { name?: unknown }).name; + // `AbortController#abort(reason)` and `signal.throwIfAborted()` both surface + // as a DOMException named "AbortError". + return name === "AbortError" || name === "TimeoutError"; +} + +export type ShutdownPhase = "idle" | "cancelling" | "stopped"; + +/** A job as it stood when shutdown was requested. */ +export interface PendingJobReport { + id: string; + operation: string; + status: JobStatus; + attempts: number; + maxAttempts: number; +} + +export interface ShutdownOptions { + /** Recorded in the report and in telemetry. Defaults to "shutdown". */ + reason?: string; +} + +export interface ShutdownReport { + readonly reason: string; + readonly phase: ShutdownPhase; + /** + * Queued and retrying jobs. They were never started, so nothing was abandoned + * and the next process can pick them up unchanged. + */ + readonly pending: readonly PendingJobReport[]; + /** + * Jobs whose handler was in flight when the abort landed. Their side effect + * may be partial, which is why they are reported separately and moved to + * `retrying` rather than settled. + */ + readonly interrupted: readonly PendingJobReport[]; + /** Jobs that had already settled before the abort. Untouched. */ + readonly settled: number; +} + +/** How long `shutdown()` waits for in-flight handlers before giving up. */ +export const SHUTDOWN_DRAIN_TIMEOUT_MS = 5_000; + export interface RetryPolicy { /** Base delay in ms; the nth retry waits `retryDelayMs * 2^(n - 1)`. */ retryDelayMs: number; @@ -114,8 +187,26 @@ export interface RunSummary { failed: number; retried: number; deadLettered: number; + /** + * Attempts abandoned because the framework was cancelled or shut down. Counted + * separately from `failed` because they are neither successes nor faults: the + * handler stopped early and the attempt is owed another run. + */ + cancelled: number; + /** Jobs still running when the summary was produced (async handlers). */ + inFlight: number; } +const EMPTY_SUMMARY = (): RunSummary => ({ + ran: 0, + succeeded: 0, + failed: 0, + retried: 0, + deadLettered: 0, + cancelled: 0, + inFlight: 0 +}); + const DEFAULT_POLICY: RetryPolicy = { retryDelayMs: 250, maxAttempts: 3, @@ -136,6 +227,37 @@ export interface WorkerFramework { idempotency(): IdempotencyStore; /** Executes every due job once and returns the lifecycle summary. */ drainDueJobs(now?: number): RunSummary; + /** + * Like `drainDueJobs`, but waits for async handlers to settle, so the summary + * is final rather than provisional (#142). + */ + drain(now?: number): Promise; + /** + * Cancel in-flight work and stop starting new work, then report what each job + * was doing when the abort landed (#142). + * + * The contract per state: + * queued / retrying — never started, left exactly as they are and reported + * as `pending`. Nothing was abandoned, so the next process re-runs them. + * running — the handler had started. Its side effect may be + * partial, so the attempt is moved to `retrying` (or `dead_lettered` if + * the budget is spent) and reported as `interrupted`. It is never marked + * `succeeded`: an abandoned attempt that claims completion is exactly the + * failure this exists to prevent. + * succeeded / dead_lettered — settled before the abort, untouched. + * + * Waits up to `timeoutMs` for in-flight handlers; one that ignores its signal + * is left running and reported as interrupted. + */ + shutdown(options?: ShutdownOptions & { timeoutMs?: number }): Promise; + /** Abort without waiting — the signal flips immediately, drain is not awaited. */ + cancel(reason?: string): void; + /** Current phase: `idle`, `cancelling` while draining, then `stopped`. */ + phase(): ShutdownPhase; + /** The signal every handler's context carries. */ + signal(): AbortSignal; + /** Jobs started but not yet settled. */ + inFlight(): number; /** Idempotent reprocess of a single job by id. */ processJob(id: string): RunSummary; /** @@ -188,6 +310,16 @@ export function createWorkerFramework( const mergedPolicy: RetryPolicy = { ...DEFAULT_POLICY, ...policy }; const registrations = new Map(); const jobs = new Map(); + /** + * Cancellation state (#142). One controller covers the whole framework, so a + * single `cancel()` reaches every handler that is in flight, and every job + * started afterwards sees an already-aborted signal. + */ + let controller = new AbortController(); + let shutdownPhase: ShutdownPhase = "idle"; + let shutdownReason: string | null = null; + /** Handlers started and not yet settled. */ + const inFlightSet = new Set>(); const idempotencyStore = createIdempotencyStore({ ttlMs: mergedPolicy.idempotencyTtlMs }); const quotaStore = options.quota ?? createQuotaStore({ storage: createMemoryQuotaStorage() }); @@ -219,7 +351,15 @@ export function createWorkerFramework( return Date.parse(job.nextAttemptAt) <= now; } - function runJob(job: JobPayload): void { + /** + * Run one job to a settled or retryable state. + * + * Returns a promise when the handler is async, so callers can wait for a final + * outcome. The handler is *awaited* on purpose (#142): it used to be invoked + * and discarded, so an async handler was marked `succeeded` before it had done + * anything and its rejection escaped as an unhandled rejection. + */ + function runJob(job: JobPayload): void | Promise { job.attempts += 1; const registration = registrations.get(job.operation); @@ -240,62 +380,150 @@ export function createWorkerFramework( } move(job, "start", "drained"); - try { - registration.handler(job.params, { - correlationId: job.correlationId, - attempt: job.attempts - }); - move(job, "succeed", "handler_returned"); - job.lastError = undefined; - emitTelemetry({ - op: "worker.run", - actorType: "worker", - result: "success", - correlationId: job.correlationId, - payload: { operation: job.operation, attempt: job.attempts } - }); - } catch (error) { - const code = errorCodeOf(error); - job.lastError = { - code, - message: error instanceof Error ? error.message : "unknown failure" - }; - // The job's own budget wins; the registration only supplies the default. - if (job.attempts >= job.maxAttempts) { - move(job, "exhaust", "attempts_exhausted"); + const settle = (error?: unknown): void => { + // Cancellation wins over both outcomes, and it is decided *only* by our own + // signal. A handler that threw because the abort reached it has not failed, + // and one that returned after the abort has not succeeded: its side effect + // is unknown, so neither `succeed` nor the failure path may claim + // otherwise. + // + // Keying off the signal rather than the error's class matters: a handler + // with its own timeout throws a cancellation-shaped error for an unrelated + // reason, and that is a failure this framework should report, not a + // shutdown it did not initiate. + if (controller.signal.aborted) { + const cancelledError = + error === undefined + ? new WorkerCancelledError( + `Cancelled during ${job.operation}: ${shutdownReason ?? "shutdown"}` + ) + : (error as Error); + job.lastError = { + code: "cancelled", + message: cancelledError instanceof Error ? cancelledError.message : "cancelled" + }; + if (job.attempts >= job.maxAttempts) { + move(job, "exhaust", "cancelled_attempts_exhausted"); + } else { + move(job, "fail", "cancelled"); + } emitTelemetry({ - op: "worker.exhausted", + op: "worker.cancelled", actorType: "worker", result: "failure", correlationId: job.correlationId, - errorCode: code, - payload: { operation: job.operation, attempts: job.attempts } + errorCode: "cancelled", + payload: { + operation: job.operation, + attempt: job.attempts, + reason: shutdownReason ?? "cancel" + } }); - } else { - move(job, "fail", "handler_threw"); - job.nextAttemptAt = new Date( - Date.now() + retryBackoff(registration.retryDelayMs, job.attempts) - ).toISOString(); + return; + } + + if (error === undefined) { + move(job, "succeed", "handler_returned"); + job.lastError = undefined; emitTelemetry({ - op: "worker.retry", + op: "worker.run", actorType: "worker", - result: "failure", + result: "success", correlationId: job.correlationId, - errorCode: code, payload: { operation: job.operation, attempt: job.attempts } }); + return; } + + handleFailure(job, registration, error); + }; + + const context: WorkerContext = { + correlationId: job.correlationId, + attempt: job.attempts, + signal: controller.signal + }; + + try { + const result = registration.handler(job.params, context); + if (result && typeof (result as Promise).then === "function") { + const promise = (result as Promise).then( + () => settle(), + (error: unknown) => settle(error) + ); + inFlightSet.add(promise); + void promise.finally(() => inFlightSet.delete(promise)); + return promise; + } + settle(); + } catch (error) { + settle(error); + } + } + + /** The failure path, shared by the sync and async routes. */ + function handleFailure( + job: JobPayload, + registration: { retryDelayMs: number }, + error: unknown + ): void { + const code = errorCodeOf(error); + job.lastError = { + code, + message: error instanceof Error ? error.message : "unknown failure" + }; + + // The job's own budget wins; the registration only supplies the default. + if (job.attempts >= job.maxAttempts) { + move(job, "exhaust", "attempts_exhausted"); + emitTelemetry({ + op: "worker.exhausted", + actorType: "worker", + result: "failure", + correlationId: job.correlationId, + errorCode: code, + payload: { operation: job.operation, attempts: job.attempts } + }); + } else { + move(job, "fail", "handler_threw"); + job.nextAttemptAt = new Date( + Date.now() + retryBackoff(registration.retryDelayMs, job.attempts) + ).toISOString(); + emitTelemetry({ + op: "worker.retry", + actorType: "worker", + result: "failure", + correlationId: job.correlationId, + errorCode: code, + payload: { operation: job.operation, attempt: job.attempts } + }); } } function bumpSummary(summary: RunSummary, job: JobPayload): void { - if (job.status === "succeeded") summary.succeeded += 1; + // A cancelled attempt ends in `retrying` like any other non-terminal state, + // so it is counted from `lastError.code` rather than from the status: it is + // not a fault, and reporting it as one is what made cancellation look like a + // broken worker. + if (job.lastError?.code === "cancelled") summary.cancelled += 1; + else if (job.status === "succeeded") summary.succeeded += 1; else if (job.status === "retrying") summary.retried += 1; else if (job.status === "dead_lettered") summary.deadLettered += 1; else summary.failed += 1; } + /** Snapshot of one job for the shutdown report. */ + function reportJob(job: JobPayload): PendingJobReport { + return { + id: job.id, + operation: job.operation, + status: job.status, + attempts: job.attempts, + maxAttempts: job.maxAttempts + }; + } + // Named so `submit()` can reach the public enqueue without re-entering the // object literal while it is being built. const framework: WorkerFramework = { @@ -395,28 +623,169 @@ export function createWorkerFramework( }, drainDueJobs(now = Date.now()): RunSummary { - const summary: RunSummary = { ran: 0, succeeded: 0, failed: 0, retried: 0, deadLettered: 0 }; + const summary = EMPTY_SUMMARY(); + // Once cancelling, no *new* work starts (#142). Work already in flight + // still settles; only the queue is frozen. + if (shutdownPhase !== "idle") { + summary.inFlight = inFlightSet.size; + return summary; + } + const due = [...jobs.values()].filter((job) => isDue(job, now)); for (const job of due) { summary.ran += 1; - runJob(job); + const settled = runJob(job); + // An async handler has not finished yet, so its outcome is counted when + // it settles; the provisional summary reports it as in flight. + if (settled) continue; bumpSummary(summary, job); } + + summary.inFlight = inFlightSet.size; return summary; }, + async drain(now = Date.now()): Promise { + const summary = EMPTY_SUMMARY(); + if (shutdownPhase !== "idle") { + summary.inFlight = inFlightSet.size; + return summary; + } + + const due = [...jobs.values()].filter((job) => isDue(job, now)); + + // Jobs an earlier synchronous drain left running. This call waits for them + // below, so they belong in the summary it returns — otherwise a caller who + // drained and then awaited would see the work complete and a summary that + // never mentions it. + const inherited = [...jobs.values()].filter((job) => job.status === "running"); + const counted = new Set(); + + for (const job of due) { + summary.ran += 1; + counted.add(job.id); + const settled = runJob(job); + if (!settled) bumpSummary(summary, job); + } + + // Await every handler still running, not only the ones this call started. + await Promise.all([...inFlightSet]); + + for (const job of due) bumpSummary(summary, job); + for (const job of inherited) { + if (counted.has(job.id)) continue; + counted.add(job.id); + summary.ran += 1; + bumpSummary(summary, job); + } + summary.inFlight = inFlightSet.size; + return summary; + }, + + cancel(reason = "cancel"): void { + if (shutdownPhase !== "idle") return; + shutdownReason = reason; + shutdownPhase = "cancelling"; + controller.abort(new WorkerCancelledError(`Worker cancelled: ${reason}`)); + }, + + async shutdown(options = {}): Promise { + const reason = options.reason ?? "shutdown"; + const timeoutMs = options.timeoutMs ?? SHUTDOWN_DRAIN_TIMEOUT_MS; + + // What each job was doing *when the abort landed*, captured before the + // drain. Classifying afterwards would miss a handler that ignores its + // signal and never settles: it is still `running`, its side effect is + // unknown, and it is exactly the job the caller most needs told about. + const atAbort = new Map( + [...jobs.values()].map((job) => [job.id, job.status]) + ); + + framework.cancel(reason); + + // Give in-flight handlers a chance to observe the signal and stop. One + // that ignores it is left running; the report below says so, because a + // job whose side effect is still in progress is precisely what a caller + // needs to know about at shutdown. + if (inFlightSet.size > 0) { + await Promise.race([ + Promise.allSettled([...inFlightSet]), + new Promise((resolve) => setTimeout(resolve, timeoutMs)) + ]); + } + + shutdownPhase = "stopped"; + + const pending: PendingJobReport[] = []; + const interrupted: PendingJobReport[] = []; + let settledCount = 0; + + for (const job of jobs.values()) { + const statusAtAbort = atAbort.get(job.id) ?? job.status; + if (statusAtAbort === "succeeded" || statusAtAbort === "dead_lettered") { + settledCount += 1; + } else if (statusAtAbort === "running" || job.lastError?.code === "cancelled") { + // Was running, or settled as cancelled while we waited. Either way the + // attempt was abandoned and is owed another run. + interrupted.push(reportJob(job)); + } else { + // Queued or retrying: never started, so nothing was abandoned. + pending.push(reportJob(job)); + } + } + + emitTelemetry({ + op: "worker.shutdown", + actorType: "worker", + result: "success", + payload: { + reason, + pending: pending.length, + interrupted: interrupted.length, + settled: settledCount + } + }); + + return { + reason, + phase: shutdownPhase, + pending, + interrupted, + settled: settledCount + }; + }, + + phase(): ShutdownPhase { + return shutdownPhase; + }, + + signal(): AbortSignal { + return controller.signal; + }, + + inFlight(): number { + return inFlightSet.size; + }, + processJob(id: string): RunSummary { - const summary: RunSummary = { ran: 0, succeeded: 0, failed: 0, retried: 0, deadLettered: 0 }; + const summary = EMPTY_SUMMARY(); const job = jobs.get(id); if (!job) return summary; // Idempotent reprocessing: a terminal job has no legal `start` event, so // it is left untouched instead of running its handler a second time. if (workerJobMachine.isTerminal(job.status)) return summary; + // Reprocessing is starting new work, so it is refused while cancelling + // for the same reason a drain is. + if (shutdownPhase !== "idle") { + summary.inFlight = inFlightSet.size; + return summary; + } summary.ran += 1; - runJob(job); - bumpSummary(summary, job); + const settled = runJob(job); + if (!settled) bumpSummary(summary, job); + summary.inFlight = inFlightSet.size; return summary; }, @@ -454,6 +823,12 @@ export function createWorkerFramework( reset(): void { jobs.clear(); registrations.clear(); + inFlightSet.clear(); + // A reset framework accepts work again: a fresh controller replaces the + // aborted one, and the phase returns to idle, so the next drain runs. + shutdownPhase = "idle"; + shutdownReason = null; + controller = new AbortController(); idempotencyStore.reset(); // A shared quota store belongs to the caller; only an owned one is cleared. if (!options.quota) quotaStore.reset();