2018-10-20 14:38:47 -05:00
|
|
|
use std::thread;
|
|
|
|
|
2018-10-15 16:44:23 -05:00
|
|
|
use crossbeam_channel::{bounded, unbounded, Receiver, Sender};
|
|
|
|
use drop_bomb::DropBomb;
|
|
|
|
|
2018-10-20 14:38:47 -05:00
|
|
|
use crate::Result;
|
2018-09-02 06:46:15 -05:00
|
|
|
|
2018-09-08 05:15:01 -05:00
|
|
|
pub struct Worker<I, O> {
|
|
|
|
pub inp: Sender<I>,
|
|
|
|
pub out: Receiver<O>,
|
|
|
|
}
|
|
|
|
|
|
|
|
impl<I, O> Worker<I, O> {
|
|
|
|
pub fn spawn<F>(name: &'static str, buf: usize, f: F) -> (Self, ThreadWatcher)
|
|
|
|
where
|
|
|
|
F: FnOnce(Receiver<I>, Sender<O>) + Send + 'static,
|
|
|
|
I: Send + 'static,
|
|
|
|
O: Send + 'static,
|
|
|
|
{
|
2018-10-16 13:08:52 -05:00
|
|
|
let (worker, inp_r, out_s) = worker_chan(buf);
|
2018-09-08 05:15:01 -05:00
|
|
|
let watcher = ThreadWatcher::spawn(name, move || f(inp_r, out_s));
|
|
|
|
(worker, watcher)
|
|
|
|
}
|
|
|
|
|
|
|
|
pub fn stop(self) -> Receiver<O> {
|
|
|
|
self.out
|
|
|
|
}
|
|
|
|
|
|
|
|
pub fn send(&self, item: I) {
|
|
|
|
self.inp.send(item)
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2018-09-02 06:46:15 -05:00
|
|
|
pub struct ThreadWatcher {
|
|
|
|
name: &'static str,
|
|
|
|
thread: thread::JoinHandle<()>,
|
|
|
|
bomb: DropBomb,
|
|
|
|
}
|
|
|
|
|
|
|
|
impl ThreadWatcher {
|
2018-09-08 05:15:01 -05:00
|
|
|
fn spawn(name: &'static str, f: impl FnOnce() + Send + 'static) -> ThreadWatcher {
|
2018-09-02 06:46:15 -05:00
|
|
|
let thread = thread::spawn(f);
|
|
|
|
ThreadWatcher {
|
|
|
|
name,
|
|
|
|
thread,
|
|
|
|
bomb: DropBomb::new(format!("ThreadWatcher {} was not stopped", name)),
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
pub fn stop(mut self) -> Result<()> {
|
|
|
|
info!("waiting for {} to finish ...", self.name);
|
|
|
|
let name = self.name;
|
|
|
|
self.bomb.defuse();
|
2018-10-15 16:44:23 -05:00
|
|
|
let res = self
|
|
|
|
.thread
|
|
|
|
.join()
|
2018-09-02 06:46:15 -05:00
|
|
|
.map_err(|_| format_err!("ThreadWatcher {} died", name));
|
|
|
|
match &res {
|
|
|
|
Ok(()) => info!("... {} terminated with ok", name),
|
2018-10-15 16:44:23 -05:00
|
|
|
Err(_) => error!("... {} terminated with err", name),
|
2018-09-02 06:46:15 -05:00
|
|
|
}
|
|
|
|
res
|
|
|
|
}
|
|
|
|
}
|
2018-09-08 04:36:02 -05:00
|
|
|
|
|
|
|
/// Sets up worker channels in a deadlock-avoind way.
|
|
|
|
/// If one sets both input and output buffers to a fixed size,
|
|
|
|
/// a worker might get stuck.
|
2018-10-16 13:08:52 -05:00
|
|
|
fn worker_chan<I, O>(buf: usize) -> (Worker<I, O>, Receiver<I>, Sender<O>) {
|
2018-09-08 04:36:02 -05:00
|
|
|
let (input_sender, input_receiver) = bounded::<I>(buf);
|
|
|
|
let (output_sender, output_receiver) = unbounded::<O>();
|
2018-10-15 16:44:23 -05:00
|
|
|
(
|
2018-10-16 13:08:52 -05:00
|
|
|
Worker {
|
|
|
|
inp: input_sender,
|
|
|
|
out: output_receiver,
|
|
|
|
},
|
2018-10-15 16:44:23 -05:00
|
|
|
input_receiver,
|
|
|
|
output_sender,
|
|
|
|
)
|
2018-09-08 04:36:02 -05:00
|
|
|
}
|