API reference

yatta/jobs

Durable queues, retries, dead letters, cron and events.

36 exported symbols and 184 members, read from src/types/job.ts.

Construct

calculateBackoff

function

Calculates 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

function

Creates an independent instance.

createCron(): CronScheduler

Returns Configured .

createEvents

function

Creates an independent instance.

createEvents(jobs?: JobQueueManager): EventBus
jobs
Optional enabling `.pipe()` background routing.

Returns Configured .

createJobs

function

Creates an independent instance.

createJobs(options?: { store?: JobStore; dbPath?: string; }): JobQueueManager
options
Storage engine or SQLite database path options.

Returns Configured .

duration

function

Converts 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

function

Standardizes an enqueue retry configuration into a concrete .

normalizeRetryPolicy(config?: EnqueueOptions["retry"]): RetryPolicy
config
Optional retry options.

Returns Normalized .

parsePriority

function

Normalizes 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

class

Standard 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`.

tsts
const cron = new CronExpression("0 9 * * 1-5", "America/New_York");const nextRun = cron.getNextDate();

10 members

  • minutesproperty
    minutes: ParsedCronField
  • hoursproperty
    hours: ParsedCronField
  • daysOfMonthproperty
    daysOfMonth: ParsedCronField
  • monthsproperty
    months: ParsedCronField
  • daysOfWeekproperty
    daysOfWeek: ParsedCronField
  • parse
    parse(expr: string)
  • parseField
    parseField(field: string, min: number, max: number, name: string): ParsedCronField
  • parseRange
    parseRange(rangeStr: string, min: number, max: number, name: string): [ number, number ]
  • getZonedParts
    getZonedParts(date: Date): { year: number; month: number; day: number; hour: number; minute: number; dow: number; }
  • getNextDate
    getNextDate(from?: Date): Date

    Computes the next matching execution timestamp after the specified reference date.

    from
    Reference starting date (defaults to `new Date()`).

    Returns Next matching .

CronScheduler

class

Recurring task scheduler supporting standard 5-field cron syntax and fixed interval durations.

tsts
const cron = createCron();cron.schedule("daily-cleanup", "0 0 * * *", async () => {  await db.cleanup();});cron.every("15m", async () => {  await syncData();});

6 members

  • tasksproperty
    tasks:
  • schedule
    schedule(name: string, cronExpression: string, handler: () => Promise<void> | void, options?: { timezone?: string; concurrency?: "skip" | "allow"; }): this

    Schedules 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.

  • every
    every(interval: Duration, handler: () => Promise<void> | void, name?: ): this

    Schedules 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.

  • armCron
    armCron(task: ScheduledTask)
  • armInterval
    armInterval(task: ScheduledTask)
  • stop
    stop(name?: string): void

    Cancels and stops a specific scheduled task or all tasks.

    name
    Optional task name to cancel. If omitted, stops all tasks.

EventBus

class

High-performance event bus featuring O(1) hash routing, prefix/suffix pattern matching,

class EventBus
tsts
const 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

  • exactListenersproperty
    exactListeners:
  • prefixListenersproperty
    prefixListeners:
  • suffixListenersproperty
    suffixListeners:
  • globalListenersproperty
    globalListeners:
  • on
    on<K extends keyof RegisteredEvents>(event: K, handler: (data: RegisteredEvents[K], name: string) => void | Promise<void>): () => void

    Subscribes 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();
  • on
    on(event: string, handler: EventHandler): () => void
  • on
    on(event: string, handler: EventHandler): () => void
  • once
    once<K extends keyof RegisteredEvents>(event: K, handler: (data: RegisteredEvents[K], name: string) => void | Promise<void>): () => void

    Subscribes 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.

  • once
    once(event: string, handler: EventHandler): () => void
  • once
    once(event: string, handler: EventHandler): () => void
  • emit
    emit<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 });
  • emit
    emit(event: string, data: unknown, options?: { throwOnError?: boolean; }): Promise<void>
  • emit
    emit(event: string, data: unknown, options?: { throwOnError?: boolean; }): Promise<void>
  • pipe
    pipe<E extends keyof RegisteredEvents, J extends keyof RegisteredJobs>(event: E, jobName: J, options?: EnqueueOptions | ((data: RegisteredEvents[E]) => EnqueueOptions), transform?: (data: RegisteredEvents[E]) => JobPayload<J>): () => void

    Pipes 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!" }),);
  • waitFor
    waitFor<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

class

Fluent builder for configuring and dispatching background jobs.

tsts
await jobs.job("process-image")  .with({ imageId: 42 })  .delay("10m")  .priority("high")  .save();

11 members

  • dataproperty
    data: TData
  • optionsproperty
    options: EnqueueOptions
  • with
    with(data: TData): this

    Sets the payload data for the job.

    data
    Payload object.

    Returns Current builder for chaining.

  • delay
    delay(delay: Duration): this

    Delays job execution by the specified duration.

    delay
    Duration string (e.g. `"5m"`) or milliseconds.

    Returns Current builder for chaining.

  • at
    at(date: Date | number): this

    Schedules job execution at an exact Date or timestamp.

    date
    Target execution Date or millisecond timestamp.

    Returns Current builder for chaining.

  • retry
    retry(attempts: number, config?: EnqueueOptions["retry"]): this

    Sets 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.

  • priority
    priority(p: EnqueueOptions["priority"]): this

    Sets the execution priority tier.

    p
    Priority level (`"low"`, `"normal"`, `"high"`, `"critical"` or number).

    Returns Current builder for chaining.

  • timeout
    timeout(limit: Duration): this

    Configures execution timeout before aborting.

    limit
    Timeout duration.

    Returns Current builder for chaining.

  • unique
    unique(key: string): this

    Assigns a unique deduplication key preventing concurrent duplicates in the queue.

    key
    Unique identifier string.

    Returns Current builder for chaining.

  • onQueue
    onQueue(queueName: string): this

    Targets a specific queue name instead of `"default"`.

    queueName
    Queue name string.

    Returns Current builder for chaining.

  • save
    save(): Promise<JobRecord<TData>>

    Finalizes configuration and enqueues the job into storage.

    Returns Promise resolving to the created .

JobQueueManager

class

Primary manager for background jobs, worker pools, queues, dead-letter queues, and operational metrics.

tsts
const 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

  • handlersproperty
    handlers:
  • poolsproperty
    pools:
  • storeproperty
    store: JobStore

    Storage backend engine powering this manager.

  • handle
    handle<K extends keyof RegisteredJobs, R = unknown>(name: K, handler: (ctx: JobContext<JobPayload<K>>) => Promise<R> | R): this

    Registers 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" };});
  • handle
    handle(name: string, handler: JobHandler): this
  • handle
    handle(name: string, handler: JobHandler): this
  • enqueue
    enqueue<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" });
  • enqueue
    enqueue(name: string, data: unknown, options?: EnqueueOptions): Promise<JobRecord>
  • enqueue
    enqueue(name: string, data: unknown, options?: EnqueueOptions): Promise<JobRecord>
  • job
    job<K extends keyof RegisteredJobs>(name: K): JobBuilder<JobPayload<K>>

    Initiates a fluent for configuring and enqueuing a job.

    name
    Job task name.
  • job
    job<T = unknown>(name: string): JobBuilder<T>
  • job
    job<T = unknown>(name: string): JobBuilder<T>
  • worker
    worker(queue?: , options?: WorkerOptions): WorkerPool

    Starts 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 });
  • stopAll
    stopAll(): Promise<void>

    Gracefully shuts down all active worker pools.

  • dlqgetter
    dlq:

    Dead-letter queue (DLQ) operations for inspecting, retrying, and purging dead jobs.

  • metrics
    metrics(queue?: string): Promise<QueueMetrics>

    Returns current operational volume counts (queued, delayed, running, completed, dead, total).

    queue
    Optional queue name filter.

MemoryJobStore

class

High-speed in-memory job store designed for unit tests and ephemeral background tasks.

class MemoryJobStore implements JobStore

16 members

  • jobsproperty
    jobs:
  • uniqueKeysproperty
    uniqueKeys:
  • init
    init(): Promise<void>

    Initializes the memory store (no-op).

  • enqueue
    enqueue(job: Omit<JobRecord, "attempts" | "state" | "createdAt" | "updatedAt">): Promise<JobRecord>

    Enqueues a job into memory.

  • claimNext
    claimNext(queue: string, workerId: string, lockDurationMs: number): Promise<JobRecord | null>

    Claims the highest priority ready job in memory.

  • heartbeat
    heartbeat(id: string, workerId: string, lockDurationMs: number): Promise<boolean>

    Renews a worker lease in memory.

  • updateProgress
    updateProgress(id: string, progress: number, message?: string): Promise<void>

    Updates progress percentage and message in memory.

  • complete
    complete(id: string, result?: unknown): Promise<void>

    Marks a job completed in memory.

  • fail
    fail(id: string, error: SerializedError, nextRunAt?: number, dead?: ): Promise<void>

    Marks a job failed or dead in memory.

  • reclaimStaleJobs
    reclaimStaleJobs(staleThresholdMs: number): Promise<number>

    Reclaims orphaned running jobs in memory.

  • getJob
    getJob(id: string): Promise<JobRecord | null>

    Fetches a job record from memory.

  • getMetrics
    getMetrics(queue?: string): Promise<QueueMetrics>

    Calculates queue counts in memory.

  • listDead
    listDead(queue?: string, limit?: ): Promise<JobRecord[]>

    Lists dead jobs from memory.

  • replayDead
    replayDead(id: string): Promise<boolean>

    Replays a dead job in memory.

  • purgeQueue
    purgeQueue(queue: string, state?: JobState): Promise<number>

    Purges jobs from memory.

  • close
    close(): Promise<void>

    Clears all jobs and keys from memory.

QueueError

class

Standard error class thrown by Yatta Job Queue, Scheduler, and Event Bus.

class QueueError extends Error

SQLiteJobStore

class

Persistent SQLite-backed job queue storage engine with WAL mode,

class SQLiteJobStore implements JobStore

17 members

  • dbproperty
    db: Database
  • pathproperty
    path: string

    Resolved filesystem database path.

  • init
    init(): Promise<void>

    Creates the `_yatta_jobs` table and performance indexes.

  • enqueue
    enqueue(job: Omit<JobRecord, "attempts" | "state" | "createdAt" | "updatedAt">): Promise<JobRecord>

    Persists a job into SQLite. Deduplicates by unique key if currently active.

  • claimNext
    claimNext(queue: string, workerId: string, lockDurationMs: number): Promise<JobRecord | null>

    Atomically claims the next eligible job in the queue using an SQLite `RETURNING` subquery.

  • heartbeat
    heartbeat(id: string, workerId: string, lockDurationMs: number): Promise<boolean>

    Renews the active worker lease in SQLite.

  • updateProgress
    updateProgress(id: string, progress: number, message?: string): Promise<void>

    Updates execution progress percentage and message for a job.

  • complete
    complete(id: string, result?: unknown): Promise<void>

    Marks a job completed and saves its return result.

  • fail
    fail(id: string, error: SerializedError, nextRunAt?: number, dead?: ): Promise<void>

    Marks a job failed, updating retry state or moving to the dead-letter state.

  • reclaimStaleJobs
    reclaimStaleJobs(staleThresholdMs: number): Promise<number>

    Reclaims running jobs whose lease expired back to queued.

  • getJob
    getJob(id: string): Promise<JobRecord | null>

    Fetches a single job record by ID.

  • getMetrics
    getMetrics(queue?: string): Promise<QueueMetrics>

    Calculates metrics for all or a specific queue.

  • listDead
    listDead(queue?: string, limit?: ): Promise<JobRecord[]>

    Lists jobs in the dead-letter queue.

  • replayDead
    replayDead(id: string): Promise<boolean>

    Replays a dead-letter job by re-queuing it.

  • purgeQueue
    purgeQueue(queue: string, state?: JobState): Promise<number>

    Purges jobs from SQLite.

  • close
    close(): Promise<void>

    Closes the SQLite database connection.

  • deserialize
    deserialize(row: any): JobRecord

WorkerPool

class

Worker execution pool managing concurrent job picking, active leases,

10 members

  • isRunningproperty
    isRunning:
  • activeWorkersproperty
    activeWorkers:
  • pollTimerproperty
    pollTimer?: ReturnType<typeof setTimeout>
  • staleTimerproperty
    staleTimer?: ReturnType<typeof setInterval>
  • workerIdproperty
    workerId:
  • inFlightJobsproperty
    inFlightJobs:
  • start
    start(): this

    Starts the worker polling loop and stale recovery timer.

    Returns Current worker pool.

  • tick
    tick()
  • executeJob
    executeJob(job: JobRecord, lockDuration: number)
  • stop
    stop(): Promise<void>

    Gracefully shuts down the worker pool, aborting in-flight tasks and waiting

Types

EnqueueOptions

interface

Options configuring how a job is enqueued and scheduled.

interface EnqueueOptions

8 members

  • queueproperty
    queue?: string

    Target queue name (defaults to `"default"`).

  • delayproperty
    delay?: Duration

    Delay duration before the job becomes eligible to run (e.g. `"5m"`, `"1h"`).

  • runAtproperty
    runAt?: Date | number

    Exact future timestamp or Date to execute the job.

  • attemptsproperty
    attempts?: number

    Maximum number of retry attempts before moving to DLQ (defaults to 3).

  • priorityproperty
    priority?: "low" | "normal" | "high" | "critical" | number

    Job execution priority (`"low"`, `"normal"`, `"high"`, `"critical"` or custom number). Defaults to `"normal"`.

  • timeoutproperty
    timeout?: Duration

    Execution timeout duration before aborting (e.g. `"30s"`, `"5m"`).

  • uniqueKeyproperty
    uniqueKey?: string

    Unique key preventing duplicate jobs in the queue while one is already active.

  • retryproperty
    retry?: { type?: "exponential" | "fixed"; delay?: Duration; factor?: number; jitter?: boolean; maxDelay?: Duration; }

    Custom retry and backoff parameters.

EventRegister

interface

Global interface for project-wide event type registration and declaration merging.

interface EventRegister
tsts
declare module "../types/job" {  interface EventRegister {    "user.registered": { userId: string; email: string };    "order.completed": { orderId: string; amount: number };  }}

JobContext

interface

Execution context supplied to a .

interface JobContext<TData = unknown>

9 members

  • idproperty
    id: string

    Unique job ID.

  • nameproperty
    name: string

    Job task name.

  • queueproperty
    queue: string

    Queue name.

  • attemptsproperty
    attempts: number

    Current execution attempt count (1-indexed).

  • maxAttemptsproperty
    maxAttempts: number

    Maximum retry attempts configured.

  • dataproperty
    data: TData

    Input job payload.

  • signalproperty
    signal: AbortSignal

    Abort signal triggered if the job times out, loses its lease, or the worker shuts down.

  • progress
    progress(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.
  • log
    log(message: string, meta?: Record<string, unknown>): void

    Logs a structured message tagged with the current job name and ID.

    message
    Log message.
    meta
    Optional metadata payload.

JobRecord

interface

Complete record of a background job in storage.

interface JobRecord<TData = unknown, TResult = unknown>

19 members

  • idproperty
    id: string

    Unique job identifier.

  • queueproperty
    queue: string

    Target queue name (e.g. `"default"`, `"emails"`).

  • nameproperty
    name: string

    Task action name.

  • dataproperty
    data: TData

    Input payload data.

  • stateproperty
    state: JobState

    Current execution lifecycle state.

  • attemptsproperty
    attempts: number

    Number of execution attempts completed so far.

  • maxAttemptsproperty
    maxAttempts: number

    Maximum attempts allowed before moving to the dead-letter queue.

  • priorityproperty
    priority: number

    Numerical priority (higher values run first).

  • runAtproperty
    runAt: number

    Scheduled execution timestamp in milliseconds since epoch.

  • timeoutproperty
    timeout?: number

    Maximum execution time allowed in milliseconds before aborting.

  • retryproperty
    retry: RetryPolicy

    Retry and backoff policy applied upon execution failure.

  • leaseproperty
    lease?: { workerId: string; acquiredAt: number; expiresAt: number; }

    Active worker lease metadata while in the `"running"` state.

  • progressproperty
    progress: number

    Execution completion percentage (0 - 100).

  • progressMessageproperty
    progressMessage?: string

    Human-readable status or progress message.

  • resultproperty
    result?: TResult

    Result returned upon successful completion.

  • errorproperty
    error?: SerializedError

    Serialized error details if the job failed or died.

  • uniqueKeyproperty
    uniqueKey?: string

    Deduplication key preventing concurrent duplicate enqueueing.

  • createdAtproperty
    createdAt: number

    Initial creation timestamp.

  • updatedAtproperty
    updatedAt: number

    Last modification timestamp.

JobRegister

interface

Global interface for project-wide job type registration and declaration merging.

interface JobRegister
tsts
declare module "../types/job" {  interface JobRegister {    "send-email": { to: string; subject: string; body: string };    "transcode-video": { videoId: string; quality: string };  }}

JobStore

interface

Storage and state-management interface for job queues.

interface JobStore

14 members

  • init
    init(): Promise<void>

    Initializes tables, schemas, and indexes.

  • enqueue
    enqueue(job: Omit<JobRecord, "attempts" | "state" | "createdAt" | "updatedAt">): Promise<JobRecord>

    Enqueues a new job into the store.

  • claimNext
    claimNext(queue: string, workerId: string, lockDurationMs: number): Promise<JobRecord | null>

    Atomically claims the next eligible job for the worker with an active lease.

  • heartbeat
    heartbeat(id: string, workerId: string, lockDurationMs: number): Promise<boolean>

    Renews an active worker lease to prevent premature reclamation.

  • updateProgress
    updateProgress(id: string, progress: number, message?: string): Promise<void>

    Updates the job's completion progress and optional status message.

  • complete
    complete(id: string, result?: unknown): Promise<void>

    Marks a job as completed and stores its result.

  • fail
    fail(id: string, error: SerializedError, nextRunAt?: number, dead?: boolean): Promise<void>

    Marks a job as failed, scheduling a retry or moving it to DLQ.

  • reclaimStaleJobs
    reclaimStaleJobs(staleThresholdMs: number): Promise<number>

    Reclaims orphaned or abandoned running jobs whose leases expired.

  • getJob
    getJob(id: string): Promise<JobRecord | null>

    Fetches a job record by ID.

  • getMetrics
    getMetrics(queue?: string): Promise<QueueMetrics>

    Returns queue operational metrics partitioned by job state.

  • listDead
    listDead(queue?: string, limit?: number): Promise<JobRecord[]>

    Retrieves jobs currently residing in the dead-letter queue.

  • replayDead
    replayDead(id: string): Promise<boolean>

    Replays a dead job by resetting its attempts and moving it back to queued.

  • purgeQueue
    purgeQueue(queue: string, state?: JobState): Promise<number>

    Purges jobs from a queue, optionally filtered by state.

  • close
    close(): Promise<void>

    Closes storage connections and releases resources.

ParsedCronField

interface

Internal descriptor representing parsed values of a single cron field.

interface ParsedCronField

2 members

  • valuesproperty
    values: Set<number>

    Set of allowable integer values for this field.

  • isWildcardproperty
    isWildcard: boolean

    Whether the field was specified as a wildcard (`*`).

QueueMetrics

interface

Real-time queue volume metrics grouped by lifecycle state.

interface QueueMetrics

6 members

  • queuedproperty
    queued: number

    Count of jobs awaiting pickup.

  • delayedproperty
    delayed: number

    Count of jobs scheduled for future execution.

  • runningproperty
    running: number

    Count of jobs currently being executed by workers.

  • completedproperty
    completed: number

    Count of successfully finished jobs.

  • deadproperty
    dead: number

    Count of jobs in the dead-letter queue.

  • totalproperty
    total: number

    Total count of all jobs.

RetryPolicy

interface

Retry and backoff configuration for failed jobs.

interface RetryPolicy

5 members

  • typeproperty
    type: "exponential" | "fixed"

    Backoff strategy:

  • delayproperty
    delay: number

    Initial delay in milliseconds before the first retry attempt.

  • factorproperty
    factor: number

    Exponential multiplier factor (e.g. `2` for doubling delay each attempt).

  • jitterproperty
    jitter: boolean

    Whether to apply random jitter to prevent worker thundering herds.

  • maxDelayproperty
    maxDelay: number

    Maximum upper bound in milliseconds for retry delays.

ScheduledTask

interface

Task registration record in the cron scheduler.

interface ScheduledTask

7 members

  • nameproperty
    name: string

    Task identifier name.

  • scheduleproperty
    schedule: string | Duration

    Cron expression string or duration string/number.

  • timezoneproperty
    timezone?: string

    Optional timezone string.

  • concurrencyproperty
    concurrency: "skip" | "allow"

    Concurrency policy when a tick fires while previous run is still in progress (`"skip"` or `"allow"`).

  • handlerproperty
    handler: () => Promise<void> | void

    Execution callback.

  • isRunningproperty
    isRunning: boolean

    Whether the task is currently executing.

  • timerproperty
    timer?: ReturnType<typeof setTimeout>

    Scheduled timer reference.

SerializedError

interface

Serialized error details saved onto failed or dead jobs for inspection.

interface SerializedError

3 members

  • messageproperty
    message: string

    Error message.

  • stackproperty
    stack?: string

    Error stack trace if available.

  • codeproperty
    code?: string

    Error classification code.

WorkerOptions

interface

Configuration options for starting a background .

interface WorkerOptions

10 members

  • queueproperty
    queue?: string

    Target queue to process (defaults to `"default"`).

  • concurrencyproperty
    concurrency?: number

    Maximum concurrent jobs executed simultaneously by this worker pool (defaults to 5).

  • pollIntervalproperty
    pollInterval?: Duration

    Polling interval when no ready jobs are found (defaults to `"1s"`).

  • lockDurationproperty
    lockDuration?: Duration

    Lease lock duration before an abandoned job can be reclaimed by other workers (defaults to `"60s"`).

  • staleRecoveryIntervalproperty
    staleRecoveryInterval?: Duration

    How frequently to scan for and recover stale/orphaned jobs (defaults to `"30s"`).

  • shutdownTimeoutproperty
    shutdownTimeout?: Duration

    Maximum grace period to wait for active jobs to finish during graceful shutdown (defaults to `"10s"`).

  • onProgressproperty
    onProgress?: (job: JobRecord) => void

    Callback fired when a job updates its progress percentage.

  • onCompletedproperty
    onCompleted?: (job: JobRecord) => void

    Callback fired when a job completes successfully.

  • onFailedproperty
    onFailed?: (job: JobRecord, err: Error) => void

    Callback fired when a job execution fails (before retries).

  • onDeadproperty
    onDead?: (job: JobRecord, err: Error) => void

    Callback fired when a job exhausts all retries and is moved to the dead-letter queue.

Duration

type

Millisecond duration represented as a literal string (e.g. `"500ms"`, `"30s"`, `"15m"`, `"2h"`, `"7d"`, `"1w"`)

type Duration = `${number}${"ms" | "s" | "m" | "h" | "d" | "w"}` | number

EventHandler

type

Event handler callback signature.

type EventHandler<T = any> = (data: T, eventName: string) => void | Promise<void>
data
Event payload.
eventName
Name of the emitted event.

EventPayload

type

Resolves 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] : any

JobHandler

type

Function handler invoked by workers to execute a queued job.

type JobHandler<TData = any, TResult = any> = (ctx: JobContext<TData>) => Promise<TResult> | TResult

JobPayload

type

Resolves 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] : any

JobState

type

Represents the execution lifecycle status of a background job:

type JobState = "queued" | "delayed" | "running" | "completed" | "dead"

RegisteredEvents

type

Type alias referring to .

RegisteredJobs

type

Type alias referring to .

Tip
Most of the types above are inferred. You rarely import AuthConfig or JobPayload — declaring your schema once is enough for the rest to follow. See Typed keys.