Queues
Durable job queues on your app's own database — send now or later, receive with a visibility timeout, delete or archive when done.
A queue holds work for later: send a message now, and a worker picks it up when it gets to it — across a deploy, a restart, or a traffic spike. Camplax queues live in your app's own database as two plain tables, so there is no extension to install and no extra service to keep alive. A cron job that posts work and a route that drains it is the whole setup.
Sending
postgresQueues takes the app's database pool — db from src/db.ts in a Plax app — and .queue(name) gives a client for one named queue:
import { postgresQueues } from "@camplax/plax";
import { db } from "./db.ts";
const queues = postgresQueues(db);
const emails = queues.queue("emails");
const { id } = await emails.send({ to: "sam@acme.co", subject: "Welcome" });
The body is any JSON value up to 64 KB — bigger is refused before it reaches the database. Queue names are 1–80 characters of letters, digits, - and _.
Dedupe a retried send by passing your own id: if a request dies after the insert but before you saw the reply, sending again with the same id is a no-op — you get { id, duplicate: true } and the queue holds one message, not two:
await emails.send({ to }, { id: `welcome-${user.id}` });
Delay a message with delaySeconds — it stays invisible until the delay passes:
await emails.send({ to, kind: "nudge" }, { delaySeconds: 86400 }); // tomorrow
Batch up to your throughput needs in one call — each entry takes its own id and delaySeconds, and the reply flags per-entry duplicates:
await emails.sendBatch(users.map((u) => ({ body: { to: u.email }, id: `welcome-${u.id}` })));
Receiving
A worker asks for messages, gets them with a receipt, and deletes each one when the work is done:
const messages = await emails.receive({ max: 10, visibilitySeconds: 60 });
for (const msg of messages) {
await sendEmail(msg.body); // msg.id, msg.body, msg.attempts
await emails.delete(msg.receipt);
}
receive returns the oldest waiting messages — at most 10 per call — and hides each for visibilitySeconds (default 30). While hidden, no other worker can see it. The receipt is a one-time token minted fresh on every receive: delete and archive only act when the receipt matches the live claim, so a stale or forged receipt deletes nothing — you get a PlaxQueueError with code receipt_mismatch.
Delivery is at-least-once, never exactly-once. If a worker dies after receive but before delete, the timeout passes and another worker receives the same message again — attempts counts up (1 on first delivery). Two consequences:
- Keep handlers idempotent — do the work in a way that is safe to repeat (upsert, dedupe on a key, check before acting).
- A message that keeps killing workers keeps coming back. Watch
attempts, andarchiveor delete the poison ones yourself — there is no automatic dead-letter move yet.
Delete or archive
delete(receipt) removes the message — the normal end of the lifecycle. archive(receipt) moves it to plax_queue_archived instead, keeping the body, attempt count and an archived_at stamp — the audit trail for "what did we process" and the starting point for replay tooling.
Counting
queues.count() sums every queue; queues.count("emails") scopes to one. Either way you get { visible, inFlight } — work waiting now versus claims still inside their timeout. A growing visible count means consumers are behind; inFlight that never drains means workers are dying mid-job.
Limits
| Limit | Value | |
|---|---|---|
| Message body | 64 KB of JSON | |
receive batch | 10 messages | |
| Visibility timeout | any visibilitySeconds you pass; default 30 | |
| Queue name | 1–80 chars: a–z, A–Z, 0–9, -, _ |
Under the hood: plax_queue_messages and plax_queue_archived are ordinary tables created on first use. receive claims rows with SELECT … FOR UPDATE SKIP LOCKED in the same statement that hides them, so two workers can never receive the same message — the guarantee is at-least-once, and idempotent handlers do the rest.