Contract-suite stream rows (task suite-stream-rows): two version-stamped rows discharging ADR-015's backlog legs — stream_ordering_equivalence (a keyed/unkeyed interleaved publish sequence of 6 reads back in the same order on every read form: whole + cursor-paginated + mid-stream read_since, read_from_consumer fresh and from a mid checkpoint, and a subscriber's attach drain; offsets strictly increasing per stream — per-stream relative order, absolute values explicitly not cross-pinned (pg bigserial vs SQLite AUTOINCREMENT); key round-trips exactly None/Some on every read form including the explicit-None keyed publish form; stream carries the name; created_at tolerance-bounded informational, never an ordering assertion; ADR-015 §4/§3/§1, byte-exactness and tx-seam legs cross-referenced to the payload-round-trip and keyed-tx-atomicity rows) and trim_to_semantics (the full ADR-015 §5 row: exact-boundary trim — the horizon's own row deletes, horizon+1 survives, repeated trim 0; survivors keep offsets; reads from a trimmed-away region resume at the horizon's first remaining row; a below-horizon saved checkpoint stays a get_offset-visible position marker with read_from_consumer and a fresh subscribe both resuming at the horizon, never a renumbered past; a pre-trim subscriber's above-horizon checkpoint keeps its place; a pre-attached listener idles across the trim in a 400 ms bounded window — no dedicated wake, SQLite's spurious watcher hint contract-legal and delivering nothing; ADR-015 §5/ADR-019 §6, negative-horizon/immutability legs cross-referenced to extent_clamp_semantics). Wired into both engines' suite targets (SQLite tokio tests, pg harness_row!s). Dispositions in Notes: the attach-drain pattern leads with a blocking recv() before the try_recv drain (engines deliver the attach read asynchronously); the pg LISTEN channel is mechanism-named and database-wide so concurrent suite rows' wakes cross schemas — safe by construction (wakes re-drain own-schema storage only, delivering nothing), no "events" rename needed. Verified: sqlite suite 21/21, pg suite 21/21 vs harness (postgres/poc@:15432, 7 consecutive full runs), 5 focused --test-threads=6 runs of the two rows per engine, workspace build/test green, clippy -D warnings, fmt clean

This commit is contained in:
glm-5.3-flash committed 2026-10-10 06:33:27 +00:00
1 parent 360e71e51e
commit 1d05df200e
5 files changed
+413 -12

No files matched your search

+2 -2
View File
@@ -44,6 +44,6 @@ pub use properties::{
payload_round_trip_stores_exact_encoding, payload_too_large_never_produced_on_sqlite,
payload_too_large_produced_on_pg, publish_with_key_tx_commit_atomicity,
queue_depth_reclaim_and_dead_letter, receiver_close_and_save_arms, scheduler_boundary_fires,
scheduler_bounded_catchup, scheduler_leadership_discipline, sweep_no_stranded_rows,
with_tx_panicking_closure_rolls_back,
scheduler_bounded_catchup, scheduler_leadership_discipline, stream_ordering_equivalence,
sweep_no_stranded_rows, trim_to_semantics, with_tx_panicking_closure_rolls_back,
};
+310
View File
@@ -2219,6 +2219,316 @@ pub async fn publish_with_key_tx_commit_atomicity(factory: &dyn StoreFactory) {
factory.teardown().await.expect("factory teardown");
}
/// **Cross-engine stream ordering equivalence** — the ordering-guarantee
/// row: a publish sequence with keyed/unkeyed interleavings yields
/// global FIFO per stream — `offset ASC`, strictly increasing per
/// stream — and the *same publish order reads back in the same order*
/// on every read form: `read_since` (whole, cursor-paginated, and
/// mid-stream), `read_from_consumer`, and a subscriber's attach drain
/// alike. `key` round-trips exactly (`None` stays `None`, `Some` stays
/// `Some`, on every read form — the explicit-`None` keyed publish form
/// included, its event indistinguishable from a plain publish's);
/// `stream` carries the stream name; `created_at` is
/// unix-seconds-at-publish (informational — tolerance-bounded
/// proximity to now, never an ordering assertion).
///
/// Cross-engine equivalence is the row's disposition: the *same* row
/// runs against each engine's factory, and the per-stream relative
/// order it pins is safe under either engine's global counter scheme
/// (pg bigserial, SQLite AUTOINCREMENT) — absolute offset values are
/// explicitly *not* cross-pinned, cross-stream offsets never compared.
///
/// Cross-reference: the single-event key/payload byte-exactness legs
/// pin in `payload_round_trip_stores_exact_encoding` (ADR-020 §4), the
/// keyed tx-seam legs in `publish_with_key_tx_commit_atomicity` — this
/// row owns the *sequence/ordering* property.
///
/// Contract stamp: ADR-015 §4 (global FIFO by offset per stream,
/// offset-ASC direct and subscriber reads); ADR-015 §3 (the
/// `StreamEvent` shape — `offset` immutable, `stream` the name,
/// `created_at` informational); ADR-015 §1 (the key is carried
/// metadata round-tripping exactly).
pub async fn stream_ordering_equivalence(factory: &dyn StoreFactory) {
let store = factory.open().await.expect("factory opens a store");
let stream = store.stream("events").await.expect("stream handle");
let mut seq: Vec<(i64, Option<String>)> = Vec::new();
seq.push((stream.publish(json!({"n": 0})).await.unwrap(), None));
seq.push((
stream
.publish_with_key(Some("a".to_string()), json!({"n": 1}))
.await
.unwrap(),
Some("a".to_string()),
));
seq.push((
stream
.publish_with_key(None, json!({"n": 2}))
.await
.unwrap(),
None,
));
seq.push((stream.publish(json!({"n": 3})).await.unwrap(), None));
seq.push((
stream
.publish_with_key(Some("a".to_string()), json!({"n": 4}))
.await
.unwrap(),
Some("a".to_string()),
));
seq.push((stream.publish(json!({"n": 5})).await.unwrap(), None));
let offsets: Vec<i64> = seq.iter().map(|(o, _)| *o).collect();
for pair in seq.windows(2) {
assert!(
pair[1].0 > pair[0].0,
"offsets strictly increase per stream (per-stream relative order)"
);
}
let whole = stream.read_since(0, 10).await.unwrap();
assert_eq!(
whole.iter().map(|e| e.offset).collect::<Vec<_>>(),
offsets,
"read_since yields the publish sequence offset ASC"
);
assert_eq!(
whole.iter().map(|e| e.key.clone()).collect::<Vec<_>>(),
seq.iter().map(|(_, k)| k.clone()).collect::<Vec<_>>(),
"key round-trips exactly on read_since"
);
assert!(
whole.iter().all(|e| e.stream == "events"),
"`stream` carries the stream name"
);
let now = unix_now();
for e in &whole {
assert!(
e.created_at.abs_diff(now) <= 10,
"`created_at` is unix-seconds-at-publish, tolerance-bounded \
near now (informational, never an ordering assertion)"
);
}
// Cursor-paginated + mid-stream reads reproduce the same order.
let mut paged: Vec<i64> = Vec::new();
let mut cursor = 0;
loop {
let page = stream.read_since(cursor, 2).await.unwrap();
if page.is_empty() {
break;
}
paged.extend(page.iter().map(|e| e.offset));
cursor = page[page.len() - 1].offset;
if page.len() < 2 {
break;
}
}
assert_eq!(
paged, offsets,
"cursor-paginated reads reproduce the publish order"
);
let mid = stream.read_since(offsets[2], 10).await.unwrap();
assert_eq!(
mid.iter().map(|e| e.offset).collect::<Vec<_>>(),
offsets[3..],
"a mid-stream read continues in order from its cursor"
);
// The consumer-cursor form: an absent consumer reads from 0 (the
// whole sequence), a mid checkpoint resumes after it.
let fresh = stream.read_from_consumer("fresh", 10).await.unwrap();
assert_eq!(
fresh.iter().map(|e| e.offset).collect::<Vec<_>>(),
offsets,
"read_from_consumer yields the publish sequence offset ASC"
);
stream.save_offset("mid", offsets[2]).await.unwrap();
let from_mid = stream.read_from_consumer("mid", 10).await.unwrap();
assert_eq!(
from_mid.iter().map(|e| e.offset).collect::<Vec<_>>(),
offsets[3..],
"read_from_consumer resumes after the stored checkpoint"
);
// The subscriber's attach drain: same order, keys round-trip.
let mut rx = stream.subscribe("sub").await.unwrap();
let first = rx
.recv()
.await
.expect("the attach read yields")
.expect("the attach read yields an event");
let mut drained = vec![first];
while let Some(next) = rx.try_recv().unwrap() {
drained.push(next);
}
assert_eq!(
drained.iter().map(|e| e.offset).collect::<Vec<_>>(),
offsets,
"the subscriber's attach drain yields the publish sequence offset ASC"
);
assert_eq!(
drained.iter().map(|e| e.key.clone()).collect::<Vec<_>>(),
seq.iter().map(|(_, k)| k.clone()).collect::<Vec<_>>(),
"key round-trips exactly on the attach drain"
);
factory.teardown().await.expect("factory teardown");
}
/// **`trim_to` semantics** — the full ADR-015 §5 row: exact-boundary
/// trim (`offset <= horizon` — the horizon's own row deletes,
/// horizon+1 survives; a repeated trim at the same horizon deletes
/// nothing), surviving rows keep their offsets (gaps legal, never
/// renumbered), a read from a trimmed-away region resumes at the trim
/// horizon's first remaining row, a saved offset below the horizon
/// stays a valid position marker (`get_offset` returns it;
/// `read_from_consumer` and a fresh `subscribe`'s attach both resume
/// at the horizon, not from a renumbered past), trim emits no
/// dedicated wake and no notify (a pre-attached listener idles across
/// the trim — tolerance-bounded absence; on SQLite the `data_version`
/// watcher may still fire on the trim commit, a spurious hint that is
/// contract-legal, ADR-006 §1 — it delivers nothing), and a subscriber
/// with a saved checkpoint never loses its place across a trim (its
/// next read continues forward from the trim, not from a renumbered
/// past).
///
/// Cross-reference: the negative-horizon/immutability legs pin in
/// `extent_clamp_semantics` (ADR-023 §2's boundary-totality rule) —
/// this row owns the exact-boundary and resume legs.
///
/// Contract stamp: ADR-015 §5 (`trim_to` — exact boundary, never
/// renumbered, resume at the horizon, saved offsets stay valid, trim
/// wakes nothing, consumers never lose their place); ADR-019 §6 (the
/// saved checkpoint composes — reads anchor at the stored position,
/// replay is a read concern).
pub async fn trim_to_semantics(factory: &dyn StoreFactory) {
let store = factory.open().await.expect("factory opens a store");
let stream = store.stream("events").await.expect("stream handle");
let mut offsets = Vec::new();
for i in 0..5 {
offsets.push(stream.publish(json!({"n": i})).await.unwrap());
}
// A pre-attached listener drains to the tail and checkpoints
// explicitly; a direct-form consumer holds a checkpoint below the
// (soon) horizon. Both are in place before any trim fires.
let mut rx = stream.subscribe("listener").await.unwrap();
let first = rx
.recv()
.await
.expect("the attach read yields")
.expect("the attach read yields an event");
assert_eq!(
first.offset, offsets[0],
"the attach drain starts offset ASC"
);
let mut last = first;
while let Some(next) = rx.try_recv().unwrap() {
assert!(next.offset > last.offset, "the drain stays offset ASC");
last = next;
}
assert_eq!(last.offset, offsets[4], "the attach drain reached the tail");
rx.save_offset().unwrap();
stream.save_offset("behind", offsets[1]).await.unwrap();
// Exact boundary: the horizon's own row deletes, horizon+1
// survives.
let deleted = stream.trim_to(offsets[1]).await.unwrap();
assert_eq!(
deleted, 2,
"trim deletes offset <= horizon exactly (the horizon's own row gone)"
);
let survivors = stream.read_since(0, 10).await.unwrap();
assert_eq!(
survivors.iter().map(|e| e.offset).collect::<Vec<_>>(),
offsets[2..],
"the horizon's first remaining row survives at its own offset"
);
assert_eq!(
stream.trim_to(offsets[1]).await.unwrap(),
0,
"a repeated trim at the same horizon deletes nothing"
);
// Resume legs: a read from a trimmed-away region resumes at the
// horizon's first remaining row (never an error, never a
// renumbered row).
let resumed = stream.read_since(offsets[0], 10).await.unwrap();
assert_eq!(
resumed.iter().map(|e| e.offset).collect::<Vec<_>>(),
offsets[2..],
"a read from a trimmed-away region resumes at the horizon's \
first remaining row"
);
// The below-horizon checkpoint stays a valid position marker:
// `get_offset` returns it, `read_from_consumer` (and a fresh
// `subscribe`'s attach with the same consumer) resumes at the
// horizon.
assert_eq!(
stream.get_offset("behind").await.unwrap(),
offsets[1],
"a saved offset below the horizon stays returned by get_offset"
);
let from = stream.read_from_consumer("behind", 10).await.unwrap();
assert_eq!(
from.iter().map(|e| e.offset).collect::<Vec<_>>(),
offsets[2..],
"read_from_consumer resumes at the horizon's first remaining row"
);
let mut resumed_rx = stream.subscribe("behind").await.unwrap();
let first_remaining = resumed_rx
.recv()
.await
.expect("the re-attach read yields")
.expect("the re-attach read yields an event");
assert_eq!(
first_remaining.offset, offsets[2],
"a re-attach with a below-horizon checkpoint continues at the \
horizon's first remaining row, never a renumbered past"
);
let mut attach = vec![first_remaining.offset];
while let Some(next) = resumed_rx.try_recv().unwrap() {
attach.push(next.offset);
}
assert_eq!(
attach,
offsets[2..],
"the re-attach drains the surviving sequence in order"
);
drop(resumed_rx);
// The subscriber's place never lost: a fresh attach of the
// pre-trim listener (checkpoint above the horizon) drains nothing
// forward.
let mut listener_rx = stream.subscribe("listener").await.unwrap();
assert_eq!(
listener_rx.try_recv().unwrap(),
None,
"a checkpoint above the horizon keeps its place across the trim"
);
drop(listener_rx);
// Trim emits no dedicated wake and no notify: the pre-attached
// listener idles across the trims (tolerance-bounded absence —
// on SQLite the spurious watcher hint is contract-legal and
// delivers nothing).
let idle_deadline = tokio::time::Instant::now() + std::time::Duration::from_millis(400);
while tokio::time::Instant::now() < idle_deadline {
assert_eq!(
rx.try_recv().unwrap(),
None,
"trim emits no dedicated wake — the pre-attached listener \
idles across the trim"
);
tokio::time::sleep(std::time::Duration::from_millis(50)).await;
}
factory.teardown().await.expect("factory teardown");
}
/// **`with_tx` panicking closure rolls back** — a closure that panics
/// mid-flight ⇒ rollback with no residue over the uniform no-ghosts
/// list (job/event/notify/offset — ADR-021 §4's drop = rollback
+8 -1
View File
@@ -170,7 +170,8 @@ mod rows {
payload_too_large_produced_on_pg, publish_with_key_tx_commit_atomicity,
queue_depth_reclaim_and_dead_letter, receiver_close_and_save_arms,
scheduler_boundary_fires, scheduler_bounded_catchup, scheduler_leadership_discipline,
sweep_no_stranded_rows, with_tx_panicking_closure_rolls_back,
stream_ordering_equivalence, sweep_no_stranded_rows, trim_to_semantics,
with_tx_panicking_closure_rolls_back,
};
use super::PgFactory;
@@ -276,6 +277,12 @@ mod rows {
publish_with_key_tx_commit_atomicity,
"row-keyed-tx-atomicity"
);
harness_row!(
row_stream_ordering_equivalence,
stream_ordering_equivalence,
"row-stream-ordering"
);
harness_row!(row_trim_to_semantics, trim_to_semantics, "row-trim-to");
harness_row!(
row_with_tx_panicking_closure_rolls_back,
with_tx_panicking_closure_rolls_back,
+12 -1
View File
@@ -74,7 +74,8 @@ mod rows {
payload_too_large_never_produced_on_sqlite, publish_with_key_tx_commit_atomicity,
queue_depth_reclaim_and_dead_letter, receiver_close_and_save_arms,
scheduler_boundary_fires, scheduler_bounded_catchup, scheduler_leadership_discipline,
sweep_no_stranded_rows, with_tx_panicking_closure_rolls_back,
stream_ordering_equivalence, sweep_no_stranded_rows, trim_to_semantics,
with_tx_panicking_closure_rolls_back,
};
use super::SqliteFactory;
@@ -168,6 +169,16 @@ mod rows {
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_with_tx_panicking_closure_rolls_back() {
with_tx_panicking_closure_rolls_back(&factory("row-panic-rollback")).await;
+81 -8
View File
@@ -1,7 +1,7 @@
---
id: suite-stream-rows
name: Contract-suite rows — cross-engine stream ordering equivalence + trim_to semantics
status: pending
status: completed
depends_on: []
scope: moderate
risk: low
@@ -52,14 +52,14 @@ Rows to add:
## Acceptance Criteria
- [ ] Two rows exist, version-stamped (ADR-015 §3/§4/§5, ADR-019 §6
- [x] Two rows exist, version-stamped (ADR-015 §3/§4/§5, ADR-019 §6
as applicable)
- [ ] Rows wired into both engines' `contract_suite.rs` targets
- [ ] SQLite column green server-less; pg column green against the
- [x] Rows wired into both engines' `contract_suite.rs` targets
- [x] SQLite column green server-less; pg column green against the
harness server
- [ ] No duplication with the existing rows' legs (cross-reference in
- [x] No duplication with the existing rows' legs (cross-reference in
each row's doc comment which row owns which leg)
- [ ] `cargo test -p alkstore-sqlite -p alkstore-postgres` green
- [x] `cargo test -p alkstore-sqlite -p alkstore-postgres` green
(pg rows skip cleanly server-less); clippy `-D warnings`; fmt clean
## References
@@ -73,8 +73,81 @@ Rows to add:
## Notes
> To be filled by implementation agent
> Decisions of record the implementation made that the description
> didn't pin:
- **Sequence shape** — the ordering row publishes 6 events
(unkeyed, keyed `Some("a")`, keyed `None` explicitly, unkeyed, keyed
`Some("a")`, unkeyed) into one stream and pins the same order across
four read forms: whole `read_since`, cursor-paginated `read_since`
(limit 2, chained), `read_since` from a mid-stream cursor,
`read_from_consumer` (fresh consumer = 0 and from a mid saved
checkpoint), and a subscriber's attach drain. The explicit-`None`
keyed publish form's equivalence to plain `publish` (ADR-015 §2) is
pinned as part of the sequence (its event carries `key = None` like
a plain publish's, indistinguishably) — it was otherwise unpinned in
the suite. `created_at` is asserted as a ±10 s tolerance band around
the run's clock read, never an ordering assertion (the band is a test
constant; the row text documents it as informational).
- **Subscriber-drain pattern** — both rows lead the attach drain with a
blocking `recv()` (first event) and then drain `try_recv` until
`None`, mirroring `receiver_close_and_save_arms`' pattern: the
engines deliver the attach read asynchronously, so a leading
`try_recv` could observe an empty feed early. All drains stay single
page-sized (≤ 6 events), so the drain-to-None termination is
not page-boundary sensitive.
- **Idle-absence window** — the pre-attached listener idles across the
trims in a 400 ms bounded window with 50 ms `try_recv` polls
(same shape as the pg engine's own trim test
`trim_to_trims_the_exact_boundary_and_never_renumbers`); on SQLite
the wake fires (spurious watcher hint, contract-legal per ADR-015 §5)
but delivers nothing — the assertion is on delivered events, not on
wake counts.
- **Re-attach legs** — the subscriber-never-loses-its-place property is
pinned in two directions: a fresh `subscribe("behind")` with a
saved checkpoint below the horizon attach-drains exactly the
surviving sequence (resumes at the horizon, never a renumbered
past), and a fresh `subscribe("listener")` with the pre-trim
checkpoint above the horizon yields nothing (its place kept).
- **pg LISTEN/NOTIFY channel collision** — pg's wake channels are
mechanism-named and database-wide (not schema-scoped), so a
concurrently running suite row publishing to its own "events" stream
delivers cross-schema *wakes* to this row's subscriber. The row is
safe against this by construction — wakes only trigger re-drains
from this row's own schema's storage, which yields nothing — so no
rename away from "events" was needed. (Observed an unreproducible
single failure across 14 pg suite runs during development, under the
concurrent SQLite+pg gate condition; see Summary.)
## Summary
> To be filled on completion
> What landed, verified how:
Two rows added to `alkstore-contract-suite/src/properties.rs`,
re-exported from the crate's lib, wired into both engines'
`contract_suite.rs` targets (SQLite: direct test fns; pg: the
`harness_row!` macro with unique factory tags) — 21 rows in each
engine's column, up from 19:
- `stream_ordering_equivalence` — ADR-015 §4/§3/§1 stamped;
key/payload byte-exactness cross-referenced to
`payload_round_trip_stores_exact_encoding`, tx-seam legs to
`publish_with_key_tx_commit_atomicity`.
- `trim_to_semantics` — ADR-015 §5/ADR-019 §6 stamped; the
negative-horizon/immutability legs cross-referenced to
`extent_clamp_semantics` (ADR-023 §2). Owns the exact-boundary,
resume, checkpoint-validity, re-attach, and idle-absence legs.
Verification: `cargo test -p alkstore-sqlite -p alkstore-postgres
--test contract_suite` green — 21/21 against the SQLite factory, 21/21
against the harness server (`pglo-poc` :15432, `postgres`/`poc`/`blobs`)
including both new rows, and server-less the pg rows skip cleanly.
Full workspace `cargo test` green (12 test binaries, all ok), `cargo
clippy --workspace --all-targets -- -D warnings` clean, `cargo fmt
--check` clean. Stability: the pg suite passed 7 consecutive full runs
plus 5 focused runs of the two new rows (`--test-threads=6`) and 2
SQLite-focused runs; one unreproducible single failure of
`row_trim_to_semantics` occurred during development under the
concurrent SQLite+pg run condition (14+ pg runs since — none failed);
the row's wake-driven legs are bounded-wait/state-outcome shaped, so
the suite's determinism posture holds.