use macp_runtime::log_store::{EntryKind, LogEntry, LogStore};
use macp_runtime::macp_core::policy::PolicyDefinition;
use macp_runtime::pb::{CommitmentPayload, Envelope, SessionStartPayload};
use macp_runtime::registry::SessionRegistry;
use macp_runtime::replay::replay_session;
use macp_runtime::runtime::Runtime;
use macp_runtime::session::{Session, SessionState};
use macp_runtime::storage::FileBackend;
use prost::Message;
use std::sync::Arc;
const MODE: &str = "macp.mode.handoff.v1";
const OWNER: &str = "agent://owner";
const TARGET: &str = "agent://target";
const POLICY_ID: &str = "handoff-auto-accept";
const TIMEOUT_MS: i64 = 60;
const SYNTHETIC_ID: &str = "implicit-accept:h1";
struct Harness {
rt: Runtime,
_dir: Option<tempfile::TempDir>,
}
fn make_harness() -> Harness {
harness_over(Arc::new(macp_runtime::storage::MemoryBackend), None)
}
fn make_durable_harness() -> Harness {
let dir = tempfile::tempdir().unwrap();
let storage = Arc::new(FileBackend::new(dir.path().to_path_buf()).unwrap());
harness_over(storage, Some(dir))
}
fn harness_over(
storage: Arc<dyn macp_runtime::storage::StorageBackend>,
dir: Option<tempfile::TempDir>,
) -> Harness {
let rt = Runtime::new(
storage,
Arc::new(SessionRegistry::new()),
Arc::new(LogStore::new()),
);
rt.register_policy(PolicyDefinition {
policy_id: POLICY_ID.into(),
mode: MODE.into(),
description: "implicit accept after a short timeout".into(),
rules: serde_json::json!({
"acceptance": { "implicit_accept_timeout_ms": TIMEOUT_MS },
"commitment": { "authority": "initiator_only" }
}),
schema_version: 1,
})
.expect("policy registers");
Harness { rt, _dir: dir }
}
fn assert_log_is_uncompacted(entries: &[LogEntry]) {
assert!(
!entries
.iter()
.any(|e| e.entry_kind == EntryKind::Checkpoint && e.compacted_incoming_ordinals > 0),
"this assertion needs the full log; use make_harness() (MemoryBackend)"
);
}
fn new_sid() -> String {
uuid::Uuid::new_v4().as_hyphenated().to_string()
}
fn env(
message_type: &str,
message_id: &str,
sid: &str,
sender: &str,
payload: Vec<u8>,
) -> Envelope {
Envelope {
macp_version: "1.0".into(),
mode: MODE.into(),
message_type: message_type.into(),
message_id: message_id.into(),
session_id: sid.into(),
sender: sender.into(),
timestamp_unix_ms: chrono::Utc::now().timestamp_millis(),
payload,
}
}
fn start_payload(max_suspend_ms: i64) -> Vec<u8> {
SessionStartPayload {
intent: "escalate".into(),
participants: vec![OWNER.into(), TARGET.into()],
mode_version: "1.0.0".into(),
configuration_version: "cfg-1".into(),
policy_version: POLICY_ID.into(),
ttl_ms: 60_000,
context_id: String::new(),
extensions: std::collections::HashMap::new(),
roots: vec![],
max_suspend_ms,
}
.encode_to_vec()
}
fn commitment(mode_version: &str) -> Vec<u8> {
CommitmentPayload {
commitment_id: "c1".into(),
action: "handoff.accepted".into(),
authority_scope: "support".into(),
reason: "bound".into(),
mode_version: mode_version.into(),
policy_version: POLICY_ID.into(),
configuration_version: "cfg-1".into(),
outcome_positive: true,
supersedes: None,
}
.encode_to_vec()
}
fn offer_payload() -> Vec<u8> {
macp_runtime::handoff_pb::HandoffOfferPayload {
handoff_id: "h1".into(),
target_participant: TARGET.into(),
scope: "support".into(),
reason: "escalate".into(),
}
.encode_to_vec()
}
fn accept_payload(implicit: bool) -> Vec<u8> {
macp_runtime::handoff_pb::HandoffAcceptPayload {
handoff_id: "h1".into(),
accepted_by: TARGET.into(),
reason: "ready".into(),
implicit,
}
.encode_to_vec()
}
fn context_payload() -> Vec<u8> {
macp_runtime::handoff_pb::HandoffContextPayload {
handoff_id: "h1".into(),
content_type: "text/plain".into(),
context: b"background".to_vec(),
}
.encode_to_vec()
}
async fn session_with_offer(rt: &Runtime) -> String {
let sid = new_sid();
rt.process(
&env("SessionStart", "start-1", &sid, OWNER, start_payload(0)),
None,
)
.await
.expect("session start");
rt.process(
&env("HandoffOffer", "offer-1", &sid, OWNER, offer_payload()),
None,
)
.await
.expect("offer");
let session = rt.get_session_checked(&sid).await.unwrap();
assert_eq!(
session.semantics_rev,
macp_runtime::macp_core::session::CURRENT_SEMANTICS_REV,
"the synthesis path is gated on rev >= 2"
);
sid
}
async fn log_of(rt: &Runtime, sid: &str) -> Vec<LogEntry> {
rt.log_store.get_log(sid).await.expect("session log")
}
fn incoming(entries: &[LogEntry]) -> Vec<&LogEntry> {
assert_log_is_uncompacted(entries);
entries
.iter()
.filter(|e| e.entry_kind == EntryKind::Incoming)
.collect()
}
fn synthetic_of(entries: &[LogEntry]) -> &LogEntry {
let found: Vec<&LogEntry> = entries
.iter()
.filter(|e| e.message_id == SYNTHETIC_ID)
.collect();
assert_eq!(
found.len(),
1,
"exactly one synthetic entry must be in the log"
);
found[0]
}
fn offer_received_at(entries: &[LogEntry]) -> i64 {
entries
.iter()
.find(|e| e.message_type == "HandoffOffer")
.expect("offer entry")
.received_at_ms
}
async fn sleep_past_the_deadline() {
tokio::time::sleep(std::time::Duration::from_millis(TIMEOUT_MS as u64 + 40)).await;
}
fn assert_replay_matches(live: &Session, replayed: &Session) {
assert_eq!(replayed.state, live.state, "state");
assert_eq!(replayed.resolution, live.resolution, "resolution");
assert_eq!(
replayed.mode_state, live.mode_state,
"mode_state must be byte-identical"
);
assert_eq!(
replayed.seen_message_ids, live.seen_message_ids,
"dedup sets must agree"
);
}
async fn replay_live_log(h: &Harness, sid: &str) -> Session {
let full = log_of(&h.rt, sid).await;
let stripped: Vec<LogEntry> = full
.iter()
.filter(|e| e.entry_kind != EntryKind::Checkpoint)
.cloned()
.collect();
let entries = if stripped.iter().any(|e| e.message_type == "SessionStart") {
stripped
} else {
full
};
replay_session(
sid,
&entries,
h.rt.mode_registry(),
Some(h.rt.policy_registry()),
)
.expect("the live log must replay")
}
#[tokio::test]
async fn lazy_synthesis_enters_history_before_the_trigger() {
let h = make_harness();
let sid = session_with_offer(&h.rt).await;
sleep_past_the_deadline().await;
let result =
h.rt.process(
&env("Commitment", "commit-1", &sid, OWNER, commitment("1.0.0")),
None,
)
.await
.expect("the commitment resolves once the accept is in history");
assert_eq!(result.session_state, SessionState::Resolved);
let entries = log_of(&h.rt, &sid).await;
let accepted = incoming(&entries);
let ids: Vec<&str> = accepted.iter().map(|e| e.message_id.as_str()).collect();
assert_eq!(
ids,
vec!["start-1", "offer-1", SYNTHETIC_ID, "commit-1"],
"the synthetic must sit immediately before the trigger"
);
let syn = synthetic_of(&entries);
assert_eq!(syn.entry_kind, EntryKind::Incoming);
assert_eq!(syn.message_type, "HandoffAccept");
assert_eq!(
syn.sender, TARGET,
"the accept is the target's, not the runtime's"
);
let payload =
macp_runtime::handoff_pb::HandoffAcceptPayload::decode(&*syn.raw_payload).unwrap();
assert_eq!(payload.handoff_id, "h1");
assert_eq!(payload.accepted_by, TARGET);
assert!(
payload.implicit,
"implicit = true is what makes it synthetic"
);
assert_eq!(payload.reason, "implicit accept (timeout)");
let live = h.rt.get_session_checked(&sid).await.unwrap();
let state: serde_json::Value = serde_json::from_slice(&live.mode_state).unwrap();
assert_eq!(state["offers"]["h1"]["disposition"], "Accepted");
assert_eq!(state["offers"]["h1"]["accepted_by"], TARGET);
assert!(live.seen_message_ids.contains(SYNTHETIC_ID));
}
#[tokio::test]
async fn synthetic_timestamp_is_the_deadline_not_observation_time() {
let h = make_harness();
let sid = session_with_offer(&h.rt).await;
tokio::time::sleep(std::time::Duration::from_millis(TIMEOUT_MS as u64 + 400)).await;
h.rt.process(
&env("Commitment", "commit-1", &sid, OWNER, commitment("1.0.0")),
None,
)
.await
.expect("resolves");
let entries = log_of(&h.rt, &sid).await;
let expected_d = offer_received_at(&entries) + TIMEOUT_MS;
let syn = synthetic_of(&entries);
assert_eq!(
syn.timestamp_unix_ms, expected_d,
"envelope clock must be D"
);
assert_eq!(syn.received_at_ms, expected_d, "entry clock must be D");
assert_eq!(syn.received_at_ms, syn.timestamp_unix_ms);
let commit = entries
.iter()
.find(|e| e.message_id == "commit-1")
.expect("commitment entry");
assert!(
commit.received_at_ms > expected_d + 100,
"the trigger must be well after D for this test to mean anything \
(trigger {} vs D {expected_d})",
commit.received_at_ms
);
}
#[tokio::test]
async fn synthetic_timestamp_excludes_a_pause_inside_the_window() {
let h = make_harness();
let sid = session_with_offer(&h.rt).await;
h.rt.suspend_session(&sid, "hold", OWNER)
.await
.expect("suspend");
tokio::time::sleep(std::time::Duration::from_millis(TIMEOUT_MS as u64 + 40)).await;
h.rt.resume_session(&sid, "go", OWNER)
.await
.expect("resume");
let session = h.rt.get_session_checked(&sid).await.unwrap();
assert!(
session.accumulated_suspended_ms >= TIMEOUT_MS,
"the pause must dominate the window for this test to bite"
);
sleep_past_the_deadline().await;
h.rt.process(
&env("Commitment", "commit-1", &sid, OWNER, commitment("1.0.0")),
None,
)
.await
.expect("resolves once enough unsuspended time has elapsed");
let entries = log_of(&h.rt, &sid).await;
let syn = synthetic_of(&entries);
let naive = offer_received_at(&entries) + TIMEOUT_MS;
assert_eq!(syn.received_at_ms, syn.timestamp_unix_ms);
assert!(
syn.timestamp_unix_ms > naive,
"the pause inside the window must push D out (D {} vs naive {naive})",
syn.timestamp_unix_ms
);
let resume_at = entries
.iter()
.find(|e| e.message_type == "SessionResume")
.expect("resume entry")
.received_at_ms;
assert!(
syn.timestamp_unix_ms >= resume_at,
"with the whole pre-pause run shorter than the timeout, D must fall after the resume"
);
let syn_idx = entries
.iter()
.position(|e| e.message_id == SYNTHETIC_ID)
.unwrap();
let resume_idx = entries
.iter()
.position(|e| e.message_type == "SessionResume")
.unwrap();
assert!(syn_idx > resume_idx);
let live = h.rt.get_session_checked(&sid).await.unwrap();
assert_replay_matches(&live, &replay_live_log(&h, &sid).await);
}
#[tokio::test]
async fn live_history_replays_byte_identically() {
let h = make_harness();
let sid = session_with_offer(&h.rt).await;
sleep_past_the_deadline().await;
h.rt.process(
&env("Commitment", "commit-1", &sid, OWNER, commitment("1.0.0")),
None,
)
.await
.expect("resolves");
let live = h.rt.get_session_checked(&sid).await.unwrap();
let replayed = replay_live_log(&h, &sid).await;
assert_eq!(live.state, SessionState::Resolved);
assert_replay_matches(&live, &replayed);
assert!(replayed.seen_message_ids.contains(SYNTHETIC_ID));
}
#[tokio::test]
async fn non_commitment_message_triggers_synthesis() {
let h = make_harness();
let sid = session_with_offer(&h.rt).await;
sleep_past_the_deadline().await;
h.rt.process(
&env("HandoffContext", "ctx-1", &sid, OWNER, context_payload()),
None,
)
.await
.expect("context is accepted at any disposition");
let entries = log_of(&h.rt, &sid).await;
let ids: Vec<&str> = incoming(&entries)
.iter()
.map(|e| e.message_id.as_str())
.collect();
assert_eq!(
ids,
vec!["start-1", "offer-1", SYNTHETIC_ID, "ctx-1"],
"a non-Commitment trigger must synthesize, and synthesize first"
);
let live = h.rt.get_session_checked(&sid).await.unwrap();
assert_eq!(
live.state,
SessionState::Open,
"the session is still open — synthesis is not resolution"
);
assert_replay_matches(&live, &replay_live_log(&h, &sid).await);
}
#[tokio::test]
async fn late_explicit_accept_loses_to_history_order() {
let h = make_harness();
let sid = session_with_offer(&h.rt).await;
sleep_past_the_deadline().await;
let err =
h.rt.process(
&env(
"HandoffAccept",
"acc-1",
&sid,
TARGET,
accept_payload(false),
),
None,
)
.await
.expect_err("the offer is no longer Offered by the time this is dispatched");
assert_eq!(err.to_string(), "InvalidPayload");
let entries = log_of(&h.rt, &sid).await;
let ids: Vec<&str> = incoming(&entries)
.iter()
.map(|e| e.message_id.as_str())
.collect();
assert_eq!(
ids,
vec!["start-1", "offer-1", SYNTHETIC_ID],
"the synthetic is in history; the losing explicit accept is not"
);
let live = h.rt.get_session_checked(&sid).await.unwrap();
let state: serde_json::Value = serde_json::from_slice(&live.mode_state).unwrap();
assert_eq!(state["offers"]["h1"]["disposition"], "Accepted");
assert_eq!(
state["offers"]["h1"]["outcome_reason"], "implicit accept (timeout)",
"the implicit accept must own the outcome, not the late explicit one"
);
assert_replay_matches(&live, &replay_live_log(&h, &sid).await);
}
#[tokio::test]
async fn squatter_cannot_block_synthesis() {
let h = make_harness();
let sid = session_with_offer(&h.rt).await;
let err =
h.rt.process(
&env(
"HandoffContext",
SYNTHETIC_ID,
&sid,
OWNER,
context_payload(),
),
None,
)
.await
.expect_err("the reserved namespace is closed to clients");
assert_eq!(err.to_string(), "InvalidEnvelope");
let live = h.rt.get_session_checked(&sid).await.unwrap();
assert!(
!live.seen_message_ids.contains(SYNTHETIC_ID),
"a rejected squat must not consume the id the runtime needs"
);
sleep_past_the_deadline().await;
let result =
h.rt.process(
&env("Commitment", "commit-1", &sid, OWNER, commitment("1.0.0")),
None,
)
.await
.expect("the squat must not have stranded the session");
assert_eq!(result.session_state, SessionState::Resolved);
let entries = log_of(&h.rt, &sid).await;
synthetic_of(&entries);
let live = h.rt.get_session_checked(&sid).await.unwrap();
assert_replay_matches(&live, &replay_live_log(&h, &sid).await);
}
#[tokio::test]
async fn rejected_trigger_leaves_dedup_intact_and_snapshot_current() {
let h = make_durable_harness();
let sid = session_with_offer(&h.rt).await;
sleep_past_the_deadline().await;
let err =
h.rt.process(
&env("Commitment", "commit-1", &sid, OWNER, commitment("9.9.9")),
None,
)
.await
.expect_err("mode_version does not match the session binding");
assert_eq!(err.to_string(), "InvalidPayload");
let live = h.rt.get_session_checked(&sid).await.unwrap();
let entries = log_of(&h.rt, &sid).await;
synthetic_of(&entries);
let ids: Vec<&str> = incoming(&entries)
.iter()
.map(|e| e.message_id.as_str())
.collect();
assert_eq!(ids, vec!["start-1", "offer-1", SYNTHETIC_ID]);
let snapshot =
h.rt.storage
.load_session(&sid)
.await
.expect("snapshot load")
.expect("a snapshot must exist");
assert_eq!(
snapshot.mode_state, live.mode_state,
"the synthesis must persist its own snapshot"
);
assert_eq!(snapshot.seen_message_ids, live.seen_message_ids);
assert!(snapshot.seen_message_ids.contains(SYNTHETIC_ID));
assert!(
!live.seen_message_ids.contains("commit-1"),
"a rejected message must not consume its dedup slot"
);
let result =
h.rt.process(
&env("Commitment", "commit-1", &sid, OWNER, commitment("1.0.0")),
None,
)
.await
.expect("the same message_id must still be usable after a rejection");
assert_eq!(result.session_state, SessionState::Resolved);
let live = h.rt.get_session_checked(&sid).await.unwrap();
assert_eq!(
live.seen_message_ids
.iter()
.filter(|id| id.as_str() == SYNTHETIC_ID)
.count(),
1
);
assert_replay_matches(&live, &replay_live_log(&h, &sid).await);
}
#[tokio::test]
async fn snapshot_is_current_immediately_after_the_rejected_trigger() {
let h = make_durable_harness();
let sid = session_with_offer(&h.rt).await;
sleep_past_the_deadline().await;
h.rt.process(
&env("Commitment", "commit-1", &sid, OWNER, commitment("9.9.9")),
None,
)
.await
.expect_err("rejected");
let live = h.rt.get_session_checked(&sid).await.unwrap();
let snapshot =
h.rt.storage
.load_session(&sid)
.await
.expect("snapshot load")
.expect("a snapshot must exist");
assert_eq!(
snapshot.mode_state, live.mode_state,
"the synthesis's own save is the only thing that can have written this"
);
assert!(snapshot.seen_message_ids.contains(SYNTHETIC_ID));
assert_eq!(snapshot.seen_message_ids, live.seen_message_ids);
let replayed = replay_live_log(&h, &sid).await;
assert_eq!(replayed.mode_state, snapshot.mode_state);
assert_replay_matches(&live, &replayed);
}
async fn seed_from_fixture(h: &Harness, sid: &str, session: &Session, log: &[LogEntry]) {
h.rt.registry
.insert_recovered_session(sid.to_string(), session.clone())
.await;
h.rt.log_store.replace_session_log(sid, log.to_vec()).await;
}
#[tokio::test]
async fn sweep_settles_the_offer_with_no_further_message() {
let h = make_harness();
let sid = session_with_offer(&h.rt).await;
sleep_past_the_deadline().await;
let before = log_of(&h.rt, &sid).await;
let ids_before: Vec<&str> = incoming(&before)
.iter()
.map(|e| e.message_id.as_str())
.collect();
assert_eq!(
ids_before,
vec!["start-1", "offer-1"],
"the deadline must still be unobserved before the sweep runs"
);
assert_eq!(
h.rt.sweep_due_synthetic_accepts().await,
1,
"the sweep must report the one entry it appended"
);
let after = log_of(&h.rt, &sid).await;
let ids_after: Vec<&str> = incoming(&after)
.iter()
.map(|e| e.message_id.as_str())
.collect();
assert_eq!(
ids_after,
vec!["start-1", "offer-1", SYNTHETIC_ID],
"the sweep must append the synthetic accept, and nothing else"
);
assert_eq!(synthetic_of(&after).entry_kind, EntryKind::Incoming);
let live = h.rt.get_session_checked(&sid).await.unwrap();
assert!(
live.seen_message_ids.contains(SYNTHETIC_ID),
"the synthetic's deterministic id must hold its dedup slot"
);
assert_eq!(
live.state,
SessionState::Open,
"synthesis is not resolution — the session stays open for its Commitment"
);
assert_replay_matches(&live, &replay_live_log(&h, &sid).await);
assert_eq!(
h.rt.sweep_due_synthetic_accepts().await,
0,
"a second sweep must be a no-op"
);
assert_eq!(log_of(&h.rt, &sid).await.len(), after.len());
h.rt.process(
&env("Commitment", "commit-1", &sid, OWNER, commitment("1.0.0")),
None,
)
.await
.expect("the swept accept must satisfy commitment_ready");
}
#[tokio::test]
async fn sweep_leaves_an_offer_alone_before_its_deadline() {
let h = make_harness();
let sid = session_with_offer(&h.rt).await;
assert_eq!(h.rt.sweep_due_synthetic_accepts().await, 0);
let entries = log_of(&h.rt, &sid).await;
assert!(
!entries.iter().any(|e| e.message_id == SYNTHETIC_ID),
"an offer inside its window must survive a sweep"
);
sleep_past_the_deadline().await;
assert_eq!(h.rt.sweep_due_synthetic_accepts().await, 1);
}
#[tokio::test]
async fn eager_and_lazy_synthesis_write_byte_identical_entries() {
let origin = make_harness();
let sid = session_with_offer(&origin.rt).await;
let fixture_session = origin.rt.get_session_checked(&sid).await.unwrap();
let fixture_log = log_of(&origin.rt, &sid).await;
assert!(
!fixture_log.iter().any(|e| e.message_id == SYNTHETIC_ID),
"the fixture must be captured before anything synthesizes"
);
sleep_past_the_deadline().await;
let eager = make_harness();
seed_from_fixture(&eager, &sid, &fixture_session, &fixture_log).await;
let observed_by_sweep = chrono::Utc::now().timestamp_millis();
assert_eq!(eager.rt.sweep_due_synthetic_accepts().await, 1);
let from_sweep = synthetic_of(&log_of(&eager.rt, &sid).await).clone();
tokio::time::sleep(std::time::Duration::from_millis(TIMEOUT_MS as u64)).await;
let lazy = make_harness();
seed_from_fixture(&lazy, &sid, &fixture_session, &fixture_log).await;
let observed_by_trigger = chrono::Utc::now().timestamp_millis();
lazy.rt
.process(
&env("Commitment", "commit-1", &sid, OWNER, commitment("9.9.9")),
None,
)
.await
.expect_err("the trigger is rejected; the synthesis ahead of it is not");
let from_lazy = synthetic_of(&log_of(&lazy.rt, &sid).await).clone();
assert!(
observed_by_trigger - observed_by_sweep >= TIMEOUT_MS,
"the two observations must be far enough apart to be distinguishable"
);
assert_eq!(
serde_json::to_string(&from_sweep).expect("entry serializes"),
serde_json::to_string(&from_lazy).expect("entry serializes"),
"the eager and lazy paths must write the same history bytes"
);
assert_eq!(from_sweep.message_id, from_lazy.message_id);
assert_eq!(from_sweep.sender, from_lazy.sender);
assert_eq!(from_sweep.raw_payload, from_lazy.raw_payload);
assert_eq!(from_sweep.entry_kind, from_lazy.entry_kind);
let expected_d = offer_received_at(&fixture_log) + TIMEOUT_MS;
assert_eq!(from_sweep.timestamp_unix_ms, expected_d);
assert_eq!(from_sweep.received_at_ms, expected_d);
assert!(
expected_d < observed_by_sweep,
"D must precede both observations, or this proves nothing"
);
let eager_session = eager.rt.get_session_checked(&sid).await.unwrap();
let lazy_session = lazy.rt.get_session_checked(&sid).await.unwrap();
assert_eq!(
eager_session.mode_state, lazy_session.mode_state,
"mode_state must be byte-identical across the two paths"
);
}
#[tokio::test]
async fn sweep_skips_a_suspended_session() {
let h = make_harness();
let sid = session_with_offer(&h.rt).await;
h.rt.suspend_session(&sid, "hold", OWNER)
.await
.expect("suspend");
sleep_past_the_deadline().await;
sleep_past_the_deadline().await;
assert_eq!(
h.rt.sweep_due_synthetic_accepts().await,
0,
"a suspended session must not be swept"
);
let entries = log_of(&h.rt, &sid).await;
assert!(
!entries.iter().any(|e| e.message_id == SYNTHETIC_ID),
"a suspended session must not gain a synthetic entry"
);
let paused = h.rt.get_session_checked(&sid).await.unwrap();
assert_eq!(paused.state, SessionState::Suspended);
assert!(!paused.seen_message_ids.contains(SYNTHETIC_ID));
h.rt.resume_session(&sid, "go", OWNER)
.await
.expect("resume");
assert_eq!(
h.rt.sweep_due_synthetic_accepts().await,
0,
"the suspended interval must not count toward the timeout"
);
sleep_past_the_deadline().await;
assert_eq!(h.rt.sweep_due_synthetic_accepts().await, 1);
let live = h.rt.get_session_checked(&sid).await.unwrap();
assert!(live.seen_message_ids.contains(SYNTHETIC_ID));
assert_replay_matches(&live, &replay_live_log(&h, &sid).await);
}
#[tokio::test]
async fn sweep_skips_a_terminal_session() {
let h = make_harness();
let sid = session_with_offer(&h.rt).await;
sleep_past_the_deadline().await;
h.rt.process(
&env("Commitment", "commit-1", &sid, OWNER, commitment("1.0.0")),
None,
)
.await
.expect("resolves");
let resolved_len = log_of(&h.rt, &sid).await.len();
assert_eq!(
h.rt.sweep_due_synthetic_accepts().await,
0,
"a resolved session must not be swept"
);
assert_eq!(
log_of(&h.rt, &sid).await.len(),
resolved_len,
"the sweep must not append to a terminal session's history"
);
let sid2 = session_with_offer(&h.rt).await;
h.rt.cancel_session(&sid2, "done", OWNER)
.await
.expect("cancel");
sleep_past_the_deadline().await;
assert_eq!(h.rt.sweep_due_synthetic_accepts().await, 0);
assert!(
!log_of(&h.rt, &sid2)
.await
.iter()
.any(|e| e.message_id == SYNTHETIC_ID),
"a cancelled session must not gain a synthetic entry"
);
}