ADR-009: scheduler collapses into queues — schedule()/unschedule()/ run_schedules (opt-in, no ambient timers), @every-only v1 grammar (dissolves honker's local-TZ cron brittleness), boundary guarantee row (at-least-once per boundary, fixed 64-cap catch-up with skip-forward, row-locked fire tx as engine-generic no-double-fire floor), __alkstore_scheduler leadership lock, InvalidSpec + LeadershipLost taxonomy additions. ADR-010: queue depth pinned engine-uniformly — three-state machine (pending/processing/dead) with delete-on-ack, get_job sees dead rows, heartbeat = renewal with late-heartbeat refusal, reclaim-consumes- an-attempt stated as contract text, equal-jitter exponential backoff (range definitionally pinned, 1 h cap), QueueOpts stamped onto job rows at enqueue (no per-queue registry), move-to-dead dead-letter with retention-sweep support and no redrive API, sweep_expired carries the no-stranded-rows property (fixes honker's expired- processing zombie hole — SQLite-side realization rides OQ-06 as a concrete fork candidate), one engine-owned pg schema, queues are rows not tables, result-storage cut-flag stands. Also: full honker-machinery and pgboss-rs reference reads persisted (docs/research/reference-*.md — the honker defect list pre-stages the OQ-06 quality read), queues.md rewritten from design-space frame to resolved-depth spec, core-contract/engines/README/overview/deployment/ ADR-002 propagated.
19 KiB
status, last_updated
| status | last_updated |
|---|---|
| draft | 2026-10-05 |
Reference read: pgboss-rs queue semantics
(pgboss-rs is design-reference only per ADR-005 — never a dependency; referenced by ADR-010.)
Reference: /workspace/pgboss-rs @ 98f7d9e ("Merge pull request
#19 ... codecov-action-7", HEAD of main — verified via git log).
Crate: pgboss v0.1.0-rc6 (Cargo.toml), edition 2024, sqlx
0.8 / tokio / chrono / serde_json. Provenance: "Inspired by, compatible
with and partially ported from pg-boss Node.js package. Heavily
influenced by ... faktory-rs" (README.md:10-12). Internal schema
version constant CURRENT_PGBOSS_APP_VERSION = 26, minimum supported
26 (src/lib.rs:90-91) — the port targets the node pg-boss v10.x
table schema and asserts on it, so its DDL is a faithful design
reference even where the Rust runtime layer is thin (the Makefile
diffs schemas against upstream node pg-boss's own dump, Makefile:46-95).
~3,600 lines, mostly SQL-string builders. Key files:
src/sql/ddl.rs (schema), src/sql/dml.rs (all job SQL),
src/job.rs, src/queue.rs, src/client/ (thin API layer).
1. Job states and schema
Seven-state enum, stored as a Postgres enum type per schema
(src/sql/ddl.rs:7-17), mapped in JobState (src/job.rs:21-42):
| state | meaning | ordinal (enum order, load-bearing) |
|---|---|---|
created |
registered, not yet consumed (default, ddl.rs:105) | 0 |
retry |
failed, awaiting re-fetch | 1 |
active |
being processed | 2 |
completed |
done | 3 |
cancelled |
cancelled via API | 4 |
failed |
terminal failure | 5 |
Enum ordering is load-bearing: fetchability/liveness guards are
ordinal comparisons — state < 'active' in the fetch CTE means
created or retry (dml.rs:138); state < 'completed' guards
cancel/fail (dml.rs:184, 281). No obsolete, no dead-letter
state — dead-lettering manifests as a row inserted on a DLQ queue.
Table layout (7 relations + 2 stored procs per schema, installed by
one transaction guarded by an advisory lock at first Client::connect,
src/sql/mod.rs:38-58):
{schema}.version—version int PK, cron_on timestamptz(ddl.rs :19-28);cron_onvestigial.{schema}.queue— registry + embedded monitor counters (ddl.rs :30-61):name, policy, retry_limit, retry_delay, retry_backoff, retry_delay_max, expire_seconds, retention_seconds, deletion_seconds, dead_letter (FK → queue.name, check (dead_letter is distinct from name)), partition, table_name, deferred_count, queued_count, warning_queued, active_count, total_count, singletons_active text[], monitor_on, maintain_on, created_on, updated_on.{schema}.schedule— cron rows (unused by port code):name FK→queue on delete cascade, key default '', cron text, timezone text, data jsonb, options jsonb, PK(name, key)(ddl.rs:63-80).{schema}.subscription— event subscription rows (unused):(event, name FK→queue), PK(event, name)(ddl.rs:82-95).{schema}.job— the jobs table, partitioned by LIST(name) with PK(name, id)(ddl.rs:97-129). Columns:id uuid default gen_random_uuid(),name,priority int default 0,data jsonb,state,retry_limit int default 2,retry_count int default 0,retry_delay int default 0(seconds),retry_backoff bool default false,retry_delay_max int nullable,expire_seconds int default 900(15 min),deletion_seconds int default 604800(7 days),singleton_key text,singleton_on timestamp(epoch-quantized bucket),start_after timestamptz default now(),created_on,started_on,completed_on(timestamptz),keep_until timestamptz default now() + interval '1209600'(14 days),output jsonb,dead_letter text,policy text.{schema}.job_common—LIKE jobDEFAULT partition with the full index set (ddl.rs:131-176): fetch indexi5 (name, start_after) INCLUDE (priority, created_on, id) WHERE state < 'active', plus four partial unique indexes implementing queue policies/throttling (i1policy short,i2singleton,i3stately,i4singleton_on-bucket,i6exclusive).- Per-partitioned-queue tables named
j || sha224(queue_name)hex, created bycreate_queue(queue_name, options jsonb)plpgsql (ddl.rs:185-275), attached as partitionsFOR VALUES IN (queue_name).
Layout: one PostgreSQL schema (default "pgboss", src/client /opts.rs:9) shared by everything; a queue is either a partition-filter
over the shared table (partition=false, default) or a dedicated
partition table (partition=true, src/queue.rs:104-106). Isolation
at partition level, only when opted in.
2. Enqueue options
Client::send_job (src/client/public/job_ops.rs:20-84) binds a
Job + JobOptions struct serialized to jsonb and merged with queue
defaults inside SQL (dml.rs:68-129). Resolution rule:
COALESCE(job-level, queue-level, hard default).
| option | resolved | default |
|---|---|---|
id |
COALESCE($1, gen_random_uuid()) |
server-generated |
queue_name |
must exist (JOIN queue), else Error::DoesNotExist |
— |
data |
jsonb passthrough | — |
priority |
COALESCE(j.priority, 0) |
0 (higher = fetched first) |
start_after / delay_for |
COALESCE(j.start_after, now()) |
immediate |
retry_limit |
COALESCE(j.retry_limit, q.retry_limit) |
2 (ddl.rs:106, 221) |
retry_delay (secs) |
COALESCE(j.retry_delay, q.retry_delay) |
0 (ddl.rs:220) |
retry_backoff |
COALESCE(j.retry_backoff, q.retry_backoff, false) |
false |
retry_delay_max |
COALESCE(...) |
NULL (uncapped); not exposed in Rust API |
expire_in (expire_seconds) |
COALESCE(j.expire_in, q.expire_seconds) |
900 s (ddl.rs:223) |
retain_for (retention_seconds) |
keep_until = COALESCE(j.start_after, now()) + COALESCE(j.retain_for, q.retention_seconds) * interval '1s' (dml.rs:103) |
14 days (ddl.rs:224) |
delete_after (deletion_seconds) |
COALESCE(j.delete_after, q.deletion_seconds) |
7 days; not exposed in Rust JobOptions (job.rs:76-114) |
singleton_for |
epoch-bucket quantization (dml.rs:98-100) | NULL |
singleton_key |
passthrough | NULL |
dead_letter |
queue-level only; copied from q.dead_letter (dml.rs:109) |
— |
policy |
queue-level only | standard |
Also consumed from options jsonb but unsettable from the Rust API:
retry_delay_max, singleton_offset, delete_after,
expire_in-as-integer.
Errors at enqueue (job_ops.rs:38-80, mapped from
unique-index/constraint names): Error::Conflict (duplicate
user-supplied id / _pkey), Error::Throttled (_i1…_i6),
Error::DoesNotExist (queue or DLQ missing via q_fkey/dlq_fkey or
empty result).
Explicit client API (src/client/public/): send_job,
send_data, fetch_job, fetch_jobs, get_job (non-consuming,
job_ops.rs:105-133), complete_job(s), fail_job(s)/ fail_job_with_details, cancel_job(s), resume_job(s) (cancelled →
created again, dml.rs:194-208), delete_job(s),
create_queue/create_standard_queue/get_queue(s)/delete_queue,
force_maintain.
3. Claim / fetch mechanics
Fetch SQL (dml.rs:133-176), one atomic CTE:
WITH next AS (
SELECT id FROM {schema}.job
WHERE name = $1 AND state < 'active' AND start_after < now()
ORDER BY priority DESC, created_on, id
LIMIT $2
FOR UPDATE SKIP LOCKED -- dml.rs:141-142
)
UPDATE {schema}.job j SET
state = 'active', started_on = now(),
retry_count = CASE WHEN started_on IS NULL THEN retry_count
ELSE retry_count + 1 END -- dml.rs:147
FROM next WHERE j.id = next.id RETURNING ...
FOR UPDATE SKIP LOCKEDunder the fetch tx; N workers poll concurrently without contention. Ordering: priority DESC, FIFO bycreated_on, id tiebreak (dml.rs:139).- Batching:
fetch_job= limit 1 (job_ops.rs:109-114);fetch_jobs(queue, batch_size)(job_ops.rs:136-150). No server-side long-poll — poll-only; visibility isstart_after < now()+ caller cadence. - Lease/visibility: no lease renewal, no
extendcall. Right-to- run bounded byexpire_seconds(returned asexpire_in, job.rs :199): the monitor sweep fails anyactivejob withstarted_on + expire_seconds < now()(dml.rs:288-303). At this revision the sweep is not scheduled by any background loop — a crashed worker's job staysactiveuntil someone callsClient::force_maintain(); the e2e test documents this explicitly (tests/e2e/maintenance.rs:49-66: "still active though time is exceeded" until forced). - Completion/ack:
complete_jobsrequiresstate = 'active'strictly; setsstate='completed', completed_on=now(), output=$3(dml.rs:223-237). Failure/nack:fail_jobs_by_jidsallowsstate < 'completed'(includes created) and runs the delete+reinsert retry/fail/dlq CTE (dml.rs:279-463); failure details land inoutput(fail_job_with_details, job_ops.rs:201-240). - Retry-count counting rule (off-by-design quirk):
retry_countincrements on re-fetch, not on failure — first delivery leavesretry_count = 0(asserted,tests/e2e/job_change_state.rs:95-98). With terminal conditionretry_count < retry_limit → retry, else failed(dml.rs:344-346),retry_limit = ngives n+1 total attempts (e2e confirmsretry_limit=1→ two attempts, job_change_state.rs:100-109).
4. Retry and backoff
- Where configured: per-job (
Job::retry_limit/retry_delay/ retry_backoff) with per-queue fallback (Queue::...), resolved at enqueue via COALESCE; queue defaults on thequeuerow (ddl.rs :36-40). - Counting: §3 — increment on re-fetch; terminal at
retry_count = retry_limit. - Strategies (dml.rs:352-361, the
fail_jobsCTE recompute ofstart_afteron reinsert):retry_backoff = false→ fixed:start_after = now() + retry_delay(retry_delay = 0 default ⇒ immediate retry).retry_backoff = true→ exponential with full jitter against the baseretry_delayplus an increment independent of base:now() + LEAST(retry_delay_max, retry_delay + (2^LEAST(16, retry_count+1) / 2 + 2^LEAST(16, retry_count+1)/2 * random())) * interval '1s'— jitter uniform ±100% of the doubling term, exponent capped at 16, capretry_delay_max(NULL = uncapped; PGLEASTignores NULLs).- Terminal failure keeps
start_afterunchanged (dml.rs:352). - No custom user-supplied backoff function; only the fixed/
exponential flag. (Node pg-boss had
retryBackoffsimilarly; upstream also hadretryDelayMax, carried in SQL but never exposed in the Rust type API.)
- Mechanics: failure runs
DELETE ... RETURNING *then re-INSERTs the same row asretry(newstart_after,completed_on=NULL) orfailed(withcompleted_on=now()— failure is semantically completion) (dml.rs:306-433).ON CONFLICT DO NOTHINGon the retried insert + afailed_jobsinsert for ids not inretried_jobshandles singleton-policy conflicts during retry. Defaults: retry_limit 2, retry_delay 0 s, backoff false. Timeout-failure (monitor sweep) writes output{ "value": { "message": "job timed out" } }— deliberate byte-compatibility with node pg-boss (dml.rs:300-302).
5. Dead-letter
- Trigger: terminal
failedwhenretry_count = retry_limitduring any fail path (workerfail_jobor monitor-timeout sweep). Thefail_jobsCTE then dead-letters (dml.rs:434-457): for eachfailedrow whosedead_letteris set (copied at enqueue from the queue row), insert a new job on the DLQ —name = r.dead_letter, data, output copied; retry_limit/retry_backoff/retry_delay from the DLQ queue row; keep_until = now() + q.retention_seconds; deletion_seconds = q.deletion_seconds— gated byJOIN queue q ON q.name = r.dead_letter(a nonexistent DLQ silently skips) and implicitlyname <> dead_letter(also enforced by thequeue.dead_letterCHECK, ddl.rs:43). - Where: the DLQ is just another regular queue (rows in the
same
jobtable with a differentname), created ahead of time (README example, README.md:27-36). No dedicated dead-letter table, no distinct state. A DLQ-consumed job'spolicyisNone(src/job.rs:204-208) so singleton indexes don't police DLQ redelivery. - Requeue surface: none. Available ops on a failed/dead job:
get_job(inspect incl. output andretry_count),delete_job(s);resume_jobdoes not coverfailed(onlycancelled → created, dml.rs:194-208). DLQ resubmission is by hand (get payload, send anew). - Retention/cleanup: governed entirely by timestamps —
keep_until(archive-at) anddeletion_seconds(delete-after-completion). At this revision nothing sweeps them: no background maintenance, no archive table. Git history shows anarchivetable +archive_jobsprocedure existed (f63d9dc) and was removed inf9db249("Adjust ddl, retire create_job procedure") during the upstream-schema convergence;queue_ops::delete_queue's doc still claims "Any jobs in the archive table are retained" (src/client/public/queue_ops.rs:55-57) — a stale reference to the removed archive design.
6. Maintenance
- What exists: exactly one op —
Client::force_maintain()(src/client/public/maintain_ops.rs:17-24) runningfail_jobs_by_timeout(expireactivejobs paststarted_on + expire_seconds; retry/fail/DLQ per §4-5; returnsMaintenanceStats { expired }only). Doc wording ("operations that are normally performed on a schedule, such as expiration, archival, and dropping", maintain_ops.rs:15-16) describes intent beyond what's wired. - Reserved but unwired schema: the
queuetable carries full monitor/maintenance state —monitor_on,maintain_oncadence markers, countersdeferred_count/queued_count/active_count/ total_count,warning_queued,singletons_active text[](ddl.rs :46-53) — none read or written by any DML at this revision.fail_jobs_by_timeouthas TODOs for queue-scoped sweeping and monitor-on coordination (dml.rs:289-293); a placeholder stub_g()for maintenance exists unpublished (dml.rs:466-471). - Cadence defaults: none — no background poller/loop in the client;
Clientis a pure pool + statement holder (src/client/mod.rs:50-55). Callers poll and maintain. - Multi-node coordination: advisory-lock-based for one thing only —
install_appwraps all DDL in a transaction holdingpg_advisory_xact_lock(hash(current_database || '.pgboss.{schema}$key'))withSET LOCAL lock_timeout/idle_in_transaction_session_timeout = 30000(src/sql/mod.rs:5-21, 38-58) — concurrent-bootstrap safe (tests/e2e/queue.rs:36-52), and lets the port coexist with a node pg-boss instance creating the same schema. No leader election, no distributed maintenance runtime;version.cron_on(ddl.rs:24, logged at connect,src/client/mod.rs:67-71) is upstream's reserved column for single-owner cron/maintenance coordination, unused here.
7. Cron / repeatable jobs
Upstream node pg-boss has schedule/unschedule/getCron with cron
expressions, every intervals, timezones, and singleton dedup. The
port has the storage shape and none of the machinery:
{schema}.scheduleexists with exactly the upstream columns:name → queue FK, key text (default '', PK (name,key) = dedup key), cron text, timezone text, data jsonb, options jsonb(ddl.rs:63-80).- No Rust API creates schedule rows; no cron parsing, no scheduler
tick, no
cron_onusage beyond the column and a connect-time log. Job::singleton_for+singleton_keyare a related-but-different mechanism: time-bucketed throttling (singleton_on= largest epoch bucket of sizesingleton_for, dml.rs:96-100) enforced by unique indexi4(ddl.rs:157-158), surfaced asError::Throttled(job_ops.rs:60-64; e2etests/e2e/job_send.rs:181-215).- Missed-boundary catch-up: N/A — no cron executor. Timezone: column exists, unimplemented.
8. What the Rust port dropped / left thinner than node pg-boss
- No background maintenance/monitor runtime — no
monitorInterval, no per-queue monitor cadence; expiration (and stalled-activerecovery) only viaforce_maintain(). The queue table's whole monitor/counters column set is dead weight at this revision. - No cron scheduler —
scheduletable andversion.cron_onare scaffolding only. - No archive stage — node pg-boss moves expired jobs to an
archivetable; the port had it, removed it (f9db249). Pipeline is active-table → (retry | failed | DLQ copy) → manual/never deletion. - No worker framework — no
work()/WorkOptions, worker pools, concurrency-per-queue, batch handlers,onCompletewiring; the library stops at fetch/complete/fail primitives (worker loops are the caller's — cf.src/bin/loadtest.rs:49-76). - No
onCompletechild jobs —output jsonbon the job itself, thesubscriptiontable an unused remnant. - No LISTEN/NOTIFY wakeup / no pub-sub — pure polling; the
subscriptiontable suggests intent, nothing more. - No maintenance-queues-as-jobs design — node pg-boss v10
self-manages by enqueuing internal jobs on
pq_<schema>queues; the port replaces that with direct SQL sweeps plus TODOs (dml.rs:289-293). - API-surface gaps vs SQL —
retry_delay_max,delete_after,singleton_offsetconsumed by the enqueue SQL but without setters; DLQ redrive absent;send_jobcannot target an unregistered queue (needs explicitcreate_queue). - Migration story — "must be >= v26, else panic; no upward
migrations" (
src/client/mod.rs:72-86). - Preserved fidelity — schema deliberately diffed against node
pg-boss @
3da860f0(Makefile:58,sql/mod.rs:4); state enum, index set (i1..i6), timeout error message, advisory-lock scheme, epoch-bucket singleton math all match upstream — schema-comparable to pg-boss v10.
9. Carry-in observations for the alkstore pg engine
- The enum-ordered state machine + ordinal comparisons are elegant but fragile (insertion order is semantic; adding a state shifts everything).
- Everything atomic and contention-free happens in one SQL
statement with
FOR UPDATE SKIP LOCKED— enqueue resolves defaults in SQL against the queue row; fetch claims in one CTE; fail runs delete+reinsert in one CTE including DLQ fan-out. The core pattern to keep or consciously redesign. - Retry counting on re-fetch (not on fail), total attempts = retry_limit + 1 — a subtle contract anyone re-deriving must decide to keep or fix.
singleton_onepoch-bucketing + partial unique indexes gives O(1) throttle correctness without locks — reusable if dedup-key semantics ever enter scope.- The port shows a polling + explicit maintenance-call architecture can ship with zero background machinery — but the crashed-worker recovery gap (no automatic expiration sweep) is exactly the hole a real store fills with either a built-in maintenance loop or upstream's maintenance-via-jobs design.