mod common;
use common::TestCall;
use std::sync::Arc;
use std::time::Duration;
use unb::{handler, Handler, HandlerError, Reply, Request, State};
use unb_client::pair;
use unb_core::Kind;
use unb_runtime::Wire;
use unb_server::Node;
use common::{ghost_establish, wait_until};
use schemars::JsonSchema;
use serde::Deserialize;
use serde_json::{json, Value};
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 a_runtime_subject_propagates_and_withdraws_across_four_hops() {
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 runtime subject reaches the far end", || {
a.reachable_names().contains(&"late.echo".to_string())
})
.await;
let reply = a.request("late.echo", json!({})).await.unwrap();
assert_eq!(reply["ok"], true);
d.remove_subject("late.echo").await.unwrap();
wait_until("the withdrawal reaches the far end", || {
!a.reachable_names().contains(&"late.echo".to_string())
})
.await;
assert!(a.request("late.echo", json!({})).await.is_err());
}
#[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 delta", || {
a.reachable_names().contains(&"delta".to_string())
})
.await;
a.request("delta", json!({})).await.unwrap();
drop(b);
wait_until("delta stays reachable through the alternate bridge", || {
a.reachable_names().contains(&"delta".to_string())
})
.await;
request_eventually(&a, "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, "delta").await;
}
#[tokio::test(flavor = "multi_thread")]
async fn conflicting_owners_fail_closed_at_the_meeting_node_and_recover_on_withdrawal() {
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 node marks the subject conflicted", || {
!middle
.reachable_names()
.contains(&"shared.subject".to_string())
})
.await;
let error = middle
.request("shared.subject", json!({}))
.await
.unwrap_err();
assert_eq!(error.code, unb_core::ErrorCode::Conflict);
assert!(
error.message.contains("owner-left") && error.message.contains("owner-right"),
"{}",
error.message
);
right.remove_subject("shared.subject").await.unwrap();
wait_until("the surviving owner becomes routable", || {
middle
.reachable_names()
.contains(&"shared.subject".to_string())
})
.await;
let reply = middle.request("shared.subject", json!({})).await.unwrap();
assert_eq!(reply["served_by"], "shared.subject");
}
#[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("job.run", json!({})).await.unwrap();
drop(owner);
wait_until("the dead owner's routes withdraw", || {
!caller.reachable_names().contains(&"job.run".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("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"
);
}