unb-server 1.0.0

unb inbound server: Node, request/subscribe handlers, catalog, relay orchestration, accept
Documentation
use std::collections::BTreeSet;
use std::sync::Arc;
use std::time::Duration;

use unb_core::{
    CoreInput, Detail, DiscoverEvent, DiscoverPlan, DiscoverWalk, Mode, Scope, StreamKey,
    WalkInput, WalkOutput, DEFAULT_HOPS,
};
use unb_runtime::{CancellationToken, ProtocolCoreHandle};
use tokio::sync::mpsc;

use crate::node::Node;

const DISCOVER_TIMEOUT: Duration = Duration::from_secs(5);

impl Node {
    pub(crate) fn query_discovery_neighbor(
        self: &Arc<Self>,
        stream: StreamKey,
        peer: String,
        plan: DiscoverPlan,
        handle: ProtocolCoreHandle,
    ) {
        let node = self.clone();
        tokio::spawn(async move {
            let (input, mut output) = mpsc::channel(64);
            let cancel = node.cancellation.child_token();
            let timeout = plan
                .timeout_ms
                .map(Duration::from_millis)
                .unwrap_or(DISCOVER_TIMEOUT);
            let leg = node.run_discovery_leg(&peer, plan, &input, &cancel, timeout);
            tokio::pin!(leg);
            loop {
                tokio::select! {
                    biased;
                    event = output.recv() => match event {
                        Some(WalkInput::NeighborEvent { event, .. }) => {
                            if handle.submit(CoreInput::DiscoveryNeighborEvent {
                                stream: stream.clone(),
                                peer: peer.clone(),
                                event,
                            }).await.is_err() {
                                return;
                            }
                        }
                        _ => return,
                    },
                    completed = &mut leg => {
                        while let Ok(WalkInput::NeighborEvent { event, .. }) = output.try_recv() {
                            if handle.submit(CoreInput::DiscoveryNeighborEvent {
                                stream: stream.clone(),
                                peer: peer.clone(),
                                event,
                            }).await.is_err() {
                                return;
                            }
                        }
                        let input = if completed {
                            CoreInput::DiscoveryNeighborDone {
                                stream,
                                peer: peer.clone(),
                            }
                        } else {
                            CoreInput::DiscoveryNeighborTimeout {
                                stream,
                                peer: peer.clone(),
                            }
                        };
                        let _ = handle.submit(input).await;
                        return;
                    }
                }
            }
        });
    }

    pub async fn discover_events(
        self: &Arc<Self>,
        detail: Detail,
        scope: Scope,
    ) -> Vec<DiscoverEvent> {
        let plan = DiscoverPlan {
            discover_id: String::new(),
            detail,
            scope,
            hops: DEFAULT_HOPS,
            visited: BTreeSet::new(),
            timeout_ms: None,
            mode: Mode::PartialOk,
        };
        let discover_id = plan.discover_id.clone();
        let snapshot = self
            .snapshot
            .load()
            .node_core
            .catalog_snapshot(plan.detail.is_full());
        let candidates = self.discovery_candidates(&plan.visited).await;
        let mut walk = DiscoverWalk::start(snapshot, plan, candidates);
        let (input_tx, mut input_rx) = mpsc::channel::<WalkInput>(64);
        let cancel = self.cancellation.child_token();
        let mut events = Vec::new();

        loop {
            while let Some(output) = walk.drain() {
                match output {
                    WalkOutput::Emit(event) => events.push(event),
                    WalkOutput::AskNeighbor { peer, plan } => self.spawn_discovery_leg(
                        peer,
                        plan,
                        input_tx.clone(),
                        DISCOVER_TIMEOUT,
                        cancel.child_token(),
                    ),
                    WalkOutput::Finish => {
                        events.push(DiscoverEvent::Done {
                            discover_id: discover_id.clone(),
                        });
                        cancel.cancel();
                        return events;
                    }
                }
            }
            tokio::select! {
                biased;
                () = cancel.cancelled() => return events,
                input = input_rx.recv() => match input {
                    Some(input) => walk.handle(input),
                    None => return events,
                }
            }
        }
    }

    pub(crate) async fn discovery_candidates(&self, visited: &BTreeSet<String>) -> Vec<String> {
        self.peers
            .read()
            .await
            .keys()
            .filter(|name| !visited.contains(*name))
            .cloned()
            .collect()
    }

    fn spawn_discovery_leg(
        self: &Arc<Self>,
        peer: String,
        plan: DiscoverPlan,
        input: mpsc::Sender<WalkInput>,
        timeout: Duration,
        cancel: CancellationToken,
    ) {
        let node = self.clone();
        tokio::spawn(async move {
            let completed = node
                .run_discovery_leg(&peer, plan, &input, &cancel, timeout)
                .await;
            let feedback = if completed {
                WalkInput::NeighborDone { peer }
            } else {
                WalkInput::NeighborTimeout { peer }
            };
            let _ = input.send(feedback).await;
        });
    }

    async fn run_discovery_leg(
        &self,
        peer: &str,
        plan: DiscoverPlan,
        input: &mpsc::Sender<WalkInput>,
        cancel: &CancellationToken,
        timeout: Duration,
    ) -> bool {
        let Ok(Ok(_permit)) =
            tokio::time::timeout(timeout, self.dispatch_slots.clone().acquire_owned()).await
        else {
            return false;
        };
        let Some(link) = self.peer(peer).await else {
            return false;
        };
        let mut stream = match link.wire.client_session().discover(plan).await {
            Ok(stream) => stream,
            Err(_) => return false,
        };
        let deadline = tokio::time::sleep(timeout);
        tokio::pin!(deadline);
        let completed = loop {
            tokio::select! {
                biased;
                () = cancel.cancelled() => break false,
                () = &mut deadline => break false,
                msg = stream.next() => match msg {
                    Ok(Some(frame)) => match frame.kind {
                        unb_core::Kind::Event => {
                            if let Ok(event) =
                                serde_json::from_slice::<DiscoverEvent>(&frame.payload)
                            {
                                if input
                                    .send(WalkInput::NeighborEvent {
                                        peer: peer.to_string(),
                                        event,
                                    })
                                    .await
                                    .is_err()
                                {
                                    break false;
                                }
                            }
                        }
                        unb_core::Kind::Response => break true,
                        _ => {}
                    },
                    Ok(None) => break false,
                    Err(_) => break false,
                }
            }
        };
        completed
    }
}