mem engine streams — coverage reconciliation against the seam task's foundation with the two unpinned engine-side pins (mem-engine-streams): the read_from_consumer cursor-read surface (absent-consumer head read, checkpoint-forward resume, tail-empty, extent clamps, name validation) and the ADR-015 trim property (saved sub-horizon checkpoints untouched — offsets never renumbered; consumer reads resume at the horizon's first remaining row). No production code changed; tests only (+2, mem suite at 99). Task file: status completed, Notes/Summary filled
This commit is contained in:
1 parent
f37fb0d17d
commit
3eca862f84
2 files changed
+128
-8
No files matched your search
@@ -644,6 +644,72 @@ mod tests {
|
||||
let _ = _a;
|
||||
}
|
||||
|
||||
#[tokio::test(flavor = "current_thread")]
|
||||
async fn read_from_consumer_reads_the_cursor_and_clamps_the_extent() {
|
||||
let store = MemStore::new();
|
||||
let handle = handle(&store, "orders");
|
||||
let a = handle.publish(serde_json::json!(1)).await.unwrap();
|
||||
let b = handle.publish(serde_json::json!(2)).await.unwrap();
|
||||
let c = handle.publish(serde_json::json!(3)).await.unwrap();
|
||||
let head = handle.read_from_consumer("c1", 10).await.unwrap();
|
||||
assert_eq!(
|
||||
head.iter().map(|e| e.offset).collect::<Vec<_>>(),
|
||||
vec![a, b, c],
|
||||
"an absent consumer reads from the head (checkpoint 0)"
|
||||
);
|
||||
handle.save_offset("c1", b).await.unwrap();
|
||||
let cursor = handle.read_from_consumer("c1", 10).await.unwrap();
|
||||
assert_eq!(
|
||||
cursor.iter().map(|e| e.offset).collect::<Vec<_>>(),
|
||||
vec![c],
|
||||
"the cursor read resumes at the saved checkpoint's forward tail"
|
||||
);
|
||||
assert!(handle.read_from_consumer("c1", 0).await.unwrap().is_empty());
|
||||
assert!(
|
||||
handle
|
||||
.read_from_consumer("c1", -2)
|
||||
.await
|
||||
.unwrap()
|
||||
.is_empty(),
|
||||
"the extent guard clamps non-positive limits to empty"
|
||||
);
|
||||
handle.save_offset("c1", c).await.unwrap();
|
||||
assert!(
|
||||
handle
|
||||
.read_from_consumer("c1", 10)
|
||||
.await
|
||||
.unwrap()
|
||||
.is_empty(),
|
||||
"a cursor at the tail reads empty"
|
||||
);
|
||||
match handle.read_from_consumer("", 10).await {
|
||||
Err(alkstore::Error::InvalidName { .. }) => {}
|
||||
_other => panic!("an empty consumer name rejects"),
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test(flavor = "current_thread")]
|
||||
async fn saved_sub_horizon_offsets_survive_trim_and_consumer_reads_resume_at_the_horizon() {
|
||||
let store = MemStore::new();
|
||||
let handle = handle(&store, "orders");
|
||||
let a = handle.publish(serde_json::json!(1)).await.unwrap();
|
||||
let b = handle.publish(serde_json::json!(2)).await.unwrap();
|
||||
let _c = handle.publish(serde_json::json!(3)).await.unwrap();
|
||||
handle.save_offset("c1", a).await.unwrap();
|
||||
assert_eq!(handle.trim_to(b - 1).await.unwrap(), 1);
|
||||
assert_eq!(
|
||||
handle.get_offset("c1").await.unwrap(),
|
||||
a,
|
||||
"the saved sub-horizon checkpoint is untouched by trim — offsets are never renumbered"
|
||||
);
|
||||
let resumed = handle.read_from_consumer("c1", 10).await.unwrap();
|
||||
assert_eq!(
|
||||
resumed.iter().map(|e| e.offset).collect::<Vec<_>>(),
|
||||
vec![b, _c],
|
||||
"reads from a trimmed-away region resume at the horizon's first remaining row"
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test(flavor = "current_thread")]
|
||||
async fn receiver_try_recv_reads_without_await() {
|
||||
let store = MemStore::new();
|
||||
|
||||
@@ -1,7 +1,7 @@
|
||||
---
|
||||
id: mem-engine-streams
|
||||
name: Mem engine — streams (in-memory log, per-consumer offsets, trim)
|
||||
status: pending
|
||||
status: completed
|
||||
depends_on: [mem-engine-seam-tx, mem-engine-notify-listen]
|
||||
scope: moderate
|
||||
risk: medium
|
||||
@@ -50,17 +50,17 @@ offset table (engine-mem.md §Mapping the contract, streams row):
|
||||
|
||||
## Acceptance Criteria
|
||||
|
||||
- [ ] Full `StreamHandle` surface over the in-memory log +
|
||||
- [x] Full `StreamHandle` surface over the in-memory log +
|
||||
per-consumer offsets
|
||||
- [ ] The ADR-015 properties unit-pinned: global FIFO, immutable
|
||||
- [x] The ADR-015 properties unit-pinned: global FIFO, immutable
|
||||
never-renumbered offsets (gaps legal post-trim), exact-boundary
|
||||
trim semantics (reads resume at horizon; saved sub-horizon
|
||||
offsets stay valid; trim wakes nothing), monotone save
|
||||
composition (direct + receiver forms interleave)
|
||||
- [ ] Key round-trips exactly (`None` vs `Some`); empty-`Some`-key
|
||||
- [x] Key round-trips exactly (`None` vs `Some`); empty-`Some`-key
|
||||
rejected at every entry point
|
||||
- [ ] Extent/offset numeric-domain guards at entry — pinned
|
||||
- [ ] `cargo check --target wasm32-unknown-unknown -p alkstore-mem`
|
||||
- [x] Extent/offset numeric-domain guards at entry — pinned
|
||||
- [x] `cargo check --target wasm32-unknown-unknown -p alkstore-mem`
|
||||
clean; workspace gates green
|
||||
|
||||
## References
|
||||
@@ -73,8 +73,62 @@ offset table (engine-mem.md §Mapping the contract, streams row):
|
||||
|
||||
## Notes
|
||||
|
||||
> To be filled by implementation agent
|
||||
> Filled by implementation agent, 2026-10-10.
|
||||
|
||||
1. **The implementation already landed in the seam task** — the same
|
||||
all-or-nothing trait-impl disposition the notify-listen, queues,
|
||||
and locks twins recorded: the `MemStreamHandle`/receiver mechanism
|
||||
in `alkstore-mem/src/stream.rs` (log + offsets + trim + subscribe
|
||||
+ the tx-staged keyed publish) was implemented with the seam's
|
||||
whole-engine scope; this task found its mechanism *implemented*.
|
||||
2. **Coverage reconciliation**: already pinned by the seam suite in
|
||||
`stream.rs` — the round-trip FIFO row (offsets monotone, offset-ASC
|
||||
reads, `created_at` from the clock, key `None` vs `Some` on the
|
||||
yielded events), the empty/whitespace-`Some`-key rejection at the
|
||||
auto-commit and `publish_with_key_tx` entry points, the read
|
||||
extent clamps (`limit <= 0` → empty), the monotone-save row
|
||||
(absent consumer 0, below/equal saves silent no-ops, forward saves
|
||||
land), the exact-boundary trim row (`offset <= horizon`,
|
||||
survivors keep offsets, negative horizon deletes nothing,
|
||||
idempotent), the subscriber row (attach replays from the saved
|
||||
checkpoint, wake-driven delivery, the receiver `save_offset()` form
|
||||
composing with the direct form, the pre-yield save no-op), the
|
||||
trim-wakes-nothing row (silent timeout with an open subscriber),
|
||||
and the engine-drop close arms (terminal `Err(Closed)` on
|
||||
`try_recv`, terminal `None` on `recv`, never reopens).
|
||||
3. **What this task added** — the two unpinned shapes: (a)
|
||||
`read_from_consumer_reads_the_cursor_and_clamps_the_extent` — the
|
||||
surface test the seam suite lacked for the consumer-cursor read
|
||||
(absent consumer reads from the head, checkpoint-forward resume,
|
||||
tail reads empty, extent clamps, empty-consumer-name
|
||||
`InvalidName`); (b) `saved_sub_horizon_offsets_survive_trim_and_
|
||||
consumer_reads_resume_at_the_horizon` — the ADR-015 trim property
|
||||
the trim row only half-pinned: the sub-horizon checkpoint is
|
||||
untouched by trim (offsets never renumbered) and consumer reads
|
||||
resume at the horizon's first remaining row.
|
||||
4. **No production code changed** — verified correct as-found; tests
|
||||
only (+2 in `stream.rs`, 99 mem tests total). One observation the
|
||||
new rows recorded: a cursor read resumes at the first *surviving*
|
||||
row after the checkpoint and continues to the tail (the pinned
|
||||
resume-at-horizon semantic is a lower-bound guarantee, not a
|
||||
one-row read).
|
||||
|
||||
## Summary
|
||||
|
||||
> To be filled on completion
|
||||
> Filled on completion, 2026-10-10.
|
||||
|
||||
Completed the mem engine's streams task against the seam task's
|
||||
foundation: reconciled the description's pin list against
|
||||
`alkstore-mem/src/stream.rs`'s existing suite (11 pre-existing
|
||||
mechanism tests), found the `StreamHandle` surface and nearly all
|
||||
ADR-015 properties already pinned, and added the two gaps — the
|
||||
`read_from_consumer` cursor-read surface (head/cursor/tail/extent/
|
||||
name-validation legs) and the sub-horizon-checkpoint-survives-trim +
|
||||
resume-at-horizon row. Acceptance criteria: full surface ✓ (seam
|
||||
task + the new cursor row), ADR-015 properties ✓ (seam rows + the new
|
||||
trim row), key round-trip/empty-key rejection ✓ (seam rows), extent
|
||||
guards ✓ (seam row + the new consumer-read clamp legs), gates ✓.
|
||||
Verified: `cargo test -p alkstore-mem` 99/99 (97 + 2 new, new tests
|
||||
green solo), workspace `cargo test` green, `cargo clippy --all-targets
|
||||
-- -D warnings` clean, `cargo fmt --check` clean,
|
||||
`cargo check --target wasm32-unknown-unknown -p alkstore-mem` clean.
|
||||
Reference in new issue
Block a user