use std::sync::Arc;
use crate::batch::{BatchStore, ItemOutcome};
use crate::case::{CaseStore, ClaimError, EventStore, TargetedDelivery, TaskStore, TimerStore};
use crate::core::{
BatchId, CaseId, CaseVersion, CorrelationKey, EffectKey, InboundEvent, Justification, OnExpiry,
Phase, Priority, RunId, Spend, StepId, StoreError, Subscription, Task, TaskId, TaskState,
Timestamp,
};
pub use super::conformance::Report;
fn ts(secs: i64) -> Timestamp {
Timestamp::from_unix_timestamp(secs).expect("representable")
}
fn keys(n: &str) -> Vec<CorrelationKey> {
vec![CorrelationKey::new("doc", n)]
}
fn effect(n: u8) -> EffectKey {
EffectKey::derive(StepId(0), Phase::Forward, u32::from(n), 1, "probe", &[n])
}
pub async fn check_cases(store: &Arc<dyn CaseStore>, r: &mut Report) {
correlating_twice_yields_one_case(store, r).await;
a_closed_case_does_not_match(store, r).await;
closing_via_set_status_also_releases_the_keys(store, r).await;
an_unmet_obligation_blocks_closure(store, r).await;
the_census_counts_every_open_case(store, r).await;
two_concurrent_messages_open_one_case(store, r).await;
a_stale_state_write_is_refused(store, r).await;
a_state_write_to_a_missing_case_is_not_found(store, r).await;
only_one_of_several_racing_writers_wins(store, r).await;
}
async fn a_stale_state_write_is_refused(store: &Arc<dyn CaseStore>, r: &mut Report) {
r.checked += 1;
let Ok(c) = store
.correlate_or_open("matter", &keys("INV-CAS"), ts(2_000))
.await
else {
r.record("case state", "correlate_or_open failed");
return;
};
let case = c.case_id();
let Ok(Some(before)) = store.case(case).await else {
r.record(
"case state",
"a case that was just opened cannot be read back",
);
return;
};
let Ok(after) = store
.put_state(case, before.version, serde_json::json!({ "by": "first" }))
.await
else {
r.record("case state", "a write at the current version was refused");
return;
};
if after <= before.version {
r.record(
"case state",
"a write did not advance the version, so no later write can tell \
whether the case moved",
);
}
match store
.put_state(case, before.version, serde_json::json!({ "by": "second" }))
.await
{
Err(StoreError::CaseConflict { .. }) => {}
Ok(_) => r.record(
"case state",
"a write made against a version the case has moved past was accepted. \
That is a lost update: the first writer's work is gone and nothing \
in the record shows it. The version check must be a predicate on the \
UPDATE, not a read followed by a write",
),
Err(e) => r.record(
"case state",
format!("a stale write must report CaseConflict, reported: {e}"),
),
}
if let Ok(Some(now)) = store.case(case).await
&& now.state != serde_json::json!({ "by": "first" })
{
r.record(
"case state",
"the refused write changed the state anyway — the check must happen \
before the row is touched",
);
}
}
async fn a_state_write_to_a_missing_case_is_not_found(store: &Arc<dyn CaseStore>, r: &mut Report) {
r.checked += 1;
let absent = CaseId::generate();
match store
.put_state(absent, CaseVersion::INITIAL, serde_json::json!({}))
.await
{
Err(StoreError::NotFound(_)) => {}
Ok(_) => r.record(
"case state",
"writing to a case that does not exist reported success. A guard whose \
result nobody reads is not a guard",
),
Err(e) => r.record(
"case state",
format!("a write to a missing case must report NotFound, reported: {e}"),
),
}
}
async fn only_one_of_several_racing_writers_wins(store: &Arc<dyn CaseStore>, r: &mut Report) {
const RACERS: usize = 8;
r.checked += 1;
let Ok(c) = store
.correlate_or_open("matter", &keys("INV-RACE-CAS"), ts(3_000))
.await
else {
r.record("case state", "correlate_or_open failed");
return;
};
let case = c.case_id();
let Ok(Some(start)) = store.case(case).await else {
r.record("case state", "cannot read back a fresh case");
return;
};
let winners = futures_util::future::join_all((0..RACERS).map(|i| {
let store = Arc::clone(store);
async move {
store
.put_state(case, start.version, serde_json::json!({ "by": i }))
.await
.is_ok()
}
}))
.await
.into_iter()
.filter(|ok| *ok)
.count();
if winners != 1 {
r.record(
"case state",
format!(
"{winners} of {RACERS} writers holding the same version succeeded; \
exactly one may. More than one means the version check is not part \
of the write, and the losers' work vanished silently"
),
);
}
}
async fn correlating_twice_yields_one_case(store: &Arc<dyn CaseStore>, r: &mut Report) {
r.checked += 1;
let k = keys("INV-1");
let Ok(first) = store.correlate_or_open("matter", &k, ts(1_000)).await else {
r.record("correlation", "correlate_or_open failed on a fresh key");
return;
};
let Ok(second) = store.correlate_or_open("matter", &k, ts(1_001)).await else {
r.record("correlation", "correlate_or_open failed on a known key");
return;
};
if first.case_id() != second.case_id() {
r.record(
"correlation",
"two messages carrying the same key opened two cases. The process then \
fragments across them and its obligations are tracked in neither",
);
}
if !matches!(second, crate::case::Correlation::Attached(_)) {
r.record(
"correlation",
"the second message must report Attached, not Opened — a caller uses \
that to decide whether this is a new matter",
);
}
}
async fn two_concurrent_messages_open_one_case(store: &Arc<dyn CaseStore>, r: &mut Report) {
const RACERS: usize = 8;
const KEYS: usize = 4;
for round in 0..KEYS {
r.checked += 1;
let k = keys(&format!("RACE-{round}"));
let mut tasks = Vec::with_capacity(RACERS);
for _ in 0..RACERS {
let store = Arc::clone(store);
let k = k.clone();
tasks.push(tokio::spawn(async move {
store.correlate_or_open("matter", &k, ts(3_000)).await
}));
}
let mut ids = std::collections::BTreeSet::new();
for t in tasks {
if let Ok(Ok(c)) = t.await {
ids.insert(c.case_id());
} else {
r.record(
"correlation",
"a concurrent correlate_or_open call failed outright",
);
return;
}
}
if ids.len() > 1 {
r.record(
"correlation",
format!(
"{RACERS} messages racing for one new matter opened {} cases. Reading \
and then inserting looks atomic when called one at a time; only the \
database can settle this, and here it did not",
ids.len()
),
);
return;
}
}
}
async fn a_closed_case_does_not_match(store: &Arc<dyn CaseStore>, r: &mut Report) {
r.checked += 1;
let k = keys("INV-2");
let Ok(opened) = store.correlate_or_open("matter", &k, ts(1_000)).await else {
return;
};
if store.close(opened.case_id()).await.is_err() {
r.record("closure", "a case with no obligations must be closable");
return;
}
let Ok(again) = store.correlate_or_open("matter", &k, ts(2_000)).await else {
r.record("closure", "a key must be reusable once its case is closed");
return;
};
if again.case_id() == opened.case_id() {
r.record(
"closure",
"a message about a settled matter reanimated the closed case. Closing \
must release the keys, or a new dispute joins an audited one",
);
}
}
async fn closing_via_set_status_also_releases_the_keys(store: &Arc<dyn CaseStore>, r: &mut Report) {
r.checked += 1;
let k = keys("INV-SS");
let Ok(opened) = store.correlate_or_open("matter", &k, ts(1_000)).await else {
return;
};
let case = opened.case_id();
let deadline = crate::core::Deadline {
case,
name: "ack".into(),
resolved_at: ts(9_000),
calendar_digest: crate::core::Digest::of(b"cal"),
warn_at: None,
state: crate::core::DeadlineState::Pending,
};
if store.register_deadline(&deadline).await.is_err() {
r.record("closure", "register_deadline failed");
return;
}
if store
.set_status(case, crate::core::CaseStatus::Closed)
.await
.is_ok()
{
r.record(
"closure",
"set_status(Closed) closed a case with a pending obligation. The agent \
path must refuse it exactly as close does",
);
}
let _ = store
.set_deadline_state(case, "ack", crate::core::DeadlineState::Met)
.await;
if store
.set_status(case, crate::core::CaseStatus::Closed)
.await
.is_err()
{
r.record(
"closure",
"a case with all obligations met must be closable",
);
return;
}
let Ok(again) = store.correlate_or_open("matter", &k, ts(2_000)).await else {
r.record("closure", "a key must be reusable once its case is closed");
return;
};
if again.case_id() == case {
r.record(
"closure",
"set_status(Closed) left the case correlatable. The status column and \
correlation-open membership are two spellings of closed and the agent \
path wrote only one",
);
}
}
async fn an_unmet_obligation_blocks_closure(store: &Arc<dyn CaseStore>, r: &mut Report) {
r.checked += 1;
let Ok(opened) = store
.correlate_or_open("matter", &keys("INV-3"), ts(1_000))
.await
else {
return;
};
let case = opened.case_id();
let deadline = crate::core::Deadline {
case,
name: "ack".into(),
resolved_at: ts(9_000),
calendar_digest: crate::core::Digest::of(b"cal"),
warn_at: None,
state: crate::core::DeadlineState::Pending,
};
if store.register_deadline(&deadline).await.is_err() {
r.record("closure", "register_deadline failed");
return;
}
if store.close(case).await.is_ok() {
r.record(
"closure",
"a case with a pending obligation was closed. Closure is when people \
stop looking, so an unmet deadline must survive it",
);
}
let _ = store
.set_deadline_state(case, "ack", crate::core::DeadlineState::Met)
.await;
if store.close(case).await.is_err() {
r.record(
"closure",
"a case whose obligations are all met must be closable",
);
}
}
async fn the_census_counts_every_open_case(store: &Arc<dyn CaseStore>, r: &mut Report) {
r.checked += 1;
let before = store.census(ts(5_000)).await.map_or(0, |c| c.open);
for i in 0..3 {
let _ = store
.correlate_or_open("bulk", &keys(&format!("C-{i}")), ts(1_000))
.await;
}
match store.census(ts(5_000)).await {
Ok(c) if c.open == before + 3 => {
if c.oldest_age_secs.is_none() {
r.record(
"census",
"an open case must report an age — a count alone cannot tell a \
healthy queue from a stuck one",
);
}
}
Ok(c) => r.record(
"census",
format!(
"census must count every open case, expected {} got {}",
before + 3,
c.open
),
),
Err(e) => r.record("census", format!("census failed: {e}")),
}
}
pub async fn check_events(store: &Arc<dyn EventStore>, r: &mut Report) {
a_repeated_event_id_is_not_buffered_twice(store, r).await;
an_event_is_claimed_by_one_waiter_only(store, r).await;
a_waiter_is_matched_by_one_event_only(store, r).await;
a_targeted_event_resumes_only_its_named_run(store, r).await;
a_claimed_event_is_never_retired(store, r).await;
}
async fn a_targeted_event_resumes_only_its_named_run(store: &Arc<dyn EventStore>, r: &mut Report) {
r.checked += 1;
let first = RunId::generate();
let target = RunId::generate();
let waiting = |run, n| Subscription {
run,
case: None,
effect: effect(n),
step: StepId(0),
phase: Phase::Forward,
kind: "continue".into(),
correlation: keys("E-TARGET"),
};
let a = waiting(first, 13);
let b = waiting(target, 14);
let _ = store.subscribe(&a, ts(1_000)).await;
let _ = store.subscribe(&b, ts(1_001)).await;
let event = InboundEvent {
source: "urn:a2a:peer-a".to_owned(),
id: "message-1".into(),
kind: "continue".into(),
correlation: keys("E-TARGET"),
payload: serde_json::json!({"answer": 42}),
};
match store.deliver_to(target, &event, ts(1_002)).await {
Ok(TargetedDelivery::Matched(sub)) if sub.run == target => {}
Ok(other) => {
r.record(
"targeted delivery",
format!("an event for {target} was not claimed by that run: {other:?}"),
);
return;
}
Err(error) => {
r.record("targeted delivery", format!("delivery failed: {error}"));
return;
}
}
if !matches!(
store.deliver_to(target, &event, ts(1_003)).await,
Ok(TargetedDelivery::Matched(_))
) {
r.record(
"targeted delivery",
"retrying a claimed event with a live subscription did not recover the prior claim",
);
}
let absent = InboundEvent {
id: "message-no-waiter".into(),
..event
};
if !matches!(
store
.deliver_to(RunId::generate(), &absent, ts(1_004))
.await,
Ok(TargetedDelivery::NotWaiting)
) {
r.record(
"targeted delivery",
"a task with no subscription did not report NotWaiting",
);
}
if !matches!(store.buffer(&absent, ts(1_005)).await, Ok(true)) {
r.record(
"targeted delivery",
"a failed targeted delivery left an orphan event in the shared buffer",
);
}
}
async fn a_claimed_event_is_never_retired(store: &Arc<dyn EventStore>, r: &mut Report) {
r.checked += 1;
let event = InboundEvent {
source: "urn:conformance".to_owned(),
id: "evt-swept".into(),
kind: "ack".into(),
correlation: keys("E-9"),
payload: serde_json::json!({}),
};
let _ = store.buffer(&event, ts(1_000)).await;
let sub = Subscription {
run: RunId::generate(),
case: None,
effect: effect(90),
step: StepId(0),
phase: Phase::Forward,
kind: "ack".into(),
correlation: keys("E-9"),
};
let _ = store.subscribe(&sub, ts(1_000)).await;
if !matches!(store.claim_for(&sub, ts(1_001)).await, Ok(Some(_))) {
r.record("sweep", "a waiting subscription did not claim its event");
return;
}
if let Err(e) = store.sweep_unclaimed(ts(9_000), "expired").await {
r.record("sweep", format!("sweep_unclaimed failed: {e}"));
return;
}
match store.dead_letters(100).await {
Ok(dead) => {
if dead.iter().any(|d| d.event.id == "evt-swept") {
r.record(
"sweep",
"an event that was claimed and delivered was retired as unclaimed. The run already resumed on it, so the dead-letter queue is now reporting a message that was in fact acted on",
);
}
}
Err(e) => r.record("sweep", format!("dead_letters failed: {e}")),
}
}
async fn a_repeated_event_id_is_not_buffered_twice(store: &Arc<dyn EventStore>, r: &mut Report) {
r.checked += 1;
let event = InboundEvent {
source: "urn:conformance".to_owned(),
id: "evt-dup".into(),
kind: "ack".into(),
correlation: keys("E-1"),
payload: serde_json::json!({}),
};
let first = store.buffer(&event, ts(1_000)).await;
let second = store.buffer(&event, ts(1_001)).await;
match (first, second) {
(Ok(true), Ok(false)) => {}
(Ok(a), Ok(b)) => r.record(
"deduplication",
format!(
"buffering one event id twice reported ({a}, {b}); it must be (true, false). \
Every counterparty retries, and a duplicate delivered twice is the message \
acted on twice"
),
),
_ => r.record("deduplication", "buffer failed"),
}
}
async fn an_event_is_claimed_by_one_waiter_only(store: &Arc<dyn EventStore>, r: &mut Report) {
r.checked += 1;
let event = InboundEvent {
source: "urn:conformance".to_owned(),
id: "evt-claim".into(),
kind: "ack".into(),
correlation: keys("E-2"),
payload: serde_json::json!({}),
};
let _ = store.buffer(&event, ts(1_000)).await;
let sub = |n: u8| Subscription {
run: RunId::generate(),
case: None,
effect: effect(n),
step: StepId(0),
phase: Phase::Forward,
kind: "ack".into(),
correlation: keys("E-2"),
};
let (a, b) = (sub(10), sub(11));
let _ = store.subscribe(&a, ts(1_000)).await;
let _ = store.subscribe(&b, ts(1_000)).await;
let first = store.claim_for(&a, ts(1_002)).await;
let second = store.claim_for(&b, ts(1_003)).await;
match (first, second) {
(Ok(Some(_)), Ok(None)) => {}
(Ok(Some(_)), Ok(Some(_))) => r.record(
"single-delivery",
"one buffered event was claimed by two waiters. Claiming is what makes \
delivery exactly-once; two runs both consuming one message is the same \
message acted on twice",
),
(Ok(None), _) => r.record(
"single-delivery",
"a waiting subscription did not claim a matching buffered event",
),
_ => r.record("single-delivery", "claim_for failed"),
}
}
async fn a_waiter_is_matched_by_one_event_only(store: &Arc<dyn EventStore>, r: &mut Report) {
r.checked += 1;
let sub = Subscription {
run: RunId::generate(),
case: None,
effect: effect(12),
step: StepId(0),
phase: Phase::Forward,
kind: "ack".into(),
correlation: keys("E-3"),
};
let _ = store.subscribe(&sub, ts(1_000)).await;
let event = InboundEvent {
source: "urn:conformance".to_owned(),
id: "evt-match".into(),
kind: "ack".into(),
correlation: keys("E-3"),
payload: serde_json::json!({}),
};
let _ = store.buffer(&event, ts(1_000)).await;
let first = store.match_waiter(&event, ts(1_001)).await;
let second = store.match_waiter(&event, ts(1_002)).await;
match (first, second) {
(Ok(Some(_)), Ok(None)) => {}
(Ok(Some(_)), Ok(Some(_))) => r.record(
"single-delivery",
"one subscription was matched twice. The arrive-before-wait direction \
must claim just as the wait-before-arrive one does",
),
(Ok(None), _) => r.record(
"single-delivery",
"an arriving event did not find the run already waiting for it",
),
_ => r.record("single-delivery", "match_waiter failed"),
}
}
pub async fn check_timers(store: &Arc<dyn TimerStore>, r: &mut Report) {
r.checked += 1;
let timer = crate::core::Timer {
run: RunId::generate(),
case: None,
effect: effect(20),
step: StepId(0),
phase: Phase::Forward,
fire_at: ts(1_000),
};
if store.arm(&timer).await.is_err() {
r.record("timers", "arm failed");
return;
}
let _ = store.arm(&timer).await;
let first = store.claim_due(ts(2_000), 10).await;
let second = store.claim_due(ts(2_000), 10).await;
match (first, second) {
(Ok(a), Ok(b)) => {
if a.len() != 1 {
r.record(
"timers",
format!(
"arming twice produced {} due timers; it must produce one, or a \
resumed run is woken twice",
a.len()
),
);
}
if !b.is_empty() {
r.record(
"single-delivery",
"a claimed timer was handed to a second sweep. Two sweepers against \
one store must not both resume the same run",
);
}
}
_ => r.record("timers", "claim_due failed"),
}
}
pub async fn check_tasks(store: &Arc<dyn TaskStore>, r: &mut Report) {
a_task_is_claimed_by_one_actor_only(store, r).await;
an_excluded_actor_cannot_claim(store, r).await;
ineligibility_outranks_contention(store, r).await;
only_the_holder_releases(store, r).await;
the_backlog_counts_work_somebody_is_holding(store, r).await;
}
async fn the_backlog_counts_work_somebody_is_holding(store: &Arc<dyn TaskStore>, r: &mut Report) {
r.checked += 1;
let t = task(70, None);
let Ok(opened) = store.open(&t).await else {
r.record("backlog", "open failed");
return;
};
let before = store.open_count().await.unwrap_or(0);
let roles = vec!["ops".to_owned()];
if store.claim(opened.id, "reviewer", &roles).await.is_err() {
r.record("backlog", "the task could not be claimed");
return;
}
let claimed = store.open_count().await.unwrap_or(0);
if claimed != before {
r.record(
"backlog",
format!(
"the backlog moved from {before} to {claimed} when a task was merely claimed. A task somebody is holding is still a decision the plane is waiting on, so this reports progress that has not happened"
),
);
}
if store
.set_state(opened.id, TaskState::Completed)
.await
.is_err()
{
r.record("backlog", "the task could not be completed");
return;
}
let done = store.open_count().await.unwrap_or(0);
if done + 1 != claimed {
r.record(
"backlog",
format!(
"the backlog went from {claimed} to {done} when a task was completed; it must fall by exactly one. A count that never moves is a dashboard that cannot show the queue draining"
),
);
}
}
fn task(id: u8, excluded: Option<&str>) -> Task {
let run = RunId::generate();
Task {
id: TaskId::derive(run, effect(id)),
run,
case: None,
kind: "approval".into(),
justification: Justification::new("needs a person", serde_json::json!({})),
candidate_roles: vec!["ops".into()],
assignee: None,
priority: Priority::Normal,
state: TaskState::Open,
on_expiry: OnExpiry::Deny,
excluded_actors: excluded.map(|a| vec![a.to_owned()]).unwrap_or_default(),
created_at: ts(1_000),
due_at: None,
}
}
async fn a_task_is_claimed_by_one_actor_only(store: &Arc<dyn TaskStore>, r: &mut Report) {
r.checked += 1;
let t = task(30, None);
if store.open(&t).await.is_err() {
r.record("tasks", "open failed");
return;
}
let roles = vec!["ops".to_owned()];
let first = store.claim(t.id, "alice", &roles).await;
let second = store.claim(t.id, "bob", &roles).await;
if first.is_err() {
r.record("tasks", "an eligible actor could not claim an open task");
}
if second.is_ok() {
r.record(
"four-eyes",
"two reviewers both hold one decision. Reservation must be atomic, or \
both believe they own it and one of them acts on a stale view",
);
}
}
async fn an_excluded_actor_cannot_claim(store: &Arc<dyn TaskStore>, r: &mut Report) {
r.checked += 1;
let t = task(31, Some("alice"));
if store.open(&t).await.is_err() {
return;
}
let roles = vec!["ops".to_owned()];
if store.claim(t.id, "alice", &roles).await.is_ok() {
r.record(
"four-eyes",
"an excluded actor claimed the task. The exclusion is the whole control: \
whoever proposed an action must not be the one who approves it",
);
}
if store.claim(t.id, "bob", &roles).await.is_err() {
r.record("four-eyes", "an eligible actor was refused");
}
}
async fn ineligibility_outranks_contention(store: &Arc<dyn TaskStore>, r: &mut Report) {
r.checked += 1;
let t = task(32, Some("alice"));
if store.open(&t).await.is_err() {
r.record("tasks", "open failed");
return;
}
let roles = vec!["ops".to_owned()];
if store.claim(t.id, "bob", &roles).await.is_err() {
r.record("tasks", "an eligible actor could not claim an open task");
return;
}
match store.claim(t.id, "alice", &roles).await {
Err(ClaimError::Excluded { .. }) => {}
Err(ClaimError::AlreadyClaimed { .. }) => r.record(
"four-eyes",
"a barred reviewer was told the task is held rather than that it is \
not theirs — they will wait for the holder to release it and be \
refused again, and meanwhile they have learnt who is reviewing what",
),
other => r.record(
"four-eyes",
format!("an excluded actor's claim was answered with {other:?}"),
),
}
let wrong = vec!["clerk".to_owned()];
match store.claim(t.id, "carol", &wrong).await {
Err(ClaimError::WrongRole { .. }) => {}
other => r.record(
"tasks",
format!("an ineligible actor's claim was answered with {other:?}"),
),
}
}
async fn only_the_holder_releases(store: &Arc<dyn TaskStore>, r: &mut Report) {
r.checked += 1;
let t = task(33, None);
if store.open(&t).await.is_err() {
r.record("tasks", "open failed");
return;
}
let roles = vec!["ops".to_owned()];
if store.claim(t.id, "bob", &roles).await.is_err() {
r.record("tasks", "an eligible actor could not claim an open task");
return;
}
match store.release(t.id, "carol").await {
Err(ClaimError::NotHeld { .. }) => {}
Ok(()) => r.record(
"tasks",
"a stranger's release reported success. Whether or not it freed the \
task, the caller now believes it did — and the holder believes they \
still have it",
),
other => r.record(
"tasks",
format!("a stranger's release was answered with {other:?}"),
),
}
match store.task(t.id).await {
Ok(Some(held)) if held.assignee.as_deref() == Some("bob") => {}
_ => r.record("tasks", "a refused release still freed the task"),
}
if store.release(t.id, "bob").await.is_err() {
r.record("tasks", "the holder could not release their own claim");
}
match store.task(t.id).await {
Ok(Some(freed)) if freed.assignee.is_none() && freed.state == TaskState::Open => {}
Ok(Some(freed)) => r.record(
"tasks",
format!(
"a released task is {:?} assigned to {:?} — it is invisible to \
the queue that must now pick it up",
freed.state, freed.assignee
),
),
_ => r.record("tasks", "a released task could not be read back"),
}
}
pub async fn check_batches(store: &Arc<dyn BatchStore>, r: &mut Report) {
r.checked += 1;
let id = BatchId::generate();
if store.open(id, "digest").await.is_err() {
r.record("batches", "open failed");
return;
}
let (first, second) = (RunId::generate(), RunId::generate());
let Ok(a) = store.reserve(id, "item-001", first).await else {
r.record("batches", "reserve failed");
return;
};
let Ok(b) = store.reserve(id, "item-001", second).await else {
r.record("batches", "the second reserve failed");
return;
};
if a.run != first || b.run != first {
r.record(
"reservation",
"reserving an item twice did not return the original run id. Overwriting \
it orphans the journal that already holds this item's effects, and they \
are performed again",
);
}
r.checked += 1;
let _ = store
.record(id, "item-001", &ItemOutcome::Succeeded, Spend::default())
.await;
let _ = store.reserve(id, "item-002", RunId::generate()).await;
match store.cursor(id).await {
Ok(c) if c.as_deref() == Some("item-001") => {}
Ok(c) => r.record(
"cursor",
format!(
"the cursor must stop before the first unfinished item, got {c:?} — a \
resume that steps over one reports the batch complete with work \
outstanding"
),
),
Err(e) => r.record("cursor", format!("cursor failed: {e}")),
}
}