--- id: sqlite-engine-seam-tx name: SQLite engine — `spawn_blocking` seam, `begin_tx`, `TxHandle`, writer-slot lease status: completed depends_on: [sqlite-engine-open-opts] scope: moderate risk: high impact: component level: implementation tags: [wave-3, sqlite-engine, seam] --- ## Description Implement the transactional seam — the highest-risk piece of the engine and the shape every other mechanism task rides: - The `spawn_blocking` bridge helper: every engine call (auto-commit and tx ops alike) round-trips a blocking thread; rusqlite connections are not `Send`-across-await. Establish the pattern once (error mapping substrate `rusqlite::Error` → `Error::Database` with the source chain preserved — ADR-008 §5's opaque fallback) and reuse it everywhere. - `begin_tx`: acquire the writer slot (substrate `Writer::acquire`), open `BEGIN IMMEDIATE`, return the boxed `TxHandle`. The handle is the **writer-slot lease** (ADR-007): holding it across `await` points parks the writer — that is the design, not a defect. - `TxHandle` impl: every `*_tx` method routes through the same connection via `spawn_blocking` (thread-affinity is handled by the slot lease — the connection moves with the handle, ops serialize on it): - Writes: `enqueue_tx` (stamps the engine's derived plain-queue defaults 300/3/5/none — no-handle-open shape, ADR-020 §3a note), `publish_tx`, `publish_with_key_tx` (key validation per ADR-015 §2), `notify_tx` (substrate's `notify()` SQL function inside the tx — delivered only at commit), `save_offset_tx` (monotone — substrate's upsert already is), `outbox_enqueue_tx` (derives the `__alkstore_outbox:{name}` backing queue name engine-side and stamps the outbox's 60/5/5 set — ADR-014 §1). - Reads (ADR-021 §1, read-your-own-writes — they run inside the caller's tx): `get_job_tx`, `get_offset_tx`, `read_since_tx`, `read_from_consumer_tx`. - `commit(self: Box)`: `COMMIT` + release the writer slot. - **Drop = rollback** (ADR-021 §4): the handle's `Drop` impl issues `ROLLBACK` and releases the writer slot — panic-safe (the drop runs across unwinds; no panic may escape drop). This is the property N-5's wave-5 probe will exercise against a real engine for the first time — get the Drop impl right here. - Name validation at every entry point via core's `validation` helpers (shared-namespace kinds vs consumer-local; ADR-008 §4 — the commit-atomic path must not bypass validation). - Payload encoding: call core's now-fallible `encode_payload` at the seam; `?` the `Codec` error (ADR-023 §1 — the engines' first consumption of the typed signature). All eleven `*_tx` methods are implemented **for real** in this task — the substrate surface is fully ported (wave 2), so nothing stubs. The opts-resolution helpers this introduces (delay-over-`run_at`, relative-`expires`, the derived-default stamp sets) are the shared engine-side arithmetic the auto-commit mechanism tasks reuse — one owner per formula (ADR-012 §2), written once here. The mechanism tasks then add what the tx path deliberately lacks: the auto-commit handles, the extent/duration guards, the backoff curve, the receiver bridges, and the scheduler runner. ## Acceptance Criteria - [x] `begin_tx` acquires the writer slot + `BEGIN IMMEDIATE`; concurrent `begin_tx` blocks (slot lease) — tested - [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 - [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 - [x] Drop is panic-safe (a panic unwinding through a scope holding the handle still rolls back and releases the slot) - [x] `with_tx` (core's provided method) works end-to-end on this engine: Ok ⇒ commit, Err ⇒ rollback — tested - [x] `encode_payload`'s `Codec` error propagates from enqueue/publish paths (typed, not silent) - [x] Substrate errors map to `Error::Database` with source chains - [x] `cargo test -p alkstore-sqlite`, clippy `-D warnings`, fmt clean ## References - docs/architecture/decisions/007-transactional-seam.md - docs/architecture/decisions/021-tx-reads-and-value-shape-fixes.md §1, §4 - docs/architecture/decisions/014-outbox-tx-enqueue.md §1 - docs/architecture/decisions/020-enqueue-opt-semantics-and-bridges.md §1–§4 - docs/architecture/decisions/023-fourth-review-round.md §1 - docs/architecture/engine-sqlite.md (Mapping the contract: begin_tx, handle ops, commit/rollback) - alkstore/src/tx.rs (the trait to implement) ## Notes - **Handle shape**: `SqliteTxHandle { conn: Option, writer: Arc, reopen: Arc Result> }` — 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)` 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` 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 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.