Documentation
¶
Overview ¶
Package queue is the PostgreSQL-native async job queue. The Enqueue and Dequeue helpers are the only public surface that touches the job_queue table; lint forbids raw INSERT INTO job_queue elsewhere.
Correlation propagation is built in: Enqueue extracts correlation_id from the caller's context (and errors if none is set); Dequeue restores the row's correlation_id onto a fresh worker context. See docs/engineering/correlation_id_propagation.md.
Spec: specs/system/job-queue.spec.yaml.
Index ¶
- Variables
- func Complete(ctx context.Context, pool *pgxpool.Pool, jobID uuid.UUID) error
- func Enqueue(ctx context.Context, pool Execer, jobType string, payload any) (uuid.UUID, error)
- func EnqueueAfter(ctx context.Context, pool Execer, jobType string, payload any, ...) (uuid.UUID, error)
- func Fail(ctx context.Context, pool *pgxpool.Pool, jobID uuid.UUID, errMsg string) error
- type Execer
- type Job
- type Status
Constants ¶
This section is empty.
Variables ¶
var ErrMissingCorrelation = errors.New("queue: enqueue requires a correlation_id on context")
ErrMissingCorrelation is returned by Enqueue when the caller's context has no correlation_id set. The queue rejects such jobs because every async unit must be traceable to its originating intent (HTTP request, cron tick, system boot).
Spec system-job-queue AC-02, C-01.
var ErrNoJob = errors.New("queue: no pending job")
ErrNoJob is returned by Dequeue when no pending job is available. Not fatal — workers poll on this.
Functions ¶
func Complete ¶
Complete marks a job as completed. Workers call this after the job runs successfully.
Spec system-job-queue AC-06.
func Enqueue ¶
Enqueue persists a new job_queue row, immediately dequeuable. The caller's ctx MUST carry a correlation_id; this is a programming-error guard per spec C-01. The row's correlation_id pins the job to the originating intent across the async boundary.
Returns the inserted job's ID. The job is in "pending" status with attempts=0 until a worker claims it via Dequeue.
Spec system-job-queue AC-01.
func EnqueueAfter ¶
func EnqueueAfter(ctx context.Context, pool Execer, jobType string, payload any, delay time.Duration) (uuid.UUID, error)
EnqueueAfter is Enqueue with a delay: the job is not dequeuable until `delay` from now (available_at = now() + delay). A non-positive delay makes it immediately available, identical to Enqueue. Used to back off and requeue a job that can't run yet (e.g. the target host is busy) without a tight re-dequeue loop, since Dequeue skips not-yet-available rows.
Spec system-job-queue AC-13 (delayed visibility).
Types ¶
type Execer ¶ added in v0.8.2
type Execer interface {
Exec(ctx context.Context, sql string, args ...any) (pgconn.CommandTag, error)
}
Execer is the write surface Enqueue needs. A *pgxpool.Pool satisfies it, and so does a pgx.Tx, which lets a caller commit the job together with the rows that describe it, or roll both back. bugs/OW-097.
type Job ¶
type Job struct {
ID uuid.UUID
JobType string
Payload []byte // JSONB; opaque to the queue
CorrelationID string
Status Status
Attempts int
LastError string
CreatedAt time.Time
LockedAt *time.Time
CompletedAt *time.Time
}
Job is the row shape exposed to consumers. Payload is opaque JSON.
func Dequeue ¶
Dequeue claims one pending job atomically via SELECT ... FOR UPDATE SKIP LOCKED. Returns (job, workerCtx, nil) on success. workerCtx is a fresh context derived from context.Background() with the row's correlation_id set — the caller's ctx is NOT used as parent so the worker-loop's own correlation_id cannot leak into per-job execution (spec C-02).
Returns ErrNoJob (sentinel) when nothing is pending. Caller should poll with a delay on this.
Spec system-job-queue AC-03, AC-04, AC-05.