11 KiB
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 |
|
moderate | high | component | implementation |
|
Description
Implement the transactional seam — the highest-risk piece of the engine and the shape every other mechanism task rides:
- The
spawn_blockingbridge helper: every engine call (auto-commit and tx ops alike) round-trips a blocking thread; rusqlite connections are notSend-across-await. Establish the pattern once (error mapping substraterusqlite::Error→Error::Databasewith the source chain preserved — ADR-008 §5's opaque fallback) and reuse it everywhere. begin_tx: acquire the writer slot (substrateWriter::acquire), openBEGIN IMMEDIATE, return the boxedTxHandle. The handle is the writer-slot lease (ADR-007): holding it acrossawaitpoints parks the writer — that is the design, not a defect.TxHandleimpl: every*_txmethod routes through the same connection viaspawn_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'snotify()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.
- Writes:
- Drop = rollback (ADR-021 §4): the handle's
Dropimpl issuesROLLBACKand 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
validationhelpers (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_payloadat the seam;?theCodecerror (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_txacquires the writer slot +BEGIN IMMEDIATE; concurrentbegin_txblocks (slot lease) — tested- All eleven
*_txmethods 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 commitcommits + 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 — testedencode_payload'sCodecerror propagates from enqueue/publish paths (typed, not silent)- Substrate errors map to
Error::Databasewith 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, theArcwriter slot it leases, and a store-supplied reopen closure (a bootstrapped open of the db path). Every*_txop routes through one privatewith_connbridge: take the conn out of the Option,spawn_blockingthe op inside an RAIIConnLeaseguard, 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 viareopen, so the next lease proceeds. The same replenish-on-consumed shape holds in the handle'sDropwhen aROLLBACKthere fails (a healthy rollback re-pools the original conn without re-opening). One-shot dispositions:committakes 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 iscommit(self: Box<Self>)consuming); the guard exists so the engine's internal state machine cannot double- release. get_job_txscopes by queue name: the substrate'sget_jobis id-global (live-then-dead lookup); the tx read filters a found row'squeuecolumn against the argument — another queue's row with the same id reads asNone(the contract'sget_job_tx(queue, job_id)is queue-scoped; ADR-019 §1'sOption<Job>shape with the queue argument validated).- Row-decode posture: the substrate's
get_job/read_sincereturn JSON strings (the fork's inherited shape, keyed to honker's function-call surface); the engine decodes them intoJob/StreamEventvia core's#[doc(hidden)] from_rowconstructors, mapping malformed rows toError::Codec(decode-side failures is exactly whatCodecis 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_attemptsis 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_txshort-circuitlimit <= 0→ emptyVecbefore any SQL (ADR-023 §2's extent class; theLIMIT -1dialect artifact dead). publish_tximplements overpublish_with_key_tx(the keyed twin is the one SQL shape; the plain form iskey = None).stream_publish/stream ops re-exported through substrate/mod.rs to the engine layer alongsidequeue_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 singleunixepoch()read),plain_queue_default_ stamps()(300/3/5/none viaQueueOpts::default()),outbox_ default_stamps()(60/5/5),outbox_backing_queue_name()(the reserved-prefix derivation), andresolve_enqueue_opts(delay- over-run_at, relativeexpires, 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.