Architecture¶
A job travels from producer to consumer through PostgreSQL: the producer inserts a row, a trigger fires a notification, and a consumer claims the row and runs the matching entrypoint.
Job flow diagram¶
Producer ──enqueue──▶ PostgreSQL ──NOTIFY──▶ EventRouter
▲ │
│ signal
│ ▼
update status QueueManager
│ │
│ dispatch
│ ▼
└───────────────── Consumer
- Producer inserts a job using
Queries.enqueue(). - A trigger emits a
table_changed_eventvia NOTIFY on the configured channel. - The EventRouter places the event in a
PGNoticeEventListenerqueue. - The QueueManager waits for events, fetches ready jobs with
FOR UPDATE SKIP LOCKED(see Row Locking & SKIP LOCKED for the mechanics), and dispatches them to registered entrypoints. - After execution, the Consumer updates job status back in PostgreSQL.
EventRouter dispatches each notification type to its typed handler: table
changes feed the listener queue, cancellations cancel the job's scope, and
health-check echoes resolve the pending probe. Because the notification arrives
over LISTEN/NOTIFY, a consumer picks up new work without waiting for its next
poll.
QueueManager processing loop¶
┌──────────────────┐
│ Wait for NOTIFY │◀─────────────────────┐
└────────┬─────────┘ │
│ │
▼ │
┌──────────────────┐ no jobs │
│ Query jobs │──────────────────────▶│
└────────┬─────────┘ │
│ found │
▼ │
┌──────────────────┐ │
│ Execute task │ │
└───┬──────────┬───┘ │
│ │ │
success│ │error │
▼ ▼ │
┌──────────┐ ┌─────────────┐ │
│successful│ │ exception │ │
└─────┬────┘ └──────┬──────┘ │
└──────────────┴──────────────────────────┘
Job status lifecycle¶
PgQueuer tracks each job's progress using a dedicated PostgreSQL ENUM type,
pgqueuer_status by default:
CREATE TYPE pgqueuer_status AS ENUM (
'queued',
'picked',
'successful',
'exception',
'canceled',
'deleted',
'failed'
);
The lifecycle of a job flows through these statuses:
queued: Newly enqueued jobs start here and wait for a worker to pick them up.picked: Set byQueueManagerwhen a worker begins processing. A heartbeat timestamp tracks active work.successful: Assigned after a job completes without errors. Details are copied to the statistics log and removed from the queue.exception: Indicates the job failed with an uncaught error. The traceback is stored for later inspection.failed: Job held in the queue for manual review whenon_failure="hold"is set. Can be re-queued or deleted manually.canceled: Jobs canceled before completion receive this status and are logged.deleted: Used when jobs are removed from the queue without running, such as during manual cleanup operations.
Status transition diagram¶
┌────────┐
│ queued │
└───┬──┬─┘
claim │ │ delete
▼ ▼
┌────────┐ ┌─────────┐
│ picked │ │ deleted │
└┬─┬──┬─┬┘ └─────────┘
│ │ │ │
success │ │ │ │ cancel
│ │ │ │
▼ │ │ ▼
┌────────────┐ │ │ ┌───────────┐
│ successful │ │ │ │ canceled │
└────────────┘ │ │ └───────────┘
error │ │ hold
▼ ▼
┌───────────┐ ┌────────┐
│ exception │ │ failed │
└───────────┘ └────────┘