use super::link::Landed;
use serde_json::Value;
use std::collections::BTreeMap;
use std::sync::mpsc::{Receiver, Sender, channel};
pub const RECEIPTS_KEPT: usize = 64;
const NO_WIRE: &str = "this window has no wire behind it";
#[derive(Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Debug)]
pub struct Ticket(u64);
pub struct Post {
acts: Sender<(Ticket, Value)>,
receipts: Receiver<(Ticket, Landed)>,
next: u64,
landed: BTreeMap<Ticket, Landed>,
}
pub struct Outbox {
acts: Receiver<(Ticket, Value)>,
receipts: Sender<(Ticket, Landed)>,
}
pub fn pair() -> (Post, Outbox) {
let (a_tx, a_rx) = channel();
let (r_tx, r_rx) = channel();
(
Post {
acts: a_tx,
receipts: r_rx,
next: 0,
landed: BTreeMap::new(),
},
Outbox {
acts: a_rx,
receipts: r_tx,
},
)
}
impl Default for Post {
fn default() -> Self {
pair().0
}
}
impl Post {
pub fn send(&mut self, act: &Value) -> Ticket {
let ticket = Ticket(self.next);
self.next = self.next.wrapping_add(1);
if self.acts.send((ticket, act.clone())).is_err() {
self.keep(ticket, Err(NO_WIRE.to_owned()));
}
ticket
}
pub fn settle(&mut self) -> Vec<Ticket> {
let arrived: Vec<(Ticket, Landed)> = self.receipts.try_iter().collect();
let tickets = arrived.iter().map(|(ticket, _)| *ticket).collect();
for (ticket, landed) in arrived {
self.keep(ticket, landed);
}
tickets
}
pub fn receipt(&mut self, ticket: Ticket) -> Option<Landed> {
self.landed.remove(&ticket)
}
fn keep(&mut self, ticket: Ticket, landed: Landed) {
self.landed.insert(ticket, landed);
while self.landed.len() > RECEIPTS_KEPT {
self.landed.pop_first();
}
}
}
impl Outbox {
pub fn next(&self) -> Option<(Ticket, Value)> {
self.acts.recv().ok()
}
pub fn try_next(&self) -> Option<(Ticket, Value)> {
self.acts.try_recv().ok()
}
pub fn publish(&self, ticket: Ticket, landed: Landed) -> bool {
self.receipts.send((ticket, landed)).is_ok()
}
}
#[cfg(test)]
mod tests;