yatta/jobs
Durable queues, retries, dead letters, cron and events.
36 exported symbols and 184 members, read from src/types/job.ts.
Construct
calculateBackoff
functionCalculates the next retry delay in milliseconds based on attempt number and policy.
calculateBackoff(attempt: number, policy: RetryPolicy): number- attempt
- Current retry attempt count.
- policy
- Configured retry policy.
Returns Milliseconds to delay before reattempting.
createCron
functionCreates an independent instance.
createCron(): CronSchedulerReturns Configured .
createEvents
functionCreates an independent instance.
createEvents(jobs?: JobQueueManager): EventBus- jobs
- Optional enabling `.pipe()` background routing.
Returns Configured .
createJobs
functionCreates an independent instance.
createJobs(options?: { store?: JobStore; dbPath?: string; }): JobQueueManager- options
- Storage engine or SQLite database path options.
Returns Configured .
duration
functionConverts a string duration (e.g. `"500ms"`, `"10s"`, `"5m"`, `"2h"`, `"1d"`, `"1w"`)
duration(value?: Duration, fallbackMs?: ): number- value
- Duration string or numeric milliseconds.
- fallbackMs
- Default fallback milliseconds if value is undefined.
Returns Milliseconds as an integer.
normalizeRetryPolicy
functionStandardizes an enqueue retry configuration into a concrete .
normalizeRetryPolicy(config?: EnqueueOptions["retry"]): RetryPolicy- config
- Optional retry options.
Returns Normalized .
parsePriority
functionNormalizes priority literals (`"critical"`, `"high"`, `"normal"`, `"low"`) or numbers into numeric values.
parsePriority(p?: EnqueueOptions["priority"]): number- p
- Priority value or string.
Returns Numerical priority (100 = critical, 50 = high, 10 = normal, 0 = low).
CronExpression
classStandard 5-field Cron parser supporting:
Standard 5-field Cron parser supporting: - Traditional Vixie Cron DOM/DOW OR-semantics. - Proper range bounds checking. - Deterministic, DST-aware time zone conversions via `Intl.DateTimeFormat`.
class CronExpressionconst cron = new CronExpression("0 9 * * 1-5", "America/New_York");const nextRun = cron.getNextDate();10 members
minutespropertyminutes: ParsedCronFieldhourspropertyhours: ParsedCronFielddaysOfMonthpropertydaysOfMonth: ParsedCronFieldmonthspropertymonths: ParsedCronFielddaysOfWeekpropertydaysOfWeek: ParsedCronFieldparseparse(expr: string)parseFieldparseField(field: string, min: number, max: number, name: string): ParsedCronFieldparseRangeparseRange(rangeStr: string, min: number, max: number, name: string): [ number, number ]getZonedPartsgetZonedParts(date: Date): { year: number; month: number; day: number; hour: number; minute: number; dow: number; }getNextDategetNextDate(from?: Date): DateComputes the next matching execution timestamp after the specified reference date.
- from
- Reference starting date (defaults to `new Date()`).
Returns Next matching .
CronScheduler
classRecurring task scheduler supporting standard 5-field cron syntax and fixed interval durations.
class CronSchedulerconst cron = createCron();cron.schedule("daily-cleanup", "0 0 * * *", async () => { await db.cleanup();});cron.every("15m", async () => { await syncData();});6 members
taskspropertytasks:scheduleschedule(name: string, cronExpression: string, handler: () => Promise<void> | void, options?: { timezone?: string; concurrency?: "skip" | "allow"; }): thisSchedules a task to run according to a standard 5-field cron expression.
- name
- Unique task name.
- cronExpression
- 5-field cron string (e.g. `"0/5 * * * *"`).
- handler
- Execution callback.
- options
- Optional timezone and concurrency settings.
Returns Current scheduler for chaining.
everyevery(interval: Duration, handler: () => Promise<void> | void, name?: ): thisSchedules a task to run repeatedly at a fixed interval duration.
- interval
- Duration string (e.g. `"10s"`, `"1h"`) or milliseconds.
- handler
- Execution callback.
- name
- Optional task name.
Returns Current scheduler for chaining.
armCronarmCron(task: ScheduledTask)armIntervalarmInterval(task: ScheduledTask)stopstop(name?: string): voidCancels and stops a specific scheduled task or all tasks.
- name
- Optional task name to cancel. If omitted, stops all tasks.
EventBus
classHigh-performance event bus featuring O(1) hash routing, prefix/suffix pattern matching,
class EventBusconst events = createEvents(jobs);events.on("user.registered", async (user) => { console.log("Welcome", user.email);});await events.emit("user.registered", { email: "alice@example.com" });15 members
exactListenerspropertyexactListeners:prefixListenerspropertyprefixListeners:suffixListenerspropertysuffixListeners:globalListenerspropertyglobalListeners:onon<K extends keyof RegisteredEvents>(event: K, handler: (data: RegisteredEvents[K], name: string) => void | Promise<void>): () => voidSubscribes a listener to a specific event pattern.
- event
- Event pattern string.
- handler
- Callback to invoke when event occurs.
Returns Unsubscribe function.
tsts const unsubscribe = events.on("order.paid", (order) => { ... });// Later:unsubscribe();onon(event: string, handler: EventHandler): () => voidonon(event: string, handler: EventHandler): () => voidonceonce<K extends keyof RegisteredEvents>(event: K, handler: (data: RegisteredEvents[K], name: string) => void | Promise<void>): () => voidSubscribes a one-time listener that automatically unsubscribes after its first invocation.
- event
- Event pattern string.
- handler
- One-time callback.
Returns Unsubscribe function to cancel before trigger.
onceonce(event: string, handler: EventHandler): () => voidonceonce(event: string, handler: EventHandler): () => voidemitemit<K extends keyof RegisteredEvents>(event: K, data: RegisteredEvents[K], options?: { throwOnError?: boolean; }): Promise<void>Emits an event, invoking all matching listeners concurrently.
- event
- Event name string.
- data
- Payload passed to listeners.
- options
- Execution options (e.g. `{ throwOnError: true }`).
tsts await events.emit("order.completed", { orderId: "123", amount: 99 });emitemit(event: string, data: unknown, options?: { throwOnError?: boolean; }): Promise<void>emitemit(event: string, data: unknown, options?: { throwOnError?: boolean; }): Promise<void>pipepipe<E extends keyof RegisteredEvents, J extends keyof RegisteredJobs>(event: E, jobName: J, options?: EnqueueOptions | ((data: RegisteredEvents[E]) => EnqueueOptions), transform?: (data: RegisteredEvents[E]) => JobPayload<J>): () => voidPipes an event directly into a background job queue!
- event
- Event to listen for.
- jobName
- Background job name to dispatch.
- options
- Enqueue options or dynamic options factory based on event data.
- transform
- Optional mapper converting the event payload into the job payload.
Returns Unsubscribe function.
tsts events.pipe("user.created", "send-welcome-email", (data) => ({ delay: "5m", priority: "high",})); // Map the event payload into a different job payload shapeevents.pipe( "user.registered", "send-email", undefined, (data) => ({ to: data.email, subject: "Welcome!" }),);waitForwaitFor<K extends keyof RegisteredEvents>(event: K, timeout?: Duration): Promise<RegisteredEvents[K]>Suspends and waits for the next occurrence of an event, resolving with its payload.
- event
- Event to await.
- timeout
- Maximum wait duration before rejecting (defaults to `"30s"`).
Returns Promise resolving to the event data.
tsts const payment = await events.waitFor("payment.confirmed", "1m");
JobBuilder
classFluent builder for configuring and dispatching background jobs.
class JobBuilderawait jobs.job("process-image") .with({ imageId: 42 }) .delay("10m") .priority("high") .save();11 members
datapropertydata: TDataoptionspropertyoptions: EnqueueOptionswithwith(data: TData): thisSets the payload data for the job.
- data
- Payload object.
Returns Current builder for chaining.
delaydelay(delay: Duration): thisDelays job execution by the specified duration.
- delay
- Duration string (e.g. `"5m"`) or milliseconds.
Returns Current builder for chaining.
atat(date: Date | number): thisSchedules job execution at an exact Date or timestamp.
- date
- Target execution Date or millisecond timestamp.
Returns Current builder for chaining.
retryretry(attempts: number, config?: EnqueueOptions["retry"]): thisSets the retry limit and optional backoff parameters for this job.
- attempts
- Maximum retry attempts.
- config
- Optional custom retry policy parameters.
Returns Current builder for chaining.
prioritypriority(p: EnqueueOptions["priority"]): thisSets the execution priority tier.
- p
- Priority level (`"low"`, `"normal"`, `"high"`, `"critical"` or number).
Returns Current builder for chaining.
timeouttimeout(limit: Duration): thisConfigures execution timeout before aborting.
- limit
- Timeout duration.
Returns Current builder for chaining.
uniqueunique(key: string): thisAssigns a unique deduplication key preventing concurrent duplicates in the queue.
- key
- Unique identifier string.
Returns Current builder for chaining.
onQueueonQueue(queueName: string): thisTargets a specific queue name instead of `"default"`.
- queueName
- Queue name string.
Returns Current builder for chaining.
savesave(): Promise<JobRecord<TData>>Finalizes configuration and enqueues the job into storage.
Returns Promise resolving to the created .
JobQueueManager
classPrimary manager for background jobs, worker pools, queues, dead-letter queues, and operational metrics.
class JobQueueManagerconst jobs = createJobs(); // 1. Register task handlerjobs.handle("send-email", async (ctx) => { await mailer.send(ctx.data);}); // 2. Start workerjobs.worker("default", { concurrency: 10 }); // 3. Enqueue jobawait jobs.enqueue("send-email", { to: "user@example.com" });16 members
handlerspropertyhandlers:poolspropertypools:storepropertystore: JobStoreStorage backend engine powering this manager.
handlehandle<K extends keyof RegisteredJobs, R = unknown>(name: K, handler: (ctx: JobContext<JobPayload<K>>) => Promise<R> | R): thisRegisters a task execution handler for a specific job name.
- name
- Job task name.
- handler
- Execution function receiving .
Returns Current manager for chaining.
tsts jobs.handle("render-video", async (ctx) => { await ctx.progress(25, "Encoding audio..."); return { file: "out.mp4" };});handlehandle(name: string, handler: JobHandler): thishandlehandle(name: string, handler: JobHandler): thisenqueueenqueue<K extends keyof RegisteredJobs>(name: K, data: JobPayload<K>, options?: EnqueueOptions): Promise<JobRecord<JobPayload<K>>>Enqueues a job into storage for background execution.
- name
- Job task name.
- data
- Payload object.
- options
- Scheduling, priority, and retry options.
Returns Created .
tsts await jobs.enqueue("send-email", { to: "alice@example.com" }, { delay: "5m" });enqueueenqueue(name: string, data: unknown, options?: EnqueueOptions): Promise<JobRecord>enqueueenqueue(name: string, data: unknown, options?: EnqueueOptions): Promise<JobRecord>jobjob<K extends keyof RegisteredJobs>(name: K): JobBuilder<JobPayload<K>>Initiates a fluent for configuring and enqueuing a job.
- name
- Job task name.
jobjob<T = unknown>(name: string): JobBuilder<T>jobjob<T = unknown>(name: string): JobBuilder<T>workerworker(queue?: , options?: WorkerOptions): WorkerPoolStarts or retrieves a managed processing the specified queue.
- queue
- Queue identifier (defaults to `"default"`).
- options
- Worker options (concurrency, intervals, callbacks).
Returns Running .
tsts const pool = jobs.worker("emails", { concurrency: 5 });stopAllstopAll(): Promise<void>Gracefully shuts down all active worker pools.
dlqgetterdlq:Dead-letter queue (DLQ) operations for inspecting, retrying, and purging dead jobs.
metricsmetrics(queue?: string): Promise<QueueMetrics>Returns current operational volume counts (queued, delayed, running, completed, dead, total).
- queue
- Optional queue name filter.
MemoryJobStore
classHigh-speed in-memory job store designed for unit tests and ephemeral background tasks.
class MemoryJobStore implements JobStore16 members
jobspropertyjobs:uniqueKeyspropertyuniqueKeys:initinit(): Promise<void>Initializes the memory store (no-op).
enqueueenqueue(job: Omit<JobRecord, "attempts" | "state" | "createdAt" | "updatedAt">): Promise<JobRecord>Enqueues a job into memory.
claimNextclaimNext(queue: string, workerId: string, lockDurationMs: number): Promise<JobRecord | null>Claims the highest priority ready job in memory.
heartbeatheartbeat(id: string, workerId: string, lockDurationMs: number): Promise<boolean>Renews a worker lease in memory.
updateProgressupdateProgress(id: string, progress: number, message?: string): Promise<void>Updates progress percentage and message in memory.
completecomplete(id: string, result?: unknown): Promise<void>Marks a job completed in memory.
failfail(id: string, error: SerializedError, nextRunAt?: number, dead?: ): Promise<void>Marks a job failed or dead in memory.
reclaimStaleJobsreclaimStaleJobs(staleThresholdMs: number): Promise<number>Reclaims orphaned running jobs in memory.
getJobgetJob(id: string): Promise<JobRecord | null>Fetches a job record from memory.
getMetricsgetMetrics(queue?: string): Promise<QueueMetrics>Calculates queue counts in memory.
listDeadlistDead(queue?: string, limit?: ): Promise<JobRecord[]>Lists dead jobs from memory.
replayDeadreplayDead(id: string): Promise<boolean>Replays a dead job in memory.
purgeQueuepurgeQueue(queue: string, state?: JobState): Promise<number>Purges jobs from memory.
closeclose(): Promise<void>Clears all jobs and keys from memory.
QueueError
classStandard error class thrown by Yatta Job Queue, Scheduler, and Event Bus.
class QueueError extends ErrorSQLiteJobStore
classPersistent SQLite-backed job queue storage engine with WAL mode,
class SQLiteJobStore implements JobStore17 members
dbpropertydb: Databasepathpropertypath: stringResolved filesystem database path.
initinit(): Promise<void>Creates the `_yatta_jobs` table and performance indexes.
enqueueenqueue(job: Omit<JobRecord, "attempts" | "state" | "createdAt" | "updatedAt">): Promise<JobRecord>Persists a job into SQLite. Deduplicates by unique key if currently active.
claimNextclaimNext(queue: string, workerId: string, lockDurationMs: number): Promise<JobRecord | null>Atomically claims the next eligible job in the queue using an SQLite `RETURNING` subquery.
heartbeatheartbeat(id: string, workerId: string, lockDurationMs: number): Promise<boolean>Renews the active worker lease in SQLite.
updateProgressupdateProgress(id: string, progress: number, message?: string): Promise<void>Updates execution progress percentage and message for a job.
completecomplete(id: string, result?: unknown): Promise<void>Marks a job completed and saves its return result.
failfail(id: string, error: SerializedError, nextRunAt?: number, dead?: ): Promise<void>Marks a job failed, updating retry state or moving to the dead-letter state.
reclaimStaleJobsreclaimStaleJobs(staleThresholdMs: number): Promise<number>Reclaims running jobs whose lease expired back to queued.
getJobgetJob(id: string): Promise<JobRecord | null>Fetches a single job record by ID.
getMetricsgetMetrics(queue?: string): Promise<QueueMetrics>Calculates metrics for all or a specific queue.
listDeadlistDead(queue?: string, limit?: ): Promise<JobRecord[]>Lists jobs in the dead-letter queue.
replayDeadreplayDead(id: string): Promise<boolean>Replays a dead-letter job by re-queuing it.
purgeQueuepurgeQueue(queue: string, state?: JobState): Promise<number>Purges jobs from SQLite.
closeclose(): Promise<void>Closes the SQLite database connection.
deserializedeserialize(row: any): JobRecord
WorkerPool
classWorker execution pool managing concurrent job picking, active leases,
class WorkerPool10 members
isRunningpropertyisRunning:activeWorkerspropertyactiveWorkers:pollTimerpropertypollTimer?: ReturnType<typeof setTimeout>staleTimerpropertystaleTimer?: ReturnType<typeof setInterval>workerIdpropertyworkerId:inFlightJobspropertyinFlightJobs:startstart(): thisStarts the worker polling loop and stale recovery timer.
Returns Current worker pool.
ticktick()executeJobexecuteJob(job: JobRecord, lockDuration: number)stopstop(): Promise<void>Gracefully shuts down the worker pool, aborting in-flight tasks and waiting
Types
EnqueueOptions
interfaceOptions configuring how a job is enqueued and scheduled.
interface EnqueueOptions8 members
queuepropertyqueue?: stringTarget queue name (defaults to `"default"`).
delaypropertydelay?: DurationDelay duration before the job becomes eligible to run (e.g. `"5m"`, `"1h"`).
runAtpropertyrunAt?: Date | numberExact future timestamp or Date to execute the job.
attemptspropertyattempts?: numberMaximum number of retry attempts before moving to DLQ (defaults to 3).
prioritypropertypriority?: "low" | "normal" | "high" | "critical" | numberJob execution priority (`"low"`, `"normal"`, `"high"`, `"critical"` or custom number). Defaults to `"normal"`.
timeoutpropertytimeout?: DurationExecution timeout duration before aborting (e.g. `"30s"`, `"5m"`).
uniqueKeypropertyuniqueKey?: stringUnique key preventing duplicate jobs in the queue while one is already active.
retrypropertyretry?: { type?: "exponential" | "fixed"; delay?: Duration; factor?: number; jitter?: boolean; maxDelay?: Duration; }Custom retry and backoff parameters.
EventRegister
interfaceGlobal interface for project-wide event type registration and declaration merging.
interface EventRegisterdeclare module "../types/job" { interface EventRegister { "user.registered": { userId: string; email: string }; "order.completed": { orderId: string; amount: number }; }}JobContext
interfaceExecution context supplied to a .
interface JobContext<TData = unknown>9 members
idpropertyid: stringUnique job ID.
namepropertyname: stringJob task name.
queuepropertyqueue: stringQueue name.
attemptspropertyattempts: numberCurrent execution attempt count (1-indexed).
maxAttemptspropertymaxAttempts: numberMaximum retry attempts configured.
datapropertydata: TDataInput job payload.
signalpropertysignal: AbortSignalAbort signal triggered if the job times out, loses its lease, or the worker shuts down.
progressprogress(percent: number, message?: string): Promise<void>Reports execution progress percentage (0 - 100) and an optional status message.
- percent
- Number between 0 and 100.
- message
- Optional human-readable message.
loglog(message: string, meta?: Record<string, unknown>): voidLogs a structured message tagged with the current job name and ID.
- message
- Log message.
- meta
- Optional metadata payload.
JobRecord
interfaceComplete record of a background job in storage.
interface JobRecord<TData = unknown, TResult = unknown>19 members
idpropertyid: stringUnique job identifier.
queuepropertyqueue: stringTarget queue name (e.g. `"default"`, `"emails"`).
namepropertyname: stringTask action name.
datapropertydata: TDataInput payload data.
statepropertystate: JobStateCurrent execution lifecycle state.
attemptspropertyattempts: numberNumber of execution attempts completed so far.
maxAttemptspropertymaxAttempts: numberMaximum attempts allowed before moving to the dead-letter queue.
prioritypropertypriority: numberNumerical priority (higher values run first).
runAtpropertyrunAt: numberScheduled execution timestamp in milliseconds since epoch.
timeoutpropertytimeout?: numberMaximum execution time allowed in milliseconds before aborting.
retrypropertyretry: RetryPolicyRetry and backoff policy applied upon execution failure.
leasepropertylease?: { workerId: string; acquiredAt: number; expiresAt: number; }Active worker lease metadata while in the `"running"` state.
progresspropertyprogress: numberExecution completion percentage (0 - 100).
progressMessagepropertyprogressMessage?: stringHuman-readable status or progress message.
resultpropertyresult?: TResultResult returned upon successful completion.
errorpropertyerror?: SerializedErrorSerialized error details if the job failed or died.
uniqueKeypropertyuniqueKey?: stringDeduplication key preventing concurrent duplicate enqueueing.
createdAtpropertycreatedAt: numberInitial creation timestamp.
updatedAtpropertyupdatedAt: numberLast modification timestamp.
JobRegister
interfaceGlobal interface for project-wide job type registration and declaration merging.
interface JobRegisterdeclare module "../types/job" { interface JobRegister { "send-email": { to: string; subject: string; body: string }; "transcode-video": { videoId: string; quality: string }; }}JobStore
interfaceStorage and state-management interface for job queues.
interface JobStore14 members
initinit(): Promise<void>Initializes tables, schemas, and indexes.
enqueueenqueue(job: Omit<JobRecord, "attempts" | "state" | "createdAt" | "updatedAt">): Promise<JobRecord>Enqueues a new job into the store.
claimNextclaimNext(queue: string, workerId: string, lockDurationMs: number): Promise<JobRecord | null>Atomically claims the next eligible job for the worker with an active lease.
heartbeatheartbeat(id: string, workerId: string, lockDurationMs: number): Promise<boolean>Renews an active worker lease to prevent premature reclamation.
updateProgressupdateProgress(id: string, progress: number, message?: string): Promise<void>Updates the job's completion progress and optional status message.
completecomplete(id: string, result?: unknown): Promise<void>Marks a job as completed and stores its result.
failfail(id: string, error: SerializedError, nextRunAt?: number, dead?: boolean): Promise<void>Marks a job as failed, scheduling a retry or moving it to DLQ.
reclaimStaleJobsreclaimStaleJobs(staleThresholdMs: number): Promise<number>Reclaims orphaned or abandoned running jobs whose leases expired.
getJobgetJob(id: string): Promise<JobRecord | null>Fetches a job record by ID.
getMetricsgetMetrics(queue?: string): Promise<QueueMetrics>Returns queue operational metrics partitioned by job state.
listDeadlistDead(queue?: string, limit?: number): Promise<JobRecord[]>Retrieves jobs currently residing in the dead-letter queue.
replayDeadreplayDead(id: string): Promise<boolean>Replays a dead job by resetting its attempts and moving it back to queued.
purgeQueuepurgeQueue(queue: string, state?: JobState): Promise<number>Purges jobs from a queue, optionally filtered by state.
closeclose(): Promise<void>Closes storage connections and releases resources.
ParsedCronField
interfaceInternal descriptor representing parsed values of a single cron field.
interface ParsedCronField2 members
valuespropertyvalues: Set<number>Set of allowable integer values for this field.
isWildcardpropertyisWildcard: booleanWhether the field was specified as a wildcard (`*`).
QueueMetrics
interfaceReal-time queue volume metrics grouped by lifecycle state.
interface QueueMetrics6 members
queuedpropertyqueued: numberCount of jobs awaiting pickup.
delayedpropertydelayed: numberCount of jobs scheduled for future execution.
runningpropertyrunning: numberCount of jobs currently being executed by workers.
completedpropertycompleted: numberCount of successfully finished jobs.
deadpropertydead: numberCount of jobs in the dead-letter queue.
totalpropertytotal: numberTotal count of all jobs.
RetryPolicy
interfaceRetry and backoff configuration for failed jobs.
interface RetryPolicy5 members
typepropertytype: "exponential" | "fixed"Backoff strategy:
delaypropertydelay: numberInitial delay in milliseconds before the first retry attempt.
factorpropertyfactor: numberExponential multiplier factor (e.g. `2` for doubling delay each attempt).
jitterpropertyjitter: booleanWhether to apply random jitter to prevent worker thundering herds.
maxDelaypropertymaxDelay: numberMaximum upper bound in milliseconds for retry delays.
ScheduledTask
interfaceTask registration record in the cron scheduler.
interface ScheduledTask7 members
namepropertyname: stringTask identifier name.
schedulepropertyschedule: string | DurationCron expression string or duration string/number.
timezonepropertytimezone?: stringOptional timezone string.
concurrencypropertyconcurrency: "skip" | "allow"Concurrency policy when a tick fires while previous run is still in progress (`"skip"` or `"allow"`).
handlerpropertyhandler: () => Promise<void> | voidExecution callback.
isRunningpropertyisRunning: booleanWhether the task is currently executing.
timerpropertytimer?: ReturnType<typeof setTimeout>Scheduled timer reference.
SerializedError
interfaceSerialized error details saved onto failed or dead jobs for inspection.
interface SerializedError3 members
messagepropertymessage: stringError message.
stackpropertystack?: stringError stack trace if available.
codepropertycode?: stringError classification code.
WorkerOptions
interfaceConfiguration options for starting a background .
interface WorkerOptions10 members
queuepropertyqueue?: stringTarget queue to process (defaults to `"default"`).
concurrencypropertyconcurrency?: numberMaximum concurrent jobs executed simultaneously by this worker pool (defaults to 5).
pollIntervalpropertypollInterval?: DurationPolling interval when no ready jobs are found (defaults to `"1s"`).
lockDurationpropertylockDuration?: DurationLease lock duration before an abandoned job can be reclaimed by other workers (defaults to `"60s"`).
staleRecoveryIntervalpropertystaleRecoveryInterval?: DurationHow frequently to scan for and recover stale/orphaned jobs (defaults to `"30s"`).
shutdownTimeoutpropertyshutdownTimeout?: DurationMaximum grace period to wait for active jobs to finish during graceful shutdown (defaults to `"10s"`).
onProgresspropertyonProgress?: (job: JobRecord) => voidCallback fired when a job updates its progress percentage.
onCompletedpropertyonCompleted?: (job: JobRecord) => voidCallback fired when a job completes successfully.
onFailedpropertyonFailed?: (job: JobRecord, err: Error) => voidCallback fired when a job execution fails (before retries).
onDeadpropertyonDead?: (job: JobRecord, err: Error) => voidCallback fired when a job exhausts all retries and is moved to the dead-letter queue.
Duration
typeMillisecond duration represented as a literal string (e.g. `"500ms"`, `"30s"`, `"15m"`, `"2h"`, `"7d"`, `"1w"`)
type Duration = `${number}${"ms" | "s" | "m" | "h" | "d" | "w"}` | numberEventHandler
typeEvent handler callback signature.
type EventHandler<T = any> = (data: T, eventName: string) => void | Promise<void>- data
- Event payload.
- eventName
- Name of the emitted event.
EventPayload
typeResolves the payload type for a specific registered event name, or falls back to `any`.
type EventPayload<K extends string> = K extends keyof RegisteredEvents ? RegisteredEvents[K] : anyJobHandler
typeFunction handler invoked by workers to execute a queued job.
type JobHandler<TData = any, TResult = any> = (ctx: JobContext<TData>) => Promise<TResult> | TResultJobPayload
typeResolves the payload type for a specific registered job name, or falls back to `any`.
type JobPayload<K extends string> = K extends keyof RegisteredJobs ? RegisteredJobs[K] : anyJobState
typeRepresents the execution lifecycle status of a background job:
type JobState = "queued" | "delayed" | "running" | "completed" | "dead"RegisteredEvents
typeType alias referring to .
type RegisteredEvents = EventRegisterRegisteredJobs
typeType alias referring to .
type RegisteredJobs = JobRegisterAuthConfig or JobPayload — declaring your schema once is enough for the rest to follow. See Typed keys.