#![allow(clippy::expect_used, clippy::unwrap_used)]
mod support;
use std::error::Error;
use std::fs;
use std::time::{Duration, Instant};
use frame_conv::{
AttachError, ConversationHandle, DepartReason, LeaveOutcome, RefusalClass, SubscriptionItem,
};
use serde::{Deserialize, Serialize};
use support::{FileStore, RunningServer, attachment, store_dir};
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
struct Ping {
note: String,
}
const OBSERVE: Duration = Duration::from_secs(12);
#[test]
fn explicit_leave_commits_and_peers_observe_the_departure() -> Result<(), Box<dyn Error>> {
let server = RunningServer::start()?;
let stores = store_dir("leave-observed")?;
let (mut observer, _observer_grant) = ConversationHandle::open(
&attachment(server.endpoint()),
FileStore::new(stores.join("observer.lpcr")),
)?;
let conversation = observer.conversation();
let (mut leaver, leaver_grant) = ConversationHandle::join(
&attachment(server.endpoint()),
conversation,
FileStore::new(stores.join("leaver.lpcr")),
)?;
let outcome = leaver.leave()?;
let LeaveOutcome::Left { seq } = outcome else {
return Err(format!("live-binding leave did not commit: {outcome:?}").into());
};
let started = Instant::now();
let mut saw_join = false;
loop {
if started.elapsed() > OBSERVE {
return Err(format!("observer never saw the departure (saw_join: {saw_join})").into());
}
match observer.next_event::<Ping>(OBSERVE)? {
Some(SubscriptionItem::PeerJoined { peer, .. }) => {
assert_eq!(peer, leaver_grant.participant);
saw_join = true;
}
Some(SubscriptionItem::PeerDeparted {
peer,
seq: observed_seq,
reason,
}) => {
assert_eq!(peer, leaver_grant.participant);
assert_eq!(reason, DepartReason::Left);
assert_eq!(
observed_seq, seq,
"the observed Left record and the leaver's commit receipt \
are two ends of one record"
);
break;
}
Some(other) => {
return Err(format!("unexpected item before the departure: {other:?}").into());
}
None => {}
}
}
server.shutdown()?;
Ok(())
}
#[test]
fn leave_is_idempotent_within_one_handle() -> Result<(), Box<dyn Error>> {
let server = RunningServer::start()?;
let stores = store_dir("leave-idempotent")?;
let (mut handle, _grant) = ConversationHandle::open(
&attachment(server.endpoint()),
FileStore::new(stores.join("participant.lpcr")),
)?;
let first = handle.leave()?;
assert!(
matches!(first, LeaveOutcome::Left { .. }),
"first leave must commit: {first:?}"
);
let second = handle.leave()?;
assert_eq!(
second,
LeaveOutcome::AlreadyLeft,
"second leave on the same handle must answer AlreadyLeft"
);
server.shutdown()?;
Ok(())
}
#[test]
fn leave_resumed_echoes_the_commit_across_restart() -> Result<(), Box<dyn Error>> {
let server = RunningServer::start()?;
let stores = store_dir("leave-echo")?;
let (mut handle, grant) = ConversationHandle::open(
&attachment(server.endpoint()),
FileStore::new(stores.join("participant.lpcr")),
)?;
let pre_leave_state = fs::read(stores.join("participant.lpcr"))?;
let LeaveOutcome::Left { seq } = handle.leave()? else {
return Err("first leave must commit".into());
};
drop(handle);
let outcome = ConversationHandle::leave_resumed(
&attachment(server.endpoint()),
&grant,
&pre_leave_state,
FileStore::new(stores.join("resumed.lpcr")),
)?;
assert_eq!(
outcome,
LeaveOutcome::Left { seq },
"the resumed leave must echo the original commit's sequence"
);
server.shutdown()?;
Ok(())
}
#[test]
fn leave_resumed_on_a_never_left_lineage_commits() -> Result<(), Box<dyn Error>> {
let server = RunningServer::start()?;
let stores = store_dir("leave-neverleft")?;
let (handle, grant) = ConversationHandle::open(
&attachment(server.endpoint()),
FileStore::new(stores.join("participant.lpcr")),
)?;
let state = fs::read(stores.join("participant.lpcr"))?;
drop(handle);
let outcome = ConversationHandle::leave_resumed(
&attachment(server.endpoint()),
&grant,
&state,
FileStore::new(stores.join("resumed.lpcr")),
)?;
assert!(
matches!(outcome, LeaveOutcome::Left { .. }),
"a never-left lineage must commit its leave: {outcome:?}"
);
server.shutdown()?;
Ok(())
}
#[test]
fn resume_after_leave_is_a_typed_identity_refusal() -> Result<(), Box<dyn Error>> {
let server = RunningServer::start()?;
let stores = store_dir("leave-resume-refused")?;
let (mut handle, grant) = ConversationHandle::open(
&attachment(server.endpoint()),
FileStore::new(stores.join("participant.lpcr")),
)?;
let pre_leave_state = fs::read(stores.join("participant.lpcr"))?;
let LeaveOutcome::Left { .. } = handle.leave()? else {
return Err("leave must commit".into());
};
drop(handle);
let refusal = ConversationHandle::resume(
&attachment(server.endpoint()),
&grant,
&pre_leave_state,
FileStore::new(stores.join("resumed.lpcr")),
)
.err();
let Some(AttachError::Refused { class, detail }) = refusal else {
return Err(format!("post-leave resume must be a typed refusal, got: {refusal:?}").into());
};
assert_eq!(class, RefusalClass::Identity, "detail: {detail}");
assert!(
detail.contains("ParticipantUnknown"),
"the exact upstream answer is preserved: {detail}"
);
server.shutdown()?;
Ok(())
}