use std::net::{IpAddr, Ipv4Addr, SocketAddr, TcpStream, ToSocketAddrs};
use std::sync::Arc;
use crate::channel::rendezvous::Pairing;
use crate::channel::rendezvous::item::{Call, Presence};
use crate::channel::rendezvous::punch::{self, Punch};
use crate::channel::{Channel, say};
use crate::dht::{Dht, Udp};
use crate::state;
type Written = fn(u64, i64, &[SocketAddr], usize) -> String;
pub(super) fn rungs(
ch: &Channel,
pairing: &Pairing,
said: &mut Vec<String>,
) -> Result<TcpStream, String> {
let bootstrap: Vec<SocketAddr> = ch
.roving
.bootstrap
.iter()
.filter_map(|name| name.to_socket_addrs().ok())
.flatten()
.collect();
if bootstrap.is_empty() {
said.push(say::no_bootstrap());
return Err("no bootstrap node resolved".to_owned());
}
let udp = Udp::bind(SocketAddr::new(IpAddr::V4(Ipv4Addr::UNSPECIFIED), 0))
.map_err(|e| format!("udp: {e}"))?;
let mut dht = Dht::new(Box::new(udp), bootstrap, ch.roving.dht.clone())?;
let cached = state::worked(&ch.key, |w| w.endpoints.clone());
if !cached.is_empty()
&& let Some(tcp) = call(ch, pairing, &mut dht, cached, say::recall, said)?
{
return Ok(tcp);
}
let got = dht.get(pairing.engine, pairing.presence_salt());
let Some(item) = got.map_err(|e| withheld(said, "presence not read", e))? else {
said.push(say::no_presence());
return Err("no presence is published under this pairing".to_owned());
};
let Some(presence) = Presence::open(&pairing.seal_key(), &item.value) else {
said.push(say::unopened(item.seq));
return Err("the presence item will not open under this pairing salt".to_owned());
};
said.push(say::presence(item.seq, &presence.endpoints));
call(ch, pairing, &mut dht, presence.endpoints, say::call, said)?
.ok_or_else(|| "nothing answered the punch inside its window".to_owned())
}
fn call(
ch: &Channel,
pairing: &Pairing,
dht: &mut Dht,
targets: Vec<SocketAddr>,
written: Written,
said: &mut Vec<String>,
) -> Result<Option<TcpStream>, String> {
let punch = punch_for(ch)?;
let mut ips = ch.roving.advertise.clone().unwrap_or_else(punch::local_ips);
for ip in seen(ch, dht) {
if !ips.contains(&ip) {
ips.push(ip);
}
}
said.extend(say::overlay(&ips));
let endpoints: Vec<SocketAddr> = ips
.into_iter()
.map(|ip| SocketAddr::new(ip, punch.port()))
.collect();
let mut nonce = [0u8; 8];
crate::dht::random(&mut nonce)?;
let call = Call {
nonce: u64::from_be_bytes(nonce),
endpoints: endpoints.clone(),
};
let sealed = call.seal(&pairing.seal_key())?;
let unix = ch.clock.unix();
let seq = state::worked(&ch.key, |w| {
w.last_seq = unix.max(w.last_seq + 1);
w.last_seq
});
let signed = pairing
.inbox_keypair()?
.sign(pairing.inbox_salt(), seq, sealed)?;
let unwritten = format!("call nonce {} not written", call.nonce);
let acks = dht.put(signed).map_err(|e| withheld(said, &unwritten, e))?;
said.push(written(call.nonce, seq, &endpoints, acks));
let what = format!("punch for nonce {}", call.nonce);
said.push(say::started(&what, &targets, ch.roving.window));
let tcp = punch.punch(targets.clone(), ch.roving.window);
said.push(match &tcp {
Some(tcp) => say::landed(&what, tcp.peer_addr().ok().map(|at| at.ip())),
None => say::expired(&what, ch.roving.window),
});
state::worked(&ch.key, |w| {
if tcp.is_some() {
w.endpoints = targets;
} else {
w.endpoints.clear();
w.observed.clear();
}
});
Ok(tcp)
}
fn seen(ch: &Channel, dht: &Dht) -> Vec<IpAddr> {
let fresh: Vec<IpAddr> = dht.observed().iter().map(SocketAddr::ip).collect();
state::worked(&ch.key, |w| {
if !fresh.is_empty() {
w.observed = fresh;
}
w.observed.clone()
})
}
fn withheld(said: &mut Vec<String>, what: &str, refusal: String) -> String {
said.push(say::withheld(what));
refusal
}
fn punch_for(ch: &Channel) -> Result<Arc<Punch>, String> {
if let Some(punch) = state::worked(&ch.key, |w| w.punch.clone()) {
return Ok(punch);
}
let bound = Arc::new(Punch::bind(0)?);
Ok(state::worked(&ch.key, |w| {
Arc::clone(w.punch.get_or_insert(bound))
}))
}