queue

Durable job queue with leases over std.sql: enqueue work, workers claim it with a lease, a worker that dies has its job reclaimed once the lease expires.

ecko get github.com/ecko-lang/queue
import queue

Pure computation: it declares no capabilities, so it cannot touch the network, the filesystem or the environment.

Version 0.5.0 - source - MIT.


init(db)

init(db) -> null.

Creates the queue's table and claim index if they do not already exist. Safe to call every time a program starts; it never touches existing rows.

enqueue(db, queue_name, payload, opts = empty_map())

enqueue(db, queue_name, payload, opts?) -> the new job's id.

payload is JSON-encoded before it is stored, so it can be any JSON-serializable value. opts.delay_ms (default 0) or an absolute opts.run_at (unix ms) delays when the job becomes claimable - run_at wins if both are given. opts.max_attempts (default 5) is how many fail calls the job survives before fail dead-letters it. opts.now overrides the clock (for tests); everything else defaults from it.

queue.enqueue(db, "emails", { to: "a@example.com" })
queue.enqueue(db, "reports", { id: 9 }, { delay_ms: 60000, max_attempts: 3 })

claim(db, queue_name, worker_id, lease_ms, opts = empty_map())

claim(db, queue_name, worker_id, lease_ms, opts?) -> a job, or null.

Atomically picks the oldest ready job for queue_name - due now (run_at has passed) and either never claimed or claimed but past its lease_expires - marks it "in_progress" under a fresh lease, and returns { id, queue_name, payload, attempts, max_attempts, worker_id, lease_token, lease_expires }. Returns null when nothing is ready. payload is the value passed to enqueue, JSON-decoded back.

lease_token changes on every claim, including a reclaim of an expired lease - it is what makes a stale worker's complete/fail/extend fail instead of racing a worker that has since taken the job over. Pass the whole job back to those calls.

The claim is one update ... returning statement, so two workers can never take the same job. Give each concurrent worker its own connection to the same database file; brief "database is locked" contention between them is retried inside claim rather than thrown.

job = queue.claim(db, "emails", "worker-1", 30000)   # 30s lease
if job == null { return }   # nothing ready right now

complete(db, job, opts = empty_map())

complete(db, job, opts?) -> null.

Marks job (as returned by claim) done. Throws { kind: "queue", reason: "lease_lost" } if another worker has since reclaimed the job, or { reason: "not_found" } if the row is gone.

fail(db, job, reason, opts = empty_map())

fail(db, job, reason, opts?) -> { status, attempts, run_at }.

Records a failed attempt on job. While job.attempts is under job.max_attempts the job goes back to "ready" after an exponential backoff (opts.backoff_ms, default 1000, doubling per attempt, capped at opts.backoff_max_ms, default 60000); once attempts reach max_attempts it is moved to "dead" instead (see dead). reason is stored as last_error. Throws the same lease errors as complete if the caller no longer holds the job.

queue.fail(db, job, "smtp timeout")
# -> { status: "ready", attempts: 1, run_at: 1700000001000 }

extend(db, job, lease_ms, opts = empty_map())

extend(db, job, lease_ms, opts?) -> job with a refreshed lease_expires.

For a job that legitimately needs more time than its original lease. Throws the same lease errors as complete if another worker has since reclaimed it.

stats(db, queue_name, opts = empty_map())

stats(db, queue_name, opts?) -> counts for queue_name.

{ ready, scheduled, in_progress, expired, done, dead, total }. ready is due and unclaimed right now; scheduled is "ready" status with run_at still in the future; in_progress is actively leased; expired is leased but past lease_expires (claimable on the next claim); done and dead are terminal.

dead(db, queue_name, opts = empty_map())

dead(db, queue_name, opts?) -> dead-lettered jobs for queue_name.

Newest first: [{ id, payload, attempts, max_attempts, last_error, created_at, updated_at }, ...]. opts.limit caps how many are returned (default 100).