Files
glm-5.3-flash 9301bb926b core: verify and pin the receiver's pre-yield offset() value (review 004 item 12, core-pin-offset-initial-value) — the three engines verified equal live, so the doc-pin branch applies: offset() reads 0 before any yield on every receiver — fresh consumer, stored positive checkpoint, stored negative seed (ADR-023 §2's totality arm, the one a checkpoint-seeded realization would surface on) — the no-yield sentinel, never the stored checkpoint: the checkpoint gates the feed's attach-read start while the receiver's offset tracks the last-yielded event only, composing with the pinned pre-yield save_offset() no-op; the attach read's prefetches sit unobserved and the first recv/try_recv observation is the first move (then the post-yield value, receiver_close_and_save_arms' pin). Verified by code read first (all three offset() bodies are self.position, all three constructors seed position: 0 unconditionally — mem stream.rs:353, sqlite :399, pg :620) then live: new suite row receiver_offset_pre_yield_is_zero (alkstore-contract-suite, stamped ADR-008 §8 + the core-contract.md pin, four arms incl. the first-observation composition leg) green on all three columns, pg against the harness server (parked pglo-poc :15432); pinned in three places: core-contract.md streams/subscribe row (dated 2026-10-11 annotation) + a verification-backlog bullet recording the verified pin, the core trait doc on offset() (core's doc comments are contract text), and the suite row itself — ADR-008 §8's shape comment deliberately left untouched (doc-only verification takes no ADR amendment, noted in the task); no engine code changed, no divergence to raise; gates green: workspace build/tests server-less, full pg column vs harness (135+26+9), clippy -D warnings, fmt --check, doc gate zero warnings, fuzz corpus replay 6/6, wasm32 mem check, taskgraph validate (96); task -> completed
2026-10-11 14:34:15 +00:00

240 lines
8.7 KiB
Rust

//! The suite-driving test target: the SQLite engine's backlog column
//! (ADR-022's option (a) — the rows live in `alkstore-contract-suite`,
//! this target is the executor that runs them against the SQLite
//! factory; wave 5 runs the same rows against pg's).
//!
//! The factory: a fresh temp file under a fresh instance directory per
//! `open` (every call gets its own never-before-used path — the
//! isolation guarantee), teardown deleting the directory (idempotent —
//! already-deleted is already-torn-down, `Ok`). Close-arm rows drop
//! the store handle (the dispose carrier) — deleting files under a
//! live store is not a close mechanism.
use std::path::PathBuf;
use std::sync::atomic::{AtomicU64, Ordering};
use alkstore::{BoxedFuture, Error, Result, Store};
use alkstore_contract_suite::StoreFactory;
use alkstore_sqlite::open;
/// The SQLite suite factory: fresh temp-file store per `open`,
/// idempotent directory-delete teardown. `Send + Sync` — instance
/// state is a path prefix and a counter.
struct SqliteFactory {
dir: PathBuf,
instance: AtomicU64,
}
impl SqliteFactory {
fn new(tag: &str) -> Self {
let dir = std::env::temp_dir().join(format!(
"alkstore-suite-{tag}-{}-{:?}",
std::process::id(),
std::thread::current().id()
));
let _ = std::fs::remove_dir_all(&dir);
std::fs::create_dir_all(&dir).unwrap();
SqliteFactory {
dir,
instance: AtomicU64::new(0),
}
}
}
impl StoreFactory for SqliteFactory {
fn open(&self) -> BoxedFuture<'_, Result<Box<dyn Store>>> {
let n = self.instance.fetch_add(1, Ordering::SeqCst);
let path = self
.dir
.join(format!("store-{n}.db"))
.to_str()
.unwrap_or_default()
.to_string();
Box::pin(async move { open(&path, Default::default()) })
}
fn teardown(&self) -> BoxedFuture<'_, Result<()>> {
let dir = self.dir.clone();
Box::pin(async move {
match std::fs::remove_dir_all(&dir) {
Ok(()) => Ok(()),
Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(()),
Err(e) => Err(Error::database(e)),
}
})
}
}
mod rows {
use alkstore_contract_suite::properties::{
backoff_curve_equivalence, concurrent_try_lock_loser_is_a_value,
drop_rollback_leaves_no_ghosts, duration_refusal_on_non_positive_ttl,
enqueue_opts_resolution, extent_clamp_semantics, in_tx_reads_see_own_writes,
job_handle_validity_predicate, lock_ttl_expiry_and_reacquisition,
name_validation_rejects_empty_and_reserved, outbox_enqueue_tx_commit_atomicity,
payload_round_trip_stores_exact_encoding, payload_too_large_never_produced_on_sqlite,
publish_with_key_tx_commit_atomicity, queue_depth_reclaim_and_dead_letter,
receiver_close_and_save_arms, receiver_offset_pre_yield_is_zero, scheduler_boundary_fires,
scheduler_bounded_catchup, scheduler_leadership_discipline, stream_ordering_equivalence,
sweep_no_stranded_rows, trim_to_semantics, wake_receiver_shapes,
with_tx_panicking_closure_rolls_back,
};
use super::SqliteFactory;
fn factory(tag: &str) -> SqliteFactory {
SqliteFactory::new(tag)
}
#[tokio::test(flavor = "multi_thread")]
async fn row_name_validation_rejects_empty_and_reserved() {
name_validation_rejects_empty_and_reserved(&factory("row-name-validation")).await;
}
#[tokio::test(flavor = "multi_thread")]
async fn row_extent_clamp_semantics() {
extent_clamp_semantics(&factory("row-extent-clamp")).await;
}
#[tokio::test(flavor = "multi_thread")]
async fn row_duration_refusal_on_non_positive_ttl() {
duration_refusal_on_non_positive_ttl(&factory("row-duration-refusal")).await;
}
#[tokio::test(flavor = "multi_thread")]
async fn row_lock_ttl_expiry_and_reacquisition() {
lock_ttl_expiry_and_reacquisition(&factory("row-lock-ttl-expiry")).await;
}
#[tokio::test(flavor = "multi_thread")]
async fn row_concurrent_try_lock_loser_is_a_value() {
concurrent_try_lock_loser_is_a_value(&factory("row-lock-loser-value")).await;
}
#[tokio::test(flavor = "multi_thread")]
async fn row_payload_round_trip_stores_exact_encoding() {
payload_round_trip_stores_exact_encoding(&factory("row-payload-round-trip")).await;
}
#[tokio::test(flavor = "multi_thread")]
async fn row_payload_too_large_never_produced_on_sqlite() {
payload_too_large_never_produced_on_sqlite(&factory("row-no-too-large")).await;
}
#[tokio::test(flavor = "multi_thread")]
async fn row_drop_rollback_leaves_no_ghosts() {
drop_rollback_leaves_no_ghosts(&factory("row-drop-rollback")).await;
}
#[tokio::test(flavor = "multi_thread")]
async fn row_in_tx_reads_see_own_writes() {
in_tx_reads_see_own_writes(&factory("row-in-tx-ryow")).await;
}
#[tokio::test(flavor = "multi_thread")]
async fn row_enqueue_opts_resolution() {
enqueue_opts_resolution(&factory("row-opts-resolution")).await;
}
#[tokio::test(flavor = "multi_thread")]
async fn row_backoff_curve_equivalence() {
backoff_curve_equivalence(&factory("row-backoff-curve")).await;
}
#[tokio::test(flavor = "multi_thread")]
async fn row_receiver_close_and_save_arms() {
receiver_close_and_save_arms(&factory("row-receiver-arms")).await;
}
#[tokio::test(flavor = "multi_thread")]
async fn row_receiver_offset_pre_yield_is_zero() {
receiver_offset_pre_yield_is_zero(&factory("row-receiver-offset-preyield")).await;
}
#[tokio::test(flavor = "multi_thread")]
async fn row_job_handle_validity_predicate() {
job_handle_validity_predicate(&factory("row-handle-predicate")).await;
}
#[tokio::test(flavor = "multi_thread")]
async fn row_queue_depth_reclaim_and_dead_letter() {
queue_depth_reclaim_and_dead_letter(&factory("row-queue-depth")).await;
}
#[tokio::test(flavor = "multi_thread")]
async fn row_sweep_no_stranded_rows() {
sweep_no_stranded_rows(&factory("row-sweep-stranded")).await;
}
#[tokio::test(flavor = "multi_thread")]
async fn row_scheduler_boundary_fires() {
scheduler_boundary_fires(&factory("row-scheduler-boundary")).await;
}
#[tokio::test(flavor = "multi_thread")]
async fn row_scheduler_bounded_catchup() {
scheduler_bounded_catchup(&factory("row-scheduler-catchup")).await;
}
#[tokio::test(flavor = "multi_thread")]
async fn row_scheduler_leadership_discipline() {
scheduler_leadership_discipline(&factory("row-scheduler-leadership")).await;
}
#[tokio::test(flavor = "multi_thread")]
async fn row_outbox_enqueue_tx_commit_atomicity() {
outbox_enqueue_tx_commit_atomicity(&factory("row-outbox-tx-atomicity")).await;
}
#[tokio::test(flavor = "multi_thread")]
async fn row_publish_with_key_tx_commit_atomicity() {
publish_with_key_tx_commit_atomicity(&factory("row-keyed-tx-atomicity")).await;
}
#[tokio::test(flavor = "multi_thread")]
async fn row_stream_ordering_equivalence() {
stream_ordering_equivalence(&factory("row-stream-ordering")).await;
}
#[tokio::test(flavor = "multi_thread")]
async fn row_trim_to_semantics() {
trim_to_semantics(&factory("row-trim-to")).await;
}
#[tokio::test(flavor = "multi_thread")]
async fn row_wake_receiver_shapes() {
wake_receiver_shapes(&factory("row-wake-shapes")).await;
}
#[tokio::test(flavor = "multi_thread")]
async fn row_with_tx_panicking_closure_rolls_back() {
with_tx_panicking_closure_rolls_back(&factory("row-panic-rollback")).await;
}
}
mod factory_shape {
use super::SqliteFactory;
use alkstore_contract_suite::StoreFactory;
/// ADR-022's factory contract: every open yield is isolated (no two
/// opens share rows) and teardown is idempotent per instance; a
/// failing assertion early-exit (the panic path) leaves teardown
/// safe against the abandoned store (the documented open/close/
/// delete order is safe at any point — the row pins the
/// teardown-side half).
#[tokio::test(flavor = "multi_thread")]
async fn opens_are_isolated_and_teardown_is_idempotent() {
let factory = SqliteFactory::new("factory-shape");
let a = factory.open().await.unwrap();
let b = factory.open().await.unwrap();
factory.teardown().await.unwrap();
factory.teardown().await.unwrap();
drop(a);
drop(b);
assert!(!factory.dir.exists(), "teardown removed the backing dir");
}
}