Files
alkstore/tasks/pg-engine-queues.md

243 lines
14 KiB
Markdown
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
---
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.