unb-server 2.0.3

unb inbound server: Node, request/subscribe handlers, catalog, relay orchestration, accept
Documentation
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"
    );
}