#![allow(clippy::expect_used, clippy::unwrap_used)]
mod support;
use std::error::Error;
use std::time::{Duration, Instant};
use frame_conv::{Anomaly, ConversationHandle, InboundRequest, RequestOutcome};
use serde::{Deserialize, Serialize};
use support::{FileStore, QUANTUM, RunningServer, attachment, store_dir};
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
struct AddContact {
name: String,
priority: u32,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
struct ContactAdded {
name: String,
total: u32,
}
#[test]
fn request_reply_round_trip_by_content() -> Result<(), Box<dyn Error>> {
let server = RunningServer::start()?;
let stores = store_dir("req-roundtrip")?;
let (mut responder, responder_grant) = ConversationHandle::open(
&attachment(server.endpoint()),
FileStore::new(stores.join("responder.lpcr")),
)?;
let conversation = responder_grant.conversation;
let responder_thread = std::thread::spawn(move || -> Result<u32, String> {
let mut handled = 0;
let inbound = responder
.next_request::<AddContact>(Duration::from_secs(25))
.map_err(|error| error.to_string())?;
if let Some(InboundRequest::Valid(request)) = inbound {
responder
.reply(
request.correlation,
&ContactAdded {
name: request.body.name,
total: 1,
},
)
.map_err(|error| error.to_string())?;
handled += 1;
} else {
return Err(format!("responder expected a valid request: {inbound:?}"));
}
Ok(handled)
});
let (mut requester, _grant) = ConversationHandle::join(
&attachment(server.endpoint()),
conversation,
FileStore::new(stores.join("requester.lpcr")),
)?;
let outcome = requester.request::<AddContact, ContactAdded>(
&AddContact {
name: "ada".to_owned(),
priority: 7,
},
Duration::from_secs(25),
)?;
let RequestOutcome::Replied {
reply, responder, ..
} = outcome
else {
return Err(format!("expected the correlated reply, observed {outcome:?}").into());
};
assert_eq!(
reply,
ContactAdded {
name: "ada".to_owned(),
total: 1,
},
"reply content diverged"
);
assert_eq!(
responder, responder_grant.participant,
"reply must carry the responder's verified identity"
);
let handled = responder_thread
.join()
.expect("responder thread must not panic")?;
assert_eq!(handled, 1, "the handler must run exactly once");
let counters = requester.anomaly_counters();
assert_eq!(counters.duplicate_replies, 0);
assert_eq!(counters.late_replies, 0);
assert_eq!(counters.gaps, 0);
std::fs::remove_dir_all(&stores)?;
server.shutdown()?;
Ok(())
}
#[test]
fn no_responder_deadline_elapses_typed() -> Result<(), Box<dyn Error>> {
let server = RunningServer::start()?;
let stores = store_dir("req-noresponder")?;
let (mut requester, _grant) = ConversationHandle::open(
&attachment(server.endpoint()),
FileStore::new(stores.join("lonely.lpcr")),
)?;
let deadline = Duration::from_secs(6);
let started = Instant::now();
let outcome = requester.request::<AddContact, ContactAdded>(
&AddContact {
name: "nobody".to_owned(),
priority: 1,
},
deadline,
)?;
let waited = started.elapsed();
let RequestOutcome::DeadlineElapsed { deadline: named } = outcome else {
return Err(format!("expected the deadline outcome, observed {outcome:?}").into());
};
assert_eq!(
named, deadline,
"the outcome must name the caller's deadline"
);
assert!(
waited >= deadline,
"returned before the deadline: {waited:?}"
);
assert!(
waited < deadline + QUANTUM + Duration::from_secs(2),
"the deadline wall drifted past one quantum of tail: {waited:?}"
);
std::fs::remove_dir_all(&stores)?;
server.shutdown()?;
Ok(())
}
#[test]
fn duplicate_reply_is_typed_anomaly_first_unaffected() -> Result<(), Box<dyn Error>> {
let server = RunningServer::start()?;
let stores = store_dir("req-duplicate")?;
let (mut responder, responder_grant) = ConversationHandle::open(
&attachment(server.endpoint()),
FileStore::new(stores.join("responder.lpcr")),
)?;
let conversation = responder_grant.conversation;
let responder_thread = std::thread::spawn(move || -> Result<(), String> {
let inbound = responder
.next_request::<AddContact>(Duration::from_secs(25))
.map_err(|error| error.to_string())?;
let Some(InboundRequest::Valid(request)) = inbound else {
return Err(format!("responder expected a valid request: {inbound:?}"));
};
responder
.reply(
request.correlation,
&ContactAdded {
name: request.body.name.clone(),
total: 1,
},
)
.map_err(|error| error.to_string())?;
responder
.reply(
request.correlation,
&ContactAdded {
name: request.body.name,
total: 2,
},
)
.map_err(|error| error.to_string())?;
Ok(())
});
let (mut requester, _grant) = ConversationHandle::join(
&attachment(server.endpoint()),
conversation,
FileStore::new(stores.join("requester.lpcr")),
)?;
let outcome = requester.request::<AddContact, ContactAdded>(
&AddContact {
name: "ada".to_owned(),
priority: 1,
},
Duration::from_secs(25),
)?;
let RequestOutcome::Replied { reply, .. } = outcome else {
return Err(format!("expected the first reply, observed {outcome:?}").into());
};
assert_eq!(reply.total, 1, "the FIRST reply must win, unaffected");
responder_thread
.join()
.expect("responder thread must not panic")?;
let mut anomalies = Vec::new();
let pump_until = Instant::now() + 2 * QUANTUM + Duration::from_secs(2);
while anomalies.is_empty() && Instant::now() < pump_until {
let _quiet = requester.next_event::<ContactAdded>(Duration::from_secs(1))?;
anomalies.extend(requester.drain_anomalies());
}
assert!(
anomalies
.iter()
.any(|anomaly| matches!(anomaly, Anomaly::DuplicateReply { .. })),
"the duplicate reply must surface on the typed anomaly queue: {anomalies:?}"
);
assert_eq!(
requester.anomaly_counters().duplicate_replies,
1,
"the named duplicate counter must read exactly one"
);
std::fs::remove_dir_all(&stores)?;
server.shutdown()?;
Ok(())
}
#[test]
fn late_reply_after_deadline_has_a_defined_observable_fate() -> Result<(), Box<dyn Error>> {
let server = RunningServer::start()?;
let stores = store_dir("req-late")?;
let (mut responder, responder_grant) = ConversationHandle::open(
&attachment(server.endpoint()),
FileStore::new(stores.join("responder.lpcr")),
)?;
let conversation = responder_grant.conversation;
let responder_thread = std::thread::spawn(move || -> Result<(), String> {
let inbound = responder
.next_request::<AddContact>(Duration::from_secs(25))
.map_err(|error| error.to_string())?;
let Some(InboundRequest::Valid(request)) = inbound else {
return Err(format!("responder expected a valid request: {inbound:?}"));
};
std::thread::sleep(Duration::from_secs(13));
responder
.reply(
request.correlation,
&ContactAdded {
name: request.body.name,
total: 1,
},
)
.map_err(|error| error.to_string())?;
Ok(())
});
let (mut requester, _grant) = ConversationHandle::join(
&attachment(server.endpoint()),
conversation,
FileStore::new(stores.join("requester.lpcr")),
)?;
let outcome = requester.request::<AddContact, ContactAdded>(
&AddContact {
name: "ada".to_owned(),
priority: 1,
},
Duration::from_secs(6),
)?;
assert!(
matches!(outcome, RequestOutcome::DeadlineElapsed { .. }),
"the deadline must win: {outcome:?}"
);
responder_thread
.join()
.expect("responder thread must not panic")?;
let mut saw_late = false;
let pump_until = Instant::now() + 2 * QUANTUM + Duration::from_secs(2);
while !saw_late && Instant::now() < pump_until {
let _quiet = requester.next_event::<ContactAdded>(Duration::from_secs(1))?;
saw_late = requester
.drain_anomalies()
.iter()
.any(|anomaly| matches!(anomaly, Anomaly::LateReply { .. }));
}
assert!(saw_late, "the late reply must surface on the anomaly queue");
assert_eq!(requester.anomaly_counters().late_replies, 1);
std::fs::remove_dir_all(&stores)?;
server.shutdown()?;
Ok(())
}
#[test]
fn reply_racing_deadline_yields_exactly_one_outcome() -> Result<(), Box<dyn Error>> {
let server = RunningServer::start()?;
let stores = store_dir("req-race")?;
let (mut responder, responder_grant) = ConversationHandle::open(
&attachment(server.endpoint()),
FileStore::new(stores.join("responder.lpcr")),
)?;
let conversation = responder_grant.conversation;
let responder_thread = std::thread::spawn(move || -> Result<(), String> {
let inbound = responder
.next_request::<AddContact>(Duration::from_secs(25))
.map_err(|error| error.to_string())?;
let Some(InboundRequest::Valid(request)) = inbound else {
return Err(format!("responder expected a valid request: {inbound:?}"));
};
std::thread::sleep(Duration::from_secs(10));
responder
.reply(
request.correlation,
&ContactAdded {
name: request.body.name,
total: 1,
},
)
.map_err(|error| error.to_string())?;
Ok(())
});
let (mut requester, _grant) = ConversationHandle::join(
&attachment(server.endpoint()),
conversation,
FileStore::new(stores.join("requester.lpcr")),
)?;
let outcome = requester.request::<AddContact, ContactAdded>(
&AddContact {
name: "ada".to_owned(),
priority: 1,
},
Duration::from_secs(6),
)?;
responder_thread
.join()
.expect("responder thread must not panic")?;
match outcome {
RequestOutcome::Replied { reply, .. } => {
assert_eq!(reply.total, 1);
let _quiet = requester.next_event::<ContactAdded>(QUANTUM)?;
assert_eq!(requester.anomaly_counters().late_replies, 0);
assert_eq!(requester.anomaly_counters().duplicate_replies, 0);
}
RequestOutcome::DeadlineElapsed { .. } => {
let mut late = 0;
let pump_until = Instant::now() + 2 * QUANTUM + Duration::from_secs(2);
while late == 0 && Instant::now() < pump_until {
let _quiet = requester.next_event::<ContactAdded>(Duration::from_secs(1))?;
late = requester.anomaly_counters().late_replies;
}
assert_eq!(late, 1, "the losing reply must land as one late anomaly");
}
RequestOutcome::ResponderFailed { .. } => {
return Err(
format!("no responder failure exists in this race (ASK-4): {outcome:?}").into(),
);
}
}
std::fs::remove_dir_all(&stores)?;
server.shutdown()?;
Ok(())
}
#[test]
fn elapsed_wait_is_benign_quiet_rearm() -> Result<(), Box<dyn Error>> {
let server = RunningServer::start()?;
let stores = store_dir("req-quiet")?;
let (mut handle, _grant) = ConversationHandle::open(
&attachment(server.endpoint()),
FileStore::new(stores.join("quiet.lpcr")),
)?;
let quiet_request = handle.next_request::<AddContact>(Duration::from_secs(1))?;
assert!(
quiet_request.is_none(),
"an elapsed request wait must be benign quiet"
);
let quiet_event = handle.next_event::<ContactAdded>(Duration::from_secs(1))?;
assert!(
quiet_event.is_none(),
"an elapsed event wait must be benign quiet"
);
let receipt = handle.publish_event(&ContactAdded {
name: "still-alive".to_owned(),
total: 1,
})?;
assert!(receipt.seq.value() > 0);
assert!(handle.attached(), "quiet waits must not detach the handle");
std::fs::remove_dir_all(&stores)?;
server.shutdown()?;
Ok(())
}
#[test]
fn schema_invalid_request_surfaces_typed_at_responder() -> Result<(), Box<dyn Error>> {
let server = RunningServer::start()?;
let stores = store_dir("req-schema")?;
let (mut responder, responder_grant) = ConversationHandle::open(
&attachment(server.endpoint()),
FileStore::new(stores.join("responder.lpcr")),
)?;
let conversation = responder_grant.conversation;
let responder_thread = std::thread::spawn(move || -> Result<bool, String> {
let inbound = responder
.next_request::<AddContact>(Duration::from_secs(25))
.map_err(|error| error.to_string())?;
let Some(InboundRequest::SchemaInvalid {
correlation,
detail,
..
}) = inbound
else {
return Err(format!("expected the typed schema refusal: {inbound:?}"));
};
if detail.is_empty() {
return Err("the schema refusal must carry its exact detail".to_owned());
}
responder
.reply(
correlation,
&ContactAdded {
name: "schema-refused".to_owned(),
total: 0,
},
)
.map_err(|error| error.to_string())?;
Ok(true)
});
let (mut requester, _grant) = ConversationHandle::join(
&attachment(server.endpoint()),
conversation,
FileStore::new(stores.join("requester.lpcr")),
)?;
let outcome = requester.request::<serde_json::Value, ContactAdded>(
&serde_json::json!("not-an-add-contact"),
Duration::from_secs(25),
)?;
let RequestOutcome::Replied { reply, .. } = outcome else {
return Err(format!("expected the typed refusal reply, observed {outcome:?}").into());
};
assert_eq!(reply.name, "schema-refused");
let saw_refusal = responder_thread
.join()
.expect("responder thread must not panic")?;
assert!(saw_refusal);
std::fs::remove_dir_all(&stores)?;
server.shutdown()?;
Ok(())
}