docs: resolve OQ-05 + OQ-09 — queue semantics depth (ADR-010) and scheduler collapse (ADR-009)

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.
This commit is contained in:
glm-5.3-flash committed 2026-10-05 03:10:52 +00:00
1 parent 7ad8ac56bc
commit 79a135c934
13 files changed
+1862 -197

No files matched your search

+18 -17
View File
@@ -24,10 +24,10 @@ pending architecture review and OQ resolution.
| Doc | Status | Purpose | Key OQs |
|---|---|---|---|
| [overview.md](overview.md) | draft | Crate family, feature surface, non-goals, evidence base | — |
| [core-contract.md](core-contract.md) | draft | The unified trait surface, delivery guarantees, tx seam | OQ-05, OQ-08, OQ-09, OQ-10 |
| [engine-sqlite.md](engine-sqlite.md) | draft | SQLite engine: honker-core/rusqlite mapping | OQ-05, OQ-06, OQ-09 |
| [engine-postgres.md](engine-postgres.md) | draft | Postgres engine: tokio-postgres/LISTEN mapping | OQ-05, OQ-08, OQ-09 |
| [queues.md](queues.md) | draft | Queue/scheduler/outbox semantics depth frame | OQ-05, OQ-09 |
| [core-contract.md](core-contract.md) | draft | The unified trait surface, delivery guarantees, tx seam | OQ-08, OQ-10 |
| [engine-sqlite.md](engine-sqlite.md) | draft | SQLite engine: honker-core/rusqlite mapping | OQ-05 (resolved), OQ-06, OQ-09 (resolved) |
| [engine-postgres.md](engine-postgres.md) | draft | Postgres engine: tokio-postgres/LISTEN mapping | OQ-05 (resolved), OQ-08, OQ-09 (resolved) |
| [queues.md](queues.md) | draft | Queue/scheduler/outbox semantics depth (resolved: ADR-009/ADR-010) | OQ-06 (rides) |
| [deployment.md](deployment.md) | draft | Host semantics, connection budgets, knobs, matrix | OQ-08 |
| [open-questions.md](open-questions.md) | draft | OQ tracker (promoted from OQ-ST register) | — |
@@ -43,29 +43,30 @@ pending architecture review and OQ resolution.
| [006](decisions/006-wake-and-delivery-contract.md) | Wake contract — opaque wake + re-read; notify-vs-streams split | Accepted |
| [007](decisions/007-transactional-seam.md) | Transactional seam — caller-held `TxHandle`, `*_tx` methods | Accepted |
| [008](decisions/008-contract-v1-pinning.md) | Contract v1 surface pinning — surface partition, `TxHandle` shape, wake type, reserved strings, error taxonomy | Accepted |
| [009](decisions/009-scheduler-collapse.md) | Scheduler collapse — queues + `schedule()`/`run_schedules`, `@every`-only v1, boundary guarantee row | Accepted |
| [010](decisions/010-queue-semantics-depth.md) | Queue semantics depth — visibility/renewal, backoff curve, dead-letter, no-stranded-rows sweep, schema layout | Accepted |
## Open Questions
Tracked in [open-questions.md](open-questions.md) (OQ-01..NN; the
Phase 0 register's OQ-ST-01..08 promote one-to-one — OQ-NN mirrors
OQ-ST-NN — with new Phase 1 questions appended after). Highlights,
in suggested resolution order (OQ-04 resolved):
in suggested resolution order (OQ-04, OQ-09, OQ-05 resolved):
- ~~**OQ-04** (high): contract pinning — exact trait shape, handle
representation, error taxonomy, reserved strings, guarantee rows
for locks/scheduler~~ — **resolved** (2026-10-05,
[ADR-008](decisions/008-contract-v1-pinning.md); the scheduler
guarantee row transferred to OQ-09).
- **OQ-05** (high): queue semantics depth — retry/backoff/dead-letter/
sweep design.
- **OQ-09** (medium): scheduler as first-class mechanism vs queues +
`schedule()`; owns the scheduler guarantee row.
- **OQ-06** (high): honker-core quality read — fork-trigger gate.
- **OQ-06** (high): honker-core quality read — fork-trigger gate;
now carries two concrete fork candidates from the queue-depth
design (the no-stranded-rows sweep fix, dead-row `get_job`
visibility — complement vs fork is the read's call).
- **OQ-08** (medium): capability-surface shape.
- **OQ-10** (medium): contract versioning across engine crates.
Resolved (Phase 0, kept with resolutions): OQ-01 (feature scope),
OQ-02 (crate split), OQ-03 (drivers), OQ-07 (extension surface cut).
Resolved (kept with resolutions): OQ-01 (feature scope), OQ-02
(crate split), OQ-03 (drivers), OQ-07 (extension surface cut),
~~OQ-04~~ (contract v1 pinning —
[ADR-008](decisions/008-contract-v1-pinning.md)), ~~OQ-09~~
(scheduler collapse — [ADR-009](decisions/009-scheduler-collapse.md);
scheduler guarantee row pinned), ~~OQ-05~~ (queue semantics depth —
[ADR-010](decisions/010-queue-semantics-depth.md)).
No deferred OQs: all open questions are actionable Phase 1 work with
complete evidence bases.
+94 -37
View File
@@ -10,10 +10,13 @@ specifies WHAT the contract is; the per-engine specs map it onto their
machinery; ADRs carry the WHY. The starting artifact is the honker-rs
surface (the Phase 0 §Interface finding) scoped to the
inventory-confirmed features — [ADR-002](decisions/002-feature-scope.md)
— and pinned against it. The exact surface is pinned by
— and pinned against it. The base surface is pinned by
[ADR-008](decisions/008-contract-v1-pinning.md) (contract v1 partition,
`TxHandle` representation, wake type, reserved strings, error
taxonomy, config split); the *obligations* are this document.
taxonomy, config split); ADR-009/ADR-010 add the first *post-v1
contract extensions* (scheduler collapse surface, `QueueOpts` depth —
versioning discipline for such extensions is OQ-10's); the
*obligations* are this document.
## Concepts
@@ -25,8 +28,8 @@ taxonomy, config split); the *obligations* are this document.
streams, queues (+outbox), locks, scheduler. Delivery guarantees are
pinned per mechanism in [ADR-006](decisions/006-wake-and-delivery-contract.md)'s
table (notify/streams/queues), extended with the locks row by
[ADR-008](decisions/008-contract-v1-pinning.md); the scheduler row
rides OQ-09's collapse decision.
[ADR-008](decisions/008-contract-v1-pinning.md), the scheduler row
and collapse by [ADR-009](decisions/009-scheduler-collapse.md).
*(Guiding principles are defined in `docs/research/phase-0.md`
§Vision and cited by number throughout this directory.)*
- **Wake** — the opaque "something changed, re-read" signal the wake
@@ -68,6 +71,9 @@ store.queue(name, opts) -> Queue handle.enqueue_tx
store.try_lock(name, owner, ttl) -> Option<Lock>
store.listen(channel) -> Box<dyn WakeReceiver>
store.outbox(name) -> Outbox
store.schedule(name, spec, queue, payload, opts) -> Result<Schedule>
store.unschedule(name) -> bool
store.run_schedules(stop) -> Result<()>
```
Payloads cross the trait as core value types (non-generic trait
@@ -140,21 +146,38 @@ business-tx shape).
### queues
Durable at-least-once work
([ADR-002](decisions/002-feature-scope.md)). The semantics-depth design
is [queues.md](queues.md)'s (OQ-05); the contract-level obligations
here:
([ADR-002](decisions/002-feature-scope.md)). Depth pinned by
[ADR-010](decisions/010-queue-semantics-depth.md); contract-level
obligations:
- `enqueue` / `enqueue_tx` — commit-atomic; `EnqueueOpts { delay,
run_at, priority, max_attempts, expires }` (the v1 skeleton,
[ADR-008](decisions/008-contract-v1-pinning.md) §1; deeper
`QueueOpts`/retry surface rides OQ-05).
run_at, priority, max_attempts, expires }`; queue-level
`QueueOpts { visibility_timeout_s, max_attempts, backoff_base_s,
dead_letter_retention_s }` ([ADR-010](decisions/010-queue-semantics-depth.md)
§3/§4).
- `claim_one` / `claim_batch` — exactly-once handout under concurrency
(POC-pinned on both engines); `cancel`, `get_job` — job management
and inspection.
- Job handle: `ack / retry / fail / heartbeat`; visibility timeouts
make at-least-once *safe* (work re-appears if a claimant dies).
- `sweep_expired` — maintenance entry point (design in
[queues.md](queues.md)).
(POC-pinned on both engines); claim ordering: priority DESC, then
ready-time, then enqueue order (FIFO under equal priority).
- Job handle: `ack / retry / fail / heartbeat`. `ack` deletes the row;
`heartbeat(extend)` is *renewal* — an absolute reset of the claim
deadline (extend is the new full deadline from now, **not additive**
to elapsed time); late heartbeat refused; **a reclaim consumes an
attempt**; `retry(err, None)` computes the queue's equal-jitter
exponential delay (range pinned by
[ADR-010](decisions/010-queue-semantics-depth.md) §3), `retry(err,
Some(d))` overrides; `fail` = immediate dead-letter. Dead letters:
move-to-dead storage, `get_job` sees dead rows (with
`last_error`/`died_at`), retention via `dead_letter_retention_s`
(default forever), no redrive API
([ADR-010](decisions/010-queue-semantics-depth.md) §1–§4).
`QueueOpts` stamp onto the job row at enqueue (§3a) — queues are
names, not config owners. The remaining v1-skeleton ops (`ack_batch`,
`cancel` — unconditional delete, not an interrupt) pin in
[ADR-010](decisions/010-queue-semantics-depth.md) §1.
- `sweep_expired(queue)` — moves *every* past-expiry row (any state)
to dead + enforces dead-letter retention — the no-stranded-rows
property ([ADR-010](decisions/010-queue-semantics-depth.md) §5);
cadence recipe in [queues.md](queues.md) (no ambient sweeper).
### named locks
@@ -182,10 +205,17 @@ Same evidence base and guarantee as queues
### scheduler
Cron/`@every` enqueueing into named queues, leader-elected where the
engine has peers (SQLite: single-host, no election needed beyond
documented posture; Postgres: advisory-lock election). Whether this is
a first-class mechanism or queues + `schedule()` is OQ-09.
Collapsed into queues per [ADR-009](decisions/009-scheduler-collapse.md):
`schedule(name, spec, queue, payload, opts)` (upsert by name),
`unschedule(name)`, and the opt-in `run_schedules(stop)` runner
(leader-elected where the engine has peers — the leadership lock is
the reserved name `__alkstore_scheduler`). Schedules never fire
without a runner (no ambient timers). v1 spec grammar:
`@every <n><unit>` only (`s|m|h|d`); cron strings are rejected
(`InvalidSpec`) pending a consumer-inventory row naming wall-clock
cron. Boundary guarantee:
[ADR-009](decisions/009-scheduler-collapse.md) §4 (at-least-once per
elapsed boundary, bounded catch-up with skip-forward past the cap).
## Cross-cutting contracts
@@ -193,12 +223,21 @@ a first-class mechanism or queues + `schedule()` is OQ-09.
The per-mechanism table in
[ADR-006](decisions/006-wake-and-delivery-contract.md) is the contract
of record for notify/streams/queues; this spec inherits it and adds
the consumer obligations:
of record for notify/streams/queues; the locks row is
[ADR-008](decisions/008-contract-v1-pinning.md) §7's and the scheduler
row is [ADR-009](decisions/009-scheduler-collapse.md) §4's. This spec
inherits them and adds the consumer obligations:
- wake idempotence (notify's at-least-once, coalescing behavior);
- explicit offset saves (streams);
- visibility-timeout budgeting (queues);
- visibility-timeout budgeting — heartbeat inside the deadline for
long work; the dual-execution window and reclaim-eats-attempt rules
are documented consumer obligations
([ADR-010](decisions/010-queue-semantics-depth.md) §2);
- running `run_schedules` for schedules to fire at all, and scheduling
sweep cadences via the collapse recipe — timeliness is the
consumer's, correctness-of-transition the engine's
([ADR-010](decisions/010-queue-semantics-depth.md) §6);
- `*_tx` operations are *only* durable after the caller's commit —
rollback drops job rows, event rows, notifications, and offset saves
together (the no-ghosts property, POC-pinned on both engines).
@@ -210,9 +249,14 @@ the consumer obligations:
`Error`, with v1 variants `PayloadTooLarge` (universal; produced
pg-side, contract-wide matchable), `ReservedName`, `InvalidName`,
`Closed`, `Codec`, and the opaque `Database` fallback (engine detail
preserved via the source chain). Pinning rule: a variant exists only
when callers can act differently on it. Queue claim "no work" and
job/lock boolean results are values, not errors.
preserved via the source chain). Post-v1 additions:
`InvalidSpec` (schedule spec grammar) and `LeadershipLost`
(`run_schedules` return on leadership loss) — both from
[ADR-009](decisions/009-scheduler-collapse.md) §6, the only
taxonomy deltas so far. Pinning rule: a variant exists only when
callers can act differently on it. Queue claim "no work" and job/lock
boolean results are values, not errors; dead-letter is observable
state (`get_job`), not an error.
### Capability surface
@@ -223,10 +267,11 @@ wake-cadence knobs) — is OQ-08's decision
### Naming / reserved namespace
Consumer-visible names (channels, streams, queues, locks) share
engine-visible namespaces on Postgres (LISTEN channel names are
server-global per database). Pinned by
[ADR-008](decisions/008-contract-v1-pinning.md) §4:
Consumer-visible names (channels, streams, queues, locks, and
schedule names) share engine-visible namespaces on Postgres (LISTEN
channel names are server-global per database). Pinned by
[ADR-008](decisions/008-contract-v1-pinning.md) §4 (schedule names
added by [ADR-009](decisions/009-scheduler-collapse.md) §1):
- Reserved prefix **`__alkstore_`**, engine-independent, applies
across all name kinds; the one v1-reserved string is
@@ -241,9 +286,10 @@ server-global per database). Pinned by
prefix (the outbox's backing queue is derived under the prefix —
honker's `_outbox:{name}` scheme is not inherited verbatim).
- SQLite: honker's `_honker_*` internal table family is storage-
internal (not consumer namespace); Postgres: queue/stream tables are
schema-scoped (layout rides OQ-05) — the *channel* namespace is
this section's.
internal (not consumer namespace); Postgres: engine tables are
schema-scoped — one engine-owned schema
([ADR-010](decisions/010-queue-semantics-depth.md) §8) — the
*channel* namespace is this section's.
## Verification backlog
@@ -268,6 +314,14 @@ before the engine specs are called `stable`:
- **Concurrent `try_lock` loser/error behavior** on SQLite (the pg
side returns cleanly; honker's busy-path under lock contention is
the thing to pin).
- **Scheduler + depth properties on Postgres** — the SQLite side's
tick/leader/catch-up machinery rides honker's test-pinned
implementation; the pg engine's re-derived tick (boundary advance,
64-cap skip-forward, leadership-loss discipline) and the
[ADR-010](decisions/010-queue-semantics-depth.md) depth properties
(visibility reclaim consuming attempts, dead-letter moves, the
no-stranded-rows sweep) pin in the contract suite at
implementation.
## Design Decisions
@@ -278,6 +332,8 @@ before the engine specs are called `stable`:
| [006](decisions/006-wake-and-delivery-contract.md) | Wake contract | opaque wake + re-read; notify-vs-streams guarantee split |
| [007](decisions/007-transactional-seam.md) | Tx seam | caller-held handle, `*_tx` methods, native commit-atomicity |
| [008](decisions/008-contract-v1-pinning.md) | Contract v1 pinning | surface partition, `TxHandle` trait shape, `Wake` type, reserved strings, error taxonomy, config split, locks guarantee row |
| [009](decisions/009-scheduler-collapse.md) | Scheduler collapse (post-v1 extension) | queues + `schedule()`/`run_schedules`, `@every`-only v1, boundary guarantee row |
| [010](decisions/010-queue-semantics-depth.md) | Queue depth (post-v1 extension) | visibility/renewal, opts stamping, backoff curve, dead-letter, sweep, layout |
## Open Questions
@@ -285,10 +341,11 @@ Open questions are tracked in
[open-questions.md](open-questions.md). Key
questions affecting this document:
- **OQ-09**: scheduler as first-class mechanism vs queues + schedule —
also owns the scheduler guarantee row ([open](open-questions.md))
- **OQ-10**: contract versioning discipline across engine crates
([open](open-questions.md))
- **OQ-08**: capability-surface shape ([open](open-questions.md))
- **OQ-05**: queue semantics depth — the v1 skeleton's extension path
([open](open-questions.md))
Resolved on this document's surface: **OQ-09** (scheduler collapse —
[ADR-009](decisions/009-scheduler-collapse.md)) and **OQ-05** (queue
semantics depth — [ADR-010](decisions/010-queue-semantics-depth.md)),
2026-10-05.
@@ -43,7 +43,12 @@ has recorded the need):
- **scheduler** — cron/`@every` enqueueing into named queues,
leader-elected. Documented-thin; the family-wide
"who sweeps/renews/reaps" need. May collapse into queues if design
shows scheduler = queues + tick (tracked in OQ-09).
shows scheduler = queues + tick (tracked in OQ-09). *(Narrowed
2026-10-05 by [ADR-009](009-scheduler-collapse.md): collapsed to
queues + `schedule()`/`run_schedules`, spec grammar `@every`-only —
the "cron" label reflected honker's menu, not any inventory row's
body; no consumer names wall-clock cron — grammar pinned by
[ADR-009](009-scheduler-collapse.md) §2.)*
**Cut-flag** (no named consumer; not silently included):
@@ -0,0 +1,244 @@
# ADR-009: Scheduler collapses into queues — `schedule()` + runner, `@every`-only v1
## Status
Accepted (2026-10-05, Phase 1 — OQ-09's resolution; resolves the
scheduler row transferred by [ADR-008](008-contract-v1-pinning.md) §7)
## Context
[ADR-002](002-feature-scope.md) took the scheduler in as
"documented-thin": the family-wide "who sweeps / renews / reaps, and
when" problem is real, but only just begins to have consumers lean on
it. OQ-09's decision rule (inherited from
[queues.md](../queues.md)'s framing of it): decide against the
inventory rows, not against the whole honker menu — the scheduler is a
first-class mechanism if consumers need inspectable/pausable schedule
*objects* (`add/pause/resume/update/list/remove`), and otherwise it
collapses into queues + a `schedule()` call with the tick machinery
absorbed.
The evidence for the decision:
- **alkblobs (sweep cadence)** — [documented:
consumer-inventory.md](../../research/consumer-inventory.md) — needs
"run the sweep every N ticks with leader election across a fleet."
An interval, not a schedule object.
- **alkfs (orphan reaping)** — [recalled: OQ-FS-07] — a periodic
orphan-cleanup pass. An interval.
- No consumer document (pinned, documented, or operator-authority)
names any of: wall-clock cron expressions, per-schedule timezones,
pausing/resuming schedules without unregistering, listing schedule
objects, or updating a schedule's expression in place.
So the collapse question itself answers cleanly. Two secondary
questions ride along and this ADR pins them with the collapse:
- **Which schedule specs does v1 accept?** Honker's parser accepts
5-field cron, seconds-field cron, and `@every <n><unit>` — but its
cron boundary arithmetic is hard-wired to *host local time*
(`chrono::Local`, honker cron.rs), which is brittleness a contract
should not inherit: a server TZ change silently re-times every
stored schedule, and the two engines would diverge if one inherited
local time and the other pinned UTC.
- **What runs the tick?** Every consumer punting "the embedder owns
cadence" is exactly the hand-rolling the substrate exists to kill —
but an *ambient* runner violates the family's no-ambient-timers
posture (alkblobs ADR-005). The line between them is opt-in.
## Decision
### 1. The scheduler collapses into queues
The contract surface is two registration methods on `Store` plus one
runner; there is no `Scheduler` handle, no schedule objects, no
inspection surface:
```text
store.schedule(name, spec, queue, payload, opts) -> Result<Schedule>
store.unschedule(name) -> bool
store.run_schedules(stop) -> Result<()> // runs until `stop`
```
- `schedule` upserts by `name` (re-register replaces the whole row) —
honker's register semantics, which makes update = re-register and
pause = unregister + re-register. Both are honest v1 substitutes;
nothing in the inventory needs finer control.
- `Schedule` is a value carrying the registered row (name, spec, queue,
opts) — read-back for confirmation, not a mechanism handle.
- **Name validation**: schedule names follow the queue-name rules —
non-empty (`InvalidName`) and reserved-prefix rejection
(`ReservedName`) at the entry point; `schedule()`/`unschedule()` are
name-bearing entry points in the
[ADR-008](008-contract-v1-pinning.md) §4 sense (their consumers-side
obligation extends the entry-point list: a reserved-prefix schedule
name could collide with the machinery's derived names). No charset
rule beyond that in v1.
- `run_schedules` is the tick: acquire the leadership lock, loop
{renew lock — on loss, *return before ticking*; fire every due
boundary; sleep until the next due boundary or `stop`}. Losing
leadership ends the runner — the consumer's recipe is to respawn it
(documented); the lock loss always precedes any stolen fire, so
respawn never double-fires a boundary.
- **Stop and return semantics**: `stop` is a cancellation token
(a cloneable handle whose flip ends the sleep early — the exact
type shapes at implementation, engine-crate docs). A leadership
loss returns `Err(LeadershipLost)` — a distinct, matchable outcome
from the clean `Ok(())` a stop produces; the respawn recipe matches
on it. (One taxonomy note: `LeadershipLost` is a *return-shape*
variant of the runner, pinned by the same act-differently rule as
ADR-008 §5 — the caller acts on it by respawning.)
- Schedules do not fire until a consumer runs `run_schedules` — the
opt-in is explicit, the no-ambient-timers posture holds,
[ADR-002](002-feature-scope.md)'s scheduler evidence is served, and
consumers who never schedule never pay for the runner.
### 2. v1 spec grammar: `@every <n><unit>` only
`spec` accepts exactly `@every <n><unit>` with `n` a positive integer
and `<unit>` in `s | m | h | d` (honker's `@every` grammar,
honker cron.rs). Everything else — 5-field cron, seconds-field cron —
parses as `InvalidSpec { spec }` and is rejected.
- **Why not cron strings**: no consumer row names one (checked against
the inventory, not against honker's menu), and accepting them forces
a timezone answer now. `@every` durations are TZ-free — the TZ
question dissolves instead of being answered.
- **Extension path** (when a consumer names wall-clock cron via a
consumer-inventory row): cron-string support with an explicit TZ
decision attached (UTC-pinned contract or a per-schedule TZ
parameter — decided by that extension, not assumed here).
- This keeps the SQLite engine off honker's cron boundary machinery
entirely (`@every` advance is `next = last + interval`) — the
local-TZ brittleness is not inherited, and a dependency requirement
is not created. *(Consumer-inventory rows change this only via the
row-first discipline.)*
### 3. Enqueue semantics per boundary
At each due boundary the tick enqueues the schedule's `payload`
(`serde_json::Value`, uniform with all payloads) into the named queue
with `ScheduleOpts { priority, max_attempts, expires }` — the
[EnqueueOpts](008-contract-v1-pinning.md) fields minus
delay/run-at (a boundary fire is never delayed by its own config; the
schedule *is* the delay). The enqueued work is ordinary queue work: it
inherits the queue mechanism's delivery guarantees
([ADR-006](006-wake-and-delivery-contract.md) queues row). The
scheduler adds no delivery machinery of its own.
### 4. Boundary guarantee — the scheduler's row in the delivery table
Extends [ADR-006](006-wake-and-delivery-contract.md)'s table,
resolving the transfer from [ADR-008](008-contract-v1-pinning.md) §7:
| Mechanism | Durability | Replay | Atomicity | Guarantee |
|---|---|---|---|---|
| scheduler (`@every`) | schedule rows + boundary fires are durable rows | n/a (fire = enqueue; the work's durability is the queues row) | fire and boundary-advance commit atomically under a row lock on the schedule row (crash mid-tick rolls back both — the boundary refires; a rogue ticker cannot double-fire) | **at-least-once per elapsed boundary while a leader runs**, with bounded catch-up: downtime fires replay boundary-by-boundary up to a fixed contract-constant per-tick cap (64, honker parity), beyond which missed boundaries skip forward to the next future boundary |
- The **bounded-catch-up + skip-forward** semantics are pinned as the
honest contract: a consumer who pauses runners for a week does not
come back to 10,000 stacked sweep jobs. (For an `@every 60s` sweep
under an outage, replays beyond the cap are useless anyway.) **The
cap is a fixed contract constant (64 per task per tick, honker
parity)** — not a knob; like the backoff cap, no consumer names one,
and the skip-forward is documented behavior, not configuration.
- **One firer per boundary** is pinned engine-generically by two
layers: (a) the leadership lock — one leader at a time, a leader
that lost the lock returns before ticking; (b) **the boundary
advance is a transactional claim on the schedule row**: the fire
transaction takes a row lock on the schedule row, re-checks
`next_fire_at <= now` under the lock, enqueues, and advances — so
even a rogue second ticker (an operator bypassing the leadership
lock) observes a boundary as already-fired and skips it, on either
engine (Postgres: `SELECT … FOR UPDATE` on the row inside the fire
tx; SQLite: the same guarantee via writer serialization —
`BEGIN IMMEDIATE`). Layer (b) is the correctness floor, layer (a)
the efficiency (one waiter instead of row-lock contention);
crash mid-tick rolls back enqueue + advance together (the boundary
refires).
- **Leader election** is a named lock under the reserved prefix —
`__alkstore_scheduler` (engine-derived name in the lock namespace per
[ADR-008](008-contract-v1-pinning.md) §4's rule; consumers can never
collide with it because every entry point rejects the reserved
prefix). Both engines: each engine's own lock machinery, engine-
internal detail. One firer per boundary across all processes sharing
the store.
### 5. Schedule storage is engine-internal
Schedule rows live in engine-owned storage (SQLite: honker's
`_honker_scheduler_tasks` — reused as-is, its `cron_expr` column
carrying `@every` specs; Postgres: the engine-owned schema per
[ADR-010](010-queue-semantics-depth.md) §8). Not consumer namespace;
schedule *names* take the same validation as queue names
(§1 — non-empty, reserved-prefix rejected). No registry
validation against queue names — enqueuing into a queue nobody has
claimed yet is fine (queues are names, not registered objects).
### 6. Error taxonomy delta
Two variants added by this surface, per ADR-008 §5's act-differently
rule: `InvalidSpec { spec }` — a schedule spec that doesn't parse
under §2's grammar (the caller fixes the spec string; rejected before
storage, no round trip) — and `LeadershipLost` — the `run_schedules`
return when leadership is lost mid-run (the caller acts by
respawning; distinct from clean stop, §1).
## Consequences
**Positive**
- The contract gains the scheduler for two registrations + one runner;
the "who sweeps" problem has a substrate answer that is one `@every`
row and one spawned task.
- The tick machinery (leader election, boundary advance, catch-up cap)
is engine-internal and mostly inherited (SQLite: honker's scheduler
machinery with `@every` specs; Postgres: re-derived, small at this
surface size — the row-locked fire tx is the one mechanism the
contract pins engine-generically, §4).
- No TZ semantics, no cron parser in the contract, no schedule-object
CRUD — the surface stays proportionate to the two inventory rows that
exist.
**Negative**
- Honker's richer scheduler (`add/update/pause/resume/list`, cron
with TZ) is deliberately not consumed at the contract level; consumers
needing inspectable schedules need an inventory row + contract
extension first.
- The opt-in runner is one more thing a deployer must remember: forgot
`run_schedules` ⇒ schedules silently never fire. The mechanism's doc
text must carry this (a silent-nothing deployment posture, the honest
twin of "no ambient timers").
- The catch-up cap is a consumer-visible loss mode for long outages —
documented (skip-forward), not silent: the schedule row's advance
makes the gap inspectable.
- The scheduler surface (and this track's `QueueOpts` depth) are
**post-v1 contract extensions** — the first exercises of the
extension path [ADR-008](008-contract-v1-pinning.md) leaves open;
their versioning discipline (how engine crates track the additions)
is OQ-10's, still open.
## References
- OQ-09 (`docs/architecture/open-questions.md`) — this ADR's
resolution; owns ADR-008 §7's transferred scheduler guarantee row.
- `docs/research/consumer-inventory.md` §scheduler — the two rows the
decision is made against.
- [ADR-002](002-feature-scope.md) — scheduler in-scope posture this
ADR completes (documented-thin → pinned thin).
- [ADR-006](006-wake-and-delivery-contract.md) — the guarantee table
§4 extends; the enqueued work inherits the queues row.
- [ADR-008](008-contract-v1-pinning.md) §1 (scheduler out of v1
pending this decision), §4 (reserved-prefix rule the leadership lock
name follows), §5 (act-differently rule the `InvalidSpec` variant
follows), §7 (the guarantee-row transfer).
- [ADR-010](010-queue-semantics-depth.md) — the queue semantics depth
this ADR's fire semantics compose with; the maintenance recipe (who
sweeps) is recorded there with the collapse machinery.
- Honker reference checkout `/workspace/honker` @ f4e53c6 — scheduler
surface (`packages/honker-rs/src/lib.rs`, `honker-core/src/honker_ops.rs`),
`@every` grammar and local-TZ brittleness (`honker-core/src/cron.rs`),
tick/leader mechanics.
- [core-contract.md](../core-contract.md) — the spec carrying this
surface.
@@ -0,0 +1,353 @@
# ADR-010: Queue semantics depth — visibility, retry/backoff, dead-letter, sweep, layout
## Status
Accepted (2026-10-05, Phase 1 — OQ-05's resolution; composes with
[ADR-009](009-scheduler-collapse.md))
## Context
Contract v1 pinned the queue *skeleton* ([ADR-008](008-contract-v1-pinning.md)
§1) and explicitly left the semantics depth to OQ-05: retry policy
shape, visibility/renewal mechanics, dead-letter move-vs-flag and
retention, sweep/maintenance design, the result-storage cut-flag's
disposition, and queue/stream/lock **table** layout (the co-tenancy
collision surface, [queues.md](../queues.md)).
The evidence base: honker's queue machinery (the SQLite-side incumbent
— `/workspace/honker` @ f4e53c6, whose functions the engine rides
directly) and the pg-boss family (the pg-side design reference —
`/workspace/pgboss-rs` @ 98f7d9e standing in for node pg-boss v10;
design-reference only per [ADR-005](005-dependency-ownership.md)).
Both models were read in full for this resolution (reference notes:
`docs/research/reference-honker-machinery.md`,
`docs/research/reference-pgboss-rs-semantics.md`). The hard
driver-coupled properties (transactional enqueue, exactly-once claim)
are already POC-pinned on both engines.
The two references disagree on several semantics points, and honker
has real gaps (a dead-letter table `get_job` can't see; expired
processing rows no path can reach — the "zombie" hole; no dead-row
retention mechanics; no backoff in the core at all). OQ-05 is the
place to fix what re-derivation should fix and inheriting should
inherit.
## Decision
### 1. The job state machine: three states, delete-on-ack
Contract job states: **`pending` → `processing` → `dead`** (+ absence:
`ack` and `cancel` delete the row).
- Honker's model, adopted whole — one live table, move-to-dead
(physical row move, savepoint-guarded), not the pg-boss seven-state
enum. `max_attempts` is frozen at enqueue; `attempts` counts
every claim.
- **`ack` deletes** (honker parity): completed work leaves no row.
The completed-job-with-output model (pg-boss's `completed` state +
`output` column) is **not** adopted — it is result storage under
another name, and result storage is a cut-flag row
([ADR-002](002-feature-scope.md)); see §7.
- **`get_job` sees dead jobs.** The returned job carries state,
payload, attempts, priority, and the timestamps `created_at`,
`run_at`, `claimed_at`, `expires_at` — plus, dead-only, `last_error`
and `died_at` (no job = `None`). This deliberately *narrows* honker's
surface (honker's `get_job` reads only live rows — post-mortem
diagnosis was SQL-only); diagnosis-by-API is the honest fix, and it
is what makes move-to-dead (vs flag-in-place) inspectable without
raw SQL.
- **`cancel` is unconditional** (honker parity): it deletes the row in
either `pending` or `processing`, regardless of which worker holds
it. Not an interrupt — the holder's next ack/heartbeat returns
false, the same shape as expiry. `ack_batch` is the batch form of
ack (count returned; non-claimed ids silently not counted).
- Claim ordering (same queue): `priority DESC`, then ready-time
(`run_at`) ascending, then enqueue order — FIFO under equal
priority, both references agree; pinned.
### 2. Visibility timeout and heartbeat: explicit renewal, late-heartbeat refusal
One answer, contract-pinned (honker's design; pg-boss's
no-renewal/expire-sweep model is not inherited):
- `visibility_timeout_s` is a per-job stamped value (see §3a's
resolution rule — default **300 s**, honker parity); each claim sets
the row's deadline from the job's stamp.
- **`heartbeat(extend)` is renewal** — an absolute reset of the
claim deadline from now. There is no progress-signal meaning; a
consumer signaling progress uses its own channels. Renewal cadence
is the consumer's obligation: handlers that might outlive the
visibility timeout must heartbeat inside it.
- **A late heartbeat is refused** (deadline already passed ⇒ returns
false): it can never steal the job back from a reclaimer. The
consequence is the honest at-least-once window — between deadline
lapse and another worker's reclaim, the original worker may still
complete: dual-execution is possible and idempotence is the
consumer's job (matches the locks row's silent-expiry honesty,
[ADR-008](008-contract-v1-pinning.md) §7).
- **A reclaim consumes an attempt** (honker's counting): a claim is
an attempt, whether fresh or a visibility reclaim. Documented
explicitly as the contract's counting rule — the footgun is stated,
not discovered: a handler that forgets to heartbeat looks like a
repeatedly-failing job and dead-letters on budget exhaustion.
- Deadline lapse itself is lazy: an expired-claim row becomes
claimable by ordinary claim (no transition fires on lapse); there
is no separate reaper for expired claims, only the budget rules
below and the no-stranded-rows sweep (§5).
### 3. Retry and backoff: explicit delay or the queue's curve
- Job handle: `ack()`, `retry(err, delay)`, `fail(err)`, and
`heartbeat(extend)` — the v1 skeleton ([ADR-008](008-contract-v1-pinning.md)
§1), with `heartbeat`'s meaning pinned by §2 and `renew` remaining
the *lock*-side term ([ADR-008](008-contract-v1-pinning.md) §8's
split).
- **`retry(err, None)`** — the deferred delay — is computed by the
engine from the queue's curve; **`retry(err, Some(d))`** overrides
it. Honker's core takes a caller-supplied delay with no curve;
pg-boss computes in-engine from a base. The contract pins the
curve so both engines compute identically.
- **The curve: equal-jitter exponential**, `base * 2^(attempt-1)`
with uniform jitter over the *lower half* of the doubling range —
`delay ∈ [base·2^(a−1)/2, base·2^(a−1)]` — capped at **1 hour** (a
pinned contract constant; no `backoff_max` knob — no consumer names
one; the explicit-delay override is the escape hatch for bespoke
policies). The range, not a jitter-label, is the definition: this
is *not* AWS-canonical full jitter (uniform over `[0, cap]`), and
the formula above is what both engines compute identically. The
attempt index `a` is the value of `attempts` on the row at the
moment the retry is scheduled (post-claim count, including
reclaims). Jitter prevents sync-failure herding across a fleet;
honker's wrappers have no jitter (a spread-out fleet of failing
workers retries in lockstep) and the pg-boss family's jittered
exponential is the design adopted.
- **`QueueOpts { visibility_timeout_s, max_attempts, backoff_base_s,
dead_letter_retention_s }`**
— queue-level defaults (respectively 300 s / 3 / 5 s / none, honker
parities except retention which is ours); resolved and stamped
per §3a. `backoff_base_s` feeds the curve. This is the
`QueueOpts` depth contract v1 deliberately withheld
([ADR-008](008-contract-v1-pinning.md) §1) — pinned here as the
extension path. **Outbox**: `outbox(name)` carries no opts of its
own; its backing queue uses a *derived* QueueOpts set as engine
defaults (honker's outbox parities: visibility 60 s, max_attempts 5,
backoff base 5 s) — engine-documented, not separately configurable
in v1.
- Explicit `fail(err)` = immediate dead-letter (honker parity; it is
the "stop retrying this" operator). `retry` at exhausted budget =
dead-letter. Exhaustion via reclaim = dead-letter (the pre-claim
sweep of already-exhausted reclaimable rows is an engine-side
laziness optimization — SQLite rides honker's; pg re-derives).
### 3a. QueueOpts resolution: stamped at enqueue, per job
Queues are names, not registered objects — so queue-level `QueueOpts`
need a defined attachment point, and the contract pins one:
- **Every job row carries the resolved opts stamped at enqueue**:
enqueue resolves `EnqueueOpts.max_attempts` over the queue handle's
`max_attempts` (the one per-job override the v1 skeleton names), and
stamps the result — `visibility_timeout_s`, `backoff_base_s`,
`dead_letter_retention_s`, `max_attempts` — onto the job row. Claims,
heartbeats, retries, and sweeps read the **job's own stamps**, never
a live registry.
- **Why stamped**: uniform across engines (no per-queue config state
to invent on either side — SQLite needs no new table beyond honker's
columns, Postgres is job-table columns), and it makes a job's
behavior immutable and inspectable (`get_job` returns the stamps
with the rest of the row) no matter which process's handle enqueued
it. Two processes opening handles with different `QueueOpts` for the
same queue name do not fight — each enqueue stamps its own opts;
the queue name is a routing key, not a configuration owner.
- The consequence, stated: opts changes apply to *future enqueues
only*; in-flight jobs keep their stamps. This is the honest,
durable-row semantics — the same reason `max_attempts` was already
frozen at enqueue in the references.
- SQLite realization note: honker's `claim_batch` takes one uniform
`timeout_s` per call, while per-job stamps make deadlines row-local —
bridging that gap (post-claim re-stamp, per-row claiming, or another
shape) is implementation work over honker's function surface and
rides the same OQ-06 assessment as §5's zombie fix. The *contract*
(deadline from the job's stamp) is engine-independent either way.
### 4. Dead-letter: move, retention = never-expire by default, no redrive API
Move-to-table (both references converge; flag-in-place is not
adopted — a dead row must not occupy the claim path's indexes):
dead rows physically move to engine-owned dead storage with
`last_error` and `died_at`. Triggers: explicit `fail`, `retry` at
budget, exhaustion-by-reclaim, and `expires` lapses (§5).
- **Retention**: dead rows live until explicitly removed by default
(honker's posture). `QueueOpts.dead_letter_retention_s` (default
`None` = forever) makes the sweep enforce per-queue dead TTL (§5) —
the retention *mechanism* is ours (honker has none); the default
stays non-magical. `last_error` source strings are contract-pinned
for the engine-triggered cases: `"max attempts exceeded"` (budget
exhaustion, any path) and `"expired"` (sweep-by-`expires`); explicit
`fail(err)`/`retry(err, …)` carry the *caller's* error string.
- **No redrive/requeue API in v1.** The recipe is pinned instead:
`get_job` (now sees dead rows, §1) + fresh `enqueue` — both
primitives already on the surface. No consumer names an API shape
for redrive; when one does, it gets a contract extension with its
own OQ.
- The engine-owned dead table/rows are storage-internal — not a
consumer namespace (no reserved-prefix question; they're not
name-addressable).
### 5. `sweep_expired`: the no-stranded-rows property
Contract semantics for `sweep_expired(queue)` (the v1 skeleton's
maintenance entry point):
- It moves to dead (`last_error = 'expired'`) **every row of the
queue past its `expires_at` in any state** — pending rows (honker
parity) *and* processing rows whose job-level expiry passed. This
is the **no-stranded-rows property**: an enqueued job with an
`expires` deadline is eventually in exactly one of pending,
processing, or dead — never stuck unreachable. Honker fails this
property: an expired *processing* row whose worker died is
unreachable by claim (predicate requires future expiry), by
pre-claim dead-lettering (same), and by `sweep_expired` (pending
only) — the zombie hole, a real defect found in the reference read.
This ADR fixes it at the contract level; the Postgres engine
enforces it directly; the SQLite engine's realization over honker's
machinery is implementation work **contingent on OQ-06** (a
complement over honker's function surface, or the fork fires and
the fix lands in owned code — the concrete fork candidate the
quality read should weigh).
- When `dead_letter_retention_s` is set, the same sweep also deletes
dead rows past their retention (the only sweeper dead rows ever
have — nothing runs without a caller).
- Sweep is single-statement atomic per queue (SQLite: writer
serialization; Postgres: row-locked `UPDATE … RETURNING`); **no
leader lock is required** for a bare `sweep_expired` call — the
multi-process recipe gets its coordination from §6's machinery, not
from the sweep itself.
### 6. Sweep/maintenance cadence: no ambient sweeper; the collapse recipe
- **The engine ships no background sweeper, no maintenance thread, no
default cadence** — the family's no-ambient-timers posture
(alkblobs ADR-005) and honker's own verified posture (nothing in
honker-core auto-runs; the only implicit maintenance is laziness in
claim paths). The engine's *responsibility* is confined to
correctness-of-transition (§1–§5); **timeliness is the
consumer's**.
- The pinned recipe for the family-wide "who sweeps" problem is the
collapse machinery ([ADR-009](009-scheduler-collapse.md)):
`store.schedule("maintenance", "@every 300s", maintenance_queue, …)`
+ `run_schedules` (leader-elected) + a worker whose handler calls
`sweep_expired` per queue. One schedule row, one worker — the
pattern every consumer was going to hand-roll, answered twice in
the contract instead of N times above it.
- Dead-letter and notifications-table hygiene (SQLite engine):
`prune_notifications`/`prune_notifications_keep_latest` remain
**out of the contract** (ADR-008 §8's disposition, resolved here by
disposition): the notifications table is the SQLite wake
mechanism's *transport* detail — consumers interact with wakes, not
rows. Its hygiene is engine-internal (capped at attach, engine-side
maintenance cadence, engine opts) — the fix for honker's unbounded
`_honker_notifications` growth must not become a consumer's chore.
### 7. Result storage: cut-flag stands
Reconsidered per OQ-05's mandate and **left cut**: no consumer row
names outcome-query-by-id; the §1 delete-on-ack decision keeps
completed jobs out of storage entirely (the pg-boss `output`/completed
-row model is that feature under another name — adopting it would
silent-include the cut flag); the documented workaround (a job writes
its own result record atomically with ack) rides the tx seam
([ADR-007](007-transactional-seam.md)) cleanly: business row + result
write + ack-equivalent in one transaction. Re-entry stays gated on a
consumer-inventory row.
### 8. Table layout: engine-owned schemas/tables; queues are rows, not tables
- **Postgres**: all engine-owned tables (job, dead, stream log,
offsets, schedule rows, internal state) live in one **engine-owned
PostgreSQL schema** — default `alkstore`, overridable per-engine
option — co-tenanted safely with consumer tables (alkblobs ADR-008's
co-tenancy precedent; schema scoping is the collision answer, and
the reserved *name* namespace ([ADR-008](008-contract-v1-pinning.md)
§4) governs the name column values inside it). Queue names become
row values in one job table — **no per-queue tables, no per-queue
schemas** (pg-boss's partition-per-queue opt-in is not inherited;
at this scale one table + partial indexes matches honker's proven
shape and keeps `queue(name)` a name, not a DDL operation — queue
creation is not registry-gated, per
[ADR-009](009-scheduler-collapse.md) §5). Per-queue config storage
is likewise unnecessary (opts stamp onto job rows, §3a).
- **SQLite**: honker's `_honker_*` table family in the caller's
database file — storage-internal per
[ADR-008](008-contract-v1-pinning.md) §4; no new tables minted by
this ADR (job-stamped opts ride honker's existing columns and the
contract's stated resolution rule, no schema change needed — this
is exactly the seam the quality read (OQ-06) assesses). The
family's exact fate rides that read: a fork re-owns the names,
nothing consumer-visible changes.
### 9. Error taxonomy: no delta
The depth adds **no new error variants** — checked against
[ADR-008](008-contract-v1-pinning.md) §5's act-differently rule: no
work claimed is a value; deadline operations are booleans; dead-letter
is observable state (`get_job`), not an error thrown; queues are
unregistered names, so a payload enqueued to a queue nobody claims
yet is not an error (it's durable work awaiting a worker).
`InvalidSpec { spec }` was the scheduler's one variant
([ADR-009](009-scheduler-collapse.md) §6) and is the only taxonomy
delta this track produces.
## Consequences
**Positive**
- The queue mechanism is fully specified end-to-end: both references'
semantics decisions are made once, pinned identically on both
engines, with the references' gaps (zombie rows, dead-job
invisibility, no retention mechanics, no jitter) fixed at the
contract level rather than inherited.
- The counting/heartbeat rules and the opts-stamping resolution are
stated as contract text — the honest at-least-once story
(dual-execution window, reclaim-eats-budget) is documented consumer
obligation, not per-engine surprise, and per-queue configuration
never fights across processes.
- One engine-owned pg schema keeps co-tenancy mechanical.
**Negative**
- Three-state + delete-on-ack means no completion history and no
built-in outcome storage — by design (the cut-flag stands, §7);
consumers needing audit trails build them on the tx seam.
- Honker's zombie fix, dead-row `get_job` visibility, and per-job
visibility stamps over the uniform-claim-timeout function surface
require either an over-machinery complement or a fork on the SQLite
side — the fork calculus gains concrete candidates to weigh
(OQ-06).
- The backoff curve's 1-hour cap and jitter formula are pinned
constants — a consumer needing a different policy uses
`retry(err, Some(d))` per attempt (correct, but manual).
## References
- OQ-05 (`docs/architecture/open-questions.md`) — this ADR's
resolution (with OQ-09, resolved by [ADR-009](009-scheduler-collapse.md)).
- `docs/research/reference-honker-machinery.md` and
`docs/research/reference-pgboss-rs-semantics.md` — the full
reference reads (schemas, state machines, defect list with
file/line cites) this resolution was made over.
- [ADR-002](002-feature-scope.md) — scope (cut-flags §7 stands),
[ADR-005](005-dependency-ownership.md) — the design-reference
postures the two reads serve,
[ADR-006](006-wake-and-delivery-contract.md) — the queues guarantee
row this ADR's §2/§3 make precise,
[ADR-007](007-transactional-seam.md) — the tx seam §7's workaround
rides,
[ADR-008](008-contract-v1-pinning.md) — the v1 skeleton this ADR
extends, §5's rule §9 applies, §8's dispositions §6 resolves.
- [ADR-009](009-scheduler-collapse.md) — the collapse machinery §6's
recipe composes with.
- OQ-06 — the SQLite-side fork candidates (§5, §3a realization note).
- [queues.md](../queues.md), [core-contract.md](../core-contract.md),
engine specs — the specs carrying this depth.
+4 -3
View File
@@ -57,9 +57,10 @@ Shared-server co-tenancy (the alkblobs ADR-008 precedent —
`/workspace/@alkdev/alkblobs/docs/architecture/decisions/` — consumer
tables co-tenant the pg instance) is supported and expected — the
[naming / reserved namespace contract](core-contract.md#naming--reserved-namespace)
protects reserved names; queue/stream/lock tables are schema-named to
avoid collisions (layout decision in [queues.md](queues.md), OQ-05's
namespace bullet).
protects reserved names; queue/stream/lock/schedule tables live in one
engine-owned PostgreSQL schema (default `alkstore`, per-engine option)
— layout per [ADR-010](decisions/010-queue-semantics-depth.md) §8
(resolved from [queues.md](queues.md)'s namespace bullet).
## Durability knobs
+14 -8
View File
@@ -26,7 +26,10 @@ obligations live in the core spec.
- Queue machinery re-derived on this driver with the pg-boss schema
family as design reference
([ADR-005](decisions/005-dependency-ownership.md)); semantics depth
is [queues.md](queues.md)'s work (OQ-05).
pinned by [ADR-010](decisions/010-queue-semantics-depth.md) — all
engine-owned tables (job, dead, stream, offsets, schedule) in one
PostgreSQL schema (default `alkstore`), queues as rows, no
per-queue tables.
## Connection architecture
@@ -62,9 +65,9 @@ obligations live in the core spec.
|---|---|
| notify / listen | `pg_notify(...)` inside the caller's tx (delivers at commit — native commit-atomicity, [ADR-007](decisions/007-transactional-seam.md)); `listen()` via LISTEN on the forwarder's connection, fanout to receivers |
| streams | durable event table + per-consumer offset cursors; `pg_notify` as the wake trigger ([ADR-006](decisions/006-wake-and-delivery-contract.md) mechanism split: durable row, LISTEN wake — the pg-boss-family shape) |
| queues | re-derived queue table + `FOR UPDATE SKIP LOCKED` claim + LISTEN-driven wake with re-poll safety net (default consumption posture, measured 5–16× vs 50 ms poll; poll-only fallback) |
| named locks | advisory-lock-semantics TTL locks (pg-boss-family design reference; depth in OQ-05's design work) |
| scheduler / outbox | pg-boss-family design reference, per [queues.md](queues.md) |
| queues | re-derived queue table + `FOR UPDATE SKIP LOCKED` claim + LISTEN-driven wake with re-poll safety net (default consumption posture, measured 5–16× vs 50 ms poll; poll-only fallback); states/dead-letter/backoff/visibility per [ADR-010](decisions/010-queue-semantics-depth.md), all in the engine-owned schema |
| named locks | advisory-lock-semantics TTL locks (pg-boss-family design reference; guarantee row pinned by [ADR-008](decisions/008-contract-v1-pinning.md) §7) |
| scheduler / outbox | collapse shape ([ADR-009](decisions/009-scheduler-collapse.md)): schedule rows in the engine-owned schema, tick re-derived (boundary advance + 64-boundary catch-up cap, honker parity), leadership via the engine's lock machinery on `__alkstore_scheduler`; outbox = helper over queues |
| begin_tx | pool checkout + `BEGIN`, returning the caller-held handle ([ADR-007](decisions/007-transactional-seam.md)) |
| handle ops | straight `.await`s through the held object; commit/rollback returns the object to the pool |
@@ -103,6 +106,8 @@ the record in [ADR-003](decisions/003-sqlite-driver.md).
| [006](decisions/006-wake-and-delivery-contract.md) | Wake contract | LISTEN push, no replay, synthetic reconnect-wake |
| [007](decisions/007-transactional-seam.md) | Tx seam | direct pooled-object handle, no bridging |
| [008](decisions/008-contract-v1-pinning.md) | Contract v1 | pinned surface; reserved reconnect-wake channel string; `PayloadTooLarge` taxonomy variant |
| [009](decisions/009-scheduler-collapse.md) | Scheduler collapse | schedule rows in the engine schema, re-derived tick, row-locked fire tx, `__alkstore_scheduler` leadership |
| [010](decisions/010-queue-semantics-depth.md) | Queue depth | job-stamped opts, equal-jitter backoff, dead-letter move, no-stranded-rows sweep, one engine-owned schema |
## Open Questions
@@ -110,9 +115,10 @@ Open questions are tracked in
[open-questions.md](open-questions.md). Key
questions affecting this document:
- **OQ-09**: scheduler collapse into queues (shared with
[queues.md](queues.md)) (open)
- **OQ-05**: queue semantics depth — the pg-boss-family design-input
work (open)
- **OQ-08**: capability surface (shared with
[deployment.md](deployment.md)) (open)
Resolved: **OQ-09** (scheduler collapse —
[ADR-009](decisions/009-scheduler-collapse.md)) and **OQ-05** (queue
semantics depth — [ADR-010](decisions/010-queue-semantics-depth.md)),
2026-10-05.
+13 -6
View File
@@ -52,9 +52,9 @@ spec and are not restated here.
|---|---|
| notify / listen | honker's notify functions inside the caller's tx; `listen()` bridges the watcher's fanout into a tokio receiver (one `spawn_blocking` thread per subscription doing `blocking_send`) |
| streams | honker's stream machinery; explicit offset saves through the tx seam |
| queues | honker's queue functions ([ADR-002](decisions/002-feature-scope.md)); semantics depth design in [queues.md](queues.md) |
| queues | honker's queue functions ([ADR-002](decisions/002-feature-scope.md)); semantics depth per [ADR-010](decisions/010-queue-semantics-depth.md) — ride + pin, except two contract properties honker's machinery doesn't provide (the no-stranded-rows `sweep_expired`, dead-row `get_job` visibility — OQ-06's concrete fork candidates) |
| named locks | honker's lock machinery |
| scheduler / outbox | honker's counterparts, per [queues.md](queues.md) |
| scheduler / outbox | collapse shape ([ADR-009](decisions/009-scheduler-collapse.md)): honker's scheduler machinery (`_honker_scheduler_tasks`) carries `@every` specs (its cron boundary machinery unused — cron strings are contract-rejected); honker's leader loop pattern (TTL lock, heartbeat, exit-before-tick-on-loss); the leadership lock is `__alkstore_scheduler`; outbox = helper over queues with the derived backing-queue name ([ADR-008](decisions/008-contract-v1-pinning.md) §4) |
| begin_tx | acquires the writer slot, opens `BEGIN IMMEDIATE`, returns the caller-held handle ([ADR-007](decisions/007-transactional-seam.md)) |
| handle ops | each `*_tx` op round-trips `spawn_blocking` to the same connection (thread-affinity note in [ADR-007](decisions/007-transactional-seam.md)) |
| commit/rollback | releases the writer slot |
@@ -102,6 +102,8 @@ consumption stands.
| [006](decisions/006-wake-and-delivery-contract.md) | Wake contract | data_version watcher, coalescing, death-closes-receivers |
| [007](decisions/007-transactional-seam.md) | Tx seam | writer-slot lease, `BEGIN IMMEDIATE`, `spawn_blocking` round-trips |
| [008](decisions/008-contract-v1-pinning.md) | Contract v1 | pinned surface; outbox backing queue derived under the reserved prefix; wake payload transport unused (non-contract) |
| [009](decisions/009-scheduler-collapse.md) | Scheduler collapse | honker scheduler machinery with `@every` specs; `__alkstore_scheduler` leadership; cron machinery unused |
| [010](decisions/010-queue-semantics-depth.md) | Queue depth | ride + pin over honker's functions; §5/§3a properties are OQ-06 fork candidates |
## Open Questions
@@ -109,7 +111,12 @@ Open questions are tracked in
[open-questions.md](open-questions.md). Key
questions affecting this document:
- **OQ-06**: honker-core quality read — fork-trigger assessment (open)
- **OQ-09**: scheduler collapse into queues (shared with
[queues.md](queues.md)) (open)
- **OQ-05**: queue semantics depth on honker's machinery (open)
- **OQ-06**: honker-core quality read — fork-trigger assessment
(open); now carrying two concrete candidates from the queue-depth
design (no-stranded-rows sweep, dead-row visibility —
complement vs fork is the read's call)
Resolved: **OQ-09** (scheduler collapse —
[ADR-009](decisions/009-scheduler-collapse.md)) and **OQ-05** (queue
semantics depth — [ADR-010](decisions/010-queue-semantics-depth.md)),
2026-10-05.
+74 -45
View File
@@ -6,17 +6,16 @@ last_updated: 2026-10-05
# alkstore — Open Questions
Centralized tracker. IDs `OQ-NN` are stable — never renumber; append.
The Phase 0 register's questions (`OQ-ST-01..08` in
`docs/research/phase-0.md`) are promoted here faithfully: **OQ-01..08
mirror OQ-ST-01..08 one-to-one**, keeping their Phase 0 statuses
(a resolved register question stays listed here as resolved, with the
ADR that carries its decision). New Phase 1 questions append from
OQ-09. Resolution order so far: **OQ-04 resolved** (2026-10-05,
Promoted, one-to-one, from Phase 0's OQ-ST register. Resolution order
so far: **OQ-04 resolved** (2026-10-05,
[ADR-008](decisions/008-contract-v1-pinning.md) — the contract surface
everything else hangs off). Next: OQ-05/OQ-09 (the queue semantics
track; OQ-09 also owns the scheduler guarantee row), OQ-06
(the SQLite dependency gate), OQ-08 (rides the now-pinned trait
shape), OQ-10 (versioning discipline for contract extensions).
everything else hangs off); **OQ-09 + OQ-05 resolved** (2026-10-05,
[ADR-009](decisions/009-scheduler-collapse.md) and
[ADR-010](decisions/010-queue-semantics-depth.md) — the queue semantics
track and the scheduler collapse/guarantee row). Next: OQ-06
(the SQLite dependency gate — now carrying two concrete fork
candidates from the queue-depth design), OQ-08 (rides the now-pinned
trait shape), OQ-10 (versioning discipline for contract extensions).
Resolved questions stay listed with their resolution; they are not
deleted.
@@ -117,43 +116,68 @@ narrowed to the pinning work its own record already scoped.)*
## Theme: Queues and scheduling
### OQ-05: Queue semantics depth — retry/backoff/dead-letter/sweep design *(== OQ-ST-05)*
### OQ-05: Queue semantics depth — retry/backoff/dead-letter/sweep design *(== OQ-ST-05)* — **RESOLVED**
- **Origin**: [queues.md](queues.md)
- **Status**: open (posture resolved: re-derived on both engines with
the pg-boss schema family as design reference and honker's queue
design as the SQLite-side one — ADR-004/ADR-005)
- **Status**: resolved (2026-10-05, Phase 1 — [ADR-010](decisions/010-queue-semantics-depth.md))
- **Priority**: high
- **Resolution**: open. The transactional and claim properties (the
hard driver-coupled part) are POC-pinned on both engines. Remaining:
the semantics-depth design — retry policy shape (attempts, backoff
curve), visibility-timeout/renewal mechanics, dead-letter
move-vs-flag and retention, sweep/maintenance design
(`sweep_expired` cadence and owner), and whether result-storage's
cut-flag gets reconsidered as part of this surface (it is
queue-adjacent tooling).
- **Resolution**: Pinned by ADR-010: (1) job state machine — three
states (`pending/processing/dead`), delete-on-ack, move-to-dead
(honker's model, pg-boss's seven-state enum not inherited),
`get_job` sees dead rows (`last_error`/`died_at` — post-mortem by
API; `cancel` unconditional, not an interrupt). (2)
Visibility/heartbeat — renewal semantics (absolute reset),
late-heartbeat refusal (dual-execution window is contract-honest),
**reclaim-consumes-an-attempt** stated as contract text. (3)
Retry/backoff — `retry(err, None)` computes the queue's
equal-jitter exponential curve (range definitionally pinned, not a
labeled distribution; base `backoff_base_s`, cap 1 h, attempt index
= row `attempts` at retry), `Some(d)` overrides; queue-level
`QueueOpts {
visibility_timeout_s, max_attempts, backoff_base_s,
dead_letter_retention_s }` **stamped onto each job row at enqueue**
(§3a — no per-queue registry; two handles with different opts for
one queue name don't fight). (4) Dead-letter — move-to-table,
retention forever by default, no redrive API (recipe: `get_job` +
fresh `enqueue`). (5) `sweep_expired` — the no-stranded-rows
property (every past-expiry row in any state moves to dead; honker's
expired-processing zombie hole fails this — fixed at contract
level, SQLite-side realization rides OQ-06). (6) Sweep cadence —
no ambient sweeper; the recipe is the collapse machinery
(schedule + runner). (7) Result-storage cut-flag reconsidered and
**stands** (delete-on-ack; pg-boss's completed-row model is result
storage under another name). (8) Table layout — pg engine-owned
schema (`alkstore` default), queues are rows not tables; SQLite
rides `_honker_*`. No new error variants (ADR-008 §5's rule).
- **Cross-references**: OQ-09, OQ-06.
### OQ-09: Is the scheduler a first-class mechanism, or queues + `schedule()`?
### OQ-09: Is the scheduler a first-class mechanism, or queues + `schedule()`? — **RESOLVED**
- **Origin**: [queues.md](queues.md) (spin-out of OQ-05's Phase 0
framing, where the scheduler's boundary question lived)
- **Status**: open
- **Status**: resolved (2026-10-05, Phase 1 — [ADR-009](decisions/009-scheduler-collapse.md))
- **Priority**: medium
- **Resolution**: open. The inventory found the scheduler need real but
thin (the family-wide "who sweeps/renews/reaps" problem), and
mechanically scheduler = cron/`@every` enqueueing into named queues
+ leader election + a tick. If design confirms the collapse, the
contract surface is queues + `store.schedule(cron, queue)` with the
scheduler mechanism absorbed; if consumers need
inspectable/pausable schedule objects (`add/pause/resume/update/
list/remove`), it stays a first-class mechanism (the full honker
shape). Decided against the inventory rows, not the whole honker
menu. **Also owns the delivery-guarantee table's scheduler row**
(transferred by [ADR-008](decisions/008-contract-v1-pinning.md) §7:
under collapse the scheduler inherits the queues row; otherwise a
dedicated row is pinned with the surface decision). The scheduler
surface is *not* part of contract v1 ([ADR-008](decisions/008-contract-v1-pinning.md) §1).
- **Resolution**: **Collapse confirmed**, decided against the
inventory rows (alkblobs sweep cadence, alkfs orphan reaping — both
intervals; no row names schedule objects or wall-clock cron): the
contract surface is `schedule(name, spec, queue, payload, opts)`
(upsert by name), `unschedule(name)`, and the opt-in
`run_schedules(stop)` runner — no `Scheduler` handle, no
schedule objects; update = re-register, pause = unregister +
re-register. **Spec grammar v1: `@every <n><unit>` only** — cron
strings rejected (`InvalidSpec`), dissolving the timezone question
honker's local-time cron brittleness poses (extension path requires
a consumer-inventory row). Boundary guarantee row pinned (the
transfer from [ADR-008](decisions/008-contract-v1-pinning.md) §7
resolved): at-least-once per elapsed boundary while a leader runs,
fire+advance commit atomically under a row lock on the schedule row
(engine-generic no-double-fire floor; the leadership lock is the
efficiency layer), bounded catch-up (**fixed contract-constant cap
64**) with skip-forward; leadership lock is the reserved
`__alkstore_scheduler` name; `run_schedules` returns
`Err(LeadershipLost)` on lock loss, `Ok(())` on clean stop. Schedule
storage engine-internal; schedules never fire without a runner (no
ambient timers).
- **Cross-references**: OQ-05, OQ-04.
## Theme: Engines and dependencies
@@ -172,8 +196,10 @@ narrowed to the pinning work its own record already scoped.)*
failure-handling, and schema-migration brittleness. Outcomes: posture
holds (no ADR change), or a fork/patch need is named (fork is normal
work per ADR-005).
- **Cross-references**: ADR-003, ADR-005, OQ-05 (if the read forces a
fork, queue-on-honker-machinery work changes shape).
- **Cross-references**: ADR-003, ADR-005, OQ-05 (the queue-depth read
handed it two concrete fork candidates: the expired-
processing-row zombie fix and dead-row `get_job` visibility —
complement-over-machinery vs fork is exactly its calculus).
## Theme: Deployment and capabilities
@@ -190,8 +216,9 @@ narrowed to the pinning work its own record already scoped.)*
(`Store::capabilities()`), a documented deployment matrix only
([deployment.md](deployment.md) carries the facts), or compile-time
knowledge only
(a consumer choosing the SQLite engine knows). Rides OQ-04: the
trait's shape constrains where capability differences can surface.
(a consumer choosing the SQLite engine knows). Rides the now-pinned
contract shape ([ADR-008](decisions/008-contract-v1-pinning.md)):
the trait constrains where capability differences can surface.
- **Cross-references**: OQ-04, [ADR-006](decisions/006-wake-and-delivery-contract.md).
### OQ-10: How do engine crates track core-contract version changes?
@@ -207,10 +234,12 @@ narrowed to the pinning work its own record already scoped.)*
test suite the engines run against the core's trait definitions?
What happens to a released engine crate when core makes a contract
breaking change?
- **Cross-references**: OQ-04 (the contract being versioned), OQ-02.
- **Cross-references**: OQ-04 (the contract being versioned — now
including its first post-v1 extensions, ADR-009/ADR-010), OQ-02.
## Deferred / Blocked
None currently. Every open OQ above is actionable Phase 1 architecture
work (contract pinning, design, quality read) with its evidence base
complete — no external arrivals are being waited on.
work (quality read, capability-surface shape, versioning discipline)
with its evidence base complete — no external arrivals are being
waited on.
+3 -1
View File
@@ -55,7 +55,7 @@ Per [ADR-002](decisions/002-feature-scope.md):
| [core-contract.md](core-contract.md) | The unified trait surface: mechanisms, delivery guarantees, tx seam, naming |
| [engine-sqlite.md](engine-sqlite.md) | SQLite engine: mapping the contract onto honker-core/rusqlite |
| [engine-postgres.md](engine-postgres.md) | Postgres engine: mapping the contract onto tokio-postgres/LISTEN |
| [queues.md](queues.md) | Queue/scheduler/outbox semantics depth (OQ-05/OQ-09) |
| [queues.md](queues.md) | Queue/scheduler/outbox semantics depth (ADR-009/ADR-010 resolved) |
| [deployment.md](deployment.md) | Host capabilities, connection budgets, deployment matrix (OQ-08) |
| [open-questions.md](open-questions.md) | OQ-01..NN tracker |
| [decisions/](decisions/) | ADRs |
@@ -72,6 +72,8 @@ Per [ADR-002](decisions/002-feature-scope.md):
| [006](decisions/006-wake-and-delivery-contract.md) | Opaque wake + re-read; notify-vs-streams guarantee split | Accepted |
| [007](decisions/007-transactional-seam.md) | Caller-held `TxHandle` with `*_tx` methods | Accepted |
| [008](decisions/008-contract-v1-pinning.md) | Contract v1 surface pinning (partition, TxHandle shape, wake type, reserved strings, error taxonomy) | Accepted |
| [009](decisions/009-scheduler-collapse.md) | Scheduler collapse (queues + `schedule()`, `@every`-only) | Accepted |
| [010](decisions/010-queue-semantics-depth.md) | Queue semantics depth (visibility, backoff, dead-letter, sweep, layout) | Accepted |
## Non-goals
+190 -79
View File
@@ -7,20 +7,21 @@ last_updated: 2026-10-05
The queue family's *posture* is decided (re-derived on both engines;
[ADR-004](decisions/004-postgres-driver.md) and
[ADR-005](decisions/005-dependency-ownership.md)) and the hard
[ADR-005](decisions/005-dependency-ownership.md)), the hard
driver-coupled properties are POC-pinned (transactional enqueue,
exactly-once claim under concurrency). What is *not* yet designed is
the semantics depth — this document's subject, and [OQ-05]'s home. It
states the design space, the decided constraints any answer must fit,
and the reference material; actual semantics decisions become ADRs
when made.
exactly-once claim under concurrency), and the semantics depth is
designed — OQ-05 and OQ-09 resolved by
[ADR-010](decisions/010-queue-semantics-depth.md) and
[ADR-009](decisions/009-scheduler-collapse.md) (2026-10-05). This
document states the resulting semantics and the constraints they were
made under; the ADRs carry the WHY.
## Decided constraints (inherited, not re-opened)
- **Mechanisms**: durable at-least-once queues; the outbox is a helper
over queues; scheduler is cron/`@every` enqueueing into named
queues ([ADR-002](decisions/002-feature-scope.md)) — pending
OQ-09's collapse decision.
over queues; the scheduler collapses into queues + `schedule()` +
`run_schedules` ([ADR-002](decisions/002-feature-scope.md),
[ADR-009](decisions/009-scheduler-collapse.md)).
- **Transactional enqueue**: commit-atomic with the caller's business
write on both engines ([ADR-007](decisions/007-transactional-seam.md));
rollback drops the job row with no ghosts.
@@ -32,104 +33,211 @@ when made.
per-subscription fanout ([ADR-006](decisions/006-wake-and-delivery-contract.md)).
- **Job options** (the contract v1 skeleton,
[ADR-008](decisions/008-contract-v1-pinning.md) §1, from the
honker-rs surface): `delay, priority, max_attempts, expires` at
enqueue; `ack / retry / fail / heartbeat` on the job handle. Deeper
`QueueOpts`/retry surface is this document's design space.
honker-rs surface): `delay, priority, max_attempts, expires, run_at`
at enqueue; `ack / retry / fail / heartbeat` on the job handle.
- **Cut-flag context**: result storage is
[ADR-002](decisions/002-feature-scope.md)'s cut-flag row; if
OQ-05's design shows queue consumers provably need
result-query-by-id, that is the channel to revisit it — with a
consumer-inventory row, not silent inclusion.
[ADR-002](decisions/002-feature-scope.md)'s cut-flag row;
[ADR-010](decisions/010-queue-semantics-depth.md) §7 reconsidered it
as part of this design and left it cut (delete-on-ack keeps
completed jobs out of storage; the pg-boss completed-row/output
model is that feature under another name). Re-entry stays gated on
a consumer-inventory row.
## Design space (what OQ-05 must pin)
## Job lifecycle (ADR-010 §1–§3a)
### Retry / backoff
```text
enqueue
│
▼
┌─────────┐ claim (attempts += 1) ┌───────────────────┐
│ pending │ ───────────────────────────► │ processing │
└─────────┘ ◄─────────────────────────── └───────────────────┘
│ retry(err, d) │ deadline = job-stamped
│ │ visibility timeout;
cancel deletes the row in EITHER │ heartbeat(extend) resets it
state (pending or processing) — ├─────────────────────────
not an interrupt; the holder's │ visibility lapse:
next ack/heartbeat returns false │ no transition fires — the
│ │ row is just reclaimable
▼ │ (a reclaim eats an attempt)
┌────────┐ fail(err) / budget │
│ dead │ ◄────── exhausted ───────────┘
└────────┘
│ ack deletes the row
│ (from processing, claim still valid)
│
▼
(gone)
- Trigger model: attempts counted from claim-completion; a job
re-appears after visibility timeout OR explicit `retry`. Does
`max_attempts` exhaustion move the job to dead-letter, flag it, or
drop it?
- Backoff curve between attempts: honker and the pg-boss family have
their own shapes (fixed, exponential); which does the contract
pin, what is configurable, what's the default?
- `heartbeat` semantics: renewal vs progress signal (or both) — and
what a missed heartbeat means (visibility re-expiry? nothing?).
expires lapse (any state) — sweep_expired(q) moves the row to
dead with last_error='expired'; also enforces dead-letter
retention. ack/cancel leave no row; get_job sees dead rows.
```
### Visibility timeouts
- **States**: `pending` → `processing` → `dead` (+ absence:
`ack`/`cancel` delete the row). `max_attempts` frozen at enqueue;
`attempts` counts every claim — fresh or reclaim.
- **Visibility**: each claim sets the deadline from the job's stamped
`visibility_timeout_s` (default 300 s, stamped per
[ADR-010](decisions/010-queue-semantics-depth.md) §3a).
`heartbeat(extend)` is renewal — an absolute reset from now (the new
full deadline, not additive); a **late heartbeat is refused** (never
steals the job back from a reclaimer). Between deadline lapse and
another worker's reclaim, the original worker may still complete —
the dual-execution window is the contract's honest at-least-once
posture; downstream idempotence is the consumer's job.
- **Reclaim consumes an attempt**: a handler that forgets to
heartbeat looks like a repeatedly-failing job and dead-letters on
budget exhaustion. Stated in the contract text
([ADR-010](decisions/010-queue-semantics-depth.md) §2), not
discovered in production.
- **Retry**: `retry(err, None)` — the engine computes the queue's
curve; `retry(err, Some(d))` overrides. `fail(err)` = immediate
dead-letter. Ordering within a queue: `priority DESC`, ready-time
ASC, enqueue order (FIFO under equal priority).
- **get_job sees dead jobs** — state, payload, attempts, timestamps,
and (dead only) `last_error`, `died_at`; post-mortem diagnosis by
API, not by SQL.
- Where the renewal lives: explicit (`heartbeat`) only, or implicit
renewal per op on the job handle? Honker's design is the SQLite-side
answer; pg-boss's is the reference. One answer, contract-pinned.
- Interaction with long business transactions (claims inside
caller txs, [ADR-007](decisions/007-transactional-seam.md)).
## Retry / backoff (ADR-010 §3, §3a)
### Dead-letter
- **The curve: equal-jitter exponential** from the job's stamped
`backoff_base_s` (default 5 s): `delay ∈ [base·2^(a−1)/2,
base·2^(a−1)]` for attempt `a` — the range is the definition (not
AWS-canonical full jitter); capped at **1 hour** (pinned contract
constant, no knob). The formula is what both engines compute
identically; the range's herding rationale lives in the ADR.
- **Where configured**: queue-level `QueueOpts
{ visibility_timeout_s, max_attempts, backoff_base_s,
dead_letter_retention_s }` (honker-parity defaults: 300 s / 3 / 5 s
/ none) — stamped onto each job row at enqueue (§3a), never read
from a live registry.
- Move-to-table (honker's `_honker_dead`, pg-boss's dead-letter
queue) vs flag-in-place. Inspection/requeue surface shape.
- Retention: does a dead-lettered job expire (an `expires`-driven
sweep — an OQ-05 sub-question) or live until explicitly cleared?
## Dead-letter (ADR-010 §4)
### Sweep / maintenance
- **Move, not flag**: dead rows physically move to engine-owned dead
storage (`last_error`, `died_at`) — never scanned by the claim path.
Triggers: explicit `fail`, `retry` at budget, exhaustion-by-reclaim,
`expires` lapses (swept).
- **Retention**: forever by default; `dead_letter_retention_s` (queue
opt) makes `sweep_expired` enforce a per-queue dead TTL. Honker's
no-retention posture kept as default; the mechanism is ours.
- **No redrive API in v1**: the recipe is `get_job` + fresh `enqueue`.
No consumer names an API shape; re-entry needs an OQ.
- `sweep_expired` is on the queue's surface; cadence ownership is
the open question. Honker's scheduler machinery (leader-elected,
missed-boundary catch-up) can run it; a consumer's own scheduler row
can too; the engine could ship a default. Decision shaped by OQ-09's
scheduler collapse — same machinery either way.
- Multi-process sweep safety on SQLite vs Postgres: named-lock
coordination (the alkblobs fleet-sweeper pattern,
[ADR-002](decisions/002-feature-scope.md)) vs native advisory locks.
## Sweep / maintenance (ADR-010 §5–§6)
### Scheduler surface
- `sweep_expired(queue)` moves *every* past-`expires_at` row (pending
**and** processing) to dead, and enforces dead-letter retention —
the **no-stranded-rows property**
([ADR-010](decisions/010-queue-semantics-depth.md) §5): a job with
an `expires` deadline is eventually in exactly one of
pending/processing/dead, never stuck unreachable. (SQLite-side
realization over honker's machinery rides OQ-06 — the zombie hole
is that read's concrete fork candidate.)
- **No ambient sweeper**: the engine ships no background maintenance,
no default cadence (the no-ambient-timers posture). Correctness of
transitions is the engine's; **timeliness is the consumer's**.
- **The "who sweeps" recipe** is the scheduler collapse's payoff
([ADR-009](decisions/009-scheduler-collapse.md)):
`store.schedule("maintenance", "@every 300s", mq, …)` +
`run_schedules(stop)` (leader-elected) + a worker handler calling
`sweep_expired` per queue. One schedule row, one worker — the
pattern every consumer would have hand-rolled.
- Multi-process sweeps need no leader lock on a bare `sweep_expired`
call (single-statement atomic per queue); the cadence coordination
comes from the scheduler machinery above. The SQLite notifications
table's hygiene is engine-internal (not a consumer chore;
`prune_notifications*` stay out of the contract, resolving ADR-008
§8's disposition).
- OQ-09's question, stated crisply: is scheduler a first-class
mechanism (own handle, methods, inspection) or queues + a
`schedule(cron_expression, queue_name, payload)` call? The
inventory's evidence is thin (documented-thin); a v1 that is queues
+ a `schedule` call with leader-election behavior documented is the
smaller contract; if a consumer needs inspectable/pausable schedule
*objects* (`add/pause/resume/update/list/remove`), that's the
full honker shape. Decide against the inventory rows, not against
the whole honker menu.
## Scheduler (ADR-009)
### Namespaces / schema layout (pg side)
Collapsed into queues: no `Scheduler` handle, no schedule objects.
- Schema-scoped DDL (the pg-boss design) vs shared-schema table
naming. Rides the reserved-namespace contract
([ADR-006](decisions/006-wake-and-delivery-contract.md),
[ADR-008](decisions/008-contract-v1-pinning.md) §4):
the reserved *channel* namespace is pinned (`__alkstore_` prefix);
queue/stream/lock **table** layout remains this document's
OQ-05 work and must be collision-proof against consumer tables in
the same database (alkblobs' ADR-008 — co-tenant-tables precedent —
at `/workspace/@alkdev/alkblobs/docs/architecture/decisions/`).
- Surface: `schedule(name, spec, queue, payload, opts)` (upsert by
name), `unschedule(name) -> bool`, `run_schedules(stop) -> Result<()>`.
Update = re-register; pause = unregister + re-register.
- **Spec grammar v1: `@every <n><unit>` only** (`s|m|h|d`). Cron
strings rejected with `InvalidSpec { spec }` (grammar rationale in
the ADR; extension path: a consumer-inventory row first).
- `run_schedules` is opt-in: schedules never fire without a runner
(the no-ambient-timers posture). The runner: leadership lock →
{renew (on loss, return before ticking — `Err(LeadershipLost)`,
distinct from clean `Ok(())` stop) → fire due boundaries → sleep}.
Leadership lock: reserved name `__alkstore_scheduler`
(engine-derived name under the reserved prefix,
[ADR-008](decisions/008-contract-v1-pinning.md) §4).
- **Boundary guarantee** ([ADR-009](decisions/009-scheduler-collapse.md)
§4): at-least-once per elapsed boundary while a leader runs; fire +
boundary-advance commit atomically under a row lock on the schedule
row (crash mid-tick refires the boundary; a rogue ticker is defeated
by the row-locked advance, not just by the leadership lock); bounded
catch-up — missed boundaries replay one-by-one up to a **fixed
contract constant cap of 64** per task per tick, beyond which
remaining boundaries skip forward. Honest loss mode, documented.
- Each boundary fire enqueues ordinary work: `payload` into the named
queue with `ScheduleOpts { priority, max_attempts, expires }` —
stamped onto the enqueued job per
[ADR-010](decisions/010-queue-semantics-depth.md) §3a; the queue
mechanism's guarantees apply, no separate delivery machinery.
- Schedule storage is engine-internal (SQLite: honker's
`_honker_scheduler_tasks`; Postgres: the engine-owned schema) — not
a consumer namespace; queues are names, not registered objects, so a
schedule may target a queue nothing has claimed yet. Schedule names
validate like queue names (non-empty, reserved-prefix rejected).
## Namespaces / schema layout (ADR-010 §8)
- **Postgres**: all engine-owned tables (job, dead, stream, offsets,
schedules, internal) in one **engine-owned schema** (default
`alkstore`, per-engine option) — co-tenancy-safe with consumer
tables (alkblobs ADR-008 precedent; the reserved *name* namespace
governs name-column values inside). **Queues are rows in one job
table** — no per-queue tables/schemas; pg-boss's partition-per-queue
opt-in is not inherited (one table + partial indexes matches honker's
proven shape and keeps `queue(name)` a name, not a DDL operation).
- **SQLite**: honker's `_honker_*` family — storage-internal
([ADR-008](decisions/008-contract-v1-pinning.md) §4); no new tables
minted by this design; the family's fate rides OQ-06 (a fork
re-owns the names; nothing consumer-visible changes).
- Partial indexes pinned identically on both engines: partial indexes
matching the claim hot path; dead rows outside it; single clock
source (second-precision timestamps); savepoint-guarded
multi-statement moves.
## Reference material
- **honker's queue design** — the SQLite-side incumbent
(`/workspace/honker`, its honker-core machinery; the engine rides it
directly, so the SQLite side's depth is largely "inherit + pin").
- **pg-boss family** — `/workspace/pgboss-rs` @ 98f7d9e (queue states,
maintenance/dead-letter behavior) and the node original (pg-boss, the
upstream of record — compare semantics the port may have dropped).
Full read: `docs/research/reference-honker-machinery.md` (schema,
claim/visibility/retry/dead-letter mechanics with file/line cites;
the defect list §8 feeds OQ-06).
- **pg-boss family** — `/workspace/pgboss-rs` @ 98f7d9e (design
reference only, [ADR-005](decisions/005-dependency-ownership.md)).
Full read: `docs/research/reference-pgboss-rs-semantics.md` (states,
backoff-with-jitter formula, dead-letter-as-queue, maintenance
postures; the node-upstream deltas §8).
- **POC ground** — `poc-pg-posture-findings.md` (the minimal queue
table + SKIP LOCKED claim measured; the property suite green);
`poc-sqlite-posture-findings.md` (the SQLite twin).
- **Consumers** — alkfs OQ-FS-14 (sync outbox), alkblobs
ops-surface/gc docs (maintenance cadences), the family-wide
"who sweeps" punt (inventory §scheduler).
- **Consumers** — alkfs OQ-FS-14 (sync outbox), alkblobs ops-surface/gc
docs (maintenance cadences), the family-wide "who sweeps" punt
(inventory §scheduler — answered by the collapse recipe above).
## Design Decisions
| ADR | Decision | Summary |
|---|---|---|
| [002](decisions/002-feature-scope.md) | Feature scope | queues + outbox in; result storage cut-flag |
| [002](decisions/002-feature-scope.md) | Feature scope | queues + outbox in; result storage cut-flag (stands, re-checked) |
| [004](decisions/004-postgres-driver.md) | pg driver | queue machinery re-derived, pg-boss as reference |
| [005](decisions/005-dependency-ownership.md) | Ownership | design-reference postures for the queue family |
| [006](decisions/006-wake-and-delivery-contract.md) | Wake contract | wake-driven consumption, guarantee table |
| [007](decisions/007-transactional-seam.md) | Tx seam | enqueue_tx commit-atomicity |
| [008](decisions/008-contract-v1-pinning.md) | Contract v1 | queue skeleton + EnqueueOpts pinned; depth is OQ-05's extension path |
| [008](decisions/008-contract-v1-pinning.md) | Contract v1 | queue skeleton + EnqueueOpts pinned; depth was this doc's design space |
| [009](decisions/009-scheduler-collapse.md) | Scheduler collapse | queues + `schedule()` + runner; `@every`-only; boundary guarantee row |
| [010](decisions/010-queue-semantics-depth.md) | Queue depth | visibility/heartbeat rules, backoff curve, dead-letter, no-stranded-rows sweep, schema layout |
## Open Questions
@@ -137,7 +245,10 @@ Open questions are tracked in
[open-questions.md](open-questions.md). Key
questions affecting this document:
- **OQ-05**: the semantics-depth pinning this document frames (open,
high priority)
- **OQ-09**: scheduler as first-class mechanism vs queues + schedule
(open)
- **OQ-06**: honker-core quality read — the SQLite-side zombie fix and
dead-row visibility are concrete fork candidates for it to weigh
(open).
Resolved on this document's surface: **OQ-05** and **OQ-09**
(2026-10-05, [ADR-010](decisions/010-queue-semantics-depth.md) /
[ADR-009](decisions/009-scheduler-collapse.md)).
+487
View File
@@ -0,0 +1,487 @@
---
status: draft
last_updated: 2026-10-05
---
# Reference read: honker queue / scheduler / outbox / lock maintenance machinery
Per-agents conventions: reference checkouts are read freely, never
wired in; checkout state noted when findings depend on code specifics.
**Checkout:** `/workspace/honker` @ `f4e53c6`
(`node-v0.5.1-10-gf4e53c6`), matching the AGENTS.md reference revision.
**Versions:** `honker-core` **0.5.0** (`honker-core/Cargo.toml`), `honker`
Rust wrapper **0.5.0** (`packages/honker-rs/Cargo.toml`). The Rust
binding is *excluded* from the root workspace; all real machinery is in
`honker-core/src/{lib.rs, honker_ops.rs, cron.rs}`.
**Architecture in one sentence:** SQLite tables (`_honker_live`,
`_honker_dead`, `_honker_locks`, `_honker_scheduler_tasks`,
`_honker_results`, `_honker_notifications`, rate-limit/stream tables) +
SQL scalar functions (`honker_*`) implemented once in Rust
(`honker_ops::attach_honker_functions`), consumed identically by
Python/Node/Rust bindings; cross-process wake is `PRAGMA data_version`
polling (1 ms default) driving an in-process channel — not
LISTEN/NOTIFY.
## 1. Queues
### 1.1 Job table schema
`BOOTSTRAP_HONKER_SQL` — `honker-core/src/lib.rs:352-432`:
**`_honker_live`** (lib.rs:353-367) — pending and processing both live
here:
| column | type | notes |
|---|---|---|
| `id` | INTEGER PK AUTOINCREMENT | |
| `queue` | TEXT NOT NULL | namespace within one table |
| `payload` | TEXT NOT NULL | stored as JSON string |
| `state` | TEXT DEFAULT `'pending'` | two states: `'pending'`, `'processing'` |
| `priority` | INTEGER DEFAULT 0 | higher = earlier (`ORDER BY priority DESC`) |
| `run_at` | INTEGER DEFAULT `(unixepoch())` | readiness deadline |
| `worker_id` | TEXT | current claim holder |
| `claim_expires_at` | INTEGER | visibility deadline; moves on heartbeat |
| `attempts` | INTEGER DEFAULT 0 | bumped on **every** claim, including reclaims |
| `max_attempts` | INTEGER DEFAULT 3 | per-row, frozen at enqueue |
| `created_at` | INTEGER DEFAULT `(unixepoch())` | |
| `expires_at` | INTEGER | job-level expiry; NULL = never |
| `claimed_at` | INTEGER, nullable | start of *current* attempt; migration added (lib.rs:486-508) |
Partial indexes tuned to the hot paths (lib.rs:368-376): claim index
`(queue, priority DESC, run_at, id) WHERE state IN
('pending','processing')`; pending index `(queue, run_at)`; processing
index `(queue, claim_expires_at)`. Dead rows never touch the claim
index.
**`_honker_dead`** (lib.rs:377-388): `id`, `queue`, `payload`,
`priority`, `run_at`, `attempts`, `max_attempts`, `last_error`,
`created_at`, `died_at`. **Move, not flag**: dead-lettering physically
DELETEs from `_honker_live` and INSERTs into `_honker_dead`. No
retention/expiry on dead rows — "Never scanned by the claim path;
retention policy is the user's problem" (lib.rs:345-347). No
`DELETE FROM _honker_dead` anywhere in the repo.
### 1.2 Enqueue and its options
Core signature (`honker_ops.rs:984-993`):
```rust
pub fn enqueue(
conn: &Connection, queue: &str, payload: &str,
run_at: Option<i64>, delay: Option<i64>,
priority: i64, max_attempts: i64, expires: Option<i64>,
) -> rusqlite::Result<i64>
```
SQL function `honker_enqueue(...)` (honker_ops.rs:593-619). Rust-level
`EnqueueOpts` (`packages/honker-rs/src/lib.rs:438-444`) — note
`max_attempts` is *not* here; it comes from queue-level `QueueOpts`
(`visibility_timeout_s: 300, max_attempts: 3`, lib.rs:421-434):
```rust
pub struct EnqueueOpts {
pub delay: Option<i64>,
pub run_at: Option<i64>,
pub priority: i64,
pub expires: Option<i64>,
}
```
Semantics (honker_ops.rs:976-1000):
- **run_at vs delay precedence**: `delay` wins — `run_at = unixepoch()
+ delay` if delay set; else literal `run_at`; else `now`.
- **expires**: relative seconds → `expires_at = unixepoch() + expires`;
NULL = never expires.
- **No synthetic wake row on enqueue** — the INSERT advancing
`data_version` *is* the wake; an earlier design writing a
`_honker_notifications` row per enqueue grew unboundedly
(honker_ops.rs:1002-1006).
`arg_i64`/`arg_opt_i64` (honker_ops.rs:40-54) accept REAL-typed whole
numbers because better-sqlite3 binds every JS number as REAL — interop
lesson (fractional values error with a diagnostic).
### 1.3 Claim mechanics
```rust
pub fn claim_batch(conn: &Connection, queue: &str, worker_id: &str,
n: i64, timeout_s: i64) -> rusqlite::Result<String>
```
(honker_ops.rs:848-914). `timeout_s` **is** the visibility timeout,
set per-claim (the Rust binding passes static `QueueOpts
.visibility_timeout_s`, honker-rs lib.rs:584-591).
Single atomic `UPDATE ... WHERE id IN (SELECT ...) RETURNING`
(honker_ops.rs:866-886):
- Claim predicate: `state IN ('pending','processing') AND attempts <
max_attempts AND (expires_at IS NULL OR expires_at > unixepoch()) AND
((state='pending' AND run_at <= now) OR (state='processing' AND
claim_expires_at < now))`.
- Claim action: `state='processing'`, `worker_id=worker`,
`claim_expires_at = now + timeout_s`, `claimed_at = now`,
`attempts += 1`.
- Ordering: `priority DESC, run_at ASC, id ASC`, `LIMIT n`.
**Reclaim = another attempt.** Claiming bumps `attempts` on reclaims
too — visibility timeouts consume the retry budget. `claimed_at` is
refreshed on reclaim but **not** on heartbeat (doc honker_ops.rs:831-847;
test honker_ops.rs:3103-3149).
`claim_one` is a thin `claim_batch(worker_id, 1)` (honker-rs
lib.rs:606-609).
**Pre-claim dead-letter sweep**: every `claim_batch` first runs
`dead_letter_exhausted_claimable` (honker_ops.rs:781-829) — rows at
`attempts >= max_attempts` that are claimable-reclaimable move to
`_honker_dead` with `last_error='max attempts exceeded'` before the
claim UPDATE (covers "worker died right after the last allowed claim").
**Idle-wake helper**: `honker_queue_next_claim_at(queue) -> unix_ts`
(honker_ops.rs:936-970) — earliest future `run_at` (pending) or
`claim_expires_at + 1` (processing with budget); 0 if nothing. Used by
the Rust `ClaimWaker::next` loop (honker-rs lib.rs:805-845): wake on
data_version or sleep until `next_at`, no busy-poll.
### 1.4 Visibility timeout / heartbeat / renewal
```rust
pub fn heartbeat(conn, job_id, worker_id, extend_s) -> Result<i64>
```
(honker_ops.rs:1306-1323); 1 = extended, 0 = refused.
- Heartbeat is a **renewal** (visibility extension), *not* a progress
signal — no progress field exists. Sets `claim_expires_at =
unixepoch() + extend_s` (absolute reset from now, not additive).
- Guard: `WHERE id = ? AND worker_id = ? AND state = 'processing' AND
claim_expires_at >= unixepoch()` — a **late heartbeat after the
visibility timeout is refused** so it cannot steal the job back from
a reclaimer (dual-execution guard, honker_ops.rs:1312-1314).
- **Missed heartbeat** ⇒ `claim_expires_at` passes ⇒ row claimable by
the normal predicate; a reclaiming claimer consumes an attempt.
Nothing transitions the row on lapse; stale `worker_id` etc. remain
until a next claim overwrites them (honker_ops.rs:843-847).
- No per-binding heartbeat thread in core; cadence is caller-
implemented.
### 1.5 Retry behavior — attempts and backoff
```rust
pub fn retry(conn, job_id, worker_id, delay_s, error) -> Result<i64>
```
(honker_ops.rs:1045-1123). Requires a still-valid processing claim;
`attempts >= max_attempts` ⇒ dead-letter (savepoint-wrapped
DELETE→INSERT, honker_ops.rs:1080-1107); else `pending` with
`run_at = unixepoch() + delay_s`, claim fields cleared.
**Backoff is NOT in the queue core.** `retry` takes a caller-chosen
`delay_s`. Curves live in wrappers:
- Rust `Outbox::retry_delay`: `base_backoff_s * 2^(attempts-1)`,
saturating (honker-rs lib.rs:543-549).
- Python `_compute_delay`: `retry_delay * backoff**(attempts-1)`
(`packages/honker/python/honker/_worker.py:108-114`); `@task` defaults
`retry_delay_s=60`, `backoff=1.0` (constant).
- **No jitter and no cap in either** (Rust caps only via `saturating_mul`).
### 1.6 Dead-letter handling
- Dedicated `_honker_dead`; **move semantics** (§1.1). Triggers:
(a) `retry()` at budget (honker_ops.rs:1080-1107), (b) explicit
`fail()` (honker_ops.rs:1135-1192), (c)
`dead_letter_exhausted_claimable` pre-claim (honker_ops.rs:781-829),
(d) `sweep_expired` expiration (honker_ops.rs:1336-1381).
- `last_error` strings: caller error, `'max attempts exceeded'`,
`'expired'`.
- **Retention: none.** No sweep/TTL/pruning of dead rows exists.
- All four move-sites are wrapped in a shared savepoint helper
`in_savepoint` (honker_ops.rs:118-254) — `DELETE ... RETURNING` +
decode + INSERT had a half-failing window; a failing decode/INSERT
previously produced `live=0, dead=0` **silent job loss** (CHANGELOG
"Core SQLite error propagation"; `UnwindUndo` undoes the frame if the
body panics). The savepoint hardening is very recent, post-0.5.1 work
in the CHANGELOG "Unreleased" section. Most instructive correctness
pattern in the codebase.
### 1.7 `sweep_expired`
```rust
pub fn sweep_expired(conn, queue) -> Result<i64>
```
(honker_ops.rs:1336-1381). **Only sweeps `state='pending'` rows with
`expires_at IS NOT NULL AND expires_at <= now`** → `_honker_dead` with
`last_error='expired'`. Does **not** touch processing rows — expired
claims are handled lazily by the claim predicate. Python doc: "The
claim path already ignores expired rows, so sweep is cleanup-only — not
correctness-critical" (`_honker.py:387-394`).
### 1.8 cancel / get_job / ack_batch
```rust
pub fn cancel(conn, job_id) -> Result<i64> // honker_ops.rs:1203-1209
pub fn get_job(conn, job_id) -> Result<String> // :1220-1295
pub fn ack(conn, job_id, worker_id) -> Result<i64> // :1028-1035
pub fn ack_batch(conn, ids_json, worker_id) -> Result<i64> // :920-934
```
- **cancel**: unconditional `DELETE ... WHERE id=? AND state IN
('pending','processing')` regardless of which worker holds it. Not
an interrupt — the holder's next `ack`/`heartbeat` returns 0
(honker-rs docs lib.rs:625-631).
- **get_job**: reads `_honker_live` only → **dead jobs are invisible**;
JSON on hit, **empty string on miss** (honker-rs maps to `None`,
lib.rs:642-650). `claimed_at: null` for never-claimed jobs.
- **ack/ack_batch**: `DELETE ... WHERE id IN (...) AND worker_id=? AND
claim_expires_at >= now RETURNING id`, returning count. Deliberately
savepoint-free (honker_ops.rs:916-919). Ack fails (0) if the claim
expired. Note: ack doesn't check `state='processing'` — keys off
`worker_id` + `claim_expires_at` only (harmless today, implicit).
## 2. Scheduler
### 2.1 Storage schema
`_honker_scheduler_tasks` (lib.rs:400-410): `name` PK, `queue`,
`cron_expr`, `payload`, `priority` (default 0), `expires_s` (nullable),
`next_fire_at` NOT NULL, `enabled` INTEGER DEFAULT 1, `max_attempts`
INTEGER DEFAULT 3. Runtime ALTER migrations tolerate the
concurrent-bootstrap "duplicate column" race (lib.rs:443-508). No
`last_fired_at`, no `tz` column — cron arithmetic is **system local
time**.
### 2.2 Surface
Core (honker_ops.rs): `scheduler_register` (:1490-1536),
`scheduler_unregister` (:1538-1551), `scheduler_tick` (:1583-1654),
`scheduler_soonest` (:1656-1664), `scheduler_pause`/`resume`
(:1669-1689, toggling `enabled`), `scheduler_list` (:1694-1752),
`scheduler_update` (:1759-1844). Plus `honker_cron_next_after(expr,
from_unix)` DETERMINISTIC SQL function (:683-692).
Rust binding (`packages/honker-rs/src/lib.rs:1365-1584`): `add`,
`add_with_max_attempts`, `remove`, `tick`, `soonest`, `pause`,
`resume`, `list`, `update` (with `ScheduleUpdate { cron_expr, payload,
priority, expires_s: Option<Option<i64>> }` — inner Option
distinguishes "clear" from "leave alone"), `update_max_attempts`, and
the blocking `run(stop, owner)` leader loop. Python mirrors
(`_scheduler.py:113-437`).
Register semantics: upsert by name, replaces entirely; `next_fire_at`
= next cron boundary **strictly after now**; clamps `max_attempts < 1`
to 1.
### 2.3 Tick and catch-up
`schema_tick(conn, now_unix) -> JSON [{name, queue, fire_at, job_id}]`
(honker_ops.rs:1583-1654):
- Selects `WHERE next_fire_at <= now AND enabled = 1`.
- Per task, loops while `next_fire_at <= now`: enqueue (run_at NULL =
claim immediately, task-level `expires_s`/`max_attempts` applied),
then `next_fire_at = cron_next_after(expr, next_fire_at)` — one fire
per missed boundary, minute-by-minute catch-up.
- `SCHEDULER_MAX_CATCHUP_FIRES = 64` cap per task per tick
(honker_ops.rs:1564-1575); beyond the cap remaining boundaries
**skip** — `next_fire_at` jumps to `next_after(now)`. Documented as
semantic: "Run the scheduler continuously, use coarser schedules, or
raise this constant."
- Final `next_fire_at` persisted in a separate UPDATE (honker_ops.rs
:1647-1651) within the caller's transaction — enqueue+advance commit
together.
- Dedup: no DB-level claim token; the advance-then-return contract
under `BEGIN IMMEDIATE` means at most one ticker observes an unfired
boundary — pinned by a 10-thread race test
(`tests/test_scheduler.py:397-470`). The leader lock is the
production gate; the SQL is "still safe if someone forgets it".
**Guarantee: at-least-once per boundary, up to 64-boundary backlog;
beyond that, gap-skip by design.** Crash mid-tick rolls back both
enqueues and the advance (boundary refires).
### 2.4 Leader loop / election
Election is a **named TTL lock in `_honker_locks`**, not SQLite
advisory APIs:
- Rust `Scheduler::run` (lib.rs:1509-1584): `LOCK_NAME="honker-scheduler"`,
TTL 60s, heartbeat every 20s (Python `_scheduler.py:135-137`, TTL 60,
heartbeat 30). Non-leaders poll every 5s.
- Leader loop (lib.rs:1539-1584): renew lock → on failure **exit
before ticking** (no dual-fire alongside the thief) → `tick()` →
`soonest` → sleep `min(heartbeat_interval, until_soonest)`.
- Python `_main_loop` additionally re-proves ownership *before every
fire* (`_scheduler.py:382-437`), runs tick+soonest in one writer tx.
Losing the lock raises `LeadershipLost`.
- Wake-on-register: register/unregister/pause/resume/update rely on
`data_version` advancing (a former synthetic notification row was
removed as unbounded-growth — `scheduler_wake` is now a no-op,
honker_ops.rs:1553-1562).
### 2.5 Cron parsing (`cron.rs`, 473 lines)
- 5-field cron (`min hour dom month dow`), 6-field with seconds, and
`@every <n><unit>` (s/m/h/d) (cron.rs:1-14, 97-129). Full syntax:
`*`, ranges, lists, steps; dow 0=Sunday; field validation with good
errors (cron.rs:131-178).
- Next-boundary search: iterative scan, bounded to 100 years
(cron.rs:207-315). DST handled: spring-forward gap skips to real
time, ambiguous fall-back picks the first instantiation
(cron.rs:330-358; tests cron.rs:456-472).
- **Brittleness — system local timezone dependence**: no TZ parameter
in the public surface (`next_after_unix` uses `chrono::Local`,
cron.rs:319-321; tz-parameterized variant is `pub(crate)` test-only).
A server TZ change silently re-times every cron task; stored
schedules are not TZ-qualified.
## 3. Outbox
Pattern: **an outbox is a regular queue named `_outbox:{name}`** — no
separate table, no separate delivery semantics (honker-rs
lib.rs:478-492; `_honker.py:865-934`).
Rust shape (`OutboxOpts`: `visibility_timeout_s: 60, max_attempts: 5,
base_backoff_s: 5`, lib.rs:454-469):
```rust
pub fn run_once<F, E>(&self, worker_id: &str, delivery: F) -> Result<bool>
where F: FnMut(serde_json::Value) -> Result<(), E>
```
(lib.rs:516-541): `claim_one` → parse payload → `delivery(payload)` →
Ok ⇒ `job.ack()` (ack failure surfaced as error, :527-529) ⇒ Err ⇒
`job.retry(base_backoff_s * 2^(attempts-1), err)` (:531-538). A **pull**
worker — app calls `run_once` in its own loop.
Transactional enqueue: `outbox.enqueue_tx(&tx, ...)` routes through
the *caller's* `Transaction` (lib.rs:506-513); single
mutex-guarded-connection means same-thread in-tx ops must use `_tx`
variants (lib.rs:241-250). Failure path: at-least-once with
exponential backoff, dead-lettering via normal budget mechanics.
**No heartbeat inside `run_once`** — delivery slower than
`visibility_timeout_s` (default 60s) can have its claim expire
mid-delivery and be redelivered (honker's own outbox is exposed to the
dual-execution window).
## 4. Locks (maintenance-relevant)
`_honker_locks (name TEXT PK, owner TEXT NOT NULL, expires_at INTEGER
NOT NULL)` (lib.rs:389-393).
```rust
pub fn lock_acquire(conn, name, owner, ttl_s) -> Result<i64> // 1/0 honker_ops.rs:1387-1415
pub fn lock_release(conn, name, owner) -> Result<i64> // :1417-1423
pub fn lock_renew(conn, name, owner, ttl_s) -> Result<i64> // 1/0 :1431-1442
```
- `lock_acquire`: opportunistic per-name expiry sweep
(`DELETE WHERE name=? AND expires_at <= now`) + `INSERT OR IGNORE` +
read-back. The PK on `name` prevents dual acquisition.
- **`INSERT OR IGNORE` does not refresh TTL on same-owner re-acquire**
— hence `lock_renew` as a separate function (honker_ops.rs:1425-1431;
lib.rs SQL comment :329-332); acquire returns 1 while silently
keeping the old expiry. `lock_renew` requires `(name, owner)` match
and positive ttl.
- `Lock::heartbeat(ttl_s)` RAII wrapper (honker-rs lib.rs:1682-1691),
with the honest caveat "holding the `Lock` value alone does not
guarantee you still own the lock".
- **No general lock-expiry sweep** — expired locks for names nobody
acquires again persist ("Expired rows ... are opportunistically
pruned on every acquire attempt", `_honker.py:962-965`). Locks also
serve as operator mutexes (skip-overlapping-runs recipe,
`_honker.py:1160-1170`).
## 5. Notifications pruning
`_honker_notifications (id, channel, payload, created_at)`
(lib.rs:308-315). **Never auto-pruned** — "no magic timer" (lib.rs
:303-305; tests pin that `notify()` never prunes, `_honker.py`
:1182-1208).
- Rust `prune_notifications(older_than_s)` (lib.rs:323-332) and
`prune_notifications_keep_latest(max_keep)` (lib.rs:336-351 —
rank-based; OFFSET trick correct across id gaps; `max_keep < 0`
errors).
- Python combines both conditions with **OR semantics** — aggressive
`older_than_s` is not blocked by `max_keep`; max_keep "is not a
floor" (`_honker.py:1182-1245`).
## 6. Result storage (brief)
`_honker_results (job_id INTEGER PK, value TEXT, created_at,
expires_at)` (lib.rs:411-416).
`result_save` (ttl ≤ 0 ⇒ NULL), `result_get` (expired reads as None),
`result_sweep` (honker_ops.rs:1850-1900). Python adds
`get_result -> (found, value)` and `wait_result` (`_honker.py:431-498`);
the Python worker saves results **before** ack and swallows
result-persistence failures so side effects aren't duplicated
(`_worker.py:75-83`).
## 7. Maintenance defaults: what runs automatically?
**Nothing. Everything is consumer-driven.** No background sweeper
thread in honker-core (only the `UpdateWatcher` data-version poller,
lib.rs:942-1026). `sweep_expired`, `result_sweep`, `rate_limit_sweep`,
`prune_notifications*`, `scheduler_tick`, `lock_renew` all invoked by
user code; the scheduler tick only **enqueues**. Implicit maintenance:
(a) pre-claim dead-lettering inside `claim_batch` (honker_ops.rs
:855-859), (b) opportunistic expired-lock deletion inside
`lock_acquire` (honker_ops.rs:1393-1397). Python workers get
idle-poll fallback (default 5s) alongside data_version wake
(`_honker.py:374-385`).
## 8. Observed defects / brittleness (feeds OQ-06)
1. **Doc drift** in the Rust binding's `sweep_expired`: "Sweep expired
processing rows back to pending" (lib.rs:652) — actual semantics is
pending → `_honker_dead` (`last_error='expired'`), never
processing→pending (honker_ops.rs:1336-1381).
2. **Zombie processing rows**: a `processing` row whose job-level
`expires_at` passed is unreachable by every path if its worker died
— claim predicate requires future expiry (honker_ops.rs:878),
pre-claim dead-lettering likewise (honker_ops.rs:792),
`sweep_expired` is pending-only (:1346). Sits in `_honker_live`
forever unless `cancel`ed or a still-live worker retries it. Narrow
but real.
3. **Silent-boundary skip after a scheduler outage** — cap 64 then
jump (honker_ops.rs:1612-1621). Deliberate and documented, but part
of the guarantee.
4. **Local-timezone cron** — `chrono::Local` hard-wired in the public
surface (cron.rs:319-321).
5. **No dead-row retention anywhere** — `_honker_dead`,
`_honker_notifications`, `_honker_stream` grow unbounded unless the
user schedules pruning (lib.rs:345-347, `_honker.py:88-90`).
6. **`lock_acquire` doesn't refresh TTL for same-owner re-acquire**
(honker_ops.rs:1398-1401); `_honker_locks` has no standalone sweep.
7. **Backoff has no jitter/cap** (Python plain `2**(attempts-1)`,
`_honker.py:930`; Python no overflow guard, Rust saturates).
8. **Reclaim consumes the attempt budget** (honker_ops.rs:870-872) —
coherent with dead-lettering but a footgun ("timeout ≠ attempt"
assumptions are wrong).
9. **Heartbeat refusal after expiry** leaves an at-least-once
dual-execution window by design (honker_ops.rs:1312-1321) —
downstream idempotency is the user's job; docs admit it
(`_scheduler.py:122-133`).
10. **`get_job` cannot see dead jobs** (honker_ops.rs:1237-1240) —
post-mortem diagnosis is SQL-only.
11. **Savepoint-class defects** (issue #133, post-0.5.1 hardening):
any `DELETE ... RETURNING` + decode + INSERT can strand rows in
neither table on a mid-flight error; fixed recently with
`in_savepoint` (honker_ops.rs:90-117, tests :2340-2700).
12. **`ack`/`ack_batch` don't check `state='processing'`**
(honker_ops.rs:922-926) — implicit rather than enforced.
## 9. Patterns worth carrying over
- One shared SQL-scalar implementation for all bindings, tested once.
- Partial indexes matching the claim hot path; dead rows out of it.
- `unixepoch()` second-precision timestamps (single clock source);
simple, multiprocess-safe.
- Wake = `PRAGMA data_version` delta, coalesced; no wake rows in the
DB (the "no synthetic notification" fix appears three times:
honker_ops.rs:1002, 1119-1120, 1553-1562).
- `queue_next_claim_at` as the sleep-computing claim counterpart
(zero-poll idle workers).
- Savepoint-safe multi-statement mutations.
@@ -0,0 +1,362 @@
---
status: draft
last_updated: 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](../architecture/decisions/010-queue-semantics-depth.md).)*
**Reference:** `/workspace/pgboss-rs` @ `98f7d9e` ("Merge pull request
#19 ... codecov-action-7", HEAD of `main` — verified via `git log`).
**Crate:** `pgboss` v**0.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_on` vestigial.
- `{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 job` DEFAULT partition with the full
index set (ddl.rs:131-176): fetch index `i5 (name, start_after)
INCLUDE (priority, created_on, id) WHERE state < 'active'`, plus four
partial unique indexes implementing queue policies/throttling (`i1`
policy *short*, `i2` *singleton*, `i3` *stately*, `i4`
singleton_on-bucket, `i6` *exclusive*).
- Per-partitioned-queue tables named `j || sha224(queue_name)` hex,
created by `create_queue(queue_name, options jsonb)` plpgsql
(ddl.rs:185-275), attached as partitions `FOR 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:
```sql
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 LOCKED` under the fetch tx; N workers poll
concurrently without contention. Ordering: priority DESC, FIFO by
`created_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 is `start_after < now()` +
caller cadence.
- **Lease/visibility**: no lease renewal, no `extend` call. Right-to-
run bounded by `expire_seconds` (returned as `expire_in`, job.rs
:199): the monitor sweep fails any `active` job with
`started_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 stays `active` until someone calls
`Client::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_jobs` requires `state = 'active'`
strictly; sets `state='completed', completed_on=now(), output=$3`
(dml.rs:223-237). **Failure/nack**: `fail_jobs_by_jids` allows
`state < 'completed'` (includes created) and runs the
delete+reinsert retry/fail/dlq CTE (dml.rs:279-463); failure details
land in `output` (`fail_job_with_details`, job_ops.rs:201-240).
- **Retry-count counting rule** (off-by-design quirk):
`retry_count` increments **on re-fetch**, not on failure — first
delivery leaves `retry_count = 0` (asserted,
`tests/e2e/job_change_state.rs:95-98`). With terminal condition
`retry_count < retry_limit → retry, else failed` (dml.rs:344-346),
`retry_limit = n` gives **n+1 total attempts** (e2e confirms
`retry_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 the `queue` row (ddl.rs
:36-40).
- **Counting**: §3 — increment on re-fetch; terminal at
`retry_count = retry_limit`.
- **Strategies** (dml.rs:352-361, the `fail_jobs` CTE recompute of
`start_after` on 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 base `retry_delay` plus 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, cap `retry_delay_max` (NULL = uncapped;
PG `LEAST` ignores NULLs).
- Terminal failure keeps `start_after` unchanged (dml.rs:352).
- **No custom user-supplied backoff function**; only the fixed/
exponential flag. (Node pg-boss had `retryBackoff` similarly;
upstream also had `retryDelayMax`, carried in SQL but never exposed
in the Rust type API.)
- **Mechanics**: failure runs `DELETE ... RETURNING *` then
**re-INSERTs the same row** as `retry` (new `start_after`,
`completed_on=NULL`) or `failed` (with `completed_on=now()` —
failure is semantically completion) (dml.rs:306-433).
`ON CONFLICT DO NOTHING` on the retried insert + a `failed_jobs`
insert for ids not in `retried_jobs` handles 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 `failed` when `retry_count = retry_limit`
during any fail path (worker `fail_job` or monitor-timeout sweep).
The `fail_jobs` CTE then dead-letters (dml.rs:434-457): for each
`failed` row whose `dead_letter` is 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 by `JOIN queue q ON
q.name = r.dead_letter` (a nonexistent DLQ **silently skips**) and
implicitly `name <> dead_letter` (also enforced by the
`queue.dead_letter` CHECK, ddl.rs:43).
- **Where**: the DLQ is **just another regular queue** (rows in the
same `job` table with a different `name`), created ahead of time
(README example, README.md:27-36). No dedicated dead-letter table,
no distinct state. A DLQ-consumed job's `policy` is `None`
(`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 and `retry_count`), `delete_job(s)`;
`resume_job` does **not** cover `failed` (only `cancelled → created`,
dml.rs:194-208). DLQ resubmission is by hand (get payload, send anew).
- **Retention/cleanup**: governed entirely by timestamps — `keep_until`
(archive-at) and `deletion_seconds` (delete-after-completion). At
this revision **nothing sweeps them**: no background maintenance, no
archive table. Git history shows an `archive` table +
`archive_jobs` procedure existed (`f63d9dc`) and was removed in
`f9db249` ("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`) running
`fail_jobs_by_timeout` (expire `active` jobs past
`started_on + expire_seconds`; retry/fail/DLQ per §4-5; returns
`MaintenanceStats { 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 `queue` table carries full
monitor/maintenance state — `monitor_on`, `maintain_on` cadence
markers, counters `deferred_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_timeout` has 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;
`Client` is 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_app` wraps all DDL in a transaction holding
`pg_advisory_xact_lock(hash(current_database || '.pgboss.{schema}$key'))`
with `SET 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}.schedule` exists 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_on` usage beyond the column and a connect-time log.
- `Job::singleton_for` + `singleton_key` are a related-but-different
mechanism: time-bucketed throttling (`singleton_on` = largest epoch
bucket of size `singleton_for`, dml.rs:96-100) enforced by unique
index `i4` (ddl.rs:157-158), surfaced as `Error::Throttled`
(job_ops.rs:60-64; e2e `tests/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
1. **No background maintenance/monitor runtime** — no `monitorInterval`,
no per-queue monitor cadence; expiration (and stalled-`active`
recovery) only via `force_maintain()`. The queue table's whole
monitor/counters column set is dead weight at this revision.
2. **No cron scheduler** — `schedule` table and `version.cron_on` are
scaffolding only.
3. **No archive stage** — node pg-boss moves expired jobs to an
`archive` table; the port had it, removed it (`f9db249`). Pipeline
is active-table → (retry | failed | DLQ copy) → manual/never
deletion.
4. **No worker framework** — no `work()`/`WorkOptions`, worker pools,
concurrency-per-queue, batch handlers, `onComplete` wiring;
the library stops at fetch/complete/fail primitives (worker loops
are the caller's — cf. `src/bin/loadtest.rs:49-76`).
5. **No `onComplete` child jobs** — `output jsonb` on the job itself,
the `subscription` table an unused remnant.
6. **No LISTEN/NOTIFY wakeup / no pub-sub** — pure polling; the
`subscription` table suggests intent, nothing more.
7. **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).
8. **API-surface gaps vs SQL** — `retry_delay_max`, `delete_after`,
`singleton_offset` consumed by the enqueue SQL but without setters;
DLQ redrive absent; `send_job` cannot target an unregistered queue
(needs explicit `create_queue`).
9. **Migration story** — "must be >= v26, else panic; no upward
migrations" (`src/client/mod.rs:72-86`).
10. **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_on` epoch-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.