use std::sync::Arc;
use std::sync::atomic::AtomicBool;
use std::sync::mpsc::{Receiver, Sender, TryRecvError, channel};
use escriba_madoguchi::errand::{Crew, Errand, Freight, Parcel};
use escriba_madoguchi::{ErrandId, Negai};
use escriba_shirube::NonEmptyAnchor;
pub struct Courier {
tx: Sender<Parcel>,
rx: Receiver<Parcel>,
crew: Crew,
next: u32,
live: Vec<(ErrandId, Arc<AtomicBool>)>,
}
impl std::fmt::Debug for Courier {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(
f,
"Courier {{ next: {}, live: {} }}",
self.next,
self.live.len()
)
}
}
impl Courier {
#[must_use]
pub fn inert() -> Self {
let (tx, rx) = channel();
Self {
tx,
rx,
crew: Crew::inert(),
next: 0,
live: Vec::new(),
}
}
pub fn hire(&mut self, crew: Crew) {
self.crew = crew;
}
pub fn send(&mut self, freight: Freight, anchor: NonEmptyAnchor) -> ErrandId {
self.next = self.next.wrapping_add(1);
let id = ErrandId(self.next);
let cancel = Arc::new(AtomicBool::new(false));
self.live.push((id, Arc::clone(&cancel)));
self.crew.get(&freight).start(
Errand {
id,
freight,
anchor,
},
cancel,
self.tx.clone(),
);
id
}
pub fn drain(&mut self, budget: usize) -> Vec<Negai> {
let mut out = Vec::new();
for _ in 0..budget {
match self.rx.try_recv() {
Ok(p) => out.push(p.slip),
Err(TryRecvError::Empty | TryRecvError::Disconnected) => break,
}
}
out
}
pub fn cancel_all(&mut self) -> usize {
let n = self.live.len();
for (_, flag) in self.live.drain(..) {
flag.store(true, std::sync::atomic::Ordering::Relaxed);
}
n
}
}
#[cfg(test)]
mod tests {
use super::Courier;
use escriba_madoguchi::Negai;
use escriba_madoguchi::errand::{Crew, Errand, Freight, Parcel, Runner};
use escriba_shirube::{Axis, NonEmptyAnchor, SessionGen, SessionKind};
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::mpsc::Sender;
fn seal() -> NonEmptyAnchor {
NonEmptyAnchor::on(Axis::Session(SessionKind::Scan, SessionGen(1)))
}
fn a_scan() -> Freight {
Freight::Scan {
raw: "x".into(),
case: escriba_search::CaseMode::Smart,
root: ".".into(),
}
}
struct Chatty(usize);
impl Runner for Chatty {
fn start(&self, e: Errand, _c: Arc<AtomicBool>, reply: Sender<Parcel>) {
for i in 0..self.0 {
let _ = reply.send(Parcel {
id: e.id,
slip: Negai::Message(i.to_string()),
});
}
}
}
struct Watcher(Arc<AtomicBool>);
impl Runner for Watcher {
fn start(&self, _e: Errand, c: Arc<AtomicBool>, _r: Sender<Parcel>) {
self.0.store(c.load(Ordering::Relaxed), Ordering::Relaxed);
let mine = Arc::clone(&c);
std::mem::forget(mine);
}
}
fn crew_of(r: impl Runner + 'static) -> Crew {
Crew {
scan: Box::new(r),
diagnostics: Box::new(escriba_madoguchi::errand::Idle("t")),
format: Box::new(escriba_madoguchi::errand::Idle("t")),
}
}
#[test]
fn an_inert_courier_still_answers_rather_than_swallowing() {
let mut c = Courier::inert();
c.send(a_scan(), seal());
let got = c.drain(16);
assert_eq!(got.len(), 1, "an unhired crew must still say something");
assert!(matches!(got[0], Negai::Message(_)));
}
#[test]
fn the_drain_is_bounded_and_the_remainder_survives() {
let mut c = Courier::inert();
c.hire(crew_of(Chatty(10)));
c.send(a_scan(), seal());
assert_eq!(c.drain(4).len(), 4, "honours the budget");
assert_eq!(c.drain(4).len(), 4);
assert_eq!(c.drain(4).len(), 2, "the rest arrives later, not never");
assert!(c.drain(4).is_empty());
}
#[test]
fn draining_with_nothing_pending_returns_empty_without_blocking() {
let mut c = Courier::inert();
assert!(c.drain(64).is_empty());
}
#[test]
fn ids_are_distinct_per_dispatch() {
let mut c = Courier::inert();
c.hire(crew_of(Chatty(0)));
let a = c.send(a_scan(), seal());
let b = c.send(a_scan(), seal());
assert_ne!(a, b);
}
#[test]
fn cancelling_sets_the_flag_and_reports_how_many_were_asked() {
let seen = Arc::new(AtomicBool::new(false));
let mut c = Courier::inert();
c.hire(crew_of(Watcher(Arc::clone(&seen))));
c.send(a_scan(), seal());
assert!(!seen.load(Ordering::Relaxed), "not cancelled at dispatch");
assert_eq!(c.cancel_all(), 1, "one errand asked to stop");
assert_eq!(c.cancel_all(), 0, "nothing left to ask");
}
}