use super::call::Answer;
use super::item::{Call, Unopened};
use super::material::Pairing;
use super::{Cadence, Ctx, Say, Stats, say};
use crate::dht::{Dht, Keypair, Udp};
use crate::ui_state::Clock;
use std::net::{IpAddr, SocketAddr, ToSocketAddrs};
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
use std::time::Instant;
mod publish;
pub(super) struct Cycle {
pairing: Pairing,
keypair: Keypair,
inbox_key: [u8; 32],
bootstrap: Vec<String>,
config: crate::dht::Config,
transport: Option<Udp>,
dht: Option<Dht>,
advertise: Vec<IpAddr>,
answer: Answer,
clock: Arc<dyn Clock>,
cadence: Cadence,
stats: Arc<Stats>,
say: Say,
last_said: Option<String>,
last_seq: i64,
last_nonce: Option<u64>,
next_publish: Instant,
next_poll: Instant,
}
impl Cycle {
pub(super) fn new(ctx: Ctx) -> Result<Cycle, String> {
let now = ctx.clock.now();
let answer = Answer {
punch: Arc::new(ctx.punch),
tls: ctx.tls,
answerer: ctx.answerer,
presence: ctx.presence,
window: ctx.cadence.window,
quiet: ctx.cadence.quiet,
stats: Arc::clone(&ctx.stats),
say: Arc::clone(&ctx.say),
};
Ok(Cycle {
keypair: ctx.pairing.keypair()?,
inbox_key: ctx.pairing.inbox_keypair()?.public(),
pairing: ctx.pairing,
bootstrap: ctx.bootstrap,
config: ctx.config,
transport: Some(ctx.transport),
dht: None,
advertise: ctx.advertise,
answer,
clock: ctx.clock,
cadence: ctx.cadence,
stats: ctx.stats,
say: ctx.say,
last_said: None,
last_seq: 0,
last_nonce: None,
next_publish: now,
next_poll: now,
})
}
pub(super) fn punch_port(&self) -> u16 {
self.answer.punch.port()
}
pub(super) fn run(mut self, stop: &AtomicBool) {
while !stop.load(Ordering::Relaxed) {
let now = self.clock.now();
if now >= self.next_poll {
let unix = self.clock.unix().max(0) as u64;
self.stats.last_poll_unix.store(unix, Ordering::Relaxed);
let outcome = self.poll();
self.tell(outcome.clone().unwrap_or_else(|_| Some(say::poll_failed())));
count(
outcome.map(drop),
&self.stats.polls,
&self.stats.poll_failures,
);
self.next_poll = now + self.cadence.poll;
}
if now >= self.next_publish {
let outcome = self.publish();
(self.say)(&outcome.clone().unwrap_or_else(|_| say::not_published()));
count(
outcome.map(drop),
&self.stats.published,
&self.stats.publish_failures,
);
self.next_publish = now + self.cadence.publish;
}
std::thread::sleep(self.cadence.tick);
}
}
fn dht(&mut self) -> Result<&mut Dht, String> {
if self.dht.is_none() {
let bootstrap: Vec<SocketAddr> = self
.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 transport = self.transport.take().ok_or("the transport is spent")?;
self.dht = Some(Dht::new(
Box::new(transport),
bootstrap,
self.config.clone(),
)?);
}
self.dht.as_mut().ok_or_else(|| "no DHT client".to_owned())
}
fn poll(&mut self) -> Result<Option<String>, String> {
let (key, salt) = (self.inbox_key, self.pairing.inbox_salt());
let Some(item) = self.dht()?.get(key, salt)? else {
return Ok(None);
};
let call = match Call::open(&self.pairing.seal_key(), &item.value) {
Ok(call) => call,
Err(Unopened::Unverified) => return Ok(Some(say::unverified(item.seq))),
Err(Unopened::NotACall) => return Ok(Some(say::unopened(item.seq))),
};
if self.last_nonce == Some(call.nonce) {
return Ok(Some(say::seen(call.nonce)));
}
self.last_nonce = Some(call.nonce);
self.stats.calls.fetch_add(1, Ordering::Relaxed);
let line = Some(say::opened(call.nonce, &call.endpoints));
self.tell(line.clone());
self.answer.clone().spawn(call);
Ok(line)
}
fn tell(&mut self, line: Option<String>) {
if line != self.last_said
&& let Some(said) = &line
{
(self.say)(said);
}
self.last_said = line;
}
}
fn count(outcome: Result<(), String>, ok: &AtomicUsize, failed: &AtomicUsize) {
match outcome {
Ok(()) => ok.fetch_add(1, Ordering::Relaxed),
Err(_) => failed.fetch_add(1, Ordering::Relaxed),
};
}