use std::net::{IpAddr, Ipv4Addr, SocketAddr, TcpStream, ToSocketAddrs};
use std::sync::Arc;
use std::time::Duration;
use super::Channel;
use super::line::{Held, Line};
use super::rendezvous::Pairing;
use super::rendezvous::item::{Call, Presence};
use super::rendezvous::punch::{self, Punch};
use crate::dht::{Config, Dht, Udp};
use crate::state;
#[derive(Clone, Debug)]
pub struct Roving {
pub direct: Duration,
pub window: Duration,
pub dht: Config,
pub bootstrap: Vec<String>,
pub advertise: Option<Vec<IpAddr>>,
}
impl Default for Roving {
fn default() -> Roving {
Roving {
direct: Duration::from_secs(5),
window: Duration::from_secs(40),
dht: Config::default(),
bootstrap: vec![
"router.bittorrent.com:6881".to_owned(),
"dht.transmissionbt.com:6881".to_owned(),
"router.utorrent.com:6881".to_owned(),
"dht.aelitis.com:6881".to_owned(),
],
advertise: None,
}
}
}
pub(crate) fn climb(ch: &Channel, held: bool) -> Result<Line, String> {
if held && let Some(line) = rung_held(ch) {
return Ok(line);
}
let direct = match rung_direct(ch) {
Ok(tcp) => return ch.line(tcp, false),
Err(refusal) => refusal,
};
let Some(pairing) = ch.pairing else {
return Err(direct);
};
if let Some(tcp) = rung_repunch(ch) {
return ch.line(tcp, true);
}
match rung_rendezvous(ch, &pairing) {
Ok(tcp) => ch.line(tcp, true),
Err(refusal) => Err(format!("{direct}; rendezvous: {refusal}")),
}
}
fn rung_held(ch: &Channel) -> Option<Line> {
let mut pool = state::worked(&ch.key, |w| std::mem::take(&mut w.held));
let now = ch.clock.now();
let mut chosen: Option<Held> = None;
while chosen.is_none()
&& let Some(mut held) = pool.pop()
{
if held.alive(now) {
chosen = Some(held);
}
}
state::worked(&ch.key, |w| w.held.append(&mut pool));
chosen.map(Line::from_held)
}
fn rung_direct(ch: &Channel) -> Result<TcpStream, String> {
let tcp = if ch.pairing.is_none() {
TcpStream::connect(&ch.address)
} else {
bounded(&ch.address, ch.roving.direct)
};
tcp.map_err(|e| format!("connect {}: {e}", ch.address))
}
fn bounded(address: &str, within: Duration) -> std::io::Result<TcpStream> {
let mut refusal = std::io::Error::other("resolved to no address");
for addr in address.to_socket_addrs()? {
match TcpStream::connect_timeout(&addr, within) {
Ok(tcp) => return Ok(tcp),
Err(e) => refusal = e,
}
}
Err(refusal)
}
fn rung_repunch(ch: &Channel) -> Option<TcpStream> {
let (endpoints, punch) = state::worked(&ch.key, |w| (w.endpoints.clone(), w.punch.clone()));
if endpoints.is_empty() {
return None;
}
punch?.punch(endpoints, ch.roving.window)
}
fn rung_rendezvous(ch: &Channel, pairing: &Pairing) -> 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() {
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 Some(item) = dht.get(pairing.engine, pairing.presence_salt())? else {
return Err("no presence is published under this pairing".to_owned());
};
let Some(presence) = Presence::open(&pairing.seal_key(), &item.value) else {
return Err("the presence item will not open under this pairing salt".to_owned());
};
let punch = punch_for(ch)?;
let mut ips = ch.roving.advertise.clone().unwrap_or_else(punch::local_ips);
for ip in dht.observed().iter().map(SocketAddr::ip) {
if !ips.contains(&ip) {
ips.push(ip);
}
}
let endpoints = 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,
};
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)?;
dht.put(signed)?;
let targets = presence.endpoints;
state::worked(&ch.key, |w| w.endpoints.clone_from(&targets));
match punch.punch(targets, ch.roving.window) {
Some(tcp) => Ok(tcp),
None => Err("nothing answered the punch inside its window".to_owned()),
}
}
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))
}))
}
#[cfg(test)]
mod tests;