use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::Arc;
use std::time::Duration;
use unb_server::{
ConnectError, ConnectionStatus, DisconnectReason, Endpoint, EndpointSet, HostConfig, Node,
PeerConnection, TcpTransport, TransportKind,
};
fn assert_clone_send_sync<T: Clone + Send + Sync>() {}
fn websocket_endpoint(address: String) -> Endpoint {
Endpoint {
kind: TransportKind::WebSocket,
address,
cert_hash: None,
}
}
async fn hosted(name: &str) -> ((Arc<Node>, unb_server::Hosting), String) {
let node = Node::builder(name)
.insecure_accept_declared_peer_identities()
.build()
.unwrap();
host_node(node).await
}
async fn host_node(node: Arc<Node>) -> ((Arc<Node>, unb_server::Hosting), String) {
let hosting = HostConfig::tcp(([127, 0, 0, 1], 0), TcpTransport::plain())
.start(&node)
.await
.unwrap();
let url = format!("ws://{}", hosting.websocket_addr().unwrap());
((node, hosting), url)
}
#[test]
fn peer_connection_is_clone_send_and_sync() {
assert_clone_send_sync::<PeerConnection>();
}
#[tokio::test(flavor = "multi_thread")]
async fn connect_returns_a_connected_peer_handle() {
let (server, url) = hosted("peer").await;
let node = Node::builder("caller")
.insecure_accept_declared_peer_identities()
.build()
.unwrap();
let connection = node
.connect(EndpointSet::from(websocket_endpoint(url)))
.await
.unwrap();
assert_eq!(connection.peer(), "peer");
assert_eq!(connection.status(), ConnectionStatus::Connected);
drop(server);
}
#[tokio::test(flavor = "multi_thread")]
async fn failed_initial_connect_does_not_poison_a_later_connection() {
let node = Node::builder("caller")
.insecure_accept_declared_peer_identities()
.build()
.unwrap();
let error = match node.connect(EndpointSet::new()).await {
Ok(_) => panic!("empty endpoint set unexpectedly connected"),
Err(error) => error,
};
assert_eq!(error, ConnectError::NoSupportedEndpoint);
let (server, url) = hosted("peer").await;
let connection = node
.connect(EndpointSet::from(websocket_endpoint(url)))
.await
.unwrap();
assert_eq!(connection.peer(), "peer");
assert_eq!(connection.status(), ConnectionStatus::Connected);
drop(server);
}
#[tokio::test(flavor = "multi_thread")]
async fn clones_observe_the_same_status_transition() {
let (server, url) = hosted("peer").await;
let node = Node::builder("caller")
.insecure_accept_declared_peer_identities()
.build()
.unwrap();
let connection = node
.connect(EndpointSet::from(websocket_endpoint(url)))
.await
.unwrap();
let observer = connection.clone();
let changed = tokio::spawn(observer.changed());
connection.disconnect();
assert_eq!(
tokio::time::timeout(Duration::from_secs(1), changed)
.await
.unwrap()
.unwrap(),
ConnectionStatus::Disconnected {
reason: DisconnectReason::ExplicitDisconnect,
}
);
assert_eq!(
connection.status(),
ConnectionStatus::Disconnected {
reason: DisconnectReason::ExplicitDisconnect,
}
);
drop(server);
}
#[tokio::test(flavor = "multi_thread")]
async fn repeated_connects_converge_on_one_logical_connection() {
let (server, url) = hosted("peer").await;
let node = Node::builder("caller")
.insecure_accept_declared_peer_identities()
.build()
.unwrap();
let endpoints = EndpointSet::from(websocket_endpoint(url));
let first = node.connect(endpoints.clone()).await.unwrap();
let second = node.connect(endpoints).await.unwrap();
first.disconnect();
assert_eq!(
second.status(),
ConnectionStatus::Disconnected {
reason: DisconnectReason::ExplicitDisconnect,
}
);
drop(server);
}
#[tokio::test(flavor = "multi_thread")]
async fn connecting_a_different_peer_keeps_connections_independent() {
let (first_server, first_url) = hosted("peer-a").await;
let (second_server, second_url) = hosted("peer-b").await;
let node = Node::builder("caller")
.insecure_accept_declared_peer_identities()
.build()
.unwrap();
let first = node
.connect(EndpointSet::from(websocket_endpoint(first_url)))
.await
.unwrap();
let second = node
.connect(EndpointSet::from(websocket_endpoint(second_url)))
.await
.unwrap();
second.disconnect();
assert_eq!(first.peer(), "peer-a");
assert_eq!(first.status(), ConnectionStatus::Connected);
assert_eq!(second.peer(), "peer-b");
assert_eq!(
second.status(),
ConnectionStatus::Disconnected {
reason: DisconnectReason::ExplicitDisconnect,
}
);
drop(first_server);
drop(second_server);
}
#[tokio::test(flavor = "multi_thread")]
async fn selected_session_retirement_starts_automatic_maintenance() {
let (server, url) = hosted("peer").await;
let node = Node::builder("caller")
.insecure_accept_declared_peer_identities()
.build()
.unwrap();
let connection = node
.connect(EndpointSet::from(websocket_endpoint(url)))
.await
.unwrap();
let changed = connection.changed();
server.0.shutdown();
let status = tokio::time::timeout(Duration::from_secs(1), changed)
.await
.unwrap();
assert!(matches!(status, ConnectionStatus::Connecting));
drop(server);
}
#[tokio::test(flavor = "multi_thread")]
async fn node_shutdown_is_visible_to_every_connection_clone() {
let (server, url) = hosted("peer").await;
let node = Node::builder("caller")
.insecure_accept_declared_peer_identities()
.build()
.unwrap();
let connection = node
.connect(EndpointSet::from(websocket_endpoint(url)))
.await
.unwrap();
let observer = connection.clone();
let changed = observer.changed();
node.shutdown();
assert_eq!(
tokio::time::timeout(Duration::from_secs(1), changed)
.await
.unwrap(),
ConnectionStatus::Disconnected {
reason: DisconnectReason::NodeShutdown,
}
);
assert_eq!(
connection.status(),
ConnectionStatus::Disconnected {
reason: DisconnectReason::NodeShutdown,
}
);
drop(server);
}
#[tokio::test(flavor = "multi_thread")]
async fn automatic_maintenance_restores_a_new_instance_of_the_same_peer() {
let admissions = Arc::new(AtomicUsize::new(0));
let release = Arc::new(tokio::sync::Notify::new());
let (first, first_url) = hosted("peer").await;
let replacement = Node::builder("peer")
.peer_layer_fn({
let admissions = admissions.clone();
let release = release.clone();
move |mut request, next| {
let admissions = admissions.clone();
let release = release.clone();
async move {
admissions.fetch_add(1, Ordering::SeqCst);
release.notified().await;
request.accept_declared();
next.admit(request).await
}
}
})
.build()
.unwrap();
let (second, second_url) = host_node(replacement).await;
let node = Node::builder("caller")
.connect_timeout(Duration::from_millis(200))
.insecure_accept_declared_peer_identities()
.build()
.unwrap();
let connection = node
.connect(EndpointSet::from([
websocket_endpoint(first_url),
websocket_endpoint(second_url),
]))
.await
.unwrap();
let connecting = connection.changed();
first.0.shutdown();
drop(first);
assert_eq!(connecting.await, ConnectionStatus::Connecting);
tokio::time::timeout(Duration::from_secs(1), async {
while admissions.load(Ordering::SeqCst) < 1 {
tokio::task::yield_now().await;
}
})
.await
.unwrap();
assert_eq!(connection.status(), ConnectionStatus::Connecting);
let connected = connection.changed();
release.notify_waiters();
assert_eq!(
tokio::time::timeout(Duration::from_secs(1), connected)
.await
.unwrap(),
ConnectionStatus::Connected
);
assert!(node.reachable_names().iter().any(|name| name == "peer"));
drop(second);
}
#[tokio::test(flavor = "multi_thread")]
async fn disconnect_rejects_a_late_automatic_candidate() {
let admissions = Arc::new(AtomicUsize::new(0));
let release = Arc::new(tokio::sync::Notify::new());
let (first, first_url) = hosted("peer").await;
let replacement = Node::builder("peer")
.peer_layer_fn({
let admissions = admissions.clone();
let release = release.clone();
move |mut request, next| {
let admissions = admissions.clone();
let release = release.clone();
async move {
admissions.fetch_add(1, Ordering::SeqCst);
release.notified().await;
request.accept_declared();
next.admit(request).await
}
}
})
.build()
.unwrap();
let (second, second_url) = host_node(replacement).await;
let node = Node::builder("caller")
.connect_timeout(Duration::from_millis(200))
.insecure_accept_declared_peer_identities()
.build()
.unwrap();
let connection = node
.connect(EndpointSet::from([
websocket_endpoint(first_url),
websocket_endpoint(second_url),
]))
.await
.unwrap();
first.0.shutdown();
drop(first);
tokio::time::timeout(Duration::from_secs(1), async {
while admissions.load(Ordering::SeqCst) < 1 {
tokio::task::yield_now().await;
}
})
.await
.unwrap();
connection.disconnect();
release.notify_waiters();
tokio::time::sleep(Duration::from_millis(50)).await;
assert_eq!(
connection.status(),
ConnectionStatus::Disconnected {
reason: DisconnectReason::ExplicitDisconnect,
}
);
drop(second);
}
#[tokio::test(flavor = "multi_thread")]
async fn deliberate_connect_after_disconnect_returns_a_new_handle_generation() {
let (peer, url) = hosted("peer").await;
let node = Node::builder("caller")
.insecure_accept_declared_peer_identities()
.build()
.unwrap();
let endpoints = EndpointSet::from(websocket_endpoint(url));
let old = node.connect(endpoints.clone()).await.unwrap();
old.disconnect();
let fresh = node.connect(endpoints).await.unwrap();
assert_eq!(fresh.status(), ConnectionStatus::Connected);
assert_eq!(
old.status(),
ConnectionStatus::Disconnected {
reason: DisconnectReason::ExplicitDisconnect,
}
);
fresh.disconnect();
assert_eq!(
old.status(),
ConnectionStatus::Disconnected {
reason: DisconnectReason::ExplicitDisconnect,
}
);
drop(peer);
}
#[tokio::test(flavor = "multi_thread")]
async fn automatic_maintenance_rejects_other_identities_and_keeps_retrying() {
let (first, first_url) = hosted("peer").await;
let admissions = Arc::new(AtomicUsize::new(0));
let other = Node::builder("other")
.peer_layer_fn({
let admissions = admissions.clone();
move |mut request, next| {
let admissions = admissions.clone();
async move {
admissions.fetch_add(1, Ordering::SeqCst);
request.accept_declared();
next.admit(request).await
}
}
})
.build()
.unwrap();
let (other, other_url) = host_node(other).await;
let node = Node::builder("caller")
.connect_timeout(Duration::from_millis(100))
.insecure_accept_declared_peer_identities()
.build()
.unwrap();
let connection = node
.connect(EndpointSet::from([
websocket_endpoint(first_url),
websocket_endpoint(other_url),
]))
.await
.unwrap();
let connecting = connection.changed();
first.0.shutdown();
drop(first);
assert_eq!(connecting.await, ConnectionStatus::Connecting);
tokio::time::timeout(Duration::from_secs(2), async {
while admissions.load(Ordering::SeqCst) < 2 {
tokio::task::yield_now().await;
}
})
.await
.unwrap();
assert_eq!(connection.status(), ConnectionStatus::Connecting);
connection.disconnect();
drop(other);
}