243 lines
14 KiB
Markdown
243 lines
14 KiB
Markdown
---
|
||
id: pg-engine-queues
|
||
name: Postgres engine — queues (`Queue`, `FOR UPDATE SKIP LOCKED` claims, job handles, backoff curve)
|
||
status: completed
|
||
depends_on: [pg-engine-seam-tx, pg-engine-notify-listen]
|
||
scope: moderate
|
||
risk: medium
|
||
impact: component
|
||
level: implementation
|
||
tags: [wave-4, postgres-engine]
|
||
---
|
||
|
||
## Description
|
||
|
||
Implement the queue mechanism's pg arm — the re-derived queue
|
||
machinery (ADR-004/005's greenfield posture, pg-boss schema family as
|
||
design reference; ADR-010's depth) over the schema task's job/dead
|
||
tables and the forwarder's wake substrate:
|
||
|
||
- **`Store::queue(name, opts)`** — validated constructor (shared-
|
||
namespace kinds; construction stays free, queues are names not
|
||
registered objects). The handle carries its `QueueOpts` as the
|
||
stamp source.
|
||
- **`enqueue`** (auto-commit): resolves `EnqueueOpts` over the
|
||
handle's stamps via the seam task's `resolve_enqueue_opts` +
|
||
`stamps_with_override` (the pg engine's own `resolution.rs` —
|
||
one clock read per enqueue, delay-over-`run_at`, relative-
|
||
`expires`, `max_attempts` the one per-job override), `encode_payload`
|
||
at the seam, `INSERT` on a pool connection. Then **wake**:
|
||
`pg_notify` on the queue's wake channel (the mechanism-name-is-
|
||
the-channel realization) — the LISTEN-driven claim path's push
|
||
half.
|
||
- **`claim_one` / `claim_batch`** — exactly-once handout under
|
||
concurrency: the `FOR UPDATE SKIP LOCKED` claim (POC-pinned under
|
||
4×4 concurrency) on a pool connection, in one statement or a
|
||
CTE-scoped transaction — implementer's choice, but the
|
||
exactly-once property is what the tests pin (two workers, one row;
|
||
N×M concurrent drain, zero double-claims). Semantics per ADR-010:
|
||
ordering `priority DESC, run_at ASC, enqueue order`; `attempts +=
|
||
1` per claim (a reclaim consumes an attempt); the claim sets the
|
||
row's deadline from the job's own stamped `visibility_timeout_s`;
|
||
the pre-claim dead-letter sweep for exhausted reclaimables (the
|
||
substrate's D-12 shape — a worker dying right after its last
|
||
allowed claim dead-letters, not strands); `worker_id` stamps the
|
||
claimant column (non-empty validated, consumer-local — ADR-019 §2).
|
||
The **extent guard** (`n <= 0` → empty `Vec` at trait-impl entry,
|
||
ADR-023 §2). Claims return boxed `JobHandle`s (ops + row value).
|
||
- **`JobHandle` impl** (ADR-019 §3, ADR-010 §2's uniform validity
|
||
predicate — `processing` + unexpired claim deadline; refusal =
|
||
`Ok(false)`, never error):
|
||
- one-shot `ack`/`retry`/`fail` as `'static` boxed futures over the
|
||
consumed handle (mirroring `TxHandle::commit`'s shape);
|
||
repeatable `heartbeat(&self, extend)` — absolute reset of the
|
||
claim deadline (not additive), late heartbeat refused.
|
||
- `ack` deletes the row; `retry(err, None)` computes the queue's
|
||
**equal-jitter exponential curve** — the pg engine's own
|
||
implementation of ADR-010 §3's pinned form
|
||
(`base * 2^(attempt-1)`, jitter over `[half, cap]` integerized
|
||
inclusive, capped 3600 s; attempt index = the row's post-claim
|
||
count). **The jitter's randomness source is this task's flagged
|
||
choice**: the SQLite engine used a std-only `RandomState`-hash
|
||
draw (no `rand` dep, ADR-005's posture) — reuse that approach
|
||
unless a pg-native source is materially better, and **flag the
|
||
choice in Notes**. `retry(err, Some(d))` overrides; `fail(None)`
|
||
pins `"failed"`; `retry`'s exhaust string pins
|
||
`"max attempts exceeded"` (the engine-default strings land
|
||
identically on both engines — the backlog row's wording).
|
||
- Dead-letter moves (retry-at-budget, fail) write the dead table
|
||
with `last_error`/`died_at`; dead rows are `get_job`-visible.
|
||
- **`ack_batch(ids)`** — the batch ack applies the uniform predicate
|
||
per id (lapsed or non-claimed ids silently not counted; count
|
||
returned; no worker filter — the wave-3 D-31 delta's shape is the
|
||
contract's, not a substrate artifact).
|
||
- **`cancel(job_id)`** — unconditional delete (pending or processing,
|
||
not an interrupt); `Ok(false)` when row-less.
|
||
- **`get_job(job_id)`** — dead-visible pure read (pool connection;
|
||
`Option<Job>` with `last_error`/`died_at` on dead rows; the
|
||
resolved values — not raw opts — are what the row shows, ADR-020).
|
||
- **`sweep_expired()`** — both-states no-stranded-rows move (any
|
||
past-expiry row → dead with `last_error='expired'`) + retention
|
||
enforcement (`dead_letter_retention_s`, `None` = forever) —
|
||
ADR-010 §5; the queue-scoped handle form.
|
||
- **LISTEN-driven claim with re-poll safety net** (the engine's
|
||
default consumption posture, ADR-004): the *contract* claim ops are
|
||
pull-shaped (callers claim in their own loops); the wake substrate
|
||
makes those loops reactive — a consumer that `listen(queue_name)`s
|
||
and claims on wake gets the measured 3–6 ms path; the re-poll
|
||
safety net (interval-tunable, `PgOpts`-carried if wired —
|
||
implementer's choice, documented either way) covers LISTEN loss
|
||
(the no-replay hole). The engine's own tests exercise the
|
||
wake-driven claim loop (commit → notify → wake → claim, the POC's
|
||
property) to pin the posture; the contract surface itself stays
|
||
pull-shaped.
|
||
- **Closed-store posture** on all queue ops (the wave-3 shape).
|
||
|
||
## Acceptance Criteria
|
||
|
||
- [x] `enqueue` resolves opts over the handle's stamps (the shared
|
||
resolution arithmetic; `max_attempts` the one override);
|
||
payload bytes exact; wake notify fires on the queue's channel
|
||
- [x] `claim_one`/`claim_batch`: exactly-once under concurrency
|
||
(two-workers-one-row; 6×40 concurrent drain, zero
|
||
double-claims); ordering row pinned; reclaim consumes an
|
||
attempt; deadline from the row's stamp; extent guard
|
||
(`n <= 0` → empty `Vec`); pre-claim exhausted-reclaimable sweep
|
||
- [x] `JobHandle`: validity predicate uniform (lapse refuses all ops,
|
||
`Ok(false)`); heartbeat absolute-reset + late-refused; ack
|
||
deletes; retry curve matches ADR-010 §3's range/cap
|
||
(statistical + cap tests, the SQLite task's shape); dead-letter
|
||
defaults (`"max attempts exceeded"`, `"failed"`, caller strings
|
||
verbatim); dead rows `get_job`-visible
|
||
- [x] `ack_batch` per-id predicate (mixed lapsed/valid ids); `cancel`
|
||
unconditional both states; `sweep_expired` both-states +
|
||
retention (incl. forever-none)
|
||
- [x] Wake-driven claim loop pinned end-to-end (enqueue → wake →
|
||
claim, cross-connection)
|
||
- [x] Validation (empty/reserved queue names, empty `worker_id`) +
|
||
closed-store fail-closed on all queue ops
|
||
- [x] `cargo test -p alkstore-postgres` (harness server), clippy
|
||
`-D warnings`, fmt clean; gates green server-less
|
||
|
||
## References
|
||
|
||
- docs/architecture/core-contract.md (queues section)
|
||
- docs/architecture/decisions/010-queue-semantics-depth.md (§1–§5, the depth)
|
||
- docs/architecture/decisions/019-mechanism-handle-surfaces.md §1–§3
|
||
- docs/architecture/decisions/020-enqueue-opt-semantics-and-bridges.md
|
||
- docs/architecture/decisions/023-fourth-review-round.md §2 (extent guard)
|
||
- docs/architecture/decisions/004-postgres-driver.md (re-derive posture, LISTEN-driven default)
|
||
- docs/architecture/queues.md (lifecycle, consumption posture)
|
||
- docs/research/reference-pgboss-rs-semantics.md (design reference)
|
||
- alkstore-sqlite/src/queue.rs (the structural twin — same semantics, own SQL)
|
||
- alkstore/src/queue.rs, job_handle.rs (the traits)
|
||
|
||
## Notes
|
||
|
||
- **The claim is one statement — the task's "one statement or a
|
||
CTE-scoped transaction — implementer's choice" resolved to the one
|
||
statement**: an `UPDATE … RETURNING` over an `id IN (SELECT … FOR
|
||
UPDATE SKIP LOCKED ORDER BY … LIMIT n)` subquery (the picked-rows
|
||
CTE inline). The exactly-once property and the ordering row are
|
||
what the tests pin; the statement shape rides.
|
||
- **The dead-letter moves are transactional — the #133 defect class
|
||
explicitly excluded on the pg arm**: the naive DELETE-then-INSERT
|
||
pair across bare auto-commit statements would strand a row in
|
||
*neither* table when the dead INSERT failed (the substrate's
|
||
"stranded in NEITHER table is the property" test row). All four
|
||
move paths (retry-at-budget, `fail`, the pre-claim exhausted-
|
||
reclaimable sweep, `sweep_expired`'s both-states move + retention)
|
||
run inside an explicit `BEGIN`…`COMMIT` frame through a pool
|
||
`in_tx` helper — the seam's pooled-object posture scoped to one op
|
||
(a failed `BEGIN`/`COMMIT`/`ROLLBACK` discards the object via
|
||
`Object::take`, the seam's unknowable-state arm verbatim; the
|
||
body's error arm rolls back and re-pools) — and the single-job ops
|
||
(`retry`/`fail`) ride one frame each: the validity-predicate row
|
||
re-read (`SELECT … FOR UPDATE`) + the dead INSERT + the live
|
||
DELETE commit atomically, and a reclaim racing the op resolves
|
||
inside the frame (the loser's predicate re-read sees 0 rows).
|
||
- **The claim frame includes the pre-claim sweep atomically**: the
|
||
at-budget reclaimables' dead-letter move and the claim statement
|
||
commit together (the sweep alone would strand on its second
|
||
statement's failure; together they're one atomic unit, and the
|
||
claim predicate skipping already-dead rows is automatic).
|
||
- **The jitter's randomness source (the flagged choice)**: the
|
||
SQLite engine's std-only `RandomState`-hash draw reused verbatim —
|
||
no `rand` dep (ADR-005's posture), the same formula body (the
|
||
exponent clamps at 30 so the doubling saturates; the draw is
|
||
uniform over the integerized `[ceil(half), cap]` inclusive), its
|
||
own `backoff_delay_s` (wave 5 pins the twins' equivalence).
|
||
- **The enqueue wake is best-effort, logged-and-swallowed** (the
|
||
stream task's twin posture): the durable row is the truth, the
|
||
re-poll safety net / the consumer's re-read covers a failed wake —
|
||
never may a wake fee fail (or stall) a committed enqueue. The wake
|
||
channel is the queue's own name (the mechanism-name-is-the-channel
|
||
realization); the wake payload is empty (`''` — the wake carries
|
||
no content, ADR-008 §3).
|
||
- **The ack_batch predicate rides `claim_expires_at >= now` with the
|
||
queue's name scope** (the seam's scoping rule at the auto-commit
|
||
shape): a queue handle acks its own rows only — matching the
|
||
queue-scoped `get_job`/`cancel`/sweep predicates (the tx seam's
|
||
scoping rule carried to every queue op).
|
||
- **Params riding SQL arithmetic need explicit `::bigint` casts
|
||
where they meet column expressions** (`$2::bigint +
|
||
visibility_timeout_s`, the run_at/heartbeat arithmetic, the
|
||
retention TTL subtraction) — the seam task's noted
|
||
extended-protocol param-typing posture (the server infers param
|
||
types from context; `$2 + col` where `$2` sits beside no literal
|
||
is the E42725 "operator is not unique" arm, observed and fixed in
|
||
this task's first harness run).
|
||
- **`enqueue_row` lifted from `PgTxHandle`'s inherent method into
|
||
the module-level function the auto-commit path shares** (one owner
|
||
of the INSERT — the tx paths bridge it through a one-line inherent
|
||
wrapper; the schema string is the parameter — the tx handle keeps
|
||
its field, the queue handle passes its own).
|
||
- **PgStore's `queue` wiring replaced the stub; the open task's
|
||
`store_trait_methods_are_wiring_stubs` test updated accordingly**
|
||
(the constructor pinned as wired alongside notify/listen/stream;
|
||
`outbox` is the remaining stub-message representative).
|
||
- **Test details**: the drain test opens sibling stores over the
|
||
same schema per task (the honest cross-instance worker shape —
|
||
queue rows are schema state, not store-instance state; the
|
||
bigserial id range is probed via a raw admin count, not a
|
||
`1..=n` loop); the expiry backdating rides admin UPDATEs
|
||
(schema-qualified); the wake-driven loop test pins
|
||
enqueue→notify→wake→claim cross-connection plus the poll-only
|
||
fallback's shape (a claim loop with no wake still lands work).
|
||
|
||
## Summary
|
||
|
||
Landed alkstore-postgres's queue mechanism: `queue.rs` (the
|
||
`PgQueueHandle` carrying its `QueueOpts` stamps; `enqueue` resolving
|
||
over the handle's stamps through the shared resolution arithmetic
|
||
with `encode_payload` at the seam and the best-effort
|
||
queue-channel wake; `claim_one`/`claim_batch` as the one-statement
|
||
`FOR UPDATE SKIP LOCKED` claim in a `BEGIN`-framed transaction with
|
||
the atomically-included pre-claim exhausted-reclaimable sweep, the
|
||
job's-own-stamp deadline, the attempt-counting, the ADR-023 §2
|
||
extent guard, and the ADR-019 §2 worker-id validation; the full
|
||
`JobHandle` impl — one-shot `ack`/`retry`/`fail` over `in_tx`
|
||
frames with the uniform validity predicate + `FOR UPDATE` row
|
||
re-read + atomic dead-letter moves, repeatable absolute-reset
|
||
`heartbeat`, the engine-side equal-jitter backoff curve (std-only
|
||
`RandomState` jitter, capped 3600 s, flagged), the pinned
|
||
dead-letter default strings, dead-visible queue-scoped `get_job`,
|
||
unconditional `cancel`, worker-less per-id `ack_batch`, and the
|
||
both-states + retention `sweep_expired`), the `Store::queue` wiring
|
||
replacing its stub, and `enqueue_row` lifted to the shared
|
||
module-level INSERT owner. neither-table state). 15 new acceptance tests + 2 server-less unit
|
||
pins (89 lib tests green ×3 runs at 8/4/1 threads against the harness
|
||
server, server-less skips clean; workspace `cargo test` green),
|
||
covering the lifecycle + stamp resolutions, the extent guard, the
|
||
two-workers-one-row + 6×40 cross-instance drain exactly-once, the
|
||
ordering row, the statistical + capped + end-to-end backoff curve,
|
||
the healthy-renewal + lapsed-refusals + reclaim-consumes-an-attempt
|
||
predicate row, the D-12 pre-claim sweep, the dead-letter defaults +
|
||
dead-row `get_job` visibility, cancel + ack_batch mixed-worker
|
||
counts, the both-states sweep + retention + forever-none, the
|
||
wake-driven claim loop end-to-end (cross-connection wake → claim,
|
||
plus the poll-only fallback), queue-name scoping, validation
|
||
(empty/reserved queue names, empty `worker_id`), and the
|
||
closed-store fail-closed on all queue ops; `cargo clippy
|
||
--all-targets -- -D warnings` and `cargo fmt --check` clean
|
||
workspace-wide. |