@sdxc/jobs
Declared background jobs dispatched over a pluggable queue backend
- Depends on
@standard-schema/spec- Used by
- uptime, auth-saas
- Source
- packages/jobs
Background jobs declared in one map, run by one dispatcher, over a queue backend of your choosing.
Installation
npm add @sdxc/jobs
Two adapters ship: @sdxc/jobs/cloudflare for Cloudflare Queues, and @sdxc/jobs/memory
for tests and for anything running without a platform behind it. Writing a third means
implementing JobQueue from @sdxc/jobs/queue and running @sdxc/jobs/conformance
against it.
@cloudflare/workers-types, remix and vitest are optional peers, needed only by the
Cloudflare adapter, the context's typed key store, and the conformance suite respectively.
Overview
A job is declared, not implemented, in the map: its name, the payload it carries, the
cron it runs on, and whatever meta its own callers read off it. The map is declaration and nothing
else — no handler, no queue — so importing it costs its schemas. The handler lives in
its own module and is loaded only when a message for it arrives.
Every job is enqueuable. A job that declares a cron is additionally schedulable, and
the two share one path: a cron delivery enqueues a message and returns, so a scheduled
job gets the same middleware, timeout, logging, retries, and dead-letter queue as any
other. Nothing runs inside the scheduled handler.
The shape is the one an HTTP router has: a map of addressable definitions, a runtime
that maps handlers onto them, and middleware that publishes values into a typed context.
Context keys come from
remix/router,
so one key serves an HTTP middleware and a job middleware alike.
Usage
Declaring the map
The key is the job's name, and the name is what the message is addressed to — a wire contract, since messages enqueued by one deploy are consumed by the next.
import { job, jobs } from "@sdxc/jobs";
import * as s from "remix/data-schema";
export default jobs({
clean: job({ cron: "0 0 * * *", meta: { monitorId: "8f1c…" } }),
checkHttp: job({ input: s.object({ monitorId: s.string() }) }),
digests: {
daily: job({ cron: "0 8 * * *" }),
weekly: job({ cron: "0 9 * * 1" }),
},
});
Nested keys are dot-joined, so digests.daily is that job's name.
Writing a handler
import { createJobHandler } from "@sdxc/jobs";
import jobs from "~/app/jobs";
export default createJobHandler(jobs.checkHttp, async (ctx) => {
let monitor = await ctx.log.time("db", () => ctx.database.find(ctx.input.monitorId));
ctx.log.set({ monitor: { id: monitor.id, region: monitor.region } });
ctx.log.note("check.started");
});
ctx.log is the run's record, emitted once when the run ends: set() a field worth
querying by, note() what is worth reading once a query has found the record, time()
anything with a duration.
Wiring the worker
import { createJobDispatcher } from "@sdxc/jobs";
import * as cloudflare from "@sdxc/jobs/cloudflare";
import { env } from "cloudflare:workers";
import { logger } from "~/bootstrap/logger";
import jobs from "~/app/jobs";
export const dispatcher = createJobDispatcher({
logger,
queue: cloudflare.queue(() => env.QUEUE),
middleware: [database()],
timeout: "5 minutes",
/**
* Runs once the ending is decided and before the delivery is settled, so a report that
* has to reach a service is finished by then.
*/
async onEnd(ctx, status) {
if (status.type !== "done") return;
let watched = ctx.of(jobs.clean);
if (watched !== null) await uptime(watched.meta.monitorId);
},
});
dispatcher.map(jobs.clean, () => import("~/app/jobs/clean"));
dispatcher.map(jobs.checkHttp, () => import("~/app/jobs/check-http"));
/** Both job handlers, bound to the dispatcher they delegate to. */
const handlers = cloudflare.worker(dispatcher, {
deadLetterQueue: "ping-dlq",
deadLetter: () => env.DLQ,
});
export default {
async scheduled(controller) {
await handlers.scheduled(controller);
},
async queue(batch) {
await handlers.queue(batch);
},
} satisfies ExportedHandler<Cloudflare.Env>;
Enqueuing a job
await dispatcher.enqueue(jobs.checkHttp, { monitorId: monitor.id });
await dispatcher.enqueueMany(
jobs.checkHttp,
monitors.map((m) => ({ monitorId: m.id })),
);
A call site that should not pull the dispatcher — and its middleware, loaders, and every handler behind them — into its module graph builds the body instead and sends it through whatever the app already writes with:
await sendQueueBatch([messageBody(jobs.checkHttp, { monitorId: monitor.id })]);
What gets logged
Every invocation the dispatcher serves emits one record through the logger it was
given, so a cron trigger, a queue batch, and each job in the batch are each one query
away. Fields are flat scalars under dotted keys; the kind says which they are.
| Kind | Opened by | Fields |
|---|---|---|
cron | dispatcher.tick() | cron.expression, cron.scheduled_at, jobs.enqueued |
queue | dispatcher.deliverBatch() | queue.name, queue.batch_size, job.count, and one of jobs.done, jobs.retried, jobs.refused, jobs.failed, jobs.timed_out, jobs.dead_lettered per message |
job | Each message, under the batch | job.name, job.id, job.attempts, job.batch_size, job.cron when declared, job.ending, and whatever the handler set() |
A job log's outcome follows its ending: ok for done, degraded for retry (with
job.delay_s when a backoff was asked for), error for refuse, timeout, and
failed, each carrying the failure under error.*. A message the dispatcher refuses
before dispatch and a message delivered on the dead-letter queue get a job log too,
ending refused (with job.refusal) or dead_letter (with job.dead_letter) and
keeping the body as job.body, so "every job that did not end well today" is one query
over kind: "job" and outcome. The batch log degrades when any of its jobs did.
{
"service": "uptime",
"kind": "job",
"job.name": "checkHttp",
"job.id": "…",
"job.attempts": 2,
"job.batch_size": 80,
"job.ending": "retry",
"job.delay_s": 300,
"monitor.id": "mon_…",
"db.count": 1,
"db.duration_ms": 12.4,
"outcome": "degraded",
"duration_ms": 412,
"notes": [{ "at": 410.2, "level": "warn", "name": "job.retry", "reason": "Rate limited" }],
}
A dispatcher built without a logger emits the same records, carrying no service,
which is how one under test runs unchanged.
API
jobs(tree: JobTree): JobMap
Builds the app's job map, naming every leaf after the key it is filed under.
Parameters:
tree: The declared jobs, keyed by the name each is known by on the wire. Groups may nest, and a nested job's name is its dot-joined path
Returns:
The same shape, with every leaf a
JobDefinitionthat knows its name, schedule, schema, andmeta
Example:
export default jobs({ clean: job({ cron: "0 0 * * *" }) });
job(options?: JobOptions): JobLeaf
Declares one job for a map. Holds no handler, so importing a map costs its schemas.
Parameters:
options.input: Schema the payload is parsed against before the handler runs. The payload travels under a key of its own, so an object schema and a bare value both serve. Absent means no payloadoptions.cron: Cron expression this job is enqueued on, spelled exactly as the matching trigger inwrangler.jsonc. Parsed here, so an expression the platform would reject throws at declaration; the type is five space-separated fields, so anything coarser than that is a compile error. Absent means it is only ever enqueued explicitlyoptions.meta: Anything this job's own callers read off it, unconstrained and inferred as written. Nothing in the package reads it — a dispatcher-level hook reaches it throughctx.of(job), which gives it back this exact type
Returns:
The leaf to file under the name this job is known by
Throws:
InvalidCronExpressionwhencronis not an expression the platform would accept, naming the offending field and its position
Example:
let checkHttp = job({
input: s.object({ monitorId: s.string() }),
meta: { monitorId: "8f1c…" } satisfies Monitored,
});
job({ cron: "0 99 * * *" }); // Throws: out-of-range in the hour field at position 2
job({ cron: "invalid" }); // Type error: not five fields
messageBody(job: JobDefinition, input?): JSONValue
Builds the body one message carries: the job it names, and the payload beside it under
body. For a call site that sends through the app's own queue helper rather than through
the dispatcher.
{ "job": "notify", "body": { "monitorId": "…" } }
A job declaring no payload carries no body. The payload keeps a namespace of its own, so
a job whose input declares a job or a type field carries it intact.
Example:
await sendQueueBatch([messageBody(jobs.notify, { monitorId: monitor.id })]);
createJobHandler(job: JobDefinition, handler: JobHandlerFunction): JobHandler
Pairs a handler with the job it runs. Passing the job is what types ctx.input, and it
is what the dispatcher checks the handler was mapped to.
Parameters:
job: The job this handler runs, from the app's maphandler: The work, receiving one context
Returns:
The handler, carrying the job it belongs to
Example:
export default createJobHandler(jobs.clean, async (ctx) => {
let removed = await ctx.log.time("db", () => purge(ctx.database));
ctx.log.set({ rows: { removed } });
});
createJobDispatcher(options?: JobDispatcherOptions): JobDispatcher
Builds the registry both worker handlers delegate to.
Parameters:
options.logger: The worker's logging configuration, fromcreateLogger(). Every cron, queue, and job log this dispatcher opens carries it; without one they carry no serviceoptions.queue: The backend this dispatcher enqueues through, from@sdxc/jobs/cloudflare,@sdxc/jobs/memory, or an adapter of your own. A dispatcher without one still runs what a backend delivers, and throws when asked to enqueueoptions.middleware: Chain every job runs inside, in the order declaredoptions.timeout: How long a job gets beforectx.signalaborts and the dispatcher stops waitingoptions.onEnd:(ctx, status) => void | Promise<void>— runs once a delivery's ending is decided and before it is settled, which is what makes it the place for work that has to reach a service before the message is acked. Its failure is never the job's: anything it throws is recorded asjob.hook_failed, and running past the settle grace is recorded asjob.hook_overran, with the ending settling as decided either wayoptions.maxAttempts: the attempts a message gets before it is dead-lettered, for a queue whoseretriesare this package's to count. Ignored by a backend counting its own
Returns:
A dispatcher with
map,enqueue,enqueueMany,deliver,deliverBatch,tick,mapped, andcrons
Example:
export const dispatcher = createJobDispatcher({ logger, middleware: [database()] });
dispatcher.map(job: JobDefinition, load): void
Registers where a job's handler comes from. Throws when that name is already mapped.
Parameters:
job: The job, from the app's mapload: A loader returning the handler's module, or the handler itself. A loader is awaited once per isolate, and only after a message has matched and parsed
Example:
dispatcher.map(jobs.clean, () => import("~/app/jobs/clean"));
dispatcher.enqueue(job: JobDefinition, input?): Promise<void>
Enqueues one message for a job. Takes exactly what the job's schema accepts, and no argument at all for a job that declares none. Loads no handler.
Example:
await dispatcher.enqueue(jobs.verifyDomain, { teamDomainId: domain.id });
await dispatcher.enqueue(jobs.clean);
dispatcher.enqueueMany(job: JobDefinition, inputs): Promise<void>
Enqueues one message per input in a single write. Enqueuing nothing writes nothing.
Example:
await dispatcher.enqueueMany(
jobs.notify,
changes.map((change) => ({ id: change.id })),
);
dispatcher.deliver(delivery: JobDelivery, batchSize?): Promise<Settlement>
Runs the job one delivery names and answers with what it ended as, for the caller's backend to apply. Names no platform, so a test drives it with a plain object.
Example:
let settlement = await dispatcher.deliver({ id: "m1", attempts: 1, body: { job: "clean" } });
dispatcher.deliverBatch(deliveries: JobDelivery[], options): Promise<void>
Runs every delivery that arrived together, inside one queue log, settling each through
options.apply as that one finishes rather than when the last one does. With
deadLettered, every delivery is recorded and acked instead of dispatched — which queue
a batch arrived on is the backend's to know, so the backend is the one that says so.
Example:
await dispatcher.deliverBatch(deliveries, {
queue: "ping",
deadLettered: batch.queue === "ping-dlq",
apply: (delivery, settlement) => settle(delivery, settlement),
});
dispatcher.tick(options: { now: Date; only?: CronExpression }): Promise<void>
Enqueues every mapped job that is due, in one write, inside one cron log that records
how many. Runs none of them, and loads no handler.
With only, a job is due when it declares that exact expression — the platform has
already decided which minute this is, so no schedule is evaluated. Without it, a job is
due when its schedule fires in now's minute, which is what drives a backend that has no
triggers behind it.
Example:
async scheduled(controller) {
await dispatcher.tick({ now: new Date(controller.scheduledTime), only: controller.cron });
}
dispatcher.mapped: JobDefinition[]
Every job a handler was mapped onto, for asserting that a map has no leaf nobody runs.
Example:
let names = new Set(dispatcher.mapped.map((job) => job.name));
expect([...names]).toEqual(expect.arrayContaining(["clean", "checkHttp"]));
dispatcher.crons: string[]
The distinct schedules the mapped jobs declare, for asserting that the code and
wrangler.jsonc agree.
Example:
expect([...dispatcher.crons].sort()).toEqual([...config.triggers.crons].sort());
JobContext
The context every middleware and handler shares for one delivery.
new JobContext(job: JobDefinition, init: JobContextInit)
Builds a context. The dispatcher does this per delivery; a test does it to call a handler directly.
Parameters:
job: The job being run, which suppliesnameandcron, and types whatctx.of()answersinit.id: The queue message's idinit.attempts: Which delivery of this message this is, counting from oneinit.input: The payload, already parsed against the job's schemainit.batchSize: How many messages share this invocation. Defaults to oneinit.log: Where this job's fields go, aLogfrom@sdxc/logger. One is created when omitted, under the current log when there is oneinit.signal: Aborts when the job's timeout expires. Never aborts when omitted
Example:
let ctx = new JobContext(jobs.clean, { id: "message-1", attempts: 1 });
Properties
ctx.input: The payload, typed by the job's schemactx.name,ctx.cron: The job's own declarationctx.of(job): This delivery's view of one job — its parsedinputand declaredmeta, both with that job's own types — ornullfor a delivery of another job. The name is checked, so it narrows rather than assertsctx.id,ctx.attempts: The delivery's identifier and delivery countctx.batchSize: How many messages share this invocationctx.log: The run's record —set()fields,note()breadcrumbs,time()durations — emitted once when the run endsctx.signal: Aborts when the timeout expires
ctx.ack(reason?: string): never
Finishes here: the delivery is acked and the run reported as completed. What returning
does, from anywhere in the call stack. Throws Ack.
ctx.retry(options?: RetryOptions): never
Gives up on this delivery and asks for another. Throws Retry.
Parameters:
options.delay: How long the platform holds the message, as a durationoptions.cause: What led here
Example:
if (response.status === 429) ctx.retry({ delay: "5 minutes" });
ctx.exit(reason?: string, options?: ErrorOptions): never
Gives up for good: the delivery is acked, because a redelivery reaches the same result,
and the run is reported as a failure. Throws NonRetriable.
Example:
if (team === null) ctx.exit("Team no longer exists");
ctx.timeout(reason?: string): never
Gives up because time ran out: the delivery is retried, and the run is reported as one
that never finished. Throws Timeout.
ctx.get(key: ContextKey): value | undefined
Reads a value some middleware published, falling back to the key's default.
ctx.require(key: ContextKey): value
The same, refusing to continue when nothing published one. This is how one middleware
reads what an earlier middleware put on the context: inside a chain the context is the
bare one, so an installed property like ctx.database is not visible there — only the
key is.
Example:
let database = ctx.require(Database);
ctx.has(key: ContextKey): boolean
Whether a value has been published for a key.
ctx.set(key: ContextKey, value, options?: { property: string }): void
Publishes a value, optionally installing it as a direct property so handlers read
ctx.database rather than ctx.get(Database).
Example:
ctx.set(Database, connect(), { property: "database" });
Endings
Each verb throws its own class. Job groups them for instanceof, and
@sdxc/jobs/errors exports each one individually, which is what a type position
needs.
| Ending | Thrown by | What the dispatcher does | job.ending | Outcome |
|---|---|---|---|---|
Job.Ack | ctx.ack() | Ack the delivery | done | ok |
Job.Retry | ctx.retry() | Retry, holding for delay; note job.retry | retry | degraded |
Job.NonRetriable | ctx.exit() | Ack; fail() the log with the error and its cause | refuse | error |
Job.Timeout | ctx.timeout() | Retry; fail() the log | timeout | error |
Returning normally ends done like Job.Ack; anything else thrown ends failed with
outcome error and is rethrown so the platform retries the invocation. An ending decided
once the timeout has aborted the signal is recorded as timeout whatever it was:
ctx.ack() still acks the delivery, and returning normally lets the message come back.
Example:
import { Job } from "@sdxc/jobs";
import type { Retry } from "@sdxc/jobs/errors";
function holdFor(error: Retry) {
return error.delay;
}
try {
await charge(invoice);
} catch (error) {
if (error instanceof Job.Ending) throw error; // Never swallow an ending.
ctx.retry({ reason: "Charge failed", delay: "1 minute", cause: error });
}
Job.Ending is the base all four share, so a broad catch re-throws whichever one the
handler chose without naming each.
Types
JobMiddleware<Effect>
type JobMiddleware<Effect extends ContextEffect = EmptyContextEffect> = (
ctx: JobContext,
next: () => Promise<void>,
) => void | Promise<void>;
The Effect names what the middleware publishes — { key, value, property } — which is
what makes an installed property visible to handlers.
JobTypes
Augmented by an app to name the context its handlers receive.
declare module "@sdxc/jobs" {
interface JobTypes {
context: JobDispatcherContext<typeof dispatcher>;
}
}
JobQueue
The backend a dispatcher enqueues through and is delivered from, at
@sdxc/jobs/queue. Every call answers with a Result, so no backend call throws.
interface JobQueue {
readonly retries: "backend" | "core";
send(messages: JobMessage[]): Promise<Result<void, JobQueueError>>;
claim?(options: ClaimOptions): Promise<Result<JobDelivery[], JobQueueError>>;
settle?(delivery: JobDelivery, settlement: Settlement): Promise<Result<void, JobQueueError>>;
}
claim and settle are present together on a backend that is pulled, and absent on one
that pushes deliveries into the dispatcher itself. retries says who counts attempts:
"backend" leaves the ceiling where the platform already states it, "core" hands it to
the dispatcher's maxAttempts.
worker(dispatcher, options) takes the two things only a backend can answer for: the
deadLetterQueue it also consumes, whose batches it has the dispatcher record and ack
rather than dispatch, and the deadLetter binding a refused body is written to. That
write is what gets a message no redelivery can fix onto that queue, since this platform
reaches a dead-letter queue only by exhausting retries. A worker given neither dispatches
every batch it receives and acks a refused body where it stands.
Two adapters ship. @sdxc/jobs/cloudflare exports queue and worker;
@sdxc/jobs/memory exports queue. Both are reached through a namespace import, since
queue collides with the first local variable holding one.
A memory queue adds what a test drives it with: drain(deliver) runs everything
claimable and answers with what each delivery ended as, messages and deadLettered
read back what it still holds, and reset() empties both.
import * as cloudflare from "@sdxc/jobs/cloudflare";
export const dispatcher = createJobDispatcher({
queue: cloudflare.queue(() => env.QUEUE),
});
// bootstrap/worker.ts
const handlers = cloudflare.worker(dispatcher);
export default {
scheduled: handlers.scheduled,
queue: handlers.queue,
};
createUptimeReporter(options)
The cron-monitor ping, at @sdxc/jobs/uptime. Wired to nothing: a dispatcher's onEnd
calls it, so a service having a bad minute cannot become a reason to redeliver work that
already succeeded. Bound to a token once, so a call site holds a monitor id and nothing
else.
import { createUptimeReporter } from "@sdxc/jobs/uptime";
const uptime = createUptimeReporter({ token: () => env.UPTIME_CRON_API_KEY });
It answers with a Result rather than throwing, which is what makes discarding the
outcome an act rather than an oversight. UptimeError.code is one of refused (the
service answered and said no), unreachable (no answer at all), or unconfigured (no
token resolved, so nothing was sent).
Either call works. Awaited, the ping finishes inside the settlement barrier and a failure
reaches that run's own record; handed to waitUntil from cloudflare:workers, the
delivery settles at once and the outcome lands nowhere, since the record is written before
the promise resolves. Await it unless a batch is large enough for the round trips to
matter.
Pattern: Middleware That Provides A Database
Middleware publishes what handlers read, so a handler names no container and a test provides its own.
import type { JobMiddleware } from "@sdxc/jobs";
import { createContextKey } from "remix/router";
export const Database = createContextKey<Database>();
export function database(): JobMiddleware<{
key: typeof Database;
value: Database;
property: "database";
}> {
return async (ctx, next) => {
ctx.set(Database, connect(), { property: "database" });
await next();
};
}
Pattern: One Body Of Work, Several Schedules
Two leaves that differ only by schedule share their work through a plain function, and
each gets its own thin handler module. A handler is paired with exactly one job — the
dispatcher refuses a handler mapped to a different one — so the sharing happens below
createJobHandler, not around it.
// app/jobs/send-team-digests.ts
export async function sendTeamDigests(ctx: JobHandlerContext<undefined>, period: "day" | "week") {
await mailDigests(ctx.database, period);
}
// app/jobs/send-team-daily-digests.ts
export default createJobHandler(jobs.sendTeamDailyDigests, (ctx) => sendTeamDigests(ctx, "day"));
// app/jobs/send-team-weekly-digests.ts
export default createJobHandler(jobs.sendTeamWeeklyDigests, (ctx) => sendTeamDigests(ctx, "week"));
Each leaf keeps its own meta and its own loader, so one schedule failing is one
monitor alerting.
Pattern: Cooperative Cancellation
A timeout aborts ctx.signal and stops the dispatcher waiting; it cannot stop a
handler. A loop that checks between iterations gives up cleanly, and a handler whose
work is already durable acks instead so a redelivery does not repeat it.
export default createJobHandler(jobs.sendDigests, async (ctx) => {
for (let team of await teamsDue(ctx.database)) {
if (ctx.signal.aborted) ctx.ack(); // The mail already sent must not be sent twice.
await sendDigest(team, { signal: ctx.signal });
}
});
Pattern: A Job That Enqueues Other Jobs
A sweep that fans work out is an ordinary cron job whose handler enqueues, which keeps the fan-out on the queue instead of inside the cron trigger's budget.
export default createJobHandler(jobs.enqueueDueChecks, async (ctx) => {
let due = await claimDue(ctx.database, Date.now());
await dispatcher.enqueueMany(
jobs.checkHttp,
due.map((m) => ({ monitorId: m.id })),
);
ctx.log.set({ checks: { enqueued: due.length } });
});
Pattern: Testing A Handler
A handler is a function over a context, so a test builds the context and calls it. No queue, no worker, no container.
import { createJobContext, Job } from "@sdxc/jobs";
import { Log } from "@sdxc/logger";
import handler from "~/app/jobs/clean";
import { Database } from "~/app/jobs/middleware/database";
test("deletes rows past the retention window", async () => {
let ctx = createJobContext(jobs.clean, { id: "message-1", attempts: 1 });
ctx.set(Database, await testDatabase(), { property: "database" });
await handler(ctx);
});
test("asks for a retry while the API is rate limiting", async () => {
let ctx = createJobContext(jobs.checkHttp, { id: "message-1", attempts: 1, input });
await expect(handler(ctx)).rejects.toBeInstanceOf(Job.Retry);
});
test("records the monitor it checked", async () => {
let records: Record<string, unknown>[] = [];
let log = new Log({ kind: "job", sink: (record) => void records.push(record) });
let ctx = createJobContext(jobs.checkHttp, { id: "message-1", attempts: 1, input, log });
await log.run(() => handler(ctx));
expect(records[0]).toMatchObject({ "monitor.id": input.monitorId });
});
createJobContext types the context the way the handler receives it — including whatever
the app declared its middleware installs. Building one skips that chain, so the test
populates what the chain would have; new JobContext(...) is the untyped equivalent the
dispatcher itself uses. Handing it a Log through init.log is how a test reads back
what the handler recorded.
Related Packages
@sdxc/logger- TheLoga run records into and thecreateLogger()configuration the dispatcher's logs carry@sdxc/duration- The duration strings a retry delay and a timeout take@sdxc/validate- Standard Schema validation, used to parse a payload@sdxc/cron- Parses the cron a job declares, and rejects one the platform would not accept@sdxc/result- TheResultevery queue call and every ping answers with@sdxc/cloudflare-mocks- Queue binding that drives a consumer in tests
Tips
Spell a cron exactly as its trigger - A valid expression still fires nothing if
wrangler.jsoncdoes not name it, andjob()cannot know that. Assertdispatcher.cronsagainst the config in a test.Treat a map key as a wire contract - Renaming one renames a message's address, and messages enqueued by the previous deploy are still in flight.
Map a loader, not a handler -
() => import(…)is what keeps job code out of the request path's module graph. For the same reason, prefermessageBody()plus the app's own queue helper when a controller enqueues, over importing the dispatcher.Never swallow an ending - A
catcharound actx.*call catches the thrown ending too. Re-throw it withif (error instanceof Job.Ending) throw error;.Write
return ctx.retry(…)- The verbs returnnever, but TypeScript only narrows on a never-returning call through aconstname, andctxis a parameter — so returning is what tells the compiler the lines below are unreachable.Reach for
ctx.exit()for bad input - Invalid data will not become valid on a redelivery, so acking and recording beats spending the retries.Pass
ctx.signalto every fetch - It is what makes a timeout cancel work rather than merely stop waiting for it.Let a monitor mean what it says - Ping from
onEndonly for adonestatus, so a run that timed out pings nothing and a monitor alerting is evidence the job really stopped completing.