use std::collections::HashMap;
use std::net::{IpAddr, SocketAddr};
use std::sync::{Arc, Mutex, OnceLock, PoisonError};
use std::time::Duration;
use crate::channel::line::Held;
use crate::channel::rendezvous::punch::Punch;
use crate::ui::{Channel, Model, Posted};
mod traffic;
pub use traffic::{Heard, Open, Said, Standing};
#[derive(Default)]
struct Shared {
heard: Vec<Heard>,
outbox: Vec<Posted>,
standing: Standing,
place: crate::place::Place,
stopped: bool,
}
#[derive(Clone)]
pub struct Link {
shared: Arc<Mutex<Shared>>,
beat: Duration,
}
impl Link {
pub fn new(beat: Duration) -> Self {
Self {
shared: Arc::new(Mutex::new(Shared::default())),
beat,
}
}
pub fn beat(&self) -> Duration {
self.beat
}
pub fn settle(&self, model: &mut Model) {
let mut shared = self.hold();
for heard in std::mem::take(&mut shared.heard) {
match heard.said {
Said::Frame(frame) => model.absorb(&heard.channel, crate::reply::read(&frame)),
Said::Live { conversation, read } => {
if model.conversation.as_deref() == Some(conversation.as_str()) {
model.absorb(&heard.channel, read);
}
}
Said::Signin { provider, read } => {
if model.following().as_deref() == Some(provider.as_str()) {
model.absorb(&heard.channel, read);
}
}
Said::Unreachable(why) => model.unreachable(&heard.channel, why),
Said::Acted { op, reach, said } => model.acted(&op, &reach, said),
Said::Receipt { op, frame } => {
model.receipt(&heard.channel, &op, crate::reply::read(&frame));
}
}
}
shared.outbox.append(&mut model.outbox);
shared.standing = Standing::of(model);
shared.place = crate::place::Place::of(model);
}
pub fn heard(&self, channel: &Channel, said: Said) {
self.hold().heard.push(Heard {
channel: channel.clone(),
said,
});
}
pub fn live(&self, channel: &Channel, conversation: &str, read: crate::reply::Read) {
let said = Said::Live {
conversation: conversation.to_owned(),
read,
};
self.heard(channel, said);
}
pub fn signing(&self, channel: &Channel, provider: &str, read: crate::reply::Read) {
let said = Said::Signin {
provider: provider.to_owned(),
read,
};
self.heard(channel, said);
}
pub fn standing(&self) -> Standing {
self.hold().standing.clone()
}
pub fn place(&self) -> crate::place::Place {
self.hold().place.clone()
}
pub fn compose(&self) -> Vec<Posted> {
std::mem::take(&mut self.hold().outbox)
}
pub fn stop(&self) {
self.hold().stopped = true;
}
pub fn stopped(&self) -> bool {
self.hold().stopped
}
fn hold(&self) -> std::sync::MutexGuard<'_, Shared> {
self.shared.lock().unwrap_or_else(PoisonError::into_inner)
}
}
#[derive(Default)]
pub(crate) struct Worked {
pub(crate) held: Vec<Held>,
pub(crate) endpoints: Vec<SocketAddr>,
pub(crate) observed: Vec<IpAddr>,
pub(crate) punch: Option<Arc<Punch>>,
pub(crate) last_seq: i64,
pub(crate) said: HashMap<&'static str, Vec<String>>,
}
static WORKED: OnceLock<Mutex<HashMap<String, Worked>>> = OnceLock::new();
pub(crate) fn worked<T>(key: &str, f: impl FnOnce(&mut Worked) -> T) -> T {
let mut table = WORKED
.get_or_init(Mutex::default)
.lock()
.unwrap_or_else(PoisonError::into_inner);
f(table.entry(key.to_owned()).or_default())
}
#[cfg(test)]
mod tests;