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

14 KiB
Raw Permalink Blame History

id, name, status, depends_on, scope, risk, impact, level, tags
id name status depends_on scope risk impact level tags
pg-engine-queues Postgres engine — queues (`Queue`, `FOR UPDATE SKIP LOCKED` claims, job handles, backoff curve) completed
pg-engine-seam-tx
pg-engine-notify-listen
moderate medium component implementation
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 JobHandles (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

  • 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
  • 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
  • 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
  • ack_batch per-id predicate (mixed lapsed/valid ids); cancel unconditional both states; sweep_expired both-states + retention (incl. forever-none)
  • Wake-driven claim loop pinned end-to-end (enqueue → wake → claim, cross-connection)
  • Validation (empty/reserved queue names, empty worker_id) + closed-store fail-closed on all queue ops
  • 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.