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
}
}