use super::accept::Toward;
use super::item::Call;
use super::punch::LINGER;
use super::{Punch, Say, Stats, say};
use crate::registry::presence::Presence;
use crate::wire::server::{Answerer, Quiet, serve};
use rustls::ServerConfig;
use std::net::{SocketAddr, TcpStream};
use std::sync::Arc;
use std::sync::atomic::Ordering;
use std::sync::mpsc::{self, Sender};
use std::time::{Duration, Instant};
#[derive(Clone)]
pub(super) struct Answer {
pub(super) punch: Arc<Punch>,
pub(super) tls: Arc<ServerConfig>,
pub(super) answerer: Arc<dyn Answerer>,
pub(super) presence: Presence,
pub(super) window: Duration,
pub(super) quiet: Quiet,
pub(super) stats: Arc<Stats>,
pub(super) say: Say,
pub(super) hints: Sender<Toward>,
}
impl Answer {
pub(super) fn spawn(self, call: Call) {
std::thread::spawn(move || {
let streams = self.window(call.endpoints);
if streams.is_empty() {
(self.say)(&say::expired(call.nonce, self.window));
return;
}
self.stats.punched.fetch_add(1, Ordering::Relaxed);
let peers: Vec<_> = streams
.iter()
.filter_map(|stream| stream.peer_addr().ok().as_ref().map(SocketAddr::ip))
.collect();
(self.say)(&say::landed(call.nonce, &peers));
for stream in streams {
self.stats.served.fetch_add(1, Ordering::Relaxed);
let answer = self.clone();
std::thread::spawn(move || {
serve(
stream,
&answer.tls,
answer.answerer.as_ref(),
&answer.presence,
answer.quiet,
);
(answer.say)(&say::ended(call.nonce));
});
}
});
}
fn window(&self, endpoints: Vec<SocketAddr>) -> Vec<TcpStream> {
let (tx, rx) = mpsc::sync_channel(0);
let _ = self.hints.send(Toward {
ips: endpoints.iter().map(SocketAddr::ip).collect(),
until: Instant::now() + self.window + LINGER,
tx,
});
self.punch
.toward(endpoints, self.window, &|| rx.try_iter().collect())
}
}