Wave-3 review gate: contract conformance clean; two findings fixed inline — substrate boundary move (open_writer_connection → seam.rs), writer-slot release on with_writer/begin/commit error arms + regression tests (task review-wave-3)
This commit is contained in:
1 parent
a82c543b40
commit
ecb8211694
8 files changed
+290
-34
No files matched your search
+14
-10
@@ -45,6 +45,17 @@ pub(crate) fn sqlite_error(e: rusqlite::Error) -> Error {
|
||||
Error::database(e)
|
||||
}
|
||||
|
||||
/// Provision a fresh bootstrapped connection into the writer slot
|
||||
/// after a cancellation-path connection was consumed. The engine's tx
|
||||
/// handle captures this as its `reopen` closure's body; the
|
||||
/// substrate-outcome mapping stays engine-side (ADR-012 §2 — the
|
||||
/// taxonomy mapping is engine territory).
|
||||
pub(crate) fn open_writer_connection(path: &str) -> Result<rusqlite::Connection> {
|
||||
crate::substrate::open_conn_bootstrapped(path).map_err(|e| match e {
|
||||
crate::substrate::SchemaError::Sqlite(inner) => sqlite_error(inner),
|
||||
})
|
||||
}
|
||||
|
||||
/// The short-lived writer-slot lease for auto-commit ops (the
|
||||
/// `notify` path's shape): acquire, run the op inside the
|
||||
/// `spawn_blocking` seam, release — the slot is free the instant the
|
||||
@@ -63,16 +74,9 @@ where
|
||||
"the store is closed: writer slot unavailable",
|
||||
))
|
||||
})?;
|
||||
match f(&conn) {
|
||||
Ok(out) => {
|
||||
writer.release(conn);
|
||||
Ok(out)
|
||||
}
|
||||
Err(e) => {
|
||||
drop(conn);
|
||||
Err(e)
|
||||
}
|
||||
}
|
||||
let out = f(&conn);
|
||||
writer.release(conn);
|
||||
out
|
||||
})
|
||||
.await
|
||||
}
|
||||
|
||||
@@ -120,7 +120,7 @@ impl Store for SqliteStore {
|
||||
let db_path = self.db_path.to_string_lossy().into_owned();
|
||||
let reopen: std::sync::Arc<
|
||||
dyn Fn() -> alkstore::Result<rusqlite::Connection> + Send + Sync,
|
||||
> = std::sync::Arc::new(move || crate::substrate::open_writer_connection(&db_path));
|
||||
> = std::sync::Arc::new(move || crate::seam::open_writer_connection(&db_path));
|
||||
crate::tx::SqliteTxHandle::begin(writer, reopen)
|
||||
}
|
||||
|
||||
|
||||
@@ -118,6 +118,47 @@ async fn notify_never_produces_payload_too_large() {
|
||||
cleanup(&dir);
|
||||
}
|
||||
|
||||
/// Acceptance: an auto-commit op's *error* arm releases the writer
|
||||
/// slot — the lease-out-both-arms property (ADR-007's seam posture):
|
||||
/// after a notify whose statement fails at the storage layer, the slot
|
||||
/// is free again and a subsequent writer op succeeds without parking.
|
||||
#[tokio::test(flavor = "multi_thread")]
|
||||
async fn failed_notify_releases_the_writer_slot() {
|
||||
let dir = temp_dir("err-arm-release");
|
||||
let path = temp_path("err-arm-release");
|
||||
let store = open_store(path.to_str().unwrap(), Default::default()).unwrap();
|
||||
|
||||
let side = rusqlite::Connection::open(&path).unwrap();
|
||||
side.execute_batch(
|
||||
"CREATE TRIGGER block_notify BEFORE INSERT ON __alkstore_notifications
|
||||
BEGIN SELECT RAISE(ABORT, 'notify insert blocked'); END",
|
||||
)
|
||||
.unwrap();
|
||||
|
||||
let err = store.notify("orders", json!({"x": 1})).await.unwrap_err();
|
||||
let source = match &err {
|
||||
Error::Database(source) => source.to_string(),
|
||||
other => panic!("the induced failure surfaces as Database, got {other:?}"),
|
||||
};
|
||||
assert!(
|
||||
source.contains("notify insert blocked"),
|
||||
"the source chain carries the trigger's message, got: {source}"
|
||||
);
|
||||
side.execute_batch("DROP TRIGGER block_notify;").unwrap();
|
||||
drop(side);
|
||||
|
||||
let deadline = tokio::time::Instant::now() + Duration::from_secs(5);
|
||||
let landed = tokio::time::timeout_at(deadline, store.notify("orders", json!({"x": 2})))
|
||||
.await
|
||||
.expect(
|
||||
"the failed op must release the writer slot — the slot is free again, \
|
||||
never parked",
|
||||
);
|
||||
landed.expect("the notify lands after a failed op released the slot");
|
||||
store.close();
|
||||
cleanup(&dir);
|
||||
}
|
||||
|
||||
/// Acceptance: `listen` returns a working receiver — wakes arrive
|
||||
/// after commits by other connections, and `Wake { channel }` carries
|
||||
/// the channel name only. The honest posture: the watcher fans out on
|
||||
|
||||
@@ -1015,6 +1015,49 @@ async fn failed_ops_inside_a_tx_keep_the_lease() {
|
||||
cleanup(&dir);
|
||||
}
|
||||
|
||||
/// Acceptance: a failed `BEGIN IMMEDIATE` does not strand the writer
|
||||
/// slot — the acquired connection returns to the slot on the begin's
|
||||
/// error arm (the lease-out-both-arms property at the seam's entry)
|
||||
/// and the next `begin_tx` succeeds without parking. Induced with a
|
||||
/// concurrent exclusive holder at a shortened busy timeout, so the
|
||||
/// begin fails fast with `SQLITE_BUSY` instead of parking.
|
||||
#[tokio::test(flavor = "multi_thread")]
|
||||
async fn failed_begin_releases_the_writer_slot() {
|
||||
let dir = temp_dir("begin-err-release");
|
||||
let concrete = open_store(dir.join("store.db").to_str().unwrap(), Default::default()).unwrap();
|
||||
|
||||
{
|
||||
let writer_conn = concrete.writer.acquire().unwrap();
|
||||
writer_conn
|
||||
.execute_batch("PRAGMA busy_timeout = 100;")
|
||||
.unwrap();
|
||||
concrete.writer.release(writer_conn);
|
||||
}
|
||||
let holder = rusqlite::Connection::open(dir.join("store.db")).unwrap();
|
||||
holder.execute_batch("BEGIN EXCLUSIVE; SELECT 1;").unwrap();
|
||||
|
||||
let begin = tokio::time::timeout(std::time::Duration::from_secs(5), concrete.begin_tx()).await;
|
||||
match begin {
|
||||
Err(_) => panic!("begin must fail fast at the 100 ms busy timeout, not park"),
|
||||
Ok(Ok(_handle)) => {
|
||||
panic!("begin must refuse under the exclusive holder, got a lease")
|
||||
}
|
||||
Ok(Err(e)) => assert!(
|
||||
matches!(e, Error::Database(_)),
|
||||
"a busy-refused begin errs opaque Database, got: {e:?}"
|
||||
),
|
||||
}
|
||||
drop(holder);
|
||||
|
||||
let next = tokio::time::timeout(std::time::Duration::from_secs(5), concrete.begin_tx())
|
||||
.await
|
||||
.expect("the slot must be grantable after the failed begin")
|
||||
.unwrap();
|
||||
drop(next);
|
||||
drop(concrete);
|
||||
cleanup(&dir);
|
||||
}
|
||||
|
||||
/// A concurrent writer's commit blocks while a lease holds the slot —
|
||||
/// the parked-writer posture (ADR-007 negative consequence, honest).
|
||||
#[tokio::test(flavor = "multi_thread")]
|
||||
|
||||
@@ -66,13 +66,3 @@ pub(crate) use schema::Error as SchemaError;
|
||||
pub(crate) use schema::open_conn_bootstrapped;
|
||||
pub(crate) use schema::{Readers, Writer};
|
||||
pub(crate) use watcher::{SharedUpdateWatcher, WatcherConfig};
|
||||
|
||||
/// Provision a fresh connection into the writer slot after a
|
||||
/// cancellation-path connection was consumed. Private helper over the
|
||||
/// substrate's bootstrapped open; the engine's tx handle captures it
|
||||
/// as its `reopen` closure's body.
|
||||
pub(crate) fn open_writer_connection(path: &str) -> Result<rusqlite::Connection, alkstore::Error> {
|
||||
open_conn_bootstrapped(path).map_err(|e| match e {
|
||||
SchemaError::Sqlite(inner) => alkstore::Error::database(inner),
|
||||
})
|
||||
}
|
||||
@@ -70,9 +70,13 @@ impl SqliteTxHandle {
|
||||
"the store is closed: writer slot unavailable",
|
||||
))
|
||||
})?;
|
||||
conn.execute_batch("BEGIN IMMEDIATE")
|
||||
.map_err(sqlite_error)
|
||||
.map(|_| conn)
|
||||
match conn.execute_batch("BEGIN IMMEDIATE").map_err(sqlite_error) {
|
||||
Ok(()) => Ok(conn),
|
||||
Err(e) => {
|
||||
writer2.release(conn);
|
||||
Err(e)
|
||||
}
|
||||
}
|
||||
})
|
||||
.await?;
|
||||
Ok(Box::new(SqliteTxHandle {
|
||||
@@ -387,12 +391,23 @@ impl TxHandle for SqliteTxHandle {
|
||||
return Box::pin(async { Err(Error::Closed) });
|
||||
};
|
||||
let writer = self.writer.clone();
|
||||
let writer_reopen = self.reopen.clone();
|
||||
Box::pin(async move {
|
||||
blocking(move || {
|
||||
conn.execute_batch("COMMIT").map_err(sqlite_error)?;
|
||||
writer.release(conn);
|
||||
Ok(())
|
||||
})
|
||||
blocking(
|
||||
move || match conn.execute_batch("COMMIT").map_err(sqlite_error) {
|
||||
Ok(()) => {
|
||||
writer.release(conn);
|
||||
Ok(())
|
||||
}
|
||||
Err(e) => {
|
||||
drop(conn);
|
||||
if let Ok(fresh) = (writer_reopen)() {
|
||||
writer.release(fresh);
|
||||
}
|
||||
Err(e)
|
||||
}
|
||||
},
|
||||
)
|
||||
.await
|
||||
})
|
||||
}
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
---
|
||||
status: draft
|
||||
last_updated: 2026-10-08 (waves 1–2 implemented + reviewed; ADR-023 follow-through folded into wave 3; wave 3 decomposed)
|
||||
last_updated: 2026-10-08 (wave 3 implemented + reviewed; review-gate fixes: substrate boundary move, writer-slot error-arm release; wave 4/5 decomposition may proceed)
|
||||
---
|
||||
|
||||
# alkstore — Implementation plan
|
||||
@@ -34,7 +34,7 @@ parallel).
|
||||
|---|---|---|---|
|
||||
| [1](#wave-1--foundations) | Workspace scaffold; core crate (errors, value types, full trait surface); contract-suite scaffold | — | implemented + reviewed (2026-10-08) |
|
||||
| [2](#wave-2--sqlite-substrate-fork) | honker-core fork into `alkstore-sqlite/src/substrate/`: port, deltas, provenance, test floor | wave 1 (workspace scaffold only) | implemented + reviewed (2026-10-08) |
|
||||
| [3](#wave-3--sqlite-engine) | SQLite engine: connection architecture, re-derived queue ops on contract v1, scheduler/outbox, tx seam, SQLite backlog column | waves 1 + 2 | decomposed (`tasks/`) |
|
||||
| [3](#wave-3--sqlite-engine) | SQLite engine: connection architecture, re-derived queue ops on contract v1, scheduler/outbox, tx seam, SQLite backlog column | waves 1 + 2 | implemented + reviewed (2026-10-08) |
|
||||
| 4 | Postgres engine: schema bootstrap, pool/open, listener/forwarder, all mechanisms, tx seam, pg backlog column | wave 1 | not yet decomposed |
|
||||
| 5 | Contract suite: the cross-engine equivalence properties (core-contract.md §Verification backlog), version-stamped per ADR-017 | waves 3 + 4 | not yet decomposed |
|
||||
| 6 | Release readiness: crate docs, deployment matrix final pass, README (written last, honestly), publish prep; mem-engine and fuzzing decisions | wave 5 | not yet decomposed |
|
||||
@@ -153,4 +153,10 @@ decomposes. Specific gates:
|
||||
`docs/reviews/2026-10-08-waves-1-2-general-review.md`) — M-1 fixed
|
||||
inline (sweep savepoint scope); M-2/N-2/N-4/N-6 resolved as ADR-023
|
||||
pre-decomposition; lint removal + suite adds folded into waves 3/5;
|
||||
N-1 recorded for wave 6.
|
||||
N-1 recorded for wave 6.
|
||||
- **Wave 3 review gate** (`review-wave-3`) — engine-vs-contract
|
||||
conformance code-read clean; two findings fixed inline (the
|
||||
`open_writer_connection` boundary move out of `substrate/mod.rs`;
|
||||
the writer-slot error-arm stranding in
|
||||
`with_writer`/`begin`/`commit`); two minor notes deferred to
|
||||
wave 5/6 (reader-pool close path, commit-error-arm coverage).
|
||||
+159
-2
@@ -1,7 +1,7 @@
|
||||
---
|
||||
id: review-wave-3
|
||||
name: Review gate — wave 3 (SQLite engine)
|
||||
status: pending
|
||||
status: completed
|
||||
depends_on: [sqlite-engine-integration]
|
||||
scope: moderate
|
||||
risk: low
|
||||
@@ -69,6 +69,163 @@ Check:
|
||||
|
||||
> To be filled by implementation agent
|
||||
|
||||
Two findings, both fixed inline (2026-10-08). Neither needed a new
|
||||
session or scope decision:
|
||||
|
||||
- **W-1 (boundary) — `alkstore::Error` referenced inside
|
||||
`src/substrate/mod.rs`**: `open_writer_connection` (the tx handle's
|
||||
`reopen` helper) lived in the substrate subtree and mapped into the
|
||||
contract taxonomy there — a violation of ADR-012 §2's
|
||||
contract-blind boundary ("no `alkstore` core types", as the
|
||||
substrate's own mod.rs doc text states; the wave-3 review
|
||||
checklist's grep for `alkstore::` under `src/substrate/` caught
|
||||
it). The machinery files (`schema.rs`/`ops.rs`/`queue_ops.rs`/
|
||||
`watcher.rs`) were always clean — the violation was a wave-3
|
||||
engine-layer helper that had been placed in the wrong file. Fixed:
|
||||
the helper moved to `alkstore-sqlite/src/seam.rs` (the substrate's
|
||||
one existing mapping site; identical body). No register entry —
|
||||
nothing lineage-side changed.
|
||||
- **W-2 (seam defect) — the writer slot stranded on three error
|
||||
arms**: the task checklist's "writer-slot lease released on every
|
||||
path" read found `Writer::acquire`d connections **dropped without
|
||||
`release()`/replenish** on (a) `seam.rs::with_writer`'s op-error
|
||||
arm, (b) `tx.rs::begin`'s failed `BEGIN IMMEDIATE` arm, and (c)
|
||||
`tx.rs::commit`'s failed `COMMIT` arm. `Writer` holds exactly one
|
||||
connection and has no self-replenish, so the first failed
|
||||
auto-commit op (e.g. a busy timeout expiry on any notify/enqueue/
|
||||
claim) would park *every* subsequent writer op forever on the
|
||||
condvar — a store-wide livelock from one transient storage failure.
|
||||
Fixed: (a) releases unconditionally on both arms (the failed op
|
||||
didn't corrupt the connection — a refused statement leaves it
|
||||
healthy; the release is safe); (b) releases on the begin-error arm
|
||||
(the no-transaction connection is reusable); (c) drops the failed
|
||||
connection (its state is unknowable — exactly the Drop path's
|
||||
posture) and replenishes the slot via the existing `reopen`
|
||||
closure, the same machinery the cancellation-path `ConnLease` and
|
||||
the failed-rollback Drop already use. Regression tests added:
|
||||
`failed_notify_releases_the_writer_slot` (notify_tests — trigger-
|
||||
induced storage failure, then the slot proves grantable),
|
||||
`failed_begin_releases_the_writer_slot` (tx_tests — SQLITE_BUSY
|
||||
begin at a shortened busy timeout, then `begin_tx` succeeds).
|
||||
|
||||
Verified clean otherwise — spot-check record:
|
||||
|
||||
- Contract conformance (code-read against the spec text, per the
|
||||
acceptance criterion): entry-point validation on every name-bearing
|
||||
method, auto-commit AND tx paths (tx.rs does all eleven
|
||||
`validate_shared_name`/`validate_local_name` at entry, before any
|
||||
round trip); `InvalidName` for empty/whitespace-only, `ReservedName`
|
||||
for prefixed, matching core's `validate_name`. Duration guards on
|
||||
**both** lock call sites (`try_lock`, `SqliteLockHandle::renew` —
|
||||
`ttl <= 0` → opaque `Database`, detail in the source chain, the
|
||||
substrate guard's message preserved there). Extent guards at
|
||||
trait-impl entry on `claim_batch` (`n <= 0` → empty `Vec`), both
|
||||
auto-commit stream reads, both tx stream reads, and
|
||||
`EventReceiver::read_since` — never reaching the substrate's
|
||||
`LIMIT ?`. Boundary args unguarded-and-total (`trim_to` passes
|
||||
through; a negative horizon deletes nothing idempotently).
|
||||
`Error::Closed` on already-consumed handles; `Codec` at every
|
||||
payload/decode seam; `recv()`'s `Err` remapped to `Database`-only
|
||||
(`database_only`, stream.rs). No-work vocabulary held:
|
||||
`Ok(None)`, `Ok(false)`, empty `Vec` where pinned.
|
||||
- Value shapes: `Job::from_row` field set matches ADR-021 §2's struct
|
||||
(claimed_at included); `StreamEvent::from_row` matches ADR-015 §3
|
||||
(the `topic` column carries `stream`, the one registered delta);
|
||||
dead-visible `get_job` carries `last_error`/`died_at`.
|
||||
- ADR-023 follow-through: all five rows present in
|
||||
`alkstore-contract-suite/src/properties.rs` with their contract
|
||||
stamps and green in `alkstore-sqlite/tests/contract_suite.rs` (10
|
||||
tests); extent guards/duration guards/`poll_interval`/plain-path
|
||||
open/`encode_payload` fallibility verified at their sites (see
|
||||
above); `encode_payload` consumed fallibly at all eight call sites,
|
||||
no silent fallback.
|
||||
- Backoff curve (queue.rs `backoff_delay_s`): equal-jitter over
|
||||
`[base·2^(a−1)/2, base·2^(a−1)]` integerized inclusive, capped at
|
||||
3600 s, attempt index = the row's post-claim count — matches
|
||||
ADR-010 §3's pinned form exactly.
|
||||
- Scheduler: tick fires + boundary advance in one `BEGIN IMMEDIATE`
|
||||
(writer serialization as the row lock); leadership loss returns
|
||||
`Err(LeadershipLost)` **before** ticking (top-of-loop renew, and
|
||||
exit-before-tick also holds at start — a refused acquire is
|
||||
`LeadershipLost`); the sleep renews the lease across itself; stop
|
||||
token releases the lock cleanly; 64-cap catch-up + skip-forward
|
||||
inherited (D-12).
|
||||
- Seam integrity code-read: drop = rollback issues `ROLLBACK` and
|
||||
releases the slot (failed rollback → connection dropped + slot
|
||||
replenished via `reopen` — the same posture commit now uses);
|
||||
panic-through-drop reaches the same path; the mid-op-future
|
||||
cancellation path (ConnLease Drop) replenishes; `spawn_blocking`
|
||||
used consistently — the only blocking-in-async exception is
|
||||
`EventReceiver::save_offset`, a sync *trait signature* taking the
|
||||
slot in-line, documented at the site (ADR-007's honest lease
|
||||
posture; not a violation).
|
||||
- Substrate boundary: no `alkstore::` reference remains under
|
||||
`src/substrate/` (post-W-1-fix grep; also `grep -r "allow("`
|
||||
clean — the wave-3 lint removals hold); PROVENANCE.md's
|
||||
D-29..D-36 entries cover the wave-3 deltas (verified against the
|
||||
code: URI flag dropped, `open_conn_bootstrapped`/D-30,
|
||||
`ack_batch` worker-less/D-31, the five cuts D-32..D-36 all
|
||||
genuinely absent); no unregistered divergence found in the
|
||||
lineage-adjacent reads.
|
||||
- Family standards: no panics/unwrap/expect in library code (only
|
||||
`#[cfg(test)]` modules); inline `//` comments exist only at the
|
||||
correctness-constraint sites the AGENTS.md exception names
|
||||
(backoff curve's range pin, the extent guard, the leader loop's
|
||||
one-owner-of-the-loss-decision note, `run_once`'s retry-string
|
||||
posture — all non-obvious behavior notes, acceptable);
|
||||
module-per-file holds.
|
||||
- One doc-vs-code nit, absorbed in the W-2 fix rather than a
|
||||
finding: `seam.rs`'s `with_writer` doc comment said "release on
|
||||
both arms" — the code did not do that; now it does.
|
||||
|
||||
Deferred (documented, not fixed — each is small but not
|
||||
inline-clean, or already owned by a later wave):
|
||||
|
||||
- **`Readers::acquire`'s `close()` path can drop a busy reader
|
||||
connection while held** (schema.rs — `close()` clears the pool
|
||||
but does not track outstanding acquisitions; a reader acquired
|
||||
before `close()` and released after it is discarded by
|
||||
`release()`'s closed check — correct, no leak; but a reader
|
||||
connection *still in use* across store-close keeps running on a
|
||||
closed file until its op completes — harmless in practice, the op
|
||||
finishes on an open connection that is then dropped). Noted for
|
||||
wave 6's close-path review; no consumer-visible defect.
|
||||
- **Coverage of the commit-error arm**: no test induces a failing
|
||||
`COMMIT` (needs SQLITE_FULL/IOERR-class injection; the code fix
|
||||
is pinned by code-read + the shared `reopen` machinery being
|
||||
exercised by the cancellation tests). Recorded here; natural
|
||||
home is wave 5's suite hardening.
|
||||
|
||||
## Summary
|
||||
|
||||
> To be filled on completion
|
||||
> What landed, verified how:
|
||||
|
||||
**Review gate passed: 2 findings, both fixed inline; all gates green.**
|
||||
|
||||
- **Fixed — substrate boundary violation**: `open_writer_connection`
|
||||
moved from `src/substrate/mod.rs` (which referenced
|
||||
`alkstore::Error`, breaking ADR-012 §2's contract-blind boundary) to
|
||||
`src/seam.rs`, the engine's one mapping site. Substrate grep for
|
||||
`alkstore::` now clean.
|
||||
- **Fixed — writer-slot stranding on error arms** (the
|
||||
store-wide-livelock defect): `with_writer` now releases the slot on
|
||||
the op-error arm; `begin` releases on a failed `BEGIN IMMEDIATE`;
|
||||
`commit` replenishes via the existing `reopen` closure on a failed
|
||||
`COMMIT`. Two regression tests pin the property end-to-end
|
||||
(`failed_notify_releases_the_writer_slot`,
|
||||
`failed_begin_releases_the_writer_slot`).
|
||||
- **Verified clean**: contract conformance (entry-point validation,
|
||||
error arms, value shapes, no-work vocabulary — full code-read
|
||||
record in Notes), ADR-023's five follow-through items + backlog
|
||||
rows (10 tests, stamped, green), backoff curve matches ADR-010 §3,
|
||||
scheduler tick atomicity + exit-before-tick-on-loss, seam
|
||||
integrity (drop = rollback incl. panic path, cancellation
|
||||
replenish), substrate boundary + PROVENANCE D-29..D-36 accurate,
|
||||
family standards.
|
||||
- **Gates**: `cargo build`, `cargo test` (25 core + 188 sqlite + 10
|
||||
suite + 3 harness), `cargo clippy --all-targets -- -D warnings`,
|
||||
`cargo fmt --check` — all green after the fixes.
|
||||
- **Recorded for later waves**: the reader-pool close-path note and
|
||||
the commit-error-arm coverage gap (Notes; wave 5/6 natural homes).
|
||||
|
||||
Wave 4/5 decomposition may proceed.
|
||||
Reference in new issue
Block a user