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();