SQLite engine integration: lint removal, contract-suite adoption, backlog column (task sqlite-engine-integration)
- Remove the wave-2 lint suppressions from substrate/mod.rs; the six genuinely dead surfaces the removal exposed are cut, not suppressed, and registered D-32..D-36 in PROVENANCE.md (arg_opt_i64, ops::now_unix, queue_next_claim_at, Writer::try_acquire, UpdateWatcher::spawn, SharedUpdateWatcher::new); test-observation items (subscriber_count, the poll-interval default re-export) are honestly #[cfg(test)]-gated - Contract suite: eight new version-stamped backlog rows (extent-clamp + boundary totality, duration-refusal, encode_payload round-trip, PayloadTooLarge-never-produced SQLite arm, drop=rollback no-ghosts, in-tx read-your-own-writes, enqueue-opts resolution, receiver close/save arms) - Fix the exemplar row's real-engine sequencing defect: the held tx handle across the with_tx leg deadlocked any single-writer factory (mock-invisible; ADR-007's parking is the pinned behavior) - SQLite factory: SqliteFactory in the new tests/contract_suite.rs target; all nine rows green against it; the factory contract (isolation + idempotent teardown) pinned - Engine lib docs: the single-host and writer-parking posture statements surfaced under # Posture - Gates: build/test/clippy -D warnings/fmt green; coverage 93.6% lines, misses confined to error arms
This commit is contained in:
1 parent
8502a51af7
commit
a82c543b40
12 files changed
+1067
-245
No files matched your search
@@ -37,4 +37,9 @@ pub mod properties;
|
||||
pub mod version_stamp;
|
||||
|
||||
pub use factory::StoreFactory;
|
||||
pub use properties::name_validation_rejects_empty_and_reserved;
|
||||
pub use properties::{
|
||||
duration_refusal_on_non_positive_ttl, enqueue_opts_resolution, extent_clamp_semantics,
|
||||
in_tx_reads_see_own_writes, name_validation_rejects_empty_and_reserved,
|
||||
payload_round_trip_stores_exact_encoding, payload_too_large_never_produced_on_sqlite,
|
||||
receiver_close_and_save_arms,
|
||||
};
|
||||
@@ -6,11 +6,22 @@
|
||||
//! and the broken expectation — a failing row is a test failure, which
|
||||
//! is the suite's purpose as a test-side artifact (relaxing the
|
||||
//! family-wide no-panic rule here only).
|
||||
//!
|
||||
//! Determinism: no wall-clock-timing assertions and no two-task race
|
||||
//! windows — orchestration that would need one (runner spawning,
|
||||
//! cross-connection wakes) stays engine-test-side; rows keep to
|
||||
//! single-store sequential drives and the factory's isolation
|
||||
//! guarantee. Rows exercise a *close/dispose* arm by dropping the
|
||||
//! store handle (the dispose carrier every engine implements — the
|
||||
//! drop path closes receivers terminally); teardown-then-fresh-reopen
|
||||
//! rows would be factory-contract tests (ADR-022), not contract rows,
|
||||
//! so none exist.
|
||||
|
||||
use serde_json::json;
|
||||
|
||||
use alkstore::{
|
||||
BoxedFuture, Delivery, Error, JobHandle, RESERVED_PREFIX, Result, validate_shared_name,
|
||||
BoxedFuture, Delivery, EnqueueOpts, Error, JobHandle, JobState, RESERVED_PREFIX, Result,
|
||||
StreamEvent, encode_payload, validate_shared_name,
|
||||
};
|
||||
|
||||
use crate::factory::StoreFactory;
|
||||
@@ -208,117 +219,126 @@ pub async fn name_validation_rejects_empty_and_reserved(factory: &dyn StoreFacto
|
||||
outbox.run_once("", &mut NoDelivery).await,
|
||||
);
|
||||
|
||||
// Tx paths — the twins validate identically (begin_tx route).
|
||||
let mut tx = store.begin_tx().await.expect("begin_tx opens");
|
||||
expect_invalid(
|
||||
"enqueue_tx(empty queue)",
|
||||
tx.enqueue_tx("", Default::default(), json!(null)).await,
|
||||
);
|
||||
expect_reserved(
|
||||
"enqueue_tx(reserved queue)",
|
||||
tx.enqueue_tx(RESERVED_PREFIX, Default::default(), json!(null))
|
||||
.await,
|
||||
);
|
||||
expect_invalid(
|
||||
"publish_tx(empty stream)",
|
||||
tx.publish_tx("", json!(null)).await,
|
||||
);
|
||||
expect_reserved(
|
||||
"publish_tx(reserved stream)",
|
||||
tx.publish_tx(RESERVED_PREFIX, json!(null)).await,
|
||||
);
|
||||
expect_invalid(
|
||||
"notify_tx(empty channel)",
|
||||
tx.notify_tx("", json!(null)).await,
|
||||
);
|
||||
expect_reserved(
|
||||
"notify_tx(reserved channel)",
|
||||
tx.notify_tx(RESERVED_PREFIX, json!(null)).await,
|
||||
);
|
||||
expect_invalid(
|
||||
"publish_with_key_tx(empty stream)",
|
||||
tx.publish_with_key_tx("", None, json!(null)).await,
|
||||
);
|
||||
expect_reserved(
|
||||
"publish_with_key_tx(reserved stream)",
|
||||
tx.publish_with_key_tx(RESERVED_PREFIX, None, json!(null))
|
||||
.await,
|
||||
);
|
||||
expect_invalid(
|
||||
"publish_with_key_tx(empty key)",
|
||||
tx.publish_with_key_tx("work", Some(String::new()), json!(null))
|
||||
.await,
|
||||
);
|
||||
expect_invalid(
|
||||
"publish_with_key_tx(whitespace key)",
|
||||
tx.publish_with_key_tx("work", Some(" ".to_string()), json!(null))
|
||||
.await,
|
||||
);
|
||||
expect_invalid(
|
||||
"save_offset_tx(empty stream)",
|
||||
tx.save_offset_tx("", "consumer", 0).await,
|
||||
);
|
||||
expect_reserved(
|
||||
"save_offset_tx(reserved stream)",
|
||||
tx.save_offset_tx(RESERVED_PREFIX, "consumer", 0).await,
|
||||
);
|
||||
expect_invalid(
|
||||
"save_offset_tx(empty consumer)",
|
||||
tx.save_offset_tx("work", "", 0).await,
|
||||
);
|
||||
expect_invalid("get_job_tx(empty queue)", tx.get_job_tx("", 1).await);
|
||||
expect_reserved(
|
||||
"get_job_tx(reserved queue)",
|
||||
tx.get_job_tx(RESERVED_PREFIX, 1).await,
|
||||
);
|
||||
expect_invalid(
|
||||
"get_offset_tx(empty stream)",
|
||||
tx.get_offset_tx("", "consumer").await,
|
||||
);
|
||||
expect_reserved(
|
||||
"get_offset_tx(reserved stream)",
|
||||
tx.get_offset_tx(RESERVED_PREFIX, "consumer").await,
|
||||
);
|
||||
expect_invalid(
|
||||
"read_since_tx(empty stream)",
|
||||
tx.read_since_tx("", 0, 10).await,
|
||||
);
|
||||
expect_reserved(
|
||||
"read_since_tx(reserved stream)",
|
||||
tx.read_since_tx(RESERVED_PREFIX, 0, 10).await,
|
||||
);
|
||||
expect_invalid(
|
||||
"read_from_consumer_tx(empty stream)",
|
||||
tx.read_from_consumer_tx("", "consumer", 10).await,
|
||||
);
|
||||
expect_reserved(
|
||||
"read_from_consumer_tx(reserved stream)",
|
||||
tx.read_from_consumer_tx(RESERVED_PREFIX, "consumer", 10)
|
||||
.await,
|
||||
);
|
||||
expect_invalid(
|
||||
"read_from_consumer_tx(empty consumer)",
|
||||
tx.read_from_consumer_tx("work", "", 10).await,
|
||||
);
|
||||
expect_invalid(
|
||||
"get_offset_tx(empty consumer)",
|
||||
tx.get_offset_tx("work", "").await,
|
||||
);
|
||||
expect_invalid(
|
||||
"outbox_enqueue_tx(empty outbox)",
|
||||
tx.outbox_enqueue_tx("", Default::default(), json!(null))
|
||||
.await,
|
||||
);
|
||||
expect_reserved(
|
||||
"outbox_enqueue_tx(reserved outbox)",
|
||||
tx.outbox_enqueue_tx(RESERVED_PREFIX, Default::default(), json!(null))
|
||||
.await,
|
||||
);
|
||||
// Tx paths — the twins validate identically (begin_tx route). The
|
||||
// battery lives in its own block so the handle drops before the
|
||||
// `with_tx` leg below: a real engine's `begin_tx` parks a
|
||||
// concurrent begin on the writer slot (ADR-007's honest model —
|
||||
// the mock harness never saw this), so a hold-over handle here
|
||||
// would deadlock the row, not test it.
|
||||
{
|
||||
let mut tx = store.begin_tx().await.expect("begin_tx opens");
|
||||
expect_invalid(
|
||||
"enqueue_tx(empty queue)",
|
||||
tx.enqueue_tx("", Default::default(), json!(null)).await,
|
||||
);
|
||||
expect_reserved(
|
||||
"enqueue_tx(reserved queue)",
|
||||
tx.enqueue_tx(RESERVED_PREFIX, Default::default(), json!(null))
|
||||
.await,
|
||||
);
|
||||
expect_invalid(
|
||||
"publish_tx(empty stream)",
|
||||
tx.publish_tx("", json!(null)).await,
|
||||
);
|
||||
expect_reserved(
|
||||
"publish_tx(reserved stream)",
|
||||
tx.publish_tx(RESERVED_PREFIX, json!(null)).await,
|
||||
);
|
||||
expect_invalid(
|
||||
"notify_tx(empty channel)",
|
||||
tx.notify_tx("", json!(null)).await,
|
||||
);
|
||||
expect_reserved(
|
||||
"notify_tx(reserved channel)",
|
||||
tx.notify_tx(RESERVED_PREFIX, json!(null)).await,
|
||||
);
|
||||
expect_invalid(
|
||||
"publish_with_key_tx(empty stream)",
|
||||
tx.publish_with_key_tx("", None, json!(null)).await,
|
||||
);
|
||||
expect_reserved(
|
||||
"publish_with_key_tx(reserved stream)",
|
||||
tx.publish_with_key_tx(RESERVED_PREFIX, None, json!(null))
|
||||
.await,
|
||||
);
|
||||
expect_invalid(
|
||||
"publish_with_key_tx(empty key)",
|
||||
tx.publish_with_key_tx("work", Some(String::new()), json!(null))
|
||||
.await,
|
||||
);
|
||||
expect_invalid(
|
||||
"publish_with_key_tx(whitespace key)",
|
||||
tx.publish_with_key_tx("work", Some(" ".to_string()), json!(null))
|
||||
.await,
|
||||
);
|
||||
expect_invalid(
|
||||
"save_offset_tx(empty stream)",
|
||||
tx.save_offset_tx("", "consumer", 0).await,
|
||||
);
|
||||
expect_reserved(
|
||||
"save_offset_tx(reserved stream)",
|
||||
tx.save_offset_tx(RESERVED_PREFIX, "consumer", 0).await,
|
||||
);
|
||||
expect_invalid(
|
||||
"save_offset_tx(empty consumer)",
|
||||
tx.save_offset_tx("work", "", 0).await,
|
||||
);
|
||||
expect_invalid("get_job_tx(empty queue)", tx.get_job_tx("", 1).await);
|
||||
expect_reserved(
|
||||
"get_job_tx(reserved queue)",
|
||||
tx.get_job_tx(RESERVED_PREFIX, 1).await,
|
||||
);
|
||||
expect_invalid(
|
||||
"get_offset_tx(empty stream)",
|
||||
tx.get_offset_tx("", "consumer").await,
|
||||
);
|
||||
expect_reserved(
|
||||
"get_offset_tx(reserved stream)",
|
||||
tx.get_offset_tx(RESERVED_PREFIX, "consumer").await,
|
||||
);
|
||||
expect_invalid(
|
||||
"read_since_tx(empty stream)",
|
||||
tx.read_since_tx("", 0, 10).await,
|
||||
);
|
||||
expect_reserved(
|
||||
"read_since_tx(reserved stream)",
|
||||
tx.read_since_tx(RESERVED_PREFIX, 0, 10).await,
|
||||
);
|
||||
expect_invalid(
|
||||
"read_from_consumer_tx(empty stream)",
|
||||
tx.read_from_consumer_tx("", "consumer", 10).await,
|
||||
);
|
||||
expect_reserved(
|
||||
"read_from_consumer_tx(reserved stream)",
|
||||
tx.read_from_consumer_tx(RESERVED_PREFIX, "consumer", 10)
|
||||
.await,
|
||||
);
|
||||
expect_invalid(
|
||||
"read_from_consumer_tx(empty consumer)",
|
||||
tx.read_from_consumer_tx("work", "", 10).await,
|
||||
);
|
||||
expect_invalid(
|
||||
"get_offset_tx(empty consumer)",
|
||||
tx.get_offset_tx("work", "").await,
|
||||
);
|
||||
expect_invalid(
|
||||
"outbox_enqueue_tx(empty outbox)",
|
||||
tx.outbox_enqueue_tx("", Default::default(), json!(null))
|
||||
.await,
|
||||
);
|
||||
expect_reserved(
|
||||
"outbox_enqueue_tx(reserved outbox)",
|
||||
tx.outbox_enqueue_tx(RESERVED_PREFIX, Default::default(), json!(null))
|
||||
.await,
|
||||
);
|
||||
}
|
||||
|
||||
// Tx path via `with_tx` (the provided wrapper): a validation error
|
||||
// inside the closure propagates as the wrapper's `Err` (and the
|
||||
// disposition rolls back) — the commit-atomic path cannot bypass
|
||||
// validation either.
|
||||
// validation either. The battery's handle dropped above (never
|
||||
// committed — validation never reached the storage layer), so the
|
||||
// writer slot is free for this wrapper's own begin.
|
||||
let result = store
|
||||
.with_tx(Box::new(|tx| {
|
||||
Box::pin(async move {
|
||||
@@ -331,3 +351,597 @@ pub async fn name_validation_rejects_empty_and_reserved(factory: &dyn StoreFacto
|
||||
|
||||
factory.teardown().await.expect("factory teardown");
|
||||
}
|
||||
|
||||
/// **Extent-clamp semantics** — extents clamp to empty, boundaries are
|
||||
/// total (the numeric-domain rule).
|
||||
///
|
||||
/// `claim_batch(n <= 0)` claims nothing: empty `Vec`, no error, no
|
||||
/// dialect reinterpretation — a non-positive extent never hands out
|
||||
/// more rows than a positive-`n` claim on the same queue state would
|
||||
/// (the SQLite `LIMIT -1` dialect artifact is dead at the trait-impl
|
||||
/// entry, before any substrate round trip; the queued rows survive
|
||||
/// untouched and remain claimable after). Both stream reads clamp
|
||||
/// identically: `read_since(limit <= 0)` and
|
||||
/// `read_from_consumer(limit <= 0)` return the empty `Vec` at the
|
||||
/// trait-impl entry. Boundaries are total without a guard:
|
||||
/// `trim_to(horizon < min surviving offset)` deletes nothing
|
||||
/// (idempotent, `Ok(0)`) and surviving rows' offsets are never
|
||||
/// renumbered.
|
||||
///
|
||||
/// Contract stamp: ADR-023 §2 (numeric domains: extents clamp to
|
||||
/// empty, boundaries total); ADR-010 §1 (the no-work-is-a-value
|
||||
/// vocabulary the clamp composes with).
|
||||
///
|
||||
/// Runnable on any engine (SQLite now, pg at wave 5 via the same row).
|
||||
pub async fn extent_clamp_semantics(factory: &dyn StoreFactory) {
|
||||
let store = factory.open().await.expect("factory opens a store");
|
||||
|
||||
let queue = store
|
||||
.queue("work", alkstore::QueueOpts::default())
|
||||
.await
|
||||
.unwrap();
|
||||
let mut ids = Vec::new();
|
||||
for i in 0..5 {
|
||||
ids.push(
|
||||
queue
|
||||
.enqueue(json!({"i": i}), EnqueueOpts::default())
|
||||
.await
|
||||
.unwrap(),
|
||||
);
|
||||
}
|
||||
|
||||
for n in [0i64, -1, i64::MIN] {
|
||||
let claimed = queue.claim_batch("w1", n).await.unwrap();
|
||||
assert!(
|
||||
claimed.is_empty(),
|
||||
"claim_batch(n = {n}) must clamp to the empty Vec"
|
||||
);
|
||||
}
|
||||
// The negative extent never became "no limit": the five rows are
|
||||
// still pending and a positive-n claim on the same state hands out
|
||||
// exactly its extent.
|
||||
let claimed = queue.claim_batch("w1", 2).await.unwrap();
|
||||
assert_eq!(claimed.len(), 2, "a positive extent claims normally");
|
||||
assert_eq!(
|
||||
claimed.iter().map(|h| h.job().id).collect::<Vec<_>>(),
|
||||
ids[..2]
|
||||
);
|
||||
assert_eq!(
|
||||
queue.get_job(ids[2]).await.unwrap().unwrap().state,
|
||||
JobState::Pending,
|
||||
"the guarded calls destroyed nothing"
|
||||
);
|
||||
|
||||
let stream = store.stream("events").await.unwrap();
|
||||
stream.publish(json!({"n": 1})).await.unwrap();
|
||||
stream.publish(json!({"n": 2})).await.unwrap();
|
||||
stream.save_offset("c", 5).await.unwrap();
|
||||
for limit in [0i64, -1, i64::MIN] {
|
||||
let read = stream.read_since(0, limit).await.unwrap();
|
||||
assert!(
|
||||
read.is_empty(),
|
||||
"read_since(limit = {limit}) must read nothing"
|
||||
);
|
||||
let read = stream.read_from_consumer("c", limit).await.unwrap();
|
||||
assert!(
|
||||
read.is_empty(),
|
||||
"read_from_consumer(limit = {limit}) must read nothing"
|
||||
);
|
||||
}
|
||||
let whole = stream.read_since(0, 10).await.unwrap();
|
||||
assert_eq!(
|
||||
whole.len(),
|
||||
2,
|
||||
"the clamped reads never handed the stream out in bulk"
|
||||
);
|
||||
|
||||
// Boundary totality: a horizon below every surviving offset deletes
|
||||
// nothing, idempotently, and survivors keep their offsets.
|
||||
let deleted = stream.trim_to(-7).await.unwrap();
|
||||
assert_eq!(deleted, 0, "a negative trim horizon deletes nothing");
|
||||
let after = stream.read_since(0, 10).await.unwrap();
|
||||
let survivor_offsets: Vec<i64> = after.iter().map(|e| e.offset).collect();
|
||||
assert_eq!(
|
||||
survivor_offsets,
|
||||
whole.iter().map(|e| e.offset).collect::<Vec<_>>(),
|
||||
"surviving offsets are immutable, never renumbered"
|
||||
);
|
||||
|
||||
factory.teardown().await.expect("factory teardown");
|
||||
}
|
||||
|
||||
/// **Duration-refusal** — durations are the other half of the
|
||||
/// numeric-domain rule: `try_lock(ttl <= 0)` and `renew(ttl <= 0)` are
|
||||
/// opaque rejections (`Error::Database`), never a granted lock, never
|
||||
/// a silent no-op grant, and not misreported as name errors. The lock
|
||||
/// state is untouched by a guarded call: a guarded `try_lock` over an
|
||||
/// untouched name creates no row and a later valid try succeeds — and
|
||||
/// a valid holder's subsequent guarded renews leave its window
|
||||
/// unshrunken.
|
||||
///
|
||||
/// The opaque shape is deliberate (no `InvalidArgument` variant was
|
||||
/// minted — the caller's remedy is the same regardless: fix the
|
||||
/// constant, ADR-008 §5's act-differently rule); the row pins the shape
|
||||
/// engines must agree on, not a message text.
|
||||
///
|
||||
/// Contract stamp: ADR-023 §2 (durations reject); ADR-008 §5 (the
|
||||
/// opaque fallback carries it, source chain in the detail); ADR-008 §7
|
||||
/// (the exclusion guarantee the refusal protects).
|
||||
pub async fn duration_refusal_on_non_positive_ttl(factory: &dyn StoreFactory) {
|
||||
let store = factory.open().await.expect("factory opens a store");
|
||||
|
||||
for ttl in [0i64, -1] {
|
||||
let out = store.try_lock("guarded", "a", ttl).await;
|
||||
match out {
|
||||
Ok(Some(_)) => panic!("try_lock(ttl = {ttl}) must not grant"),
|
||||
Ok(None) => panic!("try_lock(ttl = {ttl}) must be a rejection, not a silent miss"),
|
||||
Err(e) => assert!(
|
||||
matches!(e, Error::Database(_)),
|
||||
"opaque `Database` rejection shape, got: {e:?}"
|
||||
),
|
||||
}
|
||||
}
|
||||
let lock = store
|
||||
.try_lock("guarded", "a", 300)
|
||||
.await
|
||||
.unwrap()
|
||||
.expect("the guarded calls fired no round trip: the name is free");
|
||||
for ttl in [0i64, -1] {
|
||||
match lock.renew(ttl).await {
|
||||
Ok(landed) => panic!("renew(ttl = {ttl}) must not land: {landed}"),
|
||||
Err(e) => assert!(
|
||||
matches!(e, Error::Database(_)),
|
||||
"opaque `Database` rejection shape, got: {e:?}"
|
||||
),
|
||||
}
|
||||
}
|
||||
assert!(lock.renew(600).await.unwrap(), "a positive ttl renews");
|
||||
let _ = lock.release().await.unwrap();
|
||||
|
||||
factory.teardown().await.expect("factory teardown");
|
||||
}
|
||||
|
||||
/// **`encode_payload` typed-failure round-trip** — enqueue and publish
|
||||
/// store the exact serde_json serialization of the input `Value` and
|
||||
/// decode back to `Value` equality; the helper's failure arm is typed
|
||||
/// (`Error::Codec`), so no silent empty-byte store exists (which would
|
||||
/// decode to `Null` — the corruption ADR-023 §1 removed).
|
||||
///
|
||||
/// Row-driven probe: the round-trip property pins the byte-level pin
|
||||
/// through the real store paths (queue enqueue, keyed publish), both
|
||||
/// auto-commit; the helper's typed shape is pinned directly by
|
||||
/// exercising it against the same values (serialization of a `Value`
|
||||
/// is infallible in practice — the assert is that what *is* stored is
|
||||
/// exactly the helper's output bytes, and that the failure arm's type
|
||||
/// exists for the seam to carry). Byte equality composes with the
|
||||
/// cross-engine pin (the same `Value` serializes identically on both
|
||||
/// engines).
|
||||
///
|
||||
/// Contract stamp: ADR-023 §1 (the fallible helper — `Result<Vec<u8>,
|
||||
/// Error::Codec>`, no silent empty-Vec fallback); ADR-020 §4 (store
|
||||
/// exactly the serialized bytes; decode back to `Value` equality).
|
||||
pub async fn payload_round_trip_stores_exact_encoding(factory: &dyn StoreFactory) {
|
||||
let store = factory.open().await.expect("factory opens a store");
|
||||
|
||||
let value = json!({"task": "resize", "deep": [1, 2, {"n": null, "ok": true}], "f": -0.5});
|
||||
let exact = encode_payload(&value).expect("the serialization encodes");
|
||||
|
||||
let queue = store
|
||||
.queue("work", alkstore::QueueOpts::default())
|
||||
.await
|
||||
.unwrap();
|
||||
let id = queue
|
||||
.enqueue(value.clone(), EnqueueOpts::default())
|
||||
.await
|
||||
.unwrap();
|
||||
let job = queue.get_job(id).await.unwrap().expect("the row exists");
|
||||
assert_eq!(
|
||||
job.payload, exact,
|
||||
"the stored job bytes are the helper's exact serialization"
|
||||
);
|
||||
let decoded: serde_json::Value =
|
||||
alkstore::payload_as(&job.payload).expect("the stored bytes decode");
|
||||
assert_eq!(decoded, value, "the decode round-trips to value equality");
|
||||
|
||||
let stream = store.stream("events").await.unwrap();
|
||||
let offset = stream
|
||||
.publish_with_key(Some("k".to_string()), value.clone())
|
||||
.await
|
||||
.unwrap();
|
||||
assert!(offset > 0, "the publish got an offset");
|
||||
let page = stream.read_since(0, 10).await.unwrap();
|
||||
assert_eq!(page.len(), 1, "exactly one event landed");
|
||||
assert_eq!(
|
||||
page[0].payload, exact,
|
||||
"the stored event bytes are the helper's exact serialization"
|
||||
);
|
||||
assert_eq!(page[0].key.as_deref(), Some("k"), "the key round-trips");
|
||||
assert_eq!(
|
||||
alkstore::payload_as::<serde_json::Value>(&page[0].payload).unwrap(),
|
||||
value,
|
||||
"the event decode round-trips to value equality"
|
||||
);
|
||||
|
||||
factory.teardown().await.expect("factory teardown");
|
||||
}
|
||||
|
||||
/// **`PayloadTooLarge` occurrence asymmetry** — the SQLite engine
|
||||
/// never produces the variant at any payload size (no notification
|
||||
/// limit); engine-agnostic code writes the same match on both engines
|
||||
/// and the non-occurring arm simply never fires. The row drives an
|
||||
/// oversized payload through the store's notify path and asserts
|
||||
/// success.
|
||||
///
|
||||
/// This is the SQLite column's arm of the asymmetry row (the pg-side
|
||||
/// rejection arm pins with wave 4's adoption — the same row text, the
|
||||
/// engine-differing leg inside).
|
||||
///
|
||||
/// Contract stamp: ADR-016 §5 (the occurrence asymmetry — SQLite never
|
||||
/// produces `PayloadTooLarge`); ADR-008 §5 (universal variant, no
|
||||
/// engine-specific growth).
|
||||
pub async fn payload_too_large_never_produced_on_sqlite(factory: &dyn StoreFactory) {
|
||||
let store = factory.open().await.expect("factory opens a store");
|
||||
|
||||
let big = "x".repeat(2_000_000);
|
||||
store
|
||||
.notify("big", json!({"blob": big}))
|
||||
.await
|
||||
.expect("SQLite notify has no payload limit at any size");
|
||||
|
||||
factory.teardown().await.expect("factory teardown");
|
||||
}
|
||||
|
||||
/// **Drop = rollback, no-ghosts** — a transaction handle dropped
|
||||
/// before commit rolls its writes back through RAII: job rows, stream
|
||||
/// events, notifies, and offset saves all leave no residue the
|
||||
/// transaction's own reads can no longer observe, and post-drop reads
|
||||
/// (fresh tx and auto-commit handles alike) see the pre-transaction
|
||||
/// state exactly.
|
||||
///
|
||||
/// The row drives one tx through the full no-ghosts list (every write
|
||||
/// kind, ADR-021 §4's uniform list), drops without commit, then reads
|
||||
/// the subjects through every surface that saw them mid-tx.
|
||||
///
|
||||
/// Contract stamp: ADR-021 §4 (drop = rollback; the uniform no-ghosts
|
||||
/// list); ADR-007 (the caller-held-handle seam the RAII paths guard).
|
||||
pub async fn drop_rollback_leaves_no_ghosts(factory: &dyn StoreFactory) {
|
||||
let store = factory.open().await.expect("factory opens a store");
|
||||
|
||||
let mut tx = store.begin_tx().await.unwrap();
|
||||
let ghost_id = tx
|
||||
.enqueue_tx("work", EnqueueOpts::default(), json!({"ghost": true}))
|
||||
.await
|
||||
.unwrap();
|
||||
tx.publish_tx("events", json!({"ghost": true}))
|
||||
.await
|
||||
.unwrap();
|
||||
tx.save_offset_tx("events", "c", 7).await.unwrap();
|
||||
|
||||
// The ghost is real inside the tx (read-your-own-writes is what
|
||||
// makes this drop's rollback meaningful).
|
||||
assert!(
|
||||
tx.get_job_tx("work", ghost_id).await.unwrap().is_some(),
|
||||
"the tx sees its own enqueue pre-drop"
|
||||
);
|
||||
drop(tx);
|
||||
|
||||
// Fresh reads: the subjects are gone through every surface.
|
||||
let mut tx = store.begin_tx().await.unwrap();
|
||||
let gone = tx.get_job_tx("work", ghost_id).await.unwrap();
|
||||
let events = tx.read_since_tx("events", 0, 100).await.unwrap();
|
||||
let checkpoint = tx.get_offset_tx("events", "c").await.unwrap();
|
||||
tx.commit().await.unwrap();
|
||||
assert_eq!(gone, None, "the dropped tx's job left no row");
|
||||
assert_eq!(events, Vec::<StreamEvent>::new(), "no ghost events");
|
||||
assert_eq!(checkpoint, 0, "the dropped tx's offset save never landed");
|
||||
|
||||
let queue = store
|
||||
.queue("work", alkstore::QueueOpts::default())
|
||||
.await
|
||||
.unwrap();
|
||||
let drained = queue.claim_batch("probe", 10).await.unwrap();
|
||||
assert!(drained.is_empty(), "no ghost work became claimable");
|
||||
|
||||
let stream = store.stream("events").await.unwrap();
|
||||
let events = stream.read_since(0, 100).await.unwrap();
|
||||
assert!(events.is_empty(), "no ghost events on the auto-commit read");
|
||||
assert_eq!(
|
||||
stream.get_offset("c").await.unwrap(),
|
||||
0,
|
||||
"no ghost checkpoint"
|
||||
);
|
||||
|
||||
factory.teardown().await.expect("factory teardown");
|
||||
}
|
||||
|
||||
/// **In-tx read-your-own-writes** — reads issued on the tx handle see
|
||||
/// the transaction's own writes: `get_job_tx` a row the same tx
|
||||
/// enqueued, `read_since_tx` an event the same tx published, and the
|
||||
/// consumer-twin forms see their own saves. (Another connection
|
||||
/// seeing none of it pre-commit and the post-rollback no-ghosts leg
|
||||
/// pin in the drop-rollback and commit-visibility rows; this row pins
|
||||
/// the visibility itself.)
|
||||
///
|
||||
/// Contract stamp: ADR-021 §1 (tx-side read methods,
|
||||
/// read-your-own-writes); ADR-007 (the commit-atomic seam the reads
|
||||
/// ride).
|
||||
pub async fn in_tx_reads_see_own_writes(factory: &dyn StoreFactory) {
|
||||
let store = factory.open().await.expect("factory opens a store");
|
||||
|
||||
let mut tx = store.begin_tx().await.unwrap();
|
||||
let id = tx
|
||||
.enqueue_tx("work", EnqueueOpts::default(), json!({"n": 9}))
|
||||
.await
|
||||
.unwrap();
|
||||
let offset = tx.publish_tx("events", json!({"n": 9})).await.unwrap();
|
||||
|
||||
let job = tx
|
||||
.get_job_tx("work", id)
|
||||
.await
|
||||
.unwrap()
|
||||
.expect("own write visible");
|
||||
assert_eq!(job.id, id);
|
||||
assert_eq!(job.queue, "work");
|
||||
assert_eq!(job.state, JobState::Pending);
|
||||
assert_eq!(job.payload, encode_payload(&json!({"n": 9})).unwrap());
|
||||
|
||||
let page = tx.read_since_tx("events", 0, 10).await.unwrap();
|
||||
assert_eq!(page.len(), 1, "the tx sees its own publish");
|
||||
assert_eq!(page[0].offset, offset, "the assigned offset reads back");
|
||||
assert_eq!(page[0].stream, "events");
|
||||
|
||||
// The twin forms see their own saves too.
|
||||
tx.save_offset_tx("events", "c", offset).await.unwrap();
|
||||
let from = tx.read_from_consumer_tx("events", "c", 10).await.unwrap();
|
||||
assert!(
|
||||
from.is_empty(),
|
||||
"everything is at or below the checkpoint just saved"
|
||||
);
|
||||
let checkpoint = tx.get_offset_tx("events", "c").await.unwrap();
|
||||
assert_eq!(checkpoint, offset, "the tx sees its own save");
|
||||
tx.commit().await.unwrap();
|
||||
|
||||
factory.teardown().await.expect("factory teardown");
|
||||
}
|
||||
|
||||
/// **Enqueue-opts resolution** — opts-to-row stamping is engine-side
|
||||
/// arithmetic but contract-pinned identical across engines.
|
||||
/// Delay-over-`run_at` precedence (the row shows the resolved ready
|
||||
/// time only), relative `expires` resolution to the absolute row
|
||||
/// expiry, the plain-queue derived stamps on the no-handle-open
|
||||
/// tx shape (300/3/5/none), and the same defaults with
|
||||
/// `EnqueueOpts::max_attempts` override riding over them — `get_job`
|
||||
/// shows resolved values, never the raw opts.
|
||||
///
|
||||
/// The `run_at`-alone-literal and neither-field-resolves-to-now legs
|
||||
/// are wall-clock-adjacent; they pin with wave 5's equivalence row
|
||||
/// (this row's precedence/resolution legs are the engine-agnostic
|
||||
/// claim).
|
||||
///
|
||||
/// Contract stamp: ADR-020 §1 (delay wins over `run_at`); §2
|
||||
/// (relative `expires` → absolute row value); §3/§3a (the derived
|
||||
/// default stamps — plain queue 300/3/5/none, per-job override);
|
||||
/// ADR-012 §2 (engine-side arithmetic pinned equivalent).
|
||||
pub async fn enqueue_opts_resolution(factory: &dyn StoreFactory) {
|
||||
let store = factory.open().await.expect("factory opens a store");
|
||||
|
||||
const PLAIN_STAMPS: (i64, i64, i64, Option<i64>) = (300, 3, 5, None);
|
||||
|
||||
let mut tx = store.begin_tx().await.unwrap();
|
||||
let plain = tx
|
||||
.enqueue_tx("work", EnqueueOpts::default(), json!({"shape": "tx"}))
|
||||
.await
|
||||
.unwrap();
|
||||
let precedence = tx
|
||||
.enqueue_tx(
|
||||
"work",
|
||||
EnqueueOpts {
|
||||
delay: Some(60),
|
||||
run_at: Some(1_500_000_000),
|
||||
..Default::default()
|
||||
},
|
||||
json!({"shape": "precedence"}),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
let bounded = tx
|
||||
.enqueue_tx(
|
||||
"work",
|
||||
EnqueueOpts {
|
||||
expires: Some(90),
|
||||
..Default::default()
|
||||
},
|
||||
json!({"shape": "bounded"}),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
let overridden = tx
|
||||
.enqueue_tx(
|
||||
"work",
|
||||
EnqueueOpts {
|
||||
max_attempts: Some(7),
|
||||
..Default::default()
|
||||
},
|
||||
json!({"shape": "override"}),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
tx.commit().await.unwrap();
|
||||
|
||||
let queue = store
|
||||
.queue("work", alkstore::QueueOpts::default())
|
||||
.await
|
||||
.unwrap();
|
||||
let plain_job = queue.get_job(plain).await.unwrap().unwrap();
|
||||
let stamps = (
|
||||
plain_job.visibility_timeout_s,
|
||||
plain_job.max_attempts,
|
||||
plain_job.backoff_base_s,
|
||||
plain_job.dead_letter_retention_s,
|
||||
);
|
||||
assert_eq!(
|
||||
stamps, PLAIN_STAMPS,
|
||||
"the no-handle-open tx shape stamps the plain-queue derived defaults"
|
||||
);
|
||||
|
||||
let precedence_job = queue.get_job(precedence).await.unwrap().unwrap();
|
||||
assert_ne!(
|
||||
precedence_job.run_at, 1_500_000_000,
|
||||
"delay wins over run_at — run_at is not the row's ready time"
|
||||
);
|
||||
|
||||
let bounded_job = queue.get_job(bounded).await.unwrap().unwrap();
|
||||
let expires_at = bounded_job
|
||||
.expires_at
|
||||
.expect("relative expires resolves to an absolute row expiry");
|
||||
let now = unix_now();
|
||||
assert!(
|
||||
expires_at > now + 60 && expires_at <= now + 90,
|
||||
"expires is relative from enqueue (resolved, not raw): {expires_at} vs now {now}"
|
||||
);
|
||||
|
||||
let override_job = queue.get_job(overridden).await.unwrap().unwrap();
|
||||
assert_eq!(
|
||||
(
|
||||
override_job.visibility_timeout_s,
|
||||
override_job.max_attempts,
|
||||
override_job.backoff_base_s,
|
||||
override_job.dead_letter_retention_s
|
||||
),
|
||||
(PLAIN_STAMPS.0, 7, PLAIN_STAMPS.2, PLAIN_STAMPS.3),
|
||||
"the per-job override rides over the derived defaults, others untouched"
|
||||
);
|
||||
|
||||
factory.teardown().await.expect("factory teardown");
|
||||
}
|
||||
|
||||
/// **Receiver close and save arms** — the stream subscription
|
||||
/// receiver's pinned arms (ADR-021 §5): `recv()`'s error arm carries
|
||||
/// `Database` only, close is terminal `None` (this SQLite column:
|
||||
/// store disposal closes the receiver — the cause asymmetry is
|
||||
/// documented per engine), `try_recv` distinguishes idle (`Ok(None)`)
|
||||
/// from closed (`Err(Error::Closed)`), the receiver-form save is a
|
||||
/// no-op before any event yields, lands the last-yielded offset
|
||||
/// exactly once yielded, and regressions are silent no-ops — direct
|
||||
/// and receiver forms interleave without regressing one another
|
||||
/// (ADR-019 §6).
|
||||
///
|
||||
/// Contract stamp: ADR-021 §5 (receiver error/close arms —
|
||||
/// `Database`-only `Err`, terminal `None` close, SQLite closes on
|
||||
/// watcher death); ADR-019 §6 (save composition — monotone, forms
|
||||
/// interleave); ADR-008 §8 (explicit saves only).
|
||||
pub async fn receiver_close_and_save_arms(factory: &dyn StoreFactory) {
|
||||
let store = factory.open().await.expect("factory opens a store");
|
||||
|
||||
let stream = store.stream("events").await.unwrap();
|
||||
let o1 = stream.publish(json!({"n": 1})).await.unwrap();
|
||||
let o2 = stream.publish(json!({"n": 2})).await.unwrap();
|
||||
|
||||
let mut rx = stream.subscribe("c").await.unwrap();
|
||||
assert_eq!(rx.offset(), 0, "no event yielded yet");
|
||||
assert!(
|
||||
rx.save_offset().is_ok(),
|
||||
"the pre-yield save is a no-op: Ok on any engine"
|
||||
);
|
||||
// The no-op never wrote a checkpoint: a fresh consumer read-back
|
||||
// (from 0) is empty and the direct form's fresh inspection is 0.
|
||||
assert_eq!(
|
||||
store
|
||||
.stream("events")
|
||||
.await
|
||||
.unwrap()
|
||||
.get_offset("c")
|
||||
.await
|
||||
.unwrap(),
|
||||
0,
|
||||
"the pre-yield save touched no checkpoint"
|
||||
);
|
||||
|
||||
let mut last = rx
|
||||
.recv()
|
||||
.await
|
||||
.expect("the drained attach read yields o1")
|
||||
.unwrap();
|
||||
assert_eq!(last.offset, o1, "the attach read yields offset ASC");
|
||||
// Drain the rest of the attach read (o2 was published pre-attach —
|
||||
// the attach read drains to the tail, so it is pending in the
|
||||
// receiver's feed too; the idle assert below reads an empty tail).
|
||||
while let Some(next) = rx.try_recv().unwrap() {
|
||||
assert!(next.offset > last.offset, "the drain stays offset ASC");
|
||||
last = next;
|
||||
}
|
||||
rx.save_offset().unwrap();
|
||||
assert_eq!(
|
||||
rx.offset(),
|
||||
last.offset,
|
||||
"the save lands the last-yielded offset exactly"
|
||||
);
|
||||
assert!(
|
||||
store
|
||||
.stream("events")
|
||||
.await
|
||||
.unwrap()
|
||||
.get_offset("c")
|
||||
.await
|
||||
.unwrap()
|
||||
> 0,
|
||||
"the checkpoint is stored (the no-op save never wrote it; this one did)"
|
||||
);
|
||||
|
||||
// Idle: the feed drained to the tail — try_recv is Ok(None) (not
|
||||
// closed).
|
||||
assert_eq!(
|
||||
rx.try_recv().unwrap(),
|
||||
None,
|
||||
"an attached drained receiver idles, it is not closed"
|
||||
);
|
||||
|
||||
// Direct-then-receiver composition without regression: a direct
|
||||
// save ahead moves the cursor; the receiver's subsequent save
|
||||
// cannot rewind it (monotone).
|
||||
stream.save_offset("c", o2).await.unwrap();
|
||||
let direct = store
|
||||
.stream("events")
|
||||
.await
|
||||
.unwrap()
|
||||
.get_offset("c")
|
||||
.await
|
||||
.unwrap();
|
||||
assert!(direct >= o2, "the direct save landed");
|
||||
rx.save_offset().unwrap();
|
||||
assert!(
|
||||
store
|
||||
.stream("events")
|
||||
.await
|
||||
.unwrap()
|
||||
.get_offset("c")
|
||||
.await
|
||||
.unwrap()
|
||||
>= direct,
|
||||
"the receiver form never regresses the direct form's save"
|
||||
);
|
||||
|
||||
// Terminal close on disposal (this engine's documented close
|
||||
// cause: the watcher's death guard clears the feed — ADR-021 §5's
|
||||
// SQLite arm; the store handle dropping is the disposal).
|
||||
drop(store);
|
||||
let closed = rx.recv().await;
|
||||
assert!(
|
||||
closed.is_none(),
|
||||
"store disposal closes the receiver terminally: got {closed:?}"
|
||||
);
|
||||
assert!(
|
||||
matches!(rx.try_recv(), Err(Error::Closed)),
|
||||
"post-close try_recv is Err(Closed), not idle"
|
||||
);
|
||||
}
|
||||
|
||||
/// Shared clock helper for the suite's clock-adjacent rows: unix
|
||||
/// seconds now (rows never assert tight timings — the tolerance is
|
||||
/// the stamp's resolution).
|
||||
fn unix_now() -> i64 {
|
||||
std::time::SystemTime::now()
|
||||
.duration_since(std::time::UNIX_EPOCH)
|
||||
.map(|d| d.as_secs() as i64)
|
||||
.unwrap_or(0)
|
||||
}
|
||||
Reference in new issue
Block a user