SQLite engine: transactional seam — begin_tx, writer-slot lease, all eleven *_tx methods (task sqlite-engine-seam-tx)

This commit is contained in:
glm-5.3-flash committed 2026-10-08 11:50:16 +00:00
1 parent 7d400906f5
commit c38033db1a
8 files changed
+1988 -20

No files matched your search

+104 -11
View File
@@ -1,7 +1,7 @@
---
id: sqlite-engine-seam-tx
name: SQLite engine — `spawn_blocking` seam, `begin_tx`, `TxHandle`, writer-slot lease
status: pending
status: completed
depends_on: [sqlite-engine-open-opts]
scope: moderate
risk: high
@@ -65,23 +65,23 @@ and the scheduler runner.
## Acceptance Criteria
- [ ] `begin_tx` acquires the writer slot + `BEGIN IMMEDIATE`;
- [x] `begin_tx` acquires the writer slot + `BEGIN IMMEDIATE`;
concurrent `begin_tx` blocks (slot lease) — tested
- [ ] All eleven `*_tx` methods implemented with entry-point
- [x] All eleven `*_tx` methods implemented with entry-point
validation; writes commit-atomic, reads see the tx's own writes;
opts resolution (delay-over-run_at, relative-expires, derived
default stamp sets 300/3/5/none and 60/5/5) pinned by tests
- [ ] `commit` commits + releases the slot; drop (without commit)
- [x] `commit` commits + releases the slot; drop (without commit)
rolls back + releases the slot — no-ghosts tested for job rows,
stream events, notifications, and offset saves
- [ ] Drop is panic-safe (a panic unwinding through a scope holding
- [x] Drop is panic-safe (a panic unwinding through a scope holding
the handle still rolls back and releases the slot)
- [ ] `with_tx` (core's provided method) works end-to-end on this
- [x] `with_tx` (core's provided method) works end-to-end on this
engine: Ok ⇒ commit, Err ⇒ rollback — tested
- [ ] `encode_payload`'s `Codec` error propagates from enqueue/publish
- [x] `encode_payload`'s `Codec` error propagates from enqueue/publish
paths (typed, not silent)
- [ ] Substrate errors map to `Error::Database` with source chains
- [ ] `cargo test -p alkstore-sqlite`, clippy `-D warnings`, fmt clean
- [x] Substrate errors map to `Error::Database` with source chains
- [x] `cargo test -p alkstore-sqlite`, clippy `-D warnings`, fmt clean
## References
@@ -95,8 +95,101 @@ and the scheduler runner.
## Notes
> To be filled by implementation agent
- **Handle shape**: `SqliteTxHandle { conn:
Option<rusqlite::Connection>, writer: Arc<Writer>, reopen: Arc<dyn
Fn() -> Result<Connection>> }` — the handle holds the acquired
writer connection directly, the `Arc` writer slot it leases, and a
store-supplied reopen closure (a bootstrapped open of the db path).
Every `*_tx` op routes through one private `with_conn` bridge:
take the conn out of the Option, `spawn_blocking` the op inside an
RAII `ConnLease` guard, put the conn back on the future's
completion. The guard is the slot-leak fix: a future cancelled
mid-op (tokio abort) drops the lease instead — the conn is dropped
(its uncommitted tx rolls back at the SQLite layer) and the slot
is replenished with a fresh bootstrapped connection via `reopen`,
so the next lease proceeds. The same replenish-on-consumed shape
holds in the handle's `Drop` when a `ROLLBACK` there fails (a
healthy rollback re-pools the original conn without re-opening).
One-shot dispositions: `commit` takes the conn (re-take after
commit → `Error::Closed`, the internal fail-closed guard); no
panic escapes drop (a failed rollback is eprintln'd, then the
lease accounting completes).
- **Ops-after-commit fail closed** with `Error::Closed`: reachable
only through a mid-op-cancelled future racing a subsequent op (the
consumer-facing path is `commit(self: Box<Self>)` consuming); the
guard exists so the engine's internal state machine cannot double-
release.
- **`get_job_tx` scopes by queue name**: the substrate's `get_job` is
id-global (live-then-dead lookup); the tx read filters a found
row's `queue` column against the argument — another queue's row
with the same id reads as `None` (the contract's `get_job_tx(queue,
job_id)` is queue-scoped; ADR-019 §1's `Option<Job>` shape with the
queue argument validated).
- **Row-decode posture**: the substrate's `get_job`/`read_since`
return JSON strings (the fork's inherited shape, keyed to honker's
function-call surface); the engine decodes them into `Job`/
`StreamEvent` via core's `#[doc(hidden)] from_row` constructors,
mapping malformed rows to `Error::Codec` (decode-side failures is
exactly what `Codec` is for; the fork's writer always wrote well-
formed rows, so this is a posture guard, not a live path).
- **`stamps_with_override`**: `EnqueueOpts::max_attempts` is the one
per-job override resolved over the base stamp set — the other
stamps (visibility/backoff/retention) always take the base's
(`enqueue_tx` → plain defaults 300/3/5/none; `outbox_enqueue_tx` →
the outbox's 60/5/5). One helper both shapes share.
- **Payload bytes → storage**: serde_json serializes to UTF-8 bytes;
the substrate's rows are TEXT (`str::from_utf8` — a
serde_json-produced encode cannot fail this, `Codec`-typed as a
posture guard), notify's payload rides the SQL function as text
the same way.
- **Extent rule at the tx reads**: `read_since_tx`/`read_from_consu-
mer_tx` short-circuit `limit <= 0` → empty `Vec` before any SQL
(ADR-023 §2's extent class; the `LIMIT -1` dialect artifact dead).
- **`publish_tx` implements over `publish_with_key_tx`** (the keyed
twin is the one SQL shape; the plain form is `key = None`).
- **`stream_publish`/stream ops re-exported through substrate/mod.rs**
to the engine layer alongside `queue_ops` (wave 2's ported surface —
the `#![allow(dead_code)]` posture stays until the last wiring task;
unused re-exports stay legal).
- **resolution.rs is new (one owner per formula, ADR-012 §2)**:
`now_unix` (the single `unixepoch()` read), `plain_queue_default_
stamps()` (300/3/5/none via `QueueOpts::default()`), `outbox_
default_stamps()` (60/5/5), `outbox_backing_queue_name()` (the
reserved-prefix derivation), and `resolve_enqueue_opts` (delay-
over-`run_at`, relative `expires`, one clock read per enqueue).
The auto-commit mechanism tasks reuse these directly.
- **Test-side raw readers**: several no-ghosts/visibility assertions
read the db file with a bare `rusqlite::Connection::open` (read-
only or RW) — the honest cross-connection probe; `get_job_tx`'s id
namespace is global and starts at 1, so tests probe the ids their
own enqueues returned.
- **Stability**: the suite ran 4× green (timing-sensitive lease
tests use 150ms sleep + 5s bounded wait, no tight-race assertions).
## Summary
> To be filled on completion
Implemented the transactional seam's SQLite arm: `begin_tx` (acquire
the writer slot, `BEGIN IMMEDIATE`, return the boxed handle), the
`SqliteTxHandle` writer-slot lease with all eleven `*_tx` methods
implemented for real (enqueue/publish/keyed-publish/notify/save-
offset/get-job/get-offset/read-since/read-from-consumer/outbox-
enqueue + one-shot commit), `resolution.rs` as the one-owner engine-
side arithmetic (derived stamp sets, opts resolution, outbox
derivation), drop = rollback on panic-safe ADR-021 §4 paths, and the
entry-point validation / source-chain error mapping / `Codec`-typed
payload encoding postures from ADR-008 §4/§5, ADR-020 §1–§4,
ADR-021 §1/§4, ADR-023 §1/§2. `with_tx` (core's provided method)
works end-to-end over it. Substrate re-exports widened to the engine
layer (`mod.rs`, plus the `open_writer_connection` reopen helper);
the seam's `blocking` helper ungated (the open-opts task's
`#[cfg(test)]` gate dropped). 22 new tests (lease blocking +
release, stamps/opts-resolution parity per set, read-your-own-writes
for jobs and streams, notify-at-commit-only, monotone saves,
no-ghosts through commit/drop/panic paths, panic-safe drop,
mid-op-cancellation slot replenishment, `with_tx`
Ok⇒commit/Err⇒rollback, entry-point validation across all eleven
surfaces, extent refusal, committed-handle teardown idempotence,
cross-store lease independence, serialized fan-out through the
lease). Verified: `cargo test` (workspace: 25 + 3 suite + 133
sqlite), `cargo clippy --all-targets -- -D warnings`,
`cargo fmt --check` all clean; suite run 4× green.