queue.rs: SqliteQueueHandle (QueueOpts-carrying, stamp-resolving enqueue over the seam's shared resolution arithmetic), claim_one/ claim_batch through the writer slot with ADR-023 §2's extent guard (n <= 0 -> empty Vec at the trait-impl entry), the full JobHandle impl (one-shot ack/retry/fail as 'static boxed futures, repeatable absolute reset heartbeat, substrate's uniform validity predicate), and the engine-owned equal-jitter exponential backoff (std-only RandomState hash jitter, integerized inclusive [ceil(half), cap], 1-hour cap; no rand dep - documented). ack_batch loses its substrate worker filter (register D-31 - ADR-019 §1's worker-less batch form); job_from_json decode + stamps_with_ override moved to their one owners (queue.rs / resolution.rs); reader-pool helpers lifted from stream.rs into seam.rs. Store::queue wired; 10 acceptance tests (lifecycle/stamps, extent guard, exactly-once under concurrency, backoff range + cap + override, validity predicate + reclaim, dead-letter defaults + get_job visibility, cancel, ack_batch, sweep both-states + retention, validation + closed-store). Workspace 25+3+171 green, clippy -D warnings, fmt clean; sqlite suite 4x green.
175 lines
9.5 KiB
Markdown
175 lines
9.5 KiB
Markdown
---
|
||
id: sqlite-engine-queues
|
||
name: SQLite engine — queues (`Queue`, claims, job handles, backoff curve)
|
||
status: completed
|
||
depends_on: [sqlite-engine-seam-tx]
|
||
scope: moderate
|
||
risk: medium
|
||
impact: component
|
||
level: implementation
|
||
tags: [wave-3, sqlite-engine]
|
||
---
|
||
|
||
## Description
|
||
|
||
Implement `Store::queue` and the full `Queue` + `JobHandle` traits
|
||
over the substrate's re-derived queue ops (`enqueue`, `claim_batch`,
|
||
`ack`, `ack_batch`, `heartbeat`, `retry`, `fail`, `cancel`, `get_job`,
|
||
`sweep_expired` — wave 2's owned code):
|
||
|
||
- `queue(name, opts)` — validated constructor; the handle carries the
|
||
`QueueOpts` stamps future enqueues resolve over (ADR-010 §3a).
|
||
- `enqueue(payload, opts)` — the one handle-carrying shape: resolve
|
||
`EnqueueOpts` over **this handle's opts** (delay-over-`run_at`
|
||
precedence, relative-`expires` → absolute `expires_at`, per-call
|
||
`max_attempts` override — ADR-020 §1–§3); one clock read per
|
||
enqueue; `encode_payload` at the seam.
|
||
- `claim_one` / `claim_batch` — through the writer slot (claims are
|
||
writes); map the substrate's JSON row output into `Job` +
|
||
boxed `JobHandle`s. **Extent guard (ADR-023 §2)**: `n <= 0` yields
|
||
the empty `Vec` at the trait-impl entry, before the substrate —
|
||
`claim_batch(-1)` must never claim the whole queue (the `LIMIT -1`
|
||
dialect artifact dies here; the substrate passes `n` through
|
||
contract-blind). `worker_id` non-empty (`InvalidName`).
|
||
- `JobHandle` impl: `job()` (the row value), `ack`/`retry`/`fail`
|
||
(one-shot, `self: Box<Self>`), `heartbeat(extend)` (repeatable,
|
||
absolute reset). The uniform validity predicate is the substrate's;
|
||
refusals are `Ok(false)`, never errors.
|
||
- **The backoff curve — engine-side, this task**: `retry(err, None)`
|
||
computes the equal-jitter exponential from the job's stamped
|
||
`backoff_base_s`: `delay ∈ [base·2^(a−1)/2, base·2^(a−1)]`, capped
|
||
at 1 hour (ADR-010 §3's pinned formula — the engine owns one
|
||
implementation; wave 4 re-derives it identically and wave 5 pins
|
||
equivalence). Jitter needs a randomness source — use the substrate's
|
||
existing posture if one exists, else `std`-only (no `rand` dep
|
||
without a decision; a hash-of-id-and-time fallback is acceptable if
|
||
documented — flag the choice in Notes). `retry(err, Some(d))` passes
|
||
`d` through.
|
||
- `ack_batch(ids)` — per-id validity predicate, count returned.
|
||
- `cancel(job_id)`, `get_job(job_id)` (sees dead rows — map
|
||
`last_error`/`died_at`), `sweep_expired()` (no leader lock
|
||
required).
|
||
- Exhaustion dead-letters with the substrate's pinned default strings
|
||
(`"max attempts exceeded"`; `fail(None)` → `"failed"`).
|
||
|
||
## Acceptance Criteria
|
||
|
||
- [x] Full `Queue` + `JobHandle` impls; enqueue→claim→ack lifecycle
|
||
works; stamps visible via `get_job` match the resolution rules
|
||
- [x] Extent guard: `claim_batch(n <= 0)` claims nothing (empty
|
||
`Vec`, no error, never the whole queue) — pinned per ADR-023
|
||
§2's backlog row
|
||
- [x] Exactly-once handout under concurrent claims (two workers, one
|
||
row — the POC-pinned property, engine-tested here)
|
||
- [x] Backoff curve: `retry(err, None)` delays land in the pinned
|
||
range (statistical test over many retries, tolerance-bounded);
|
||
cap at 1 hour; `Some(d)` overrides
|
||
- [x] Validity predicate: lapsed-deadline ops refuse with `Ok(false)`
|
||
(heartbeat, ack, retry, fail); reclaim consumes an attempt
|
||
- [x] Dead-letter: exhaustion moves to dead with the pinned default
|
||
strings; `get_job` sees dead rows with `last_error`/`died_at`
|
||
- [x] `sweep_expired` moves expired rows (both states) + applies
|
||
retention; returns the moved count
|
||
- [x] `cargo test -p alkstore-sqlite`, clippy `-D warnings`, fmt clean
|
||
|
||
## References
|
||
|
||
- docs/architecture/core-contract.md (queues section)
|
||
- docs/architecture/decisions/010-queue-semantics-depth.md §1–§5
|
||
- docs/architecture/decisions/019-mechanism-handle-surfaces.md §2, §3
|
||
- docs/architecture/decisions/020-enqueue-opt-semantics-and-bridges.md §1–§3
|
||
- docs/architecture/decisions/023-fourth-review-round.md §2
|
||
- docs/architecture/queues.md (Retry/backoff, Dead-letter)
|
||
- alkstore-sqlite/src/substrate/queue_ops.rs
|
||
|
||
## Notes
|
||
|
||
- **`queue.rs` is new** (`SqliteQueueHandle` + `SqliteJobHandle`);
|
||
`Store::queue` wires it (validation + closed-store check at the
|
||
constructor — construction stays free, queues are names not
|
||
registered objects). The handle carries its `QueueOpts` as the stamp
|
||
source; `enqueue` resolves over the handle's stamps via the seam
|
||
task's `resolve_enqueue_opts` + a moved `stamps_with_override`
|
||
(now in `resolution.rs` — one owner, the tx paths share it).
|
||
- **Backoff jitter posture (the flagged choice)**: std-only — a fresh
|
||
`RandomState` hasher (per-process-random seed) over
|
||
`(row_id, attempts, subsecond-nanos)` draws uniformly over the
|
||
integerized range. No `rand` dependency was added (the ADR-005
|
||
dependency posture; a randomness-source change later is engine-opts
|
||
scope, not contract). The integerization is second-precision
|
||
honest: the draw is uniform over `[ceil(cap/2), cap]` **inclusive**
|
||
(never below the pinned lower half, never above the cap; the
|
||
attempt exponent is clamped `u32::try_from`-safe at 30 so the cap
|
||
arithmetic saturates rather than overflows). Statistical acceptance
|
||
test pins range membership + spread over 400 draws per curve point;
|
||
cap test pins `[1800, 3600]` at a=40.
|
||
- **Substrate delta (register D-31)**: `ack_batch(ids_json)` dropped
|
||
its `worker_id` filter — the contract's batch ack carries no worker
|
||
identity (ADR-019 §1: the batch form of ack, per-id independent
|
||
outcomes); the single-row `ack(job_id, worker_id)` keeps its
|
||
conjunct. The engine layer was the only caller. Recorded in
|
||
`PROVENANCE.md` per ADR-012 §4's delta discipline.
|
||
- **Decode owner moved**: `job_from_json` moved from `tx.rs` to
|
||
`queue.rs` (pub(crate)) so the claim pages and the tx reads share
|
||
one decoder; `stream.rs`'s reader helpers (`with_reader`,
|
||
`blocking_read`, closed-error posture) lifted into `seam.rs` as the
|
||
engine-wide reader seam (the queue's `get_job` is a reader-pool
|
||
read; claims/writes stay writer-slot).
|
||
- **One-shot op shapes**: `ack`/`retry`/`fail` destructure
|
||
`Box<Self>` and move the writer `Arc` into the `'static` future
|
||
(the consumed-handle trait shape, mirroring `TxHandle::commit`);
|
||
`heartbeat` borrows `&self` (repeatable). Refusals are `Ok(false)`
|
||
mapped from the substrate's 0-rows (the substrate owns the uniform
|
||
validity predicate engine-blind).
|
||
- **`retry(err, None)`'s exhaust string**: `None` err means the pinned
|
||
`"max attempts exceeded"` (the exhaustion path's owned string wins
|
||
over the caller-None case, per ADR-019 §3's annotation); a
|
||
caller-supplied `Some(err)` rides the dead row verbatim on both the
|
||
reschedule refusals and the dead-letter move. `fail(None)` pins
|
||
`"failed"`.
|
||
- **`get_job` is a reader-pool read** (auto-commit pure read, ADR-021
|
||
§1's read posture at the auto-commit convenience shape); the tx
|
||
twin (`get_job_tx`) stays a writer-conn read as landed in the seam
|
||
task. Claims/ack-batch/cancel/sweep/enqueue ride the writer slot.
|
||
- **Tests** (`store/queue_tests.rs`, 10): lifecycle + stamps
|
||
resolution, extent guard (`-1`, `0`), two-workers-one-row +
|
||
6×40 concurrent drain exactly-once, backoff statistical + cap +
|
||
end-to-end `run_at` + `Some(d)` override, healthy renewal + lapsed
|
||
refusals per op (one claim per one-shot op — consumed handles) +
|
||
reclaim-consumes-an-attempt + ack-loses-to-reclaim,
|
||
dead-letter defaults (`retry(None)`-exhausted, `fail(None)`,
|
||
caller strings) + dead-row `get_job` visibility, cancel
|
||
unconditional + ack_batch mixed-worker count + empty-batch,
|
||
sweep both-states move + retention count, retention vs
|
||
forever-none, validation (empty/reserved queue names, empty
|
||
`worker_id`) + closed-store fail-closed on all queue ops. Tests
|
||
probe the db file with a bare `rusqlite::Connection::open` where
|
||
stamps need honest cross-connection reads. Suite stable 4× green.
|
||
|
||
## Summary
|
||
|
||
Implemented the queue mechanism's SQLite arm: `Store::queue` (validated
|
||
constructor) → `SqliteQueueHandle` carrying its `QueueOpts` stamps,
|
||
`enqueue` resolving `EnqueueOpts` over the handle's opts through the
|
||
seam's shared resolution arithmetic (`stamps_with_override` lifted to
|
||
`resolution.rs`; one clock read per enqueue, delay-over-`run_at`,
|
||
relative-`expires`, `encode_payload` at the seam), `claim_one`/
|
||
`claim_batch` through the writer slot with the ADR-023 §2 extent guard
|
||
(`n <= 0` → empty `Vec` at the trait-impl entry, before the
|
||
contract-blind substrate), the full `JobHandle` impl (one-shot
|
||
`ack`/`retry`/`fail` as `'static` boxed futures over the consumed
|
||
handle, repeatable absolute-reset `heartbeat`; substrate's uniform
|
||
validity predicate, refusals `Ok(false)`), the engine-side backoff
|
||
curve (equal-jitter exponential from the row's stamped
|
||
`backoff_base_s`, integerized inclusive over `[ceil(half), cap]`,
|
||
capped 3600 s; std-only `RandomState` hash jitter, no `rand` dep —
|
||
documented), `ack_batch` over the worker-less per-id predicate
|
||
(substrate D-31 delta, register recorded), unconditional `cancel`,
|
||
dead-visible `get_job` via the shared `job_from_json` decode owner
|
||
(moved from `tx.rs`), and both-states+sweep with retention.
|
||
`reader-pool` helpers (`with_reader`/`blocking_read`/closed-error)
|
||
lifted from `stream.rs` into `seam.rs` as the engine-wide reader seam.
|
||
10 new acceptance tests + the substrate ack_batch test updated
|
||
(171 sqlite tests total). Verified: `cargo test` (workspace: 25 + 3 +
|
||
171 sqlite), `cargo clippy --all-targets -- -D warnings`,
|
||
`cargo fmt --check` all clean; sqlite suite 4× green. |