use crate::boundary::reply::Reply;
use serde_json::Value;
use std::collections::{BTreeMap, BTreeSet};
use std::sync::mpsc::{Receiver, Sender, channel};
pub type Landed = Result<Reply, String>;
pub struct Link {
questions: Sender<Vec<Value>>,
answers: Receiver<(String, Landed)>,
standing: BTreeMap<String, Value>,
asked: BTreeSet<String>,
landed: BTreeMap<String, Landed>,
}
pub struct LinkEnd {
questions: Receiver<Vec<Value>>,
answers: Sender<(String, Landed)>,
standing: Vec<Value>,
}
pub fn pair() -> (Link, LinkEnd) {
let (q_tx, q_rx) = channel();
let (a_tx, a_rx) = channel();
(
Link {
questions: q_tx,
answers: a_rx,
standing: BTreeMap::new(),
asked: BTreeSet::new(),
landed: BTreeMap::new(),
},
LinkEnd {
questions: q_rx,
answers: a_tx,
standing: Vec::new(),
},
)
}
impl Default for Link {
fn default() -> Self {
pair().0
}
}
impl Link {
pub fn ask(&mut self, question: &Value) -> Option<Landed> {
let key = question.to_string();
let landed = self.landed.get(&key).cloned();
self.standing.insert(key, question.clone());
landed
}
#[cfg(test)]
pub fn awaiting(&self) -> bool {
self.asked.iter().any(|key| !self.landed.contains_key(key))
}
pub fn settle(&mut self) {
for (key, answer) in self.answers.try_iter() {
self.landed.insert(key, answer);
}
let wanted: BTreeSet<String> = self.standing.keys().cloned().collect();
if wanted != self.asked {
let _ = self
.questions
.send(self.standing.values().cloned().collect());
self.asked = wanted;
self.landed.retain(|key, _| self.asked.contains(key));
}
self.standing.clear();
}
}
impl LinkEnd {
pub fn standing(&mut self) -> Vec<Value> {
if let Some(newest) = self.questions.try_iter().last() {
self.standing = newest;
}
self.standing.clone()
}
pub fn publish(&self, question: &Value, landed: Landed) -> bool {
self.answers.send((question.to_string(), landed)).is_ok()
}
}
#[cfg(test)]
mod tests;