use std::time::Duration;
use super::super::doubt::in_doubt;
use super::{Invocation, MailCell, Mailbox, lock_mail};
pub(super) const HOLD_WAITS: u32 = 240;
pub(super) const HOLD_TICK: Duration = Duration::from_millis(125);
pub(super) const HAND_OFFS: u32 = 3;
impl Mailbox {
pub fn take(&self, client: &str) -> Result<Vec<Invocation>, String> {
let _reading = self.reading(client)?;
let mut ack = true;
for _ in 0..self.waits {
let taken = self.drain(client, ack);
if !taken.is_empty() {
return Ok(taken);
}
ack = false;
std::thread::sleep(self.tick);
}
Ok(self.drain(client, ack))
}
pub(crate) fn reading(&self, client: &str) -> Result<Reading, String> {
if !lock_mail(&self.cell).reading.insert(client.to_owned()) {
return Err(format!(
"invocations: {client:?} is already holding this engine's follow-class \
read — one machine's queue has one reader, because a second would take \
work the first is parked for and neither end would learn it. Something \
else is presenting this certificate: stop it, or stop this"
));
}
Ok(Reading {
cell: self.cell.clone(),
name: client.to_owned(),
})
}
pub fn serving(&self, client: &str) -> bool {
lock_mail(&self.cell).reading.contains(client)
}
fn drain(&self, client: &str, ack: bool) -> Vec<Invocation> {
let mut slots = lock_mail(&self.cell);
let mut out = Vec::new();
for slot in slots.live.values_mut() {
if slot.client != client || slot.capture.is_some() {
continue;
}
if slot.handed >= HAND_OFFS {
slot.capture = Some(in_doubt(client, slot.handed));
} else if ack || slot.handed == 0 {
slot.handed += 1;
out.push(slot.invocation.clone());
}
}
out
}
}
pub(crate) struct Reading {
cell: MailCell,
name: String,
}
impl Drop for Reading {
fn drop(&mut self) {
lock_mail(&self.cell).reading.remove(&self.name);
}
}