mod common;
use common::TestCall;
use std::sync::Arc;
use std::time::Duration;
use common::{ghost_establish, wait_until};
use schemars::JsonSchema;
use serde::Deserialize;
use serde_json::{json, Value};
use unb::{handler, Handler, HandlerError, Reply, Request, State};
use unb_client::pair;
use unb_core::Kind;
use unb_runtime::Wire;
use unb_server::Node;
struct Epoch(u64);
#[derive(Deserialize, JsonSchema)]
struct Probe {}
#[handler]
async fn served_by(request: Request<Probe>) -> Result<Reply<Value>, HandlerError> {
Ok(Reply::new(json!({ "served_by": request.subject() })))
}
#[handler]
async fn ok(_request: Request<Probe>) -> Result<Reply<Value>, HandlerError> {
Ok(Reply::new(json!({ "ok": true })))
}
#[handler]
async fn job_run(
epoch: State<Epoch>,
_request: Request<Probe>,
) -> Result<Reply<Value>, HandlerError> {
Ok(Reply::new(json!({ "epoch": epoch.0 })))
}
fn service(name: &str, subject: &'static str) -> Arc<Node> {
Node::builder(name)
.service(served_by.at_subject(subject))
.insecure_accept_declared_peer_identities()
.build()
.unwrap()
}
async fn request_eventually(node: &Arc<Node>, subject: &str) -> Value {
tokio::time::timeout(Duration::from_secs(5), async {
loop {
if let Ok(Ok(reply)) =
tokio::time::timeout(Duration::from_millis(250), node.request(subject, json!({})))
.await
{
break reply;
}
tokio::time::sleep(Duration::from_millis(20)).await;
}
})
.await
.unwrap_or_else(|_| panic!("{subject} did not become serviceable"))
}
#[tokio::test(flavor = "multi_thread")]
async fn runtime_feature_mutation_does_not_change_a_four_hop_node_route() {
let a = Node::builder("hop-a")
.insecure_accept_declared_peer_identities()
.build()
.unwrap();
let b = Node::builder("hop-b")
.insecure_accept_declared_peer_identities()
.build()
.unwrap();
let c = Node::builder("hop-c")
.insecure_accept_declared_peer_identities()
.build()
.unwrap();
let d = Node::builder("hop-d")
.insecure_accept_declared_peer_identities()
.build()
.unwrap();
a.link(&b).await.unwrap();
b.link(&c).await.unwrap();
c.link(&d).await.unwrap();
d.add_service(ok.at_subject("late.echo")).await.unwrap();
wait_until("the destination node reaches the far end", || {
a.reachable_names().contains(&"hop-d".to_string())
})
.await;
let reply = a.request("/hop-d/late.echo", json!({})).await.unwrap();
assert_eq!(reply["ok"], true);
d.remove_subject("late.echo").await.unwrap();
assert!(a.reachable_names().contains(&"hop-d".to_string()));
let error = a.request("/hop-d/late.echo", json!({})).await.unwrap_err();
assert_eq!(error.code, unb_core::ErrorCode::UnknownSubject);
}
#[tokio::test(flavor = "multi_thread")]
async fn losing_the_selected_bridge_keeps_the_subject_reachable_through_the_alternate() {
let a = Node::builder("dia-a")
.insecure_accept_declared_peer_identities()
.build()
.unwrap();
let b = Node::builder("dia-b")
.insecure_accept_declared_peer_identities()
.build()
.unwrap();
let c = Node::builder("dia-c")
.insecure_accept_declared_peer_identities()
.build()
.unwrap();
let d = service("dia-d", "delta");
a.link(&b).await.unwrap();
a.link(&c).await.unwrap();
b.link(&d).await.unwrap();
c.link(&d).await.unwrap();
wait_until("a reaches dia-d", || {
a.reachable_names().contains(&"dia-d".to_string())
})
.await;
a.request("/dia-d/delta", json!({})).await.unwrap();
drop(b);
wait_until("dia-d stays reachable through the alternate bridge", || {
a.reachable_names().contains(&"dia-d".to_string())
})
.await;
request_eventually(&a, "/dia-d/delta").await;
let b2 = Node::builder("dia-b2")
.insecure_accept_declared_peer_identities()
.build()
.unwrap();
a.link(&b2).await.unwrap();
b2.link(&d).await.unwrap();
request_eventually(&a, "/dia-d/delta").await;
}
#[tokio::test(flavor = "multi_thread")]
async fn duplicate_feature_names_remain_addressable_by_owning_node() {
let left = service("owner-left", "shared.subject");
let right = service("owner-right", "shared.subject");
let middle = Node::builder("middle")
.insecure_accept_declared_peer_identities()
.build()
.unwrap();
left.link(&middle).await.unwrap();
right.link(&middle).await.unwrap();
wait_until("the middle reaches both owners", || {
let names = middle.reachable_names();
names.contains(&"owner-left".to_string()) && names.contains(&"owner-right".to_string())
})
.await;
assert_eq!(
middle
.request("/owner-left/shared.subject", json!({}))
.await
.unwrap()["served_by"],
"shared.subject"
);
assert_eq!(
middle
.request("/owner-right/shared.subject", json!({}))
.await
.unwrap()["served_by"],
"shared.subject"
);
right.remove_subject("shared.subject").await.unwrap();
assert!(middle
.reachable_names()
.contains(&"owner-right".to_string()));
let error = middle
.request("/owner-right/shared.subject", json!({}))
.await
.unwrap_err();
assert_eq!(error.code, unb_core::ErrorCode::UnknownSubject);
middle
.request("/owner-left/shared.subject", json!({}))
.await
.unwrap();
}
#[tokio::test(flavor = "multi_thread")]
async fn an_owner_restart_with_a_newer_epoch_resumes_service() {
let caller = Node::builder("caller")
.insecure_accept_declared_peer_identities()
.build()
.unwrap();
let owner = Node::builder("owner-x")
.epoch(1)
.state(Epoch(1))
.service(job_run.at_subject("job.run"))
.insecure_accept_declared_peer_identities()
.build()
.unwrap();
caller.link(&owner).await.unwrap();
caller.request("/owner-x/job.run", json!({})).await.unwrap();
drop(owner);
wait_until("the dead owner's routes withdraw", || {
!caller.reachable_names().contains(&"owner-x".to_string())
})
.await;
let restarted = Node::builder("owner-x")
.epoch(2)
.state(Epoch(2))
.service(job_run.at_subject("job.run"))
.insecure_accept_declared_peer_identities()
.build()
.unwrap();
caller.link(&restarted).await.unwrap();
let reply = caller.request("/owner-x/job.run", json!({})).await.unwrap();
assert_eq!(reply["epoch"], 2);
}
#[tokio::test(flavor = "multi_thread")]
async fn a_stable_topology_stops_emitting_route_updates() {
let a = service("quiet-a", "alpha");
let b = service("quiet-b", "beta");
a.link(&b).await.unwrap();
let (dial_side, ghost_side) = pair();
let ghost = Wire::open(ghost_side);
let (connected, ()) = tokio::join!(
a.connect_transport("watcher", dial_side),
ghost_establish(&ghost, "watcher")
);
connected.unwrap();
let mut observer = ghost.observe();
tokio::time::sleep(Duration::from_millis(300)).await;
let mut drained = 0;
while let Ok(Ok(envelope)) =
tokio::time::timeout(Duration::from_millis(150), observer.recv()).await
{
if envelope.kind == Kind::RouteSnapshot {
drained += 1;
}
}
a.remove_subject("never.existed").await.unwrap();
let extra = tokio::time::timeout(Duration::from_millis(300), async {
loop {
match observer.recv().await {
Ok(envelope) if envelope.kind == Kind::RouteSnapshot => {
return true;
}
Ok(_) => {}
Err(_) => return false,
}
}
})
.await;
assert!(
extra.is_err(),
"a no-op mutation must not emit a route update (drained {drained} settling frames first)"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn a_generation_gap_from_a_peer_answers_resync_required() {
use unb_core::{Envelope, RouteAck, RouteAckStatus, RouteDelta};
let node = service("resync-node", "alpha");
let (dial_side, ghost_side) = pair();
let ghost = Wire::open(ghost_side);
let (connected, ()) = tokio::join!(
node.connect_transport("gapper", dial_side),
ghost_establish(&ghost, "gapper")
);
connected.unwrap();
let mut observer = ghost.observe();
let gapped = RouteDelta {
generation: 40,
upsert: Vec::new(),
withdraw: Vec::new(),
};
ghost
.control(
Kind::RouteDelta,
Envelope::encode_payload(&serde_json::to_value(&gapped).unwrap()),
)
.await
.unwrap();
let deadline = tokio::time::Instant::now() + Duration::from_secs(3);
loop {
let remaining = deadline - tokio::time::Instant::now();
let envelope = tokio::time::timeout(remaining, observer.recv())
.await
.expect("expected a resync request before the deadline")
.expect("ghost session closed instead of requesting resync");
if envelope.kind == Kind::RouteAck {
let ack = envelope.parse_payload::<RouteAck>().unwrap();
if ack.status == RouteAckStatus::ResyncRequired {
return;
}
}
}
}
#[tokio::test(flavor = "multi_thread")]
async fn a_resync_request_ahead_of_the_emitted_generation_closes_the_session() {
use unb_core::{Envelope, RouteAck, RouteAckStatus};
let node = service("future-resync-node", "alpha");
let (dial_side, ghost_side) = pair();
let ghost = Wire::open(ghost_side);
let (connected, ()) = tokio::join!(
node.connect_transport("future-resync-peer", dial_side),
ghost_establish(&ghost, "future-resync-peer")
);
connected.unwrap();
let ack = RouteAck {
generation: 2,
status: RouteAckStatus::ResyncRequired,
};
ghost
.control(
Kind::RouteAck,
Envelope::encode_payload(&serde_json::to_value(&ack).unwrap()),
)
.await
.unwrap();
tokio::time::timeout(Duration::from_secs(3), ghost.closed())
.await
.expect("future resync request must close the session");
}
#[tokio::test(flavor = "multi_thread")]
async fn an_accepted_resync_emits_one_recovery_snapshot_at_the_next_generation() {
use unb_core::{Envelope, RouteAck, RouteAckStatus, RouteSnapshot};
let node = service("recovery-node", "alpha");
let (dial_side, ghost_side) = pair();
let ghost = Wire::open(ghost_side);
let (connected, ()) = tokio::join!(
node.connect_transport("recovery-peer", dial_side),
ghost_establish(&ghost, "recovery-peer")
);
connected.unwrap();
let mut observer = ghost.observe();
let ack = RouteAck {
generation: 1,
status: RouteAckStatus::ResyncRequired,
};
let payload = Envelope::encode_payload(&serde_json::to_value(&ack).unwrap());
ghost
.control(Kind::RouteAck, payload.clone())
.await
.unwrap();
let snapshot = tokio::time::timeout(Duration::from_secs(3), async {
loop {
match observer.recv().await {
Ok(envelope) if envelope.kind == Kind::RouteSnapshot => {
return envelope.parse_payload::<RouteSnapshot>().unwrap();
}
Ok(_) => {}
Err(_) => panic!("session closed before recovery snapshot"),
}
}
})
.await
.expect("accepted resync must emit a recovery snapshot");
assert_eq!(snapshot.generation, 2);
ghost.control(Kind::RouteAck, payload).await.unwrap();
let replayed_snapshot = tokio::time::timeout(Duration::from_millis(300), async {
loop {
match observer.recv().await {
Ok(envelope) if envelope.kind == Kind::RouteSnapshot => {
return;
}
Ok(_) => {}
Err(_) => panic!("session closed after replayed resync"),
}
}
})
.await;
assert!(
replayed_snapshot.is_err(),
"replayed resync emitted a second recovery snapshot"
);
}