mod common;
use std::time::Duration;
use common::{connect_nodes, ghost_establish, start};
use futures_util::StreamExt;
use schemars::JsonSchema;
use serde::Deserialize;
use serde_json::{json, Value};
use unb::{handler, Handler, HandlerError, Reply, Request};
use unb_client::{dial_transport, pair};
use unb_core::{Detail, DiscoverEvent, Kind, Scope};
use unb_runtime::{ClientError, Wire};
use unb_server::Node;
#[derive(Deserialize, JsonSchema)]
struct Probe {}
#[handler]
async fn chess(_request: Request<Probe>) -> Result<Reply<Value>, HandlerError> {
Ok(Reply::new(json!({})))
}
#[handler]
async fn weather(_request: Request<Probe>) -> Result<Reply<Value>, HandlerError> {
Ok(Reply::new(json!({})))
}
#[tokio::test(flavor = "multi_thread")]
async fn discover_streams_a_local_node_catalog_then_a_done_response() {
let node =
Node::builder("hub")
.service(chess.describe(
json!({ "one_line": "a chess game", "input_schema": {"from": "string"} }),
))
.insecure_accept_declared_peer_identities()
.build()
.unwrap();
let url = start(node).await;
let wire = Wire::open(dial_transport(&url).await.unwrap());
common::ready_client(&wire).await;
let mut stream = wire
.client_session()
.start(
"/hub",
Kind::Discover,
unb_core::Envelope::encode_payload(&json!({ "detail": "index", "scope": "local" })),
None,
Default::default(),
)
.await
.unwrap();
let catalog = stream
.next()
.await
.unwrap()
.expect("expected a node_catalog event");
assert_eq!(catalog.kind, Kind::Event);
let payload = catalog.payload_json();
assert_eq!(payload["type"], "node_catalog");
assert_eq!(payload["node"], "hub");
assert_eq!(payload["subjects"][0]["subject"], "chess");
assert_eq!(payload["subjects"][0]["target_path"], "/hub/chess");
assert_eq!(payload["subjects"][0]["one_line"], "a chess game");
assert!(
payload["subjects"][0].get("input_schema").is_none(),
"index detail omits schemas"
);
let done = stream
.next()
.await
.unwrap()
.expect("expected a terminal done");
assert_eq!(done.kind, Kind::Response);
assert_eq!(done.payload_json()["type"], "done");
let mut stream = wire
.client_session()
.start(
"/hub",
Kind::Discover,
unb_core::Envelope::encode_payload(&json!({ "detail": "full", "scope": "local" })),
None,
Default::default(),
)
.await
.unwrap();
let full = stream.next().await.unwrap().unwrap();
assert_eq!(
full.payload_json()["subjects"][0]["metadata"]["input_schema"]["from"],
"string",
"full detail carries schemas"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn discover_catalog_yields_typed_events_and_ends_on_done() {
let node = Node::builder("hub")
.service(chess.describe(json!({ "one_line": "a chess game" })))
.insecure_accept_declared_peer_identities()
.build()
.unwrap();
let url = start(node).await;
let wire = Wire::open(dial_transport(&url).await.unwrap());
common::ready_client(&wire).await;
let mut catalog = wire
.discover_catalog("/hub", Detail::Index, Scope::Local)
.await
.unwrap();
let mut events = Vec::new();
while let Some(event) = catalog.next().await {
events.push(event);
}
assert!(
matches!(events.first(), Some(DiscoverEvent::NodeCatalog { node, .. }) if node.as_str() == "hub"),
"first event is the local node catalog: {events:?}"
);
assert!(
matches!(events.last(), Some(DiscoverEvent::Done { .. })),
"the stream ends with the done marker: {events:?}"
);
let catalogs = events
.iter()
.filter(|event| matches!(event, DiscoverEvent::NodeCatalog { .. }))
.count();
assert_eq!(catalogs, 1, "local scope yields exactly one node catalog");
}
#[tokio::test(flavor = "multi_thread")]
async fn connected_discovery_observes_exact_atomic_catalog_revisions() {
let owner = Node::builder("mutable-owner")
.insecure_accept_declared_peer_identities()
.build()
.unwrap();
let observer = Node::builder("observer")
.insecure_accept_declared_peer_identities()
.build()
.unwrap();
connect_nodes(&observer, "mutable-owner", &owner).await;
owner.add_service(chess).await.unwrap();
assert!(observer
.reachable_names()
.contains(&"mutable-owner".to_string()));
let added = observer
.discover_events(Detail::Full, Scope::Reachable)
.await;
assert!(added.iter().any(|event| matches!(
event,
DiscoverEvent::NodeCatalog { node, revision: 1, subjects, .. }
if node == "mutable-owner"
&& subjects.len() == 1
&& subjects[0]["subject"] == "chess"
&& subjects[0]["target_path"] == "/mutable-owner/chess"
&& subjects[0]["operations"]["unary"].is_object()
)));
owner.remove_subject("chess").await.unwrap();
assert!(observer
.reachable_names()
.contains(&"mutable-owner".to_string()));
let removed = observer
.discover_events(Detail::Full, Scope::Reachable)
.await;
assert!(removed.iter().any(|event| matches!(
event,
DiscoverEvent::NodeCatalog { node, revision: 2, subjects, .. }
if node == "mutable-owner" && subjects.is_empty()
)));
}
#[tokio::test(flavor = "multi_thread")]
async fn discover_streams_node_scoped_catalogs_with_no_cross_node_dedup() {
let owner_a = Node::builder("owner-a")
.service(weather.describe(json!({ "one_line": "forecasts a" })))
.insecure_accept_declared_peer_identities()
.build()
.unwrap();
let owner_b = Node::builder("owner-b")
.service(weather.describe(json!({ "one_line": "forecasts b" })))
.insecure_accept_declared_peer_identities()
.build()
.unwrap();
let relay = Node::builder("relay-1")
.insecure_accept_declared_peer_identities()
.build()
.unwrap();
connect_nodes(&relay, "owner-a", &owner_a).await;
connect_nodes(&relay, "owner-b", &owner_b).await;
let (client_side, relay_server) = pair();
relay.serve_transport(relay_server).await;
let client = Wire::open(client_side);
common::ready_client(&client).await;
let mut stream = common::stream(
&client,
"/relay-1",
Kind::Discover,
json!({ "detail": "index" }),
)
.await;
let mut catalogs: Vec<(String, Vec<String>)> = Vec::new();
loop {
let frame = stream
.next()
.await
.unwrap()
.expect("stream ended without a terminal done");
if frame.kind == Kind::Response {
assert_eq!(frame.payload_json()["type"], "done");
break;
}
let payload = frame.payload_json();
if payload["type"] == "node_catalog" {
let node = payload["node"].as_str().unwrap_or_default().to_string();
let subjects: Vec<String> = payload["subjects"]
.as_array()
.into_iter()
.flatten()
.filter_map(|entry| entry["subject"].as_str().map(str::to_string))
.collect();
catalogs.push((node, subjects));
}
}
let nodes: Vec<&String> = catalogs.iter().map(|(node, _)| node).collect();
assert!(nodes.contains(&&"relay-1".to_string()), "{catalogs:?}");
assert!(nodes.contains(&&"owner-a".to_string()), "{catalogs:?}");
assert!(nodes.contains(&&"owner-b".to_string()), "{catalogs:?}");
let weather_owners = catalogs
.iter()
.filter(|(_, subjects)| subjects.iter().any(|s| s == "weather"))
.count();
assert_eq!(
weather_owners, 2,
"two nodes own weather; discovery is node-scoped and emits both (no dedup): {catalogs:?}"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn discover_partial_ok_warns_on_a_timed_out_neighbor_and_still_finishes() {
let relay = Node::builder("relay")
.insecure_accept_declared_peer_identities()
.max_activations(0)
.build()
.unwrap();
let (relay_side, ghost_side) = pair();
let ghost = Wire::open(ghost_side);
let (connected, ()) = tokio::join!(
relay.connect_transport("ghost", relay_side),
ghost_establish(&ghost, "ghost")
);
connected.unwrap();
let _ghost = ghost;
let (client_side, relay_server) = pair();
relay.serve_transport(relay_server).await;
let client = Wire::open(client_side);
common::ready_client(&client).await;
let mut stream = client
.client_session()
.start(
"/relay",
Kind::Discover,
unb_core::Envelope::encode_payload(&json!({ "timeout_ms": 150, "mode": "partial_ok" })),
None,
Default::default(),
)
.await
.unwrap();
let mut saw_warning = false;
loop {
let frame = stream
.next()
.await
.unwrap()
.expect("stream ended without a done");
let payload = frame.payload_json();
if frame.kind == Kind::Response {
assert_eq!(payload["type"], "done");
break;
}
if payload["type"] == "warning" {
assert_eq!(payload["node"], "ghost");
saw_warning = true;
}
}
assert!(
saw_warning,
"partial_ok surfaces a warning for the timed-out neighbor"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn discover_strict_fails_the_stream_on_a_timed_out_neighbor() {
let relay = Node::builder("relay")
.insecure_accept_declared_peer_identities()
.max_activations(0)
.build()
.unwrap();
let (relay_side, ghost_side) = pair();
let ghost = Wire::open(ghost_side);
let (connected, ()) = tokio::join!(
relay.connect_transport("ghost", relay_side),
ghost_establish(&ghost, "ghost")
);
connected.unwrap();
let _ghost = ghost;
let (client_side, relay_server) = pair();
relay.serve_transport(relay_server).await;
let client = Wire::open(client_side);
common::ready_client(&client).await;
let mut stream = client
.client_session()
.start(
"/relay",
Kind::Discover,
unb_core::Envelope::encode_payload(&json!({ "timeout_ms": 150, "mode": "strict" })),
None,
Default::default(),
)
.await
.unwrap();
loop {
let frame = match stream.next().await {
Err(ClientError::Protocol { code, .. }) => {
assert_eq!(code, unb_core::ErrorCode::PeerUnreachable);
return;
}
Ok(Some(frame)) => frame,
other => panic!("stream ended without a terminal error: {other:?}"),
};
if frame.kind == Kind::Error {
assert_eq!(frame.payload_json()["code"], "PEER_UNREACHABLE");
return;
}
assert_eq!(
frame.kind,
Kind::Event,
"only events precede the strict error"
);
}
}
#[tokio::test(flavor = "multi_thread")]
async fn discover_does_not_deadlock_when_the_fan_out_exceeds_the_input_channel() {
let relay = Node::builder("relay")
.insecure_accept_declared_peer_identities()
.build()
.unwrap();
let mut ghosts = Vec::new();
for index in 0..80 {
let name = format!("ghost-{index}");
let ghost = Node::builder(&name)
.insecure_accept_declared_peer_identities()
.build()
.unwrap();
connect_nodes(&relay, &name, &ghost).await;
ghosts.push(ghost);
}
let (client_side, relay_server) = pair();
relay.serve_transport(relay_server).await;
let client = Wire::open(client_side);
common::ready_client(&client).await;
let mut stream = client
.client_session()
.start(
"/relay",
Kind::Discover,
unb_core::Envelope::encode_payload(&json!({ "timeout_ms": 150, "mode": "partial_ok" })),
None,
Default::default(),
)
.await
.unwrap();
let done = tokio::time::timeout(Duration::from_secs(10), async {
loop {
let frame = stream
.next()
.await
.unwrap()
.expect("stream ended without a done");
if frame.kind == Kind::Response {
return frame.payload_json()["type"] == "done";
}
}
})
.await
.expect("a fan-out wider than the input channel must still terminate, not deadlock");
assert!(done);
}