use bytes::Bytes;
use serde_json::json;
use std::collections::{BTreeMap, BTreeSet};
use std::time::Duration;
use unb_core::{
ApplicationFailure, ApplicationFrame, ApplicationOrigin, ApplicationResponse,
ApplicationResult, BodyId, CapacityResult, ClientDelivery, CoreEffect, CoreError, CoreInput,
CorrelationId, Detail, DiscoverPlan, EffectId, Envelope, ErrorCode, Kind, Mode, NodeIdentity,
OperationInput, PeerAdmission, ProtocolCore, RelayOpenResult, Resolution, RetirementReason,
RouteAck, RouteAckStatus, RouteAdvertisement, RouteDelta, RouteSnapshot, Scope, SendResult,
SessionClass, SessionId, StreamKey, PROTOCOL_VERSION,
};
use web_time::Instant;
fn hello(id: &str) -> Envelope {
Envelope {
v: PROTOCOL_VERSION,
id: id.into(),
target: String::new(),
subject: String::new(),
kind: Kind::Hello,
corr: None,
seq: None,
hops: None,
body_token: None,
payload: Envelope::encode_payload(&json!({ "versions": [PROTOCOL_VERSION] })),
path: Vec::new(),
headers: Default::default(),
}
}
fn frame(id: &str, kind: Kind) -> Envelope {
Envelope {
v: PROTOCOL_VERSION,
id: id.into(),
target: if kind.is_application_request() {
"node".into()
} else {
String::new()
},
subject: String::new(),
kind,
corr: kind.is_application_request().then(|| format!("{id}-corr")),
seq: None,
hops: None,
body_token: None,
payload: Bytes::new(),
path: Vec::new(),
headers: Default::default(),
}
}
fn application_response(response: http::Response<Bytes>) -> ApplicationResponse {
let (parts, body) = response.into_parts();
ApplicationResponse {
head: http::Response::from_parts(parts, ()),
body: (!body.is_empty()).then(|| BodyId::from("test-body")),
}
}
fn drain(core: &mut ProtocolCore) -> Vec<CoreEffect> {
std::iter::from_fn(|| core.poll_effect()).collect()
}
fn without_route_state(effects: Vec<CoreEffect>) -> Vec<CoreEffect> {
effects
.into_iter()
.filter(|effect| {
!matches!(
effect,
CoreEffect::RouteSnapshotApplied { .. }
| CoreEffect::RouteDeltaApplied { .. }
| CoreEffect::RouteSessionWithdrawn { .. }
)
})
.collect()
}
fn frames(effects: &[CoreEffect]) -> Vec<Envelope> {
effects
.iter()
.filter_map(|effect| match effect {
CoreEffect::SendFrame { envelope, .. } | CoreEffect::SendProtocol { envelope, .. } => {
Some(envelope.clone())
}
CoreEffect::Send { frame, .. } => Some(frame.clone().into_envelope()),
_ => None,
})
.collect()
}
fn decisions(effects: Vec<CoreEffect>) -> Vec<CoreEffect> {
effects
.into_iter()
.filter(|effect| !matches!(effect, CoreEffect::ScheduleSessionDeadline { .. }))
.collect()
}
fn established_pair() -> (ProtocolCore, SessionId, SessionId, Instant) {
let now = Instant::now();
let dialer = SessionId::from("dialer");
let listener = SessionId::from("listener");
let mut core = ProtocolCore::new("node");
core.handle(
now,
CoreInput::SessionOpened {
session: dialer.clone(),
initiator: true,
establish_peer: true,
expected_peer: None,
},
)
.unwrap();
core.handle(
now,
CoreInput::SessionOpened {
session: listener.clone(),
initiator: false,
establish_peer: true,
expected_peer: None,
},
)
.unwrap();
let hello = frames(&drain(&mut core)).remove(0);
core.handle(
now,
CoreInput::FrameReceived {
session: listener.clone(),
envelope: hello,
},
)
.unwrap();
let listener_effects = drain(&mut core);
let welcome = frames(&listener_effects).remove(0);
assert!(listener_effects.iter().any(|effect| matches!(
effect,
CoreEffect::HandshakeEstablished { session, version: PROTOCOL_VERSION } if session == &listener
)));
core.handle(
now,
CoreInput::FrameReceived {
session: dialer.clone(),
envelope: welcome,
},
)
.unwrap();
assert!(drain(&mut core).iter().any(|effect| matches!(
effect,
CoreEffect::HandshakeEstablished { session, version: PROTOCOL_VERSION } if session == &dialer
)));
(core, dialer, listener, now)
}
fn established_listener() -> (ProtocolCore, SessionId, Instant) {
let now = Instant::now();
let session = SessionId::from("listener");
let mut core = ProtocolCore::new("node");
core.handle(
now,
CoreInput::SessionOpened {
session: session.clone(),
initiator: false,
establish_peer: true,
expected_peer: None,
},
)
.unwrap();
core.handle(
now,
CoreInput::FrameReceived {
session: session.clone(),
envelope: hello("hello"),
},
)
.unwrap();
drain(&mut core);
(core, session, now)
}
#[test]
fn identity_aware_construction_preserves_the_complete_local_identity() {
let identity = NodeIdentity {
node_id: "node".into(),
instance_id: "instance".into(),
epoch: 7,
proof: json!({"token": "proof"}),
};
let core = ProtocolCore::with_identity(identity.clone());
assert_eq!(core.node().identity(), identity);
}
fn receive(core: &mut ProtocolCore, session: &SessionId, now: Instant, mut envelope: Envelope) {
let input = if (envelope.kind.is_application_request() && envelope.kind != Kind::Discover)
|| envelope.kind.is_application_response()
|| envelope.kind == Kind::Cancel
{
if envelope.kind != Kind::Error && !envelope.payload.is_empty() {
envelope.body_token = Some("test-body".into());
envelope.payload = Bytes::new();
}
CoreInput::ApplicationFrameReceived {
session: session.clone(),
frame: ApplicationFrame::from_envelope(&envelope).unwrap(),
}
} else {
CoreInput::FrameReceived {
session: session.clone(),
envelope,
}
};
core.handle(now, input).unwrap();
}
fn add_client(core: &mut ProtocolCore, now: Instant, session: &str) -> SessionId {
let session = SessionId::from(session);
core.handle(
now,
CoreInput::SessionOpened {
session: session.clone(),
initiator: false,
establish_peer: true,
expected_peer: None,
},
)
.unwrap();
receive(core, &session, now, hello(&format!("hello-{session}")));
drain(core);
session
}
fn opening(id: &str, target: &str, hops: Option<u8>) -> Envelope {
opening_kind(id, target, Kind::Request, hops)
}
fn opening_kind(id: &str, target: &str, kind: Kind, hops: Option<u8>) -> Envelope {
let mut envelope = frame(id, kind);
envelope.target = target.into();
envelope.subject = "service".into();
envelope.hops = hops;
envelope
}
fn paired_relay() -> (ProtocolCore, StreamKey, StreamKey, Instant) {
let now = Instant::now();
let mut core = ProtocolCore::new("relay");
let peer = add_ready_peer(&mut core, now, "peer-session", "peer");
import_route(&mut core, &peer, now, 2, route("owner", "peer", "owner", 1));
let client = add_client(&mut core, now, "client");
receive(&mut core, &client, now, opening("source", "owner", None));
let CoreEffect::OpenRelay { effect, source, .. } = decisions(drain(&mut core)).remove(0) else {
panic!("routed opening must request a relay")
};
let corr = core
.open_stream_body(
&peer,
"/target-node/service",
Kind::Request,
None,
None,
Default::default(),
)
.unwrap();
drain(&mut core);
let target = StreamKey {
session: peer,
corr: CorrelationId::from(corr),
};
core.handle(
now,
CoreInput::RelayOpenCompleted {
effect,
result: RelayOpenResult::Opened(target.clone()),
},
)
.unwrap();
(core, source, target, now)
}
fn assert_protocol_close(effects: &[CoreEffect]) {
assert!(effects.iter().any(|effect| matches!(
effect,
CoreEffect::CloseTransport {
code: ErrorCode::Protocol,
..
}
)));
}
#[test]
fn first_application_opening_permanently_classifies_client() {
let (mut core, session, now) = established_listener();
receive(&mut core, &session, now, frame("request", Kind::Request));
assert_eq!(core.session_class(&session).unwrap(), SessionClass::Client);
drain(&mut core);
receive(&mut core, &session, now, frame("identify", Kind::Identify));
assert_eq!(core.session_class(&session).unwrap(), SessionClass::Client);
assert_protocol_close(&drain(&mut core));
}
#[test]
fn identify_permanently_classifies_node_candidate() {
let (mut core, session, now) = established_listener();
receive(&mut core, &session, now, frame("identify", Kind::Identify));
assert_eq!(
core.session_class(&session).unwrap(),
SessionClass::NodeCandidate
);
}
#[test]
fn ping_and_pong_leave_session_unclassified() {
let (mut core, session, now) = established_listener();
receive(&mut core, &session, now, frame("ping", Kind::Ping));
receive(&mut core, &session, now, frame("pong", Kind::Pong));
assert_eq!(
core.session_class(&session).unwrap(),
SessionClass::Unclassified
);
}
#[test]
fn control_after_client_classification_is_a_protocol_close() {
let (mut core, session, now) = established_listener();
receive(&mut core, &session, now, frame("request", Kind::Request));
drain(&mut core);
receive(
&mut core,
&session,
now,
frame("snapshot", Kind::RouteSnapshot),
);
assert_eq!(core.session_class(&session).unwrap(), SessionClass::Client);
assert_protocol_close(&drain(&mut core));
}
#[test]
fn application_opening_while_node_candidate_is_not_ready_is_a_protocol_close() {
let (mut core, session, now) = established_listener();
receive(&mut core, &session, now, frame("identify", Kind::Identify));
drain(&mut core);
receive(&mut core, &session, now, frame("request", Kind::Request));
assert_eq!(
core.session_class(&session).unwrap(),
SessionClass::NodeCandidate
);
assert_protocol_close(&drain(&mut core));
}
#[test]
fn cancelling_an_in_flight_dispatch_aborts_it_and_ignores_the_late_result() {
let (mut core, session, now) = established_listener();
receive(&mut core, &session, now, frame("request", Kind::Request));
let effects = drain(&mut core);
let CoreEffect::CheckDispatchCapacity {
effect: capacity, ..
} = &effects[0]
else {
panic!("application opening must request capacity")
};
core.handle(
now,
CoreInput::CapacityChecked {
effect: *capacity,
result: CapacityResult::Available,
},
)
.unwrap();
let effects = drain(&mut core);
let CoreEffect::InvokeApplication { effect, .. } = &effects[0] else {
panic!("available capacity must request application invocation")
};
let dispatch = *effect;
let mut cancel = frame("cancel", Kind::Cancel);
cancel.corr = Some("request-corr".into());
receive(&mut core, &session, now, cancel);
let effects = drain(&mut core);
assert!(
effects.iter().any(|effect| matches!(
effect,
CoreEffect::AbortDispatch { effect, .. } if *effect == dispatch
)),
"cancelling the stream must abort its in-flight dispatch: {effects:?}"
);
core.handle(
now,
CoreInput::DispatchCompleted {
effect: dispatch,
result: Ok(ApplicationResult::Response(application_response(
http::Response::builder()
.body(Bytes::from_static(b"late"))
.unwrap(),
))),
},
)
.unwrap();
assert!(
drain(&mut core).is_empty(),
"a late result for an aborted dispatch must be ignored"
);
}
#[test]
fn session_loss_aborts_every_in_flight_dispatch() {
let (mut core, session, now) = established_listener();
receive(&mut core, &session, now, frame("request", Kind::Request));
let effects = drain(&mut core);
let CoreEffect::CheckDispatchCapacity {
effect: capacity, ..
} = &effects[0]
else {
panic!("application opening must request capacity")
};
core.handle(
now,
CoreInput::CapacityChecked {
effect: *capacity,
result: CapacityResult::Available,
},
)
.unwrap();
let effects = drain(&mut core);
let CoreEffect::InvokeApplication { effect, .. } = &effects[0] else {
panic!("available capacity must request application invocation")
};
let dispatch = *effect;
core.handle(now, CoreInput::SessionClosed { session })
.unwrap();
let effects = drain(&mut core);
assert!(
effects.iter().any(|effect| matches!(
effect,
CoreEffect::AbortDispatch { effect, .. } if *effect == dispatch
)),
"session loss must abort its in-flight dispatches: {effects:?}"
);
}
#[test]
fn application_policy_and_send_outcomes_are_correlated_core_inputs_and_effects() {
let (mut core, session, now) = established_listener();
receive(&mut core, &session, now, frame("request", Kind::Request));
let effects = drain(&mut core);
let CoreEffect::CheckDispatchCapacity {
effect: capacity,
stream,
frame,
} = &effects[0]
else {
panic!("application opening must request capacity")
};
assert_eq!(stream.session, session);
assert_eq!(stream.corr.as_str(), "request-corr");
assert_eq!(frame.head.kind, Kind::Request);
core.handle(
now,
CoreInput::CapacityChecked {
effect: *capacity,
result: CapacityResult::Available,
},
)
.unwrap();
let effects = drain(&mut core);
let CoreEffect::InvokeApplication { effect, invocation } = &effects[0] else {
panic!("available capacity must request application invocation")
};
assert_eq!(invocation.stream, *stream);
assert_eq!(invocation.reservation, *capacity);
assert_eq!(
invocation.origin,
ApplicationOrigin::Client {
session: session.clone()
}
);
let dispatch = *effect;
core.handle(
now,
CoreInput::DispatchCompleted {
effect: dispatch,
result: Ok(ApplicationResult::Response(application_response(
http::Response::builder()
.header("x-result", "preserved")
.body(Bytes::from_static(b"ok"))
.unwrap(),
))),
},
)
.unwrap();
let effects = drain(&mut core);
let CoreEffect::Send {
effect: send,
session: target,
frame,
} = &effects[0]
else {
panic!("application response must request a correlated send")
};
assert_eq!(target, &session);
assert_eq!(frame.head.kind, Kind::Response);
assert_eq!(frame.body.as_ref().map(BodyId::as_str), Some("test-body"));
assert_eq!(frame.head.headers["x-result"], "preserved");
let send = *send;
core.handle(
now,
CoreInput::SendCompleted {
effect: send,
result: SendResult::Reserved,
},
)
.unwrap();
assert_eq!(core.poll_effect(), None);
core.handle(
now,
CoreInput::SendCompleted {
effect: send,
result: SendResult::Written,
},
)
.unwrap();
assert_eq!(core.poll_effect(), None);
core.handle(
now,
CoreInput::SendCompleted {
effect: send,
result: SendResult::Written,
},
)
.unwrap();
assert_eq!(core.poll_effect(), None);
}
#[test]
fn send_reservation_timeout_returns_busy_without_retiring_the_session() {
let (mut core, session, now) = established_listener();
receive(&mut core, &session, now, frame("request", Kind::Request));
let CoreEffect::CheckDispatchCapacity { effect, .. } = drain(&mut core).remove(0) else {
panic!("request must check capacity")
};
core.handle(
now,
CoreInput::CapacityChecked {
effect,
result: CapacityResult::Available,
},
)
.unwrap();
let CoreEffect::InvokeApplication { effect, .. } = drain(&mut core).remove(0) else {
panic!("capacity must invoke")
};
core.handle(
now,
CoreInput::DispatchCompleted {
effect,
result: Ok(ApplicationResult::Event(application_response(
http::Response::new(Bytes::new()),
))),
},
)
.unwrap();
let CoreEffect::Send { effect, .. } = drain(&mut core).remove(0) else {
panic!("event must send")
};
core.handle(
now,
CoreInput::SendCompleted {
effect,
result: SendResult::ReservationTimedOut,
},
)
.unwrap();
let effects = drain(&mut core);
assert!(frames(&effects).iter().any(|envelope| {
envelope.kind == Kind::Error && envelope.payload_json()["code"] == json!(ErrorCode::Busy)
}));
assert!(!effects
.iter()
.any(|effect| matches!(effect, CoreEffect::SessionRetired { .. })));
}
#[test]
fn capacity_and_application_failures_are_constructed_by_core() {
let (mut core, session, now) = established_listener();
receive(&mut core, &session, now, frame("busy", Kind::Request));
let CoreEffect::CheckDispatchCapacity { effect, .. } = drain(&mut core).remove(0) else {
panic!("application opening must request capacity")
};
core.handle(
now,
CoreInput::CapacityChecked {
effect,
result: CapacityResult::Busy,
},
)
.unwrap();
assert!(frames(&drain(&mut core))
.iter()
.any(|frame| { frame.kind == Kind::Error && frame.payload_json()["code"] == "BUSY" }));
let (mut core, session, now) = established_listener();
receive(&mut core, &session, now, frame("failed", Kind::Request));
let CoreEffect::CheckDispatchCapacity { effect, .. } = drain(&mut core).remove(0) else {
panic!("application opening must request capacity")
};
core.handle(
now,
CoreInput::CapacityChecked {
effect,
result: CapacityResult::Available,
},
)
.unwrap();
let CoreEffect::InvokeApplication { effect, .. } = drain(&mut core).remove(0) else {
panic!("available capacity must request invocation")
};
core.handle(
now,
CoreInput::DispatchCompleted {
effect,
result: Err(ApplicationFailure {
code: ErrorCode::Unauthorized,
message: "denied".into(),
}),
},
)
.unwrap();
assert!(frames(&drain(&mut core)).iter().any(|frame| {
frame.kind == Kind::Error && frame.payload_json()["code"] == "UNAUTHORIZED"
}));
}
#[test]
fn repeated_application_events_keep_dispatch_pending_until_terminal() {
let (mut core, session, now) = established_listener();
receive(&mut core, &session, now, frame("stream", Kind::Subscribe));
let CoreEffect::CheckDispatchCapacity { effect, .. } = drain(&mut core).remove(0) else {
panic!("stream opening must request capacity")
};
core.handle(
now,
CoreInput::CapacityChecked {
effect,
result: CapacityResult::Available,
},
)
.unwrap();
let CoreEffect::InvokeApplication { effect, .. } = drain(&mut core).remove(0) else {
panic!("capacity must produce invocation")
};
for body in ["one", "two"] {
core.handle(
now,
CoreInput::DispatchCompleted {
effect,
result: Ok(ApplicationResult::Event(application_response(
http::Response::builder()
.header("x-event", body)
.body(Bytes::copy_from_slice(body.as_bytes()))
.unwrap(),
))),
},
)
.unwrap();
let CoreEffect::Send { frame, .. } = drain(&mut core).remove(0) else {
panic!("event must produce send")
};
assert_eq!(frame.head.kind, Kind::Event);
assert_eq!(frame.head.headers["x-event"], body);
}
core.handle(
now,
CoreInput::DispatchCompleted {
effect,
result: Ok(ApplicationResult::Finished(application_response(
http::Response::builder()
.header("x-finished", "yes")
.body(Bytes::from_static(b"done"))
.unwrap(),
))),
},
)
.unwrap();
let CoreEffect::Send { frame, .. } = drain(&mut core).remove(0) else {
panic!("finished must produce terminal send")
};
assert_eq!(frame.head.kind, Kind::Response);
assert_eq!(frame.head.headers["x-finished"], "yes");
}
#[test]
fn send_failure_uses_pending_stream_context_and_retires_deterministically() {
let (mut core, session, now) = established_listener();
receive(&mut core, &session, now, frame("request", Kind::Request));
let CoreEffect::CheckDispatchCapacity { effect, .. } = drain(&mut core).remove(0) else {
panic!("request must check capacity")
};
core.handle(
now,
CoreInput::CapacityChecked {
effect,
result: CapacityResult::Available,
},
)
.unwrap();
let CoreEffect::InvokeApplication { effect, .. } = drain(&mut core).remove(0) else {
panic!("capacity must invoke")
};
core.handle(
now,
CoreInput::DispatchCompleted {
effect,
result: Ok(ApplicationResult::Response(application_response(
http::Response::new(Bytes::new()),
))),
},
)
.unwrap();
let CoreEffect::Send { effect, .. } = drain(&mut core).remove(0) else {
panic!("response must send")
};
core.handle(
now,
CoreInput::SendCompleted {
effect,
result: SendResult::Closed,
},
)
.unwrap();
assert!(matches!(
without_route_state(drain(&mut core)).as_slice(),
[CoreEffect::SessionRetired {
reason: RetirementReason::SendClosed,
..
}]
));
}
#[test]
fn identity_requests_typed_admission_and_late_results_do_nothing() {
let (mut core, session, now) = established_listener();
let identity = NodeIdentity {
node_id: "peer".into(),
instance_id: "peer-instance".into(),
epoch: 1,
proof: serde_json::Value::Null,
};
let mut identify = frame("identify", Kind::Identify);
identify.payload = serde_json::to_vec(&identity).unwrap().into();
receive(&mut core, &session, now, identify);
let CoreEffect::RequestPeerAdmission {
effect,
session: target,
remote,
} = drain(&mut core).remove(0)
else {
panic!("identity must request peer admission")
};
assert_eq!(target, session);
assert_eq!(remote, identity);
core.handle(
now,
CoreInput::PeerAdmissionCompleted {
effect,
result: PeerAdmission::Admitted(identity),
},
)
.unwrap();
drain(&mut core);
core.handle(
now,
CoreInput::CapacityChecked {
effect: EffectId::new(u64::MAX),
result: CapacityResult::Busy,
},
)
.unwrap();
assert_eq!(core.poll_effect(), None);
}
fn control(id: &str, kind: Kind, value: serde_json::Value) -> Envelope {
let mut envelope = frame(id, kind);
envelope.payload = Envelope::encode_payload(&value);
envelope
}
fn establish_peer(
core: &mut ProtocolCore,
session: &SessionId,
now: Instant,
identity: NodeIdentity,
) -> Vec<CoreEffect> {
receive(
core,
session,
now,
control(
"identify",
Kind::Identify,
serde_json::to_value(&identity).unwrap(),
),
);
let CoreEffect::RequestPeerAdmission { effect, .. } = drain(core).remove(0) else {
panic!("identity must request admission")
};
core.handle(
now,
CoreInput::PeerAdmissionCompleted {
effect,
result: PeerAdmission::Admitted(identity),
},
)
.unwrap();
drain(core);
receive(
core,
session,
now,
control("accepted", Kind::IdentityAccepted, serde_json::Value::Null),
);
drain(core);
let snapshot = unb_core::RouteSnapshot::canonical(1, Vec::new());
receive(
core,
session,
now,
control(
"snapshot",
Kind::RouteSnapshot,
serde_json::to_value(snapshot).unwrap(),
),
);
drain(core);
receive(
core,
session,
now,
control(
"ack",
Kind::RouteAck,
serde_json::json!({ "generation": 1, "status": "applied" }),
),
);
drain(core)
}
fn add_ready_peer(core: &mut ProtocolCore, now: Instant, session: &str, peer: &str) -> SessionId {
let session = SessionId::from(session);
core.handle(
now,
CoreInput::SessionOpened {
session: session.clone(),
initiator: false,
establish_peer: true,
expected_peer: None,
},
)
.unwrap();
receive(core, &session, now, hello(&format!("hello-{peer}")));
drain(core);
establish_peer(
core,
&session,
now,
NodeIdentity {
node_id: peer.into(),
instance_id: format!("{peer}-instance"),
epoch: 1,
proof: serde_json::Value::Null,
},
);
session
}
fn import_route(
core: &mut ProtocolCore,
session: &SessionId,
now: Instant,
generation: u64,
route: RouteAdvertisement,
) {
receive(
core,
session,
now,
control(
&format!("snapshot-{generation}"),
Kind::RouteSnapshot,
serde_json::to_value(RouteSnapshot::canonical(generation, vec![route])).unwrap(),
),
);
drain(core);
}
fn imported_route(
advertiser: &str,
owner: &str,
instance: &str,
epoch: u64,
revision: u64,
) -> RouteAdvertisement {
RouteAdvertisement {
destination: owner.into(),
owner: owner.into(),
owner_instance: instance.into(),
owner_epoch: epoch,
owner_revision: revision,
distance: 1,
path: vec![owner.into(), advertiser.into()],
}
}
fn route(destination: &str, advertiser: &str, owner: &str, distance: u32) -> RouteAdvertisement {
let mut path = vec![owner.into()];
if distance > 0 {
path.extend((1..distance).map(|hop| format!("hop-{hop}")));
path.push(advertiser.into());
}
RouteAdvertisement {
destination: destination.into(),
owner: owner.into(),
owner_instance: format!("{owner}-instance"),
owner_epoch: 1,
owner_revision: 1,
distance,
path,
}
}
fn route_deltas(effects: &[CoreEffect]) -> BTreeMap<String, RouteDelta> {
effects
.iter()
.filter_map(|effect| match effect {
CoreEffect::SendFrame { session, envelope } if envelope.kind == Kind::RouteDelta => {
Some((session.to_string(), envelope.parse_payload().unwrap()))
}
_ => None,
})
.collect()
}
#[test]
fn relay_opening_resolves_local_route_conflict_and_unknown_in_core() {
let now = Instant::now();
let mut local = ProtocolCore::new("local");
local
.handle(
now,
CoreInput::LocalCapabilitiesInstalled {
capabilities: BTreeMap::from([("service".into(), json!({}))]),
},
)
.unwrap();
let client = add_client(&mut local, now, "client");
receive(&mut local, &client, now, opening("local", "local", None));
assert!(matches!(
decisions(drain(&mut local)).as_slice(),
[CoreEffect::CheckDispatchCapacity { .. }]
));
let mut routed = ProtocolCore::new("relay");
let peer = add_ready_peer(&mut routed, now, "peer-session", "peer");
import_route(
&mut routed,
&peer,
now,
2,
route("owner", "peer", "owner", 1),
);
let client = add_client(&mut routed, now, "client");
receive(
&mut routed,
&client,
now,
opening("route", "owner", Some(3)),
);
assert!(matches!(
decisions(drain(&mut routed)).as_slice(),
[CoreEffect::OpenRelay { peer, frame, .. }]
if peer == "peer" && frame.head.hops == Some(2)
));
let mut conflicted = ProtocolCore::new("relay");
let left = add_ready_peer(&mut conflicted, now, "left-session", "left");
let right = add_ready_peer(&mut conflicted, now, "right-session", "right");
import_route(
&mut conflicted,
&left,
now,
2,
imported_route("left", "owner", "owner-left", 1, 1),
);
import_route(
&mut conflicted,
&right,
now,
2,
imported_route("right", "owner", "owner-right", 1, 1),
);
let client = add_client(&mut conflicted, now, "client");
receive(
&mut conflicted,
&client,
now,
opening("conflict", "owner", None),
);
assert!(frames(&drain(&mut conflicted)).iter().any(|envelope| {
envelope.kind == Kind::Error
&& envelope.payload_json()["code"] == json!(ErrorCode::PeerUnreachable)
}));
let mut unknown = ProtocolCore::new("relay");
let client = add_client(&mut unknown, now, "client");
receive(
&mut unknown,
&client,
now,
opening("unknown", "missing", None),
);
let effects = drain(&mut unknown);
let CoreEffect::AwaitTargetReadiness { effect, .. } = effects
.iter()
.find(|effect| matches!(effect, CoreEffect::AwaitTargetReadiness { .. }))
.cloned()
.expect("an unavailable target must ask the server about relevant recovery")
else {
unreachable!()
};
unknown
.handle(
now,
CoreInput::TargetReadinessCompleted {
effect,
result: unb_core::TargetReadinessResult::Unavailable {
message: "no recovery source".into(),
},
},
)
.unwrap();
assert!(frames(&drain(&mut unknown)).iter().any(|envelope| {
envelope.kind == Kind::Error
&& envelope.payload_json()["code"] == json!(ErrorCode::PeerUnreachable)
}));
}
#[test]
fn relayed_opening_re_resolves_after_target_readiness_without_replaying() {
for kind in [Kind::Request, Kind::Subscribe, Kind::Channel] {
let now = Instant::now();
let mut core = ProtocolCore::new("relay");
let peer = add_ready_peer(&mut core, now, "peer-session", "peer");
let client = add_client(&mut core, now, "client");
receive(
&mut core,
&client,
now,
opening_kind("waiting", "owner", kind, None),
);
let effects = drain(&mut core);
let CoreEffect::AwaitTargetReadiness { effect, target, .. } = effects
.iter()
.find(|effect| matches!(effect, CoreEffect::AwaitTargetReadiness { .. }))
.cloned()
.unwrap()
else {
unreachable!()
};
assert_eq!(target, "owner");
import_route(&mut core, &peer, now, 2, route("owner", "peer", "owner", 1));
core.handle(
now,
CoreInput::TargetReadinessCompleted {
effect,
result: unb_core::TargetReadinessResult::Ready,
},
)
.unwrap();
assert!(matches!(
decisions(drain(&mut core)).as_slice(),
[CoreEffect::OpenRelay { peer, frame, .. }]
if peer == "peer"
&& frame.head.kind == kind
&& frame.head.target == "owner"
&& frame.head.subject == "service"
));
}
}
#[test]
fn closing_a_source_stream_cancels_target_wait_and_releases_its_body() {
let now = Instant::now();
let mut core = ProtocolCore::new("relay");
let client = add_client(&mut core, now, "client");
let mut request = opening("waiting-body", "owner", None);
request.payload = Bytes::from_static(b"held");
receive(&mut core, &client, now, request);
let effects = drain(&mut core);
assert!(effects
.iter()
.any(|effect| matches!(effect, CoreEffect::AwaitTargetReadiness { .. })));
core.handle(
now,
CoreInput::SessionClosed {
session: client.clone(),
},
)
.unwrap();
let effects = drain(&mut core);
assert!(effects.iter().any(|effect| matches!(
effect,
CoreEffect::AbortDispatch { session, .. } if session == &client
)));
assert!(effects.iter().any(|effect| matches!(
effect,
CoreEffect::ReleaseBody { session, body }
if session == &client && body.as_str() == "test-body"
)));
}
#[test]
fn closing_a_source_stream_cancels_a_not_yet_admitted_relay_body() {
let now = Instant::now();
let mut core = ProtocolCore::new("relay");
let peer = add_ready_peer(&mut core, now, "peer-session", "peer");
import_route(&mut core, &peer, now, 2, route("owner", "peer", "owner", 1));
let client = add_client(&mut core, now, "client");
let mut request = opening("waiting-relay-body", "owner", None);
request.payload = Bytes::from_static(b"held");
receive(&mut core, &client, now, request);
assert!(decisions(drain(&mut core))
.iter()
.any(|effect| matches!(effect, CoreEffect::OpenRelay { .. })));
core.handle(
now,
CoreInput::SessionClosed {
session: client.clone(),
},
)
.unwrap();
let effects = drain(&mut core);
assert!(effects.iter().any(|effect| matches!(
effect,
CoreEffect::AbortDispatch { session, .. } if session == &client
)));
assert!(effects.iter().any(|effect| matches!(
effect,
CoreEffect::ReleaseBody { session, body }
if session == &client && body.as_str() == "test-body"
)));
}
#[test]
fn relay_open_completion_records_paired_stream_keys_and_rejects_duplicate_pairing() {
let now = Instant::now();
let mut core = ProtocolCore::new("relay");
let peer = add_ready_peer(&mut core, now, "peer-session", "peer");
import_route(&mut core, &peer, now, 2, route("owner", "peer", "owner", 1));
let client = add_client(&mut core, now, "client");
receive(&mut core, &client, now, opening("first", "owner", None));
let CoreEffect::OpenRelay { effect, source, .. } = decisions(drain(&mut core)).remove(0) else {
panic!("routed opening must request a relay")
};
let target = StreamKey {
session: peer.clone(),
corr: CorrelationId::from("outbound"),
};
assert_ne!(
source.corr, target.corr,
"each relay leg must use a fresh hop-local correlation"
);
core.handle(
now,
CoreInput::RelayOpenCompleted {
effect,
result: RelayOpenResult::Opened(target.clone()),
},
)
.unwrap();
assert_eq!(core.relay_peer(&source), Some(&target));
assert_eq!(core.relay_peer(&target), Some(&source));
receive(&mut core, &client, now, opening("second", "owner", None));
let CoreEffect::OpenRelay {
effect,
source: duplicate,
..
} = decisions(drain(&mut core)).remove(0)
else {
panic!("second routed opening must request a relay")
};
core.handle(
now,
CoreInput::RelayOpenCompleted {
effect,
result: RelayOpenResult::Opened(target),
},
)
.unwrap();
assert!(frames(&drain(&mut core)).iter().any(|envelope| {
envelope.corr.as_deref() == Some(duplicate.corr.as_str())
&& envelope.kind == Kind::Error
&& envelope.payload_json()["code"] == json!(ErrorCode::Conflict)
}));
}
#[test]
fn relay_opening_hop_handling_is_deterministic() {
let now = Instant::now();
let mut core = ProtocolCore::new("relay");
let peer = add_ready_peer(&mut core, now, "peer-session", "peer");
import_route(&mut core, &peer, now, 2, route("owner", "peer", "owner", 1));
let client = add_client(&mut core, now, "client");
receive(&mut core, &client, now, opening("default", "owner", None));
assert!(matches!(
decisions(drain(&mut core)).as_slice(),
[CoreEffect::OpenRelay { frame, .. }]
if frame.head.hops == Some(unb_core::DEFAULT_HOPS - 1)
));
receive(
&mut core,
&client,
now,
opening("exhausted", "owner", Some(0)),
);
assert!(frames(&drain(&mut core)).iter().any(|envelope| {
envelope.kind == Kind::Error
&& envelope.payload_json()["code"] == json!(ErrorCode::HopLimitExceeded)
}));
}
#[test]
fn relay_opening_shares_opaque_payload_and_preserves_non_protocol_metadata() {
let now = Instant::now();
let mut core = ProtocolCore::new("relay");
let peer = add_ready_peer(&mut core, now, "peer-session", "peer");
import_route(&mut core, &peer, now, 2, route("owner", "peer", "owner", 1));
let client = add_client(&mut core, now, "client");
let mut original = opening("opaque", "owner", Some(5));
original.subject = "opaque.subject".into();
original.kind = Kind::Subscribe;
original.seq = Some(41);
original.path = vec!["custom-path".into()];
original.payload = Bytes::from(vec![0x5a; 1024 * 1024]);
original.headers = serde_json::Map::from_iter([
("x-custom".into(), json!("preserve")),
("x-status".into(), json!(207)),
]);
let expected = original.clone();
receive(&mut core, &client, now, original);
let CoreEffect::OpenRelay { frame, .. } = decisions(drain(&mut core)).remove(0) else {
panic!("routed opening must request a relay")
};
assert_eq!(frame.body.as_ref().map(BodyId::as_str), Some("test-body"));
assert_eq!(frame.head.id, expected.id);
assert_eq!(frame.head.subject, expected.subject);
assert_eq!(frame.head.kind, expected.kind);
assert_eq!(frame.head.corr, expected.corr);
assert_eq!(frame.head.seq, expected.seq);
assert_eq!(frame.head.hops, Some(4));
assert_eq!(frame.head.path, expected.path);
assert_eq!(frame.head.headers, expected.headers);
}
#[test]
fn relay_terminal_and_closure_transitions_use_normal_core_inputs() {
for kind in [Kind::Response, Kind::Error] {
let (mut core, source, target, now) = paired_relay();
let mut terminal = frame("terminal", kind);
terminal.corr = Some(target.corr.as_str().into());
receive(&mut core, &target.session, now, terminal);
let CoreEffect::ForwardRelay {
effect,
source: forwarded_source,
target: forwarded_target,
terminal: true,
..
} = decisions(drain(&mut core)).remove(0)
else {
panic!("terminal frame must forward through core")
};
assert_eq!(forwarded_source, target);
assert_eq!(forwarded_target, source);
assert_eq!(core.relay_peer(&forwarded_source), Some(&forwarded_target));
core.handle(
now,
CoreInput::RelayForwardCompleted {
effect,
result: SendResult::Written,
},
)
.unwrap();
assert_eq!(core.relay_peer(&forwarded_source), None);
assert_eq!(core.poll_effect(), None);
}
for from_source in [true, false] {
let (mut core, source, target, now) = paired_relay();
let sender = if from_source { &source } else { &target };
let mut cancel = frame("cancel", Kind::Cancel);
cancel.corr = Some(sender.corr.as_str().into());
receive(&mut core, &sender.session, now, cancel);
let CoreEffect::ForwardRelay {
effect,
terminal: true,
..
} = decisions(drain(&mut core)).remove(0)
else {
panic!("cancel must forward through core")
};
assert!(core.relay_peer(sender).is_some());
core.handle(
now,
CoreInput::RelayForwardCompleted {
effect,
result: SendResult::Written,
},
)
.unwrap();
assert_eq!(core.relay_peer(sender), None);
}
{
let (mut core, _source, target, now) = paired_relay();
let mut event = frame("evt", Kind::Event);
event.corr = Some(target.corr.as_str().into());
event.seq = Some(0);
receive(&mut core, &target.session, now, event);
let CoreEffect::ForwardRelay { effect, .. } = decisions(drain(&mut core)).remove(0) else {
panic!("event frame must forward through core")
};
core.handle(
now,
CoreInput::RelayForwardCompleted {
effect,
result: SendResult::Refused {
code: ErrorCode::PayloadTooLarge,
message: "raise the WS collect ceiling or reach this node over WebTransport"
.into(),
},
},
)
.unwrap();
let mut saw_refusal = false;
for effect in drain(&mut core) {
assert!(
!matches!(effect, CoreEffect::SessionRetired { .. }),
"a refusal must not retire sessions"
);
if let CoreEffect::Send { frame, .. } = &effect {
if frame.head.kind == Kind::Error {
let error = frame.head.error.as_ref().unwrap();
assert_eq!(error.code, ErrorCode::PayloadTooLarge);
assert!(error.message.contains("WebTransport"));
saw_refusal = true;
}
}
}
assert!(saw_refusal, "the refusal code must reach the source leg");
}
for result in [
SendResult::ReservationTimedOut,
SendResult::Closed,
SendResult::Cancelled,
SendResult::WriteFailed("write failed".into()),
] {
let (mut core, source, target, now) = paired_relay();
let mut terminal = frame("terminal", Kind::Response);
terminal.corr = Some(target.corr.as_str().into());
receive(&mut core, &target.session, now, terminal);
let CoreEffect::ForwardRelay {
effect,
source: forwarded_source,
target: forwarded_target,
terminal: true,
..
} = decisions(drain(&mut core)).remove(0)
else {
panic!("terminal frame must forward through core")
};
assert_eq!(core.relay_peer(&forwarded_source), Some(&forwarded_target));
core.handle(
now,
CoreInput::RelayForwardCompleted {
effect,
result: result.clone(),
},
)
.unwrap();
assert_eq!(core.relay_peer(&source), None);
assert_eq!(core.relay_peer(&target), None);
assert!(!drain(&mut core).is_empty());
core.handle(
now,
CoreInput::RelayForwardCompleted {
effect,
result: SendResult::Written,
},
)
.unwrap();
assert_eq!(core.poll_effect(), None);
}
for close_source in [true, false] {
let (mut core, source, target, now) = paired_relay();
let closed = if close_source { &source } else { &target };
core.handle(
now,
CoreInput::SessionClosed {
session: closed.session.clone(),
},
)
.unwrap();
assert_eq!(core.relay_peer(&source), None);
assert_eq!(core.relay_peer(&target), None);
assert!(!drain(&mut core).is_empty());
}
}
#[test]
fn session_loss_activates_the_selected_backup_and_exports_it_in_the_same_cycle() {
let now = Instant::now();
let mut core = ProtocolCore::new("node");
let primary = add_ready_peer(&mut core, now, "primary-session", "primary");
let backup = add_ready_peer(&mut core, now, "backup-session", "backup");
let observer = add_ready_peer(&mut core, now, "observer-session", "observer");
import_route(
&mut core,
&primary,
now,
2,
route("owner", "primary", "owner", 1),
);
import_route(
&mut core,
&backup,
now,
2,
route("owner", "backup", "owner", 2),
);
core.handle(
now,
CoreInput::SessionClosed {
session: primary.clone(),
},
)
.unwrap();
let effects = drain(&mut core);
assert_eq!(
core.node().resolve("owner"),
Resolution::Route("backup".into())
);
let deltas = route_deltas(&effects);
assert_eq!(
deltas[observer.as_str()].upsert[0].path,
["owner", "hop-1", "backup", "node"]
);
assert_eq!(deltas[backup.as_str()].withdraw[0].destination, "owner");
assert!(effects.iter().any(|effect| matches!(
effect,
CoreEffect::RouteSessionWithdrawn { session, changed: true } if session == &primary
)));
}
#[test]
fn final_path_loss_withdraws_the_selected_route_and_exports_in_the_same_cycle() {
let now = Instant::now();
let mut core = ProtocolCore::new("node");
let owner = add_ready_peer(&mut core, now, "owner-session", "owner");
let first = add_ready_peer(&mut core, now, "first-session", "first");
let second = add_ready_peer(&mut core, now, "second-session", "second");
import_route(
&mut core,
&owner,
now,
2,
route("owner", "owner", "owner", 0),
);
core.handle(now, CoreInput::SessionClosed { session: owner })
.unwrap();
let effects = drain(&mut core);
assert_eq!(core.node().resolve("owner"), Resolution::Unknown);
let deltas = route_deltas(&effects);
for observer in [first, second] {
assert_eq!(deltas[observer.as_str()].withdraw[0].destination, "owner");
}
}
#[test]
fn conflicting_owner_loss_recovers_selection_and_exports_in_the_same_cycle() {
let now = Instant::now();
let mut core = ProtocolCore::new("node");
let alpha = add_ready_peer(&mut core, now, "alpha-session", "alpha-hop");
let beta = add_ready_peer(&mut core, now, "beta-session", "beta-hop");
let observer = add_ready_peer(&mut core, now, "observer-session", "observer");
import_route(
&mut core,
&alpha,
now,
2,
imported_route("alpha-hop", "owner", "owner-alpha", 1, 1),
);
import_route(
&mut core,
&beta,
now,
2,
imported_route("beta-hop", "owner", "owner-beta", 1, 1),
);
assert!(matches!(
core.node().resolve("owner"),
Resolution::Conflicted { .. }
));
core.handle(now, CoreInput::SessionClosed { session: beta })
.unwrap();
let effects = drain(&mut core);
assert_eq!(
core.node().resolve("owner"),
Resolution::Route("alpha-hop".into())
);
let deltas = route_deltas(&effects);
assert_eq!(deltas[observer.as_str()].upsert[0].owner, "owner");
assert!(!deltas.contains_key(alpha.as_str()));
}
#[test]
fn session_loss_emits_every_split_horizon_export_before_the_handle_cycle_drains() {
let now = Instant::now();
let mut core = ProtocolCore::new("node");
let lost = add_ready_peer(&mut core, now, "lost-session", "lost");
let backup = add_ready_peer(&mut core, now, "backup-session", "backup");
let first = add_ready_peer(&mut core, now, "first-session", "first");
let second = add_ready_peer(&mut core, now, "second-session", "second");
import_route(&mut core, &lost, now, 2, route("owner", "lost", "owner", 1));
import_route(
&mut core,
&backup,
now,
2,
route("owner", "backup", "owner", 2),
);
core.handle(now, CoreInput::SessionClosed { session: lost })
.unwrap();
let effects = drain(&mut core);
let deltas = route_deltas(&effects);
assert_eq!(
deltas.keys().cloned().collect::<Vec<_>>(),
[backup.to_string(), first.to_string(), second.to_string()]
);
assert_eq!(deltas[backup.as_str()].withdraw[0].destination, "owner");
for observer in [first, second] {
assert_eq!(deltas[observer.as_str()].upsert.len(), 1);
assert!(deltas[observer.as_str()].withdraw.is_empty());
}
assert_eq!(core.poll_effect(), None);
}
fn route_ack(
core: &mut ProtocolCore,
session: &SessionId,
now: Instant,
generation: u64,
status: RouteAckStatus,
) -> Vec<CoreEffect> {
receive(
core,
session,
now,
control(
"route-ack",
Kind::RouteAck,
serde_json::to_value(RouteAck { generation, status }).unwrap(),
),
);
drain(core)
}
#[test]
fn route_export_initial_snapshot_is_generation_one() {
let (mut core, session, now) = established_listener();
let identity = NodeIdentity {
node_id: "peer".into(),
instance_id: "peer-instance".into(),
epoch: 1,
proof: serde_json::Value::Null,
};
receive(
&mut core,
&session,
now,
control(
"identify",
Kind::Identify,
serde_json::to_value(&identity).unwrap(),
),
);
let CoreEffect::RequestPeerAdmission { effect, .. } = drain(&mut core).remove(0) else {
panic!("identity must request admission")
};
core.handle(
now,
CoreInput::PeerAdmissionCompleted {
effect,
result: PeerAdmission::Admitted(identity),
},
)
.unwrap();
drain(&mut core);
receive(
&mut core,
&session,
now,
control("accepted", Kind::IdentityAccepted, serde_json::Value::Null),
);
let snapshot = frames(&drain(&mut core))
.into_iter()
.find(|frame| frame.kind == Kind::RouteSnapshot)
.unwrap()
.parse_payload::<RouteSnapshot>()
.unwrap();
assert_eq!(snapshot.generation, 1);
assert_eq!(snapshot.routes.len(), 1);
assert_eq!(snapshot.routes[0].destination, "node");
assert_eq!(snapshot.routes[0].owner, "node");
assert_eq!(snapshot.routes[0].distance, 0);
assert_eq!(snapshot.routes[0].path, ["node"]);
}
#[test]
fn catalog_changes_do_not_churn_node_route_exports() {
let now = Instant::now();
let mut core = ProtocolCore::new("node");
let _session = add_ready_peer(&mut core, now, "peer-session", "peer");
core.handle(
now,
CoreInput::LocalCapabilitiesInstalled {
capabilities: BTreeMap::new(),
},
)
.unwrap();
assert_eq!(core.poll_effect(), None);
let capabilities = BTreeMap::from([("alpha".into(), serde_json::json!({}))]);
core.handle(
now,
CoreInput::LocalCapabilitiesInstalled {
capabilities: capabilities.clone(),
},
)
.unwrap();
assert_eq!(core.poll_effect(), None);
core.handle(now, CoreInput::LocalCapabilitiesInstalled { capabilities })
.unwrap();
assert_eq!(core.poll_effect(), None);
core.handle(
now,
CoreInput::LocalCapabilitiesInstalled {
capabilities: BTreeMap::new(),
},
)
.unwrap();
assert_eq!(core.poll_effect(), None);
assert_eq!(core.node().resolve("alpha"), Resolution::Unknown);
}
#[test]
fn imported_selection_recomputes_split_horizon_exports_for_ready_peers() {
let now = Instant::now();
let mut core = ProtocolCore::new("observer");
let source = add_ready_peer(&mut core, now, "source-session", "source");
let target = add_ready_peer(&mut core, now, "target-session", "target");
let route = RouteAdvertisement {
destination: "owner".into(),
owner: "owner".into(),
owner_instance: "owner-instance".into(),
owner_epoch: 1,
owner_revision: 4,
distance: 1,
path: vec!["owner".into(), "source".into()],
};
receive(
&mut core,
&source,
now,
control(
"snapshot-2",
Kind::RouteSnapshot,
serde_json::to_value(RouteSnapshot::canonical(2, vec![route])).unwrap(),
),
);
let effects = drain(&mut core);
assert!(!effects.iter().any(|effect| matches!(
effect,
CoreEffect::SendFrame { session, envelope }
if session == &source && envelope.kind == Kind::RouteDelta
)));
let exported = effects
.iter()
.find_map(|effect| match effect {
CoreEffect::SendFrame { session, envelope }
if session == &target && envelope.kind == Kind::RouteDelta =>
{
Some(envelope.parse_payload::<RouteDelta>().unwrap())
}
_ => None,
})
.unwrap();
assert_eq!(exported.generation, 2);
assert_eq!(exported.upsert.len(), 1);
assert_eq!(exported.upsert[0].destination, "owner");
assert_eq!(exported.upsert[0].distance, 2);
assert_eq!(exported.upsert[0].path, vec!["owner", "source", "observer"]);
assert!(exported.withdraw.is_empty());
}
#[test]
fn route_export_applied_ack_replay_and_behind_ack_preserve_resync_state() {
let now = Instant::now();
let mut core = ProtocolCore::new("node");
let session = add_ready_peer(&mut core, now, "peer-session", "peer");
core.handle(
now,
CoreInput::LocalCapabilitiesInstalled {
capabilities: BTreeMap::from([("alpha".into(), serde_json::json!({}))]),
},
)
.unwrap();
drain(&mut core);
assert!(frames(&route_ack(
&mut core,
&session,
now,
1,
RouteAckStatus::Applied
))
.is_empty());
assert!(frames(&route_ack(
&mut core,
&session,
now,
0,
RouteAckStatus::Applied
))
.is_empty());
assert!(frames(&route_ack(
&mut core,
&session,
now,
0,
RouteAckStatus::ResyncRequired
))
.is_empty());
let recovery = route_ack(&mut core, &session, now, 1, RouteAckStatus::ResyncRequired);
let snapshot = frames(&recovery)
.remove(0)
.parse_payload::<RouteSnapshot>()
.unwrap();
assert_eq!(snapshot.generation, 2);
assert_eq!(snapshot.routes.len(), 1);
assert_eq!(snapshot.routes[0].destination, "node");
assert!(frames(&route_ack(
&mut core,
&session,
now,
1,
RouteAckStatus::ResyncRequired
))
.is_empty());
}
#[test]
fn route_export_future_ack_retires_and_accepted_resync_advances_generation() {
let now = Instant::now();
let mut core = ProtocolCore::new("node");
let session = add_ready_peer(&mut core, now, "peer-session", "peer");
let effects = route_ack(&mut core, &session, now, 1, RouteAckStatus::ResyncRequired);
let snapshot = frames(&effects)
.remove(0)
.parse_payload::<RouteSnapshot>()
.unwrap();
assert_eq!(snapshot.generation, 2);
assert!(frames(&route_ack(
&mut core,
&session,
now,
1,
RouteAckStatus::ResyncRequired
))
.is_empty());
let effects = route_ack(&mut core, &session, now, 3, RouteAckStatus::Applied);
assert!(effects.iter().any(|effect| matches!(
effect,
CoreEffect::SessionRetired {
session: retired,
reason: RetirementReason::AcknowledgementOrder,
} if retired == &session
)));
}
#[test]
fn ready_peer_imports_apply_owner_incarnation_freshness_before_conflicts() {
let now = Instant::now();
let mut core = ProtocolCore::new("observer");
let stale = add_ready_peer(&mut core, now, "stale", "path-a");
let fresh = add_ready_peer(&mut core, now, "fresh", "path-b");
import_route(
&mut core,
&stale,
now,
2,
imported_route("path-a", "owner", "owner-old", 4, 90),
);
import_route(
&mut core,
&fresh,
now,
2,
imported_route("path-b", "owner", "owner-new", 5, 1),
);
assert_eq!(
core.node().resolve("owner"),
Resolution::Route("path-b".into())
);
}
#[test]
fn ready_peer_imports_distinct_node_destinations_without_conflict() {
let now = Instant::now();
let mut core = ProtocolCore::new("observer");
let alpha = add_ready_peer(&mut core, now, "alpha", "path-a");
let beta = add_ready_peer(&mut core, now, "beta", "path-b");
import_route(
&mut core,
&alpha,
now,
2,
imported_route("path-a", "alpha", "alpha-instance", 2, 40),
);
import_route(
&mut core,
&beta,
now,
2,
imported_route("path-b", "beta", "beta-instance", 90, 3),
);
assert_eq!(
core.node().resolve("alpha"),
Resolution::Route("path-a".into())
);
assert_eq!(
core.node().resolve("beta"),
Resolution::Route("path-b".into())
);
}
#[test]
fn ready_peer_equal_freshness_path_selection_ignores_arrival_order() {
let selected = |order: [(&str, &str); 2]| {
let now = Instant::now();
let mut core = ProtocolCore::new("observer");
let first = add_ready_peer(&mut core, now, "first", order[0].0);
let second = add_ready_peer(&mut core, now, "second", order[1].0);
import_route(
&mut core,
&first,
now,
2,
imported_route(order[0].0, "owner", "owner-instance", 7, 11),
);
import_route(
&mut core,
&second,
now,
2,
imported_route(order[1].0, "owner", "owner-instance", 7, 11),
);
core.node().resolve("owner")
};
assert_eq!(
selected([("path-b", "b"), ("path-a", "a")]),
Resolution::Route("path-a".into())
);
assert_eq!(
selected([("path-a", "a"), ("path-b", "b")]),
Resolution::Route("path-a".into())
);
}
#[test]
fn ready_peer_stale_withdrawal_preserves_new_owner_incarnation() {
let now = Instant::now();
let mut core = ProtocolCore::new("observer");
let session = add_ready_peer(&mut core, now, "path", "path-a");
import_route(
&mut core,
&session,
now,
2,
imported_route("path-a", "owner", "owner-new", 8, 2),
);
receive(
&mut core,
&session,
now,
control(
"delta-3",
Kind::RouteDelta,
serde_json::to_value(RouteDelta {
generation: 3,
upsert: Vec::new(),
withdraw: vec![unb_core::RouteWithdrawal {
destination: "owner".into(),
owner: "owner".into(),
owner_instance: "owner-old".into(),
owner_epoch: 7,
owner_revision: 90,
}],
})
.unwrap(),
),
);
drain(&mut core);
assert_eq!(
core.node().resolve("owner"),
Resolution::Route("path-a".into())
);
}
#[test]
fn identity_snapshot_and_ack_sequence_promotes_only_when_complete() {
let (mut core, session, now) = established_listener();
let identity = NodeIdentity {
node_id: "peer".into(),
instance_id: "peer-instance".into(),
epoch: 1,
proof: serde_json::Value::Null,
};
let effects = establish_peer(&mut core, &session, now, identity.clone());
assert!(core.session_ready(&session).unwrap());
assert!(effects.iter().any(|effect| matches!(
effect,
CoreEffect::SessionEstablished { session: established, peer, .. }
if established == &session && peer == &identity
)));
}
#[test]
fn post_ready_route_snapshot_and_delta_are_validated_and_applied_in_core() {
let (mut core, session, now) = established_listener();
establish_peer(
&mut core,
&session,
now,
NodeIdentity {
node_id: "peer".into(),
instance_id: "peer-instance".into(),
epoch: 1,
proof: serde_json::Value::Null,
},
);
let route = RouteAdvertisement {
destination: "peer".into(),
owner: "peer".into(),
owner_instance: "peer-instance".into(),
owner_epoch: 1,
owner_revision: 1,
distance: 0,
path: vec!["peer".into()],
};
receive(
&mut core,
&session,
now,
control(
"snapshot-2",
Kind::RouteSnapshot,
serde_json::to_value(RouteSnapshot::canonical(2, vec![route.clone()])).unwrap(),
),
);
let effects = drain(&mut core);
assert_eq!(
core.node().resolve("peer"),
Resolution::Route("peer".into())
);
assert!(effects.iter().any(|effect| matches!(
effect,
CoreEffect::RouteSnapshotApplied { session: target, peer, snapshot, changed: true }
if target == &session && peer == "peer" && snapshot.generation == 2
)));
assert!(frames(&effects).iter().any(|envelope| {
envelope.kind == Kind::RouteAck
&& envelope.parse_payload::<RouteAck>().unwrap()
== RouteAck {
generation: 2,
status: RouteAckStatus::Applied,
}
}));
let delta = RouteDelta {
generation: 3,
upsert: Vec::new(),
withdraw: vec![unb_core::RouteWithdrawal {
destination: route.destination,
owner: route.owner,
owner_instance: route.owner_instance,
owner_epoch: route.owner_epoch,
owner_revision: route.owner_revision,
}],
};
receive(
&mut core,
&session,
now,
control(
"delta-3",
Kind::RouteDelta,
serde_json::to_value(delta).unwrap(),
),
);
let effects = drain(&mut core);
assert_eq!(core.node().resolve("peer"), Resolution::Unknown);
assert!(effects.iter().any(|effect| matches!(
effect,
CoreEffect::RouteDeltaApplied { delta, changed: true, .. } if delta.generation == 3
)));
}
#[test]
fn post_ready_route_gap_ack_and_malformed_control_decisions_are_core_owned() {
let (mut core, session, now) = established_listener();
establish_peer(
&mut core,
&session,
now,
NodeIdentity {
node_id: "peer".into(),
instance_id: "peer-instance".into(),
epoch: 1,
proof: serde_json::Value::Null,
},
);
receive(
&mut core,
&session,
now,
control(
"gap",
Kind::RouteDelta,
serde_json::to_value(RouteDelta {
generation: 4,
upsert: Vec::new(),
withdraw: Vec::new(),
})
.unwrap(),
),
);
assert!(frames(&drain(&mut core)).iter().any(|envelope| {
envelope.parse_payload::<RouteAck>().is_ok_and(|ack| {
ack == RouteAck {
generation: 1,
status: RouteAckStatus::ResyncRequired,
}
})
}));
receive(
&mut core,
&session,
now,
control(
"ack",
Kind::RouteAck,
json!({ "generation": 1, "status": "resync_required" }),
),
);
let recovery = frames(&drain(&mut core))
.into_iter()
.find(|envelope| envelope.kind == Kind::RouteSnapshot)
.unwrap()
.parse_payload::<RouteSnapshot>()
.unwrap();
assert_eq!(recovery.generation, 2);
assert_eq!(recovery.routes.len(), 1);
assert_eq!(recovery.routes[0].destination, "node");
let mut malformed = frame("malformed", Kind::RouteSnapshot);
malformed.payload = Bytes::from_static(b"{");
receive(&mut core, &session, now, malformed);
assert!(drain(&mut core).iter().any(|effect| matches!(
effect,
CoreEffect::SessionRetired {
session: target,
reason: RetirementReason::MalformedRouteControl,
} if target == &session
)));
}
#[test]
fn rejected_mismatched_and_timed_out_candidates_retire_in_core() {
let (mut core, session, now) = established_listener();
let identity = NodeIdentity {
node_id: "peer".into(),
instance_id: "peer-instance".into(),
epoch: 1,
proof: serde_json::Value::Null,
};
receive(
&mut core,
&session,
now,
control(
"identify",
Kind::Identify,
serde_json::to_value(&identity).unwrap(),
),
);
let CoreEffect::RequestPeerAdmission { effect, .. } = drain(&mut core).remove(0) else {
panic!("identity must request admission")
};
core.handle(
now,
CoreInput::PeerAdmissionCompleted {
effect,
result: PeerAdmission::Rejected("denied".into()),
},
)
.unwrap();
assert!(drain(&mut core)
.iter()
.any(|effect| matches!(effect, CoreEffect::SessionRetired { session: retired, .. } if retired == &session)));
core.handle(
now,
CoreInput::PeerAdmissionCompleted {
effect,
result: PeerAdmission::Admitted(identity.clone()),
},
)
.unwrap();
assert_eq!(core.poll_effect(), None);
let mismatch = SessionId::from("mismatch");
core.handle(
now,
CoreInput::SessionOpened {
session: mismatch.clone(),
initiator: false,
establish_peer: true,
expected_peer: Some("expected".into()),
},
)
.unwrap();
receive(&mut core, &mismatch, now, hello("mismatch-hello"));
drain(&mut core);
receive(
&mut core,
&mismatch,
now,
control(
"mismatch-identify",
Kind::Identify,
serde_json::to_value(&identity).unwrap(),
),
);
let CoreEffect::RequestPeerAdmission { effect, .. } = drain(&mut core).remove(0) else {
panic!("identity must request admission")
};
core.handle(
now,
CoreInput::PeerAdmissionCompleted {
effect,
result: PeerAdmission::Admitted(identity),
},
)
.unwrap();
assert!(drain(&mut core)
.iter()
.any(|effect| matches!(effect, CoreEffect::SessionRetired { session: retired, .. } if retired == &mismatch)));
let timeout = SessionId::from("timeout");
core.handle(
now,
CoreInput::SessionOpened {
session: timeout.clone(),
initiator: false,
establish_peer: true,
expected_peer: Some("expected".into()),
},
)
.unwrap();
core.handle(
now,
CoreInput::EstablishmentTimeout {
session: timeout.clone(),
},
)
.unwrap();
let effects = drain(&mut core);
assert!(effects
.iter()
.any(|effect| matches!(effect, CoreEffect::SessionRetired { session: retired, .. } if retired == &timeout)));
core.handle(
now,
CoreInput::EstablishmentTimeout {
session: timeout.clone(),
},
)
.unwrap();
assert_eq!(core.poll_effect(), None);
}
#[test]
fn duplicate_peer_arbitration_keeps_the_deterministic_direction() {
let now = Instant::now();
let inbound = SessionId::from("inbound");
let outbound = SessionId::from("outbound");
let mut core = ProtocolCore::new("alpha");
for (session, initiator) in [(&inbound, false), (&outbound, true)] {
core.handle(
now,
CoreInput::SessionOpened {
session: session.clone(),
initiator,
establish_peer: true,
expected_peer: None,
},
)
.unwrap();
if initiator {
let welcome = control(
"welcome",
Kind::Welcome,
serde_json::json!({ "version": PROTOCOL_VERSION }),
);
receive(&mut core, session, now, welcome);
} else {
receive(&mut core, session, now, hello("hello"));
}
drain(&mut core);
}
let identity = NodeIdentity {
node_id: "zeta".into(),
instance_id: "zeta-instance".into(),
epoch: 1,
proof: serde_json::Value::Null,
};
establish_peer(&mut core, &inbound, now, identity.clone());
let effects = establish_peer(&mut core, &outbound, now, identity);
assert!(core.session_ready(&outbound).unwrap());
assert!(!core.session_ready(&inbound).unwrap());
assert!(effects.iter().any(|effect| matches!(
effect,
CoreEffect::SessionRetired { session, .. } if session == &inbound
)));
}
fn discover_frame(id: &str, target: &str, scope: Scope, mode: Mode) -> Envelope {
let plan = DiscoverPlan {
discover_id: id.into(),
detail: Detail::Index,
scope,
hops: 8,
visited: BTreeSet::new(),
timeout_ms: Some(10),
mode,
};
let mut envelope = frame(id, Kind::Discover);
envelope.target = target.into();
envelope.hops = Some(plan.hops);
envelope.payload = serde_json::to_vec(&plan).unwrap().into();
envelope
}
#[test]
fn local_discovery_emits_catalog_and_finishes_in_core() {
let (mut core, session, now) = established_listener();
receive(
&mut core,
&session,
now,
discover_frame("local", "node", Scope::Local, Mode::PartialOk),
);
let effects = drain(&mut core);
assert!(effects.iter().any(|effect| matches!(
effect,
CoreEffect::SendProtocol { envelope, .. }
if envelope.kind == Kind::Event
&& envelope.payload_json()["type"] == "node_catalog"
&& envelope.payload_json()["node"] == "node"
)));
assert!(effects.iter().any(|effect| matches!(
effect,
CoreEffect::SendProtocol { envelope, .. } if envelope.kind == Kind::Response
)));
}
#[test]
fn targeted_discovery_routes_before_starting_the_destination_walk() {
let now = Instant::now();
let mut core = ProtocolCore::new("relay");
let peer = add_ready_peer(&mut core, now, "peer-session", "peer");
import_route(
&mut core,
&peer,
now,
2,
route("target", "peer", "target", 1),
);
let client = add_client(&mut core, now, "client");
receive(
&mut core,
&client,
now,
discover_frame("targeted", "target", Scope::Local, Mode::PartialOk),
);
let effects = drain(&mut core);
let CoreEffect::QueryDiscoveryTarget {
stream,
peer,
target_path,
plan,
} = effects
.iter()
.find(|effect| matches!(effect, CoreEffect::QueryDiscoveryTarget { .. }))
.cloned()
.expect("remote discovery must route toward its final target")
else {
unreachable!()
};
assert_eq!(peer, "peer");
assert_eq!(target_path, "/target");
assert_eq!(plan.scope, Scope::Local);
assert_eq!(plan.hops, 7);
assert!(frames(&effects).is_empty());
core.handle(
now,
CoreInput::DiscoveryTargetEvent {
stream: stream.clone(),
event: unb_core::DiscoverEvent::NodeCatalog {
node: "target".into(),
instance_id: "target-instance".into(),
revision: 1,
fingerprint: "fingerprint".into(),
subjects: Vec::new(),
},
},
)
.unwrap();
assert!(frames(&drain(&mut core)).iter().any(|envelope| {
envelope.kind == Kind::Event && envelope.payload_json()["node"] == "target"
}));
core.handle(
now,
CoreInput::DiscoveryTargetEvent {
stream: stream.clone(),
event: unb_core::DiscoverEvent::Done {
discover_id: "targeted".into(),
},
},
)
.unwrap();
assert!(frames(&drain(&mut core))
.iter()
.any(|envelope| envelope.kind == Kind::Response));
core.handle(
now,
CoreInput::DiscoveryTargetFailed {
stream,
message: "late".into(),
},
)
.unwrap();
assert_eq!(core.poll_effect(), None);
}
#[test]
fn discovery_branch_timeout_and_session_close_are_core_owned() {
let now = Instant::now();
let peer_session = SessionId::from("peer-session");
let client_session = SessionId::from("client-session");
let mut core = ProtocolCore::new("alpha");
for session in [&peer_session, &client_session] {
core.handle(
now,
CoreInput::SessionOpened {
session: session.clone(),
initiator: false,
establish_peer: true,
expected_peer: None,
},
)
.unwrap();
receive(&mut core, session, now, hello("hello"));
drain(&mut core);
}
establish_peer(
&mut core,
&peer_session,
now,
NodeIdentity {
node_id: "beta".into(),
instance_id: "beta-instance".into(),
epoch: 1,
proof: serde_json::Value::Null,
},
);
receive(
&mut core,
&client_session,
now,
discover_frame("walk", "alpha", Scope::Reachable, Mode::PartialOk),
);
let mut effects = drain(&mut core);
if let Some(CoreEffect::ContinueDiscovery { stream }) = effects
.iter()
.find(|effect| matches!(effect, CoreEffect::ContinueDiscovery { .. }))
.cloned()
{
core.handle(now, CoreInput::ContinueDiscovery { stream })
.unwrap();
effects.extend(drain(&mut core));
}
let CoreEffect::QueryDiscoveryNeighbor { stream, peer, .. } = effects
.iter()
.find(|effect| matches!(effect, CoreEffect::QueryDiscoveryNeighbor { .. }))
.cloned()
.expect("reachable discovery must query the active peer")
else {
unreachable!()
};
assert_eq!(peer, "beta");
core.handle(
now,
CoreInput::DiscoveryNeighborTimeout {
stream: stream.clone(),
peer,
},
)
.unwrap();
let effects = drain(&mut core);
assert!(effects.iter().any(|effect| matches!(
effect,
CoreEffect::SendProtocol { envelope, .. }
if envelope.kind == Kind::Event
&& envelope.payload_json()["type"] == "warning"
&& envelope.payload_json()["node"] == "beta"
)));
assert!(effects.iter().any(|effect| matches!(
effect,
CoreEffect::SendProtocol { envelope, .. } if envelope.kind == Kind::Response
)));
receive(
&mut core,
&client_session,
now,
discover_frame("close", "alpha", Scope::Reachable, Mode::Strict),
);
let effects = drain(&mut core);
let (stream, peer) = effects
.iter()
.find_map(|effect| match effect {
CoreEffect::QueryDiscoveryNeighbor { stream, peer, .. } => {
Some((stream.clone(), peer.clone()))
}
_ => None,
})
.expect("strict discovery must query the active peer");
core.handle(now, CoreInput::DiscoveryNeighborTimeout { stream, peer })
.unwrap();
assert!(drain(&mut core).iter().any(|effect| matches!(
effect,
CoreEffect::Send { frame, .. }
if frame.head.kind == Kind::Error
&& frame.head.error.as_ref().is_some_and(|error| {
error.code == ErrorCode::PeerUnreachable
})
)));
core.handle(
now,
CoreInput::SessionClosed {
session: client_session.clone(),
},
)
.unwrap();
assert!(drain(&mut core).iter().any(|effect| matches!(
effect,
CoreEffect::SessionRetired {
reason: RetirementReason::SessionClosed,
..
}
)));
}
fn transfer(
core: &mut ProtocolCore,
from: &SessionId,
to: &SessionId,
now: Instant,
) -> Vec<CoreEffect> {
let outgoing = drain(core);
assert!(outgoing.iter().all(|effect| effect.session() == from));
for effect in &outgoing {
let input = match effect {
CoreEffect::SendFrame { envelope, .. } | CoreEffect::SendProtocol { envelope, .. } => {
CoreInput::FrameReceived {
session: to.clone(),
envelope: envelope.clone(),
}
}
CoreEffect::Send { frame, .. } => CoreInput::ApplicationFrameReceived {
session: to.clone(),
frame: frame.clone(),
},
_ => continue,
};
core.handle(now, input).unwrap();
}
let mut effects = outgoing;
effects.extend(drain(core));
effects
}
#[test]
fn complete_client_and_peer_establishment_simulation_is_runtime_free() {
let now = Instant::now();
let client = SessionId::from("client");
let peer = SessionId::from("peer");
let mut core = ProtocolCore::new("alpha");
for session in [&client, &peer] {
core.handle(
now,
CoreInput::SessionOpened {
session: session.clone(),
initiator: false,
establish_peer: true,
expected_peer: None,
},
)
.unwrap();
receive(&mut core, session, now, hello("hello"));
drain(&mut core);
}
let mut request = frame("request", Kind::Request);
request.target = "alpha".into();
receive(&mut core, &client, now, request);
let CoreEffect::CheckDispatchCapacity { effect, .. } = drain(&mut core).remove(0) else {
panic!("client request must check capacity")
};
core.handle(
now,
CoreInput::CapacityChecked {
effect,
result: CapacityResult::Available,
},
)
.unwrap();
let CoreEffect::InvokeApplication { effect, invocation } = drain(&mut core).remove(0) else {
panic!("available client request must invoke application")
};
assert_eq!(
invocation.origin,
ApplicationOrigin::Client {
session: client.clone()
}
);
core.handle(
now,
CoreInput::DispatchCompleted {
effect,
result: Ok(ApplicationResult::Response(application_response(
http::Response::new(Bytes::from_static(b"ok")),
))),
},
)
.unwrap();
assert!(drain(&mut core).iter().any(|effect| matches!(
effect,
CoreEffect::Send { session, frame, .. }
if session == &client
&& frame.body.as_ref().map(BodyId::as_str) == Some("test-body")
)));
assert_eq!(core.session_class(&client).unwrap(), SessionClass::Client);
let identity = NodeIdentity {
node_id: "beta".into(),
instance_id: "beta-instance".into(),
epoch: 1,
proof: serde_json::Value::Null,
};
let effects = establish_peer(&mut core, &peer, now, identity.clone());
assert!(effects.iter().any(|effect| matches!(
effect,
CoreEffect::SessionEstablished { session, peer: established, .. }
if session == &peer && established == &identity
)));
assert!(core.session_ready(&peer).unwrap());
let mut peer_request = frame("peer-request", Kind::Request);
peer_request.target = "alpha".into();
receive(&mut core, &peer, now, peer_request);
let CoreEffect::CheckDispatchCapacity { effect, .. } = drain(&mut core).remove(0) else {
panic!("peer request must check capacity")
};
core.handle(
now,
CoreInput::CapacityChecked {
effect,
result: CapacityResult::Available,
},
)
.unwrap();
let CoreEffect::InvokeApplication { invocation, .. } = drain(&mut core).remove(0) else {
panic!("ready peer request must invoke application")
};
assert_eq!(
invocation.origin,
ApplicationOrigin::Peer {
session: peer,
peer: identity,
}
);
}
#[test]
fn complete_failed_establishment_simulation_covers_rejection_order_and_timeout() {
let now = Instant::now();
let mut core = ProtocolCore::new("alpha");
let identity = NodeIdentity {
node_id: "beta".into(),
instance_id: "beta-instance".into(),
epoch: 1,
proof: serde_json::Value::Null,
};
let rejected = SessionId::from("rejected");
core.handle(
now,
CoreInput::SessionOpened {
session: rejected.clone(),
initiator: false,
establish_peer: true,
expected_peer: None,
},
)
.unwrap();
receive(&mut core, &rejected, now, hello("hello"));
drain(&mut core);
receive(
&mut core,
&rejected,
now,
control(
"identify",
Kind::Identify,
serde_json::to_value(&identity).unwrap(),
),
);
let CoreEffect::RequestPeerAdmission { effect, .. } = drain(&mut core).remove(0) else {
panic!("identity must request admission")
};
core.handle(
now,
CoreInput::PeerAdmissionCompleted {
effect,
result: PeerAdmission::Rejected("denied".into()),
},
)
.unwrap();
assert!(matches!(
without_route_state(drain(&mut core)).as_slice(),
[CoreEffect::SessionRetired {
session,
reason: RetirementReason::AdmissionRejected,
}] if session == &rejected
));
core.handle(
now,
CoreInput::PeerAdmissionCompleted {
effect,
result: PeerAdmission::Admitted(identity),
},
)
.unwrap();
assert_eq!(core.poll_effect(), None);
let malformed = SessionId::from("malformed-order");
core.handle(
now,
CoreInput::SessionOpened {
session: malformed.clone(),
initiator: false,
establish_peer: true,
expected_peer: None,
},
)
.unwrap();
receive(&mut core, &malformed, now, hello("hello"));
drain(&mut core);
receive(
&mut core,
&malformed,
now,
control("accepted", Kind::IdentityAccepted, serde_json::Value::Null),
);
assert!(drain(&mut core).iter().any(|effect| matches!(
effect,
CoreEffect::SessionRetired {
session,
reason: RetirementReason::IdentityAcceptanceOrder,
} if session == &malformed
)));
let timeout = SessionId::from("timeout");
core.handle(
now,
CoreInput::SessionOpened {
session: timeout.clone(),
initiator: false,
establish_peer: true,
expected_peer: None,
},
)
.unwrap();
core.handle(
now + Duration::from_secs(10),
CoreInput::EstablishmentTimeout {
session: timeout.clone(),
},
)
.unwrap();
assert!(matches!(
without_route_state(drain(&mut core)).as_slice(),
[CoreEffect::SessionRetired {
session,
reason: RetirementReason::EstablishmentTimeout,
}] if session == &timeout
));
}
#[test]
fn complete_duplicate_peer_and_transport_loss_simulation_preserves_one_winner() {
let now = Instant::now();
let inbound = SessionId::from("inbound");
let outbound = SessionId::from("outbound");
let mut core = ProtocolCore::new("alpha");
for (session, initiator) in [(&inbound, false), (&outbound, true)] {
core.handle(
now,
CoreInput::SessionOpened {
session: session.clone(),
initiator,
establish_peer: true,
expected_peer: None,
},
)
.unwrap();
if initiator {
receive(
&mut core,
session,
now,
control(
"welcome",
Kind::Welcome,
json!({ "version": PROTOCOL_VERSION }),
),
);
} else {
receive(&mut core, session, now, hello("hello"));
}
drain(&mut core);
}
let identity = NodeIdentity {
node_id: "zeta".into(),
instance_id: "zeta-instance".into(),
epoch: 1,
proof: serde_json::Value::Null,
};
establish_peer(&mut core, &inbound, now, identity.clone());
let effects = establish_peer(&mut core, &outbound, now, identity);
assert!(effects.iter().any(|effect| matches!(
effect,
CoreEffect::SessionRetired {
session,
reason: RetirementReason::DuplicateSessionReplaced,
} if session == &inbound
)));
assert!(core.session_ready(&outbound).unwrap());
core.handle(
now,
CoreInput::SessionClosed {
session: outbound.clone(),
},
)
.unwrap();
assert!(matches!(
without_route_state(drain(&mut core)).as_slice(),
[CoreEffect::SessionRetired {
session,
reason: RetirementReason::SessionClosed,
}] if session == &outbound
));
assert!(!core.session_ready(&outbound).unwrap());
core.handle(now, CoreInput::SessionClosed { session: outbound })
.unwrap();
assert_eq!(core.poll_effect(), None);
}
#[test]
fn complete_discovery_and_session_cleanup_simulation_ignores_all_late_results() {
let now = Instant::now();
let peer = SessionId::from("peer");
let client = SessionId::from("client");
let dispatch_client = SessionId::from("dispatch-client");
let mut core = ProtocolCore::new("alpha");
for session in [&peer, &client, &dispatch_client] {
core.handle(
now,
CoreInput::SessionOpened {
session: session.clone(),
initiator: false,
establish_peer: true,
expected_peer: None,
},
)
.unwrap();
receive(&mut core, session, now, hello("hello"));
drain(&mut core);
}
establish_peer(
&mut core,
&peer,
now,
NodeIdentity {
node_id: "beta".into(),
instance_id: "beta-instance".into(),
epoch: 1,
proof: serde_json::Value::Null,
},
);
receive(
&mut core,
&client,
now,
discover_frame("walk", "alpha", Scope::Reachable, Mode::PartialOk),
);
let mut effects = drain(&mut core);
if let Some(continuation) = effects.iter().find_map(|effect| match effect {
CoreEffect::ContinueDiscovery { stream } => Some(stream.clone()),
_ => None,
}) {
core.handle(
now,
CoreInput::ContinueDiscovery {
stream: continuation,
},
)
.unwrap();
effects.extend(drain(&mut core));
}
let (stream, neighbor) = effects
.iter()
.find_map(|effect| match effect {
CoreEffect::QueryDiscoveryNeighbor { stream, peer, .. } => {
Some((stream.clone(), peer.clone()))
}
_ => None,
})
.expect("reachable discovery must query its established peer");
core.handle(
now + Duration::from_millis(10),
CoreInput::DiscoveryNeighborTimeout {
stream: stream.clone(),
peer: neighbor.clone(),
},
)
.unwrap();
let effects = drain(&mut core);
assert!(effects.iter().any(|effect| matches!(
effect,
CoreEffect::SendProtocol { envelope, .. }
if envelope.kind == Kind::Event && envelope.payload_json()["type"] == "warning"
)));
assert!(effects.iter().any(|effect| matches!(
effect,
CoreEffect::SendProtocol { envelope, .. } if envelope.kind == Kind::Response
)));
let mut request = frame("request", Kind::Request);
request.target = "alpha".into();
receive(&mut core, &dispatch_client, now, request);
let CoreEffect::CheckDispatchCapacity {
effect: capacity, ..
} = drain(&mut core).remove(0)
else {
panic!("request must check capacity")
};
core.handle(
now,
CoreInput::CapacityChecked {
effect: capacity,
result: CapacityResult::Available,
},
)
.unwrap();
let CoreEffect::InvokeApplication {
effect: dispatch, ..
} = drain(&mut core).remove(0)
else {
panic!("available request must invoke application")
};
core.handle(
now,
CoreInput::DispatchCompleted {
effect: dispatch,
result: Ok(ApplicationResult::Response(application_response(
http::Response::new(Bytes::new()),
))),
},
)
.unwrap();
let CoreEffect::Send { effect: send, .. } = drain(&mut core).remove(0) else {
panic!("application response must send")
};
core.handle(
now,
CoreInput::SessionClosed {
session: dispatch_client.clone(),
},
)
.unwrap();
assert!(matches!(
drain(&mut core).as_slice(),
[CoreEffect::SessionRetired {
session,
reason: RetirementReason::SessionClosed,
}] if session == &dispatch_client
));
for input in [
CoreInput::CapacityChecked {
effect: capacity,
result: CapacityResult::Busy,
},
CoreInput::DispatchCompleted {
effect: dispatch,
result: Err(ApplicationFailure {
code: ErrorCode::Internal,
message: "late".into(),
}),
},
CoreInput::SendCompleted {
effect: send,
result: SendResult::WriteFailed("late".into()),
},
CoreInput::DiscoveryNeighborTimeout {
stream: stream.clone(),
peer: neighbor.clone(),
},
CoreInput::DiscoveryNeighborDone {
stream,
peer: neighbor,
},
] {
core.handle(now, input).unwrap();
assert_eq!(core.poll_effect(), None);
}
core.handle(
now,
CoreInput::SessionClosed {
session: client.clone(),
},
)
.unwrap();
drain(&mut core);
core.handle(
now,
CoreInput::ContinueDiscovery {
stream: StreamKey {
session: client,
corr: CorrelationId::from("walk-corr"),
},
},
)
.unwrap();
assert_eq!(core.poll_effect(), None);
}
#[test]
fn sessions_have_typed_key_spaces_and_route_only_their_frames_and_effects() {
let now = Instant::now();
let first = SessionId::from("first");
let second = SessionId::from("second");
let mut core = ProtocolCore::new("node");
assert_ne!(
StreamKey {
session: first.clone(),
corr: CorrelationId::from("s1"),
},
StreamKey {
session: second.clone(),
corr: CorrelationId::from("s1"),
}
);
core.handle(
now,
CoreInput::SessionOpened {
session: first.clone(),
initiator: false,
establish_peer: true,
expected_peer: None,
},
)
.unwrap();
core.handle(
now,
CoreInput::SessionOpened {
session: second.clone(),
initiator: false,
establish_peer: true,
expected_peer: None,
},
)
.unwrap();
assert_eq!(core.poll_effect(), None);
core.handle(
now,
CoreInput::FrameReceived {
session: first.clone(),
envelope: hello("first-hello"),
},
)
.unwrap();
let effects = without_route_state(std::iter::from_fn(|| core.poll_effect()).collect());
assert!(effects.iter().all(|effect| effect.session() == &first));
assert!(effects.iter().any(|effect| matches!(
effect,
CoreEffect::SendFrame { envelope, .. } if envelope.kind == Kind::Welcome
)));
assert!(effects.iter().any(|effect| matches!(
effect,
CoreEffect::HandshakeEstablished {
version: PROTOCOL_VERSION,
..
}
)));
assert_eq!(core.poll_effect(), None);
core.handle(
now,
CoreInput::FrameReceived {
session: second.clone(),
envelope: hello("second-hello"),
},
)
.unwrap();
while let Some(effect) = core.poll_effect() {
assert_eq!(effect.session(), &second);
}
}
#[test]
fn handshake_pair_keepalive_deadlines_and_corr_parity_are_session_owned() {
let (mut core, dialer, listener, now) = established_pair();
assert_eq!(
core.session_deadline(&dialer).unwrap(),
Some(now + Duration::from_secs(15))
);
assert_eq!(
core.session_deadline(&listener).unwrap(),
Some(now + Duration::from_secs(15))
);
let odd = core
.open_stream_body(
&dialer,
"/node/chess",
Kind::Request,
None,
None,
Default::default(),
)
.unwrap();
assert_eq!(odd, "s1");
transfer(&mut core, &dialer, &listener, now);
let even = core
.open_stream_body(
&listener,
"/node/browser-ui",
Kind::Request,
None,
None,
Default::default(),
)
.unwrap();
assert_eq!(even, "s2");
transfer(&mut core, &listener, &dialer, now);
core.handle(
now + Duration::from_secs(15),
CoreInput::SessionTimeout {
session: dialer.clone(),
},
)
.unwrap();
let effects = without_route_state(drain(&mut core));
assert!(effects.iter().any(|effect| matches!(
effect,
CoreEffect::SendFrame { envelope, .. } if envelope.kind == Kind::Ping
)));
assert!(effects
.iter()
.any(|effect| matches!(effect, CoreEffect::ScheduleSessionDeadline { .. })));
assert_eq!(
core.session_deadline(&dialer).unwrap(),
Some(now + Duration::from_secs(30))
);
assert_eq!(
core.session_deadline(&listener).unwrap(),
Some(now + Duration::from_secs(15))
);
}
#[test]
fn correlated_delivery_requires_a_registered_client_operation() {
let (mut core, dialer, _listener, now) = established_pair();
let corr = core
.open_stream_body(
&dialer,
"/node/todo",
Kind::Request,
None,
None,
Default::default(),
)
.unwrap();
drain(&mut core);
let mut response = frame("response", Kind::Response);
response.corr = Some(corr.clone());
response.payload = Bytes::from_static(b"ok");
receive(&mut core, &dialer, now, response);
let effects = drain(&mut core);
assert!(effects.iter().any(|effect| matches!(
effect,
CoreEffect::Deliver { session, envelope }
if session == &dialer && envelope.corr.as_deref() == Some(corr.as_str())
)));
assert!(!effects.iter().any(|effect| matches!(
effect,
CoreEffect::DeliverClient { session, operation, .. }
if session == &dialer && operation.as_str() == corr
)));
assert_eq!(
effects
.iter()
.filter(|effect| matches!(
effect,
CoreEffect::StreamClosed { session, operation }
if session == &dialer && operation.as_str() == corr
))
.count(),
1
);
}
#[test]
fn registered_client_operation_delivery_remains_client_owned() {
let (mut core, dialer, _listener, now) = established_pair();
core.handle(
now,
CoreInput::StartClientOperation {
session: dialer.clone(),
target_path: "/node/todo".into(),
kind: Kind::Request,
input: OperationInput::Body(None),
hops: None,
headers: Default::default(),
timeout: None,
},
)
.unwrap();
let outgoing = drain(&mut core);
let envelope = outgoing
.iter()
.find_map(|effect| match effect {
CoreEffect::Send { frame, .. } => Some(frame.clone()),
_ => None,
})
.unwrap();
let corr = envelope.head.corr.clone().unwrap();
let mut response = frame("response", Kind::Response);
response.corr = Some(corr.clone());
response.payload = Bytes::from_static(b"ok");
receive(&mut core, &dialer, now, response);
let effects = drain(&mut core);
assert!(effects.iter().any(|effect| matches!(
effect,
CoreEffect::DeliverClient {
session,
operation,
delivery: ClientDelivery::Terminal(frame),
} if session == &dialer
&& operation.as_str() == corr
&& frame.body.as_ref().map(BodyId::as_str) == Some("test-body")
)));
}
#[test]
fn event_sequence_and_terminal_cancel_are_isolated_across_sessions() {
let (mut core, dialer, listener, now) = established_pair();
core.handle(
now,
CoreInput::LocalCapabilitiesInstalled {
capabilities: BTreeMap::from([("todo".into(), json!({}))]),
},
)
.unwrap();
drain(&mut core);
let corr = core
.open_stream_body(
&dialer,
"/node/todo",
Kind::Subscribe,
None,
None,
Default::default(),
)
.unwrap();
transfer(&mut core, &dialer, &listener, now);
core.send_body(&listener, &corr, Some(BodyId::from("event-1")))
.unwrap();
core.send_body(&listener, &corr, Some(BodyId::from("event-2")))
.unwrap();
core.respond_body(
&listener,
&corr,
Some(BodyId::from("terminal")),
Default::default(),
)
.unwrap();
let effects = transfer(&mut core, &listener, &dialer, now);
let delivered: Vec<_> = effects
.iter()
.filter_map(|effect| match effect {
CoreEffect::Deliver { session, envelope }
if session == &dialer && envelope.corr.as_deref() == Some(corr.as_str()) =>
{
Some(envelope)
}
_ => None,
})
.collect();
assert_eq!((delivered[0].seq, delivered[1].seq), (Some(1), Some(2)));
assert_eq!(delivered[2].kind, Kind::Response);
assert!(matches!(
core.send_body(&dialer, &corr, None),
Err(CoreError::UnknownStream(_))
));
assert!(matches!(
core.send_body(&listener, &corr, None),
Err(CoreError::UnknownStream(_))
));
let cancelled = core
.open_stream_body(
&dialer,
"/node/todo",
Kind::Subscribe,
None,
None,
Default::default(),
)
.unwrap();
transfer(&mut core, &dialer, &listener, now);
core.cancel(&dialer, &cancelled).unwrap();
let effects = transfer(&mut core, &dialer, &listener, now);
assert!(effects.iter().any(|effect| matches!(
effect,
CoreEffect::Deliver { session, envelope }
if session == &listener
&& envelope.corr.as_deref() == Some(cancelled.as_str())
&& envelope.kind == Kind::Cancel
)));
assert!(matches!(
core.send_body(&dialer, &cancelled, None),
Err(CoreError::UnknownStream(_))
));
assert!(matches!(
core.send_body(&listener, &cancelled, None),
Err(CoreError::UnknownStream(_))
));
}