Files
alkstore/tasks/sqlite-engine-seam-tx.md

11 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
sqlite-engine-seam-tx SQLite engine — `spawn_blocking` seam, `begin_tx`, `TxHandle`, writer-slot lease completed
sqlite-engine-open-opts
moderate high component implementation
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<Self>): 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

  • begin_tx acquires the writer slot + BEGIN IMMEDIATE; concurrent begin_tx blocks (slot lease) — tested
  • 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) 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 the handle still rolls back and releases the slot)
  • 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 paths (typed, not silent)
  • Substrate errors map to Error::Database with source chains
  • 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<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

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.