use crate::queue::Queue;
use defer_heavy::defer;
use std::io;
use std::io::{ErrorKind, Write};
use std::sync::{Arc, OnceLock};
#[derive(Debug, Default)]
struct WritePipeInner {
queue: Arc<Queue>,
error: OnceLock<ErrorKind>,
}
impl WritePipeInner {
fn handle<T: Write + Send>(&self, mut write: T) {
defer! {
self.queue.kill();
}
loop {
let pop = match self.queue.pop() {
Ok(guard) => guard,
Err(e) => {
_ = self.error.set(e.kind());
return;
}
};
if let Err(err) = write.write_all(pop.as_slice()) {
_ = self.error.set(err.kind());
return;
}
}
}
}
#[derive(Debug)]
pub struct WritePipe(Arc<WritePipeInner>);
impl Drop for WritePipe {
fn drop(&mut self) {
self.0.queue.kill();
}
}
impl WritePipe {
pub fn new<W: Write + Send + 'static, T: FnMut(Box<dyn FnOnce() + Send>) -> io::Result<()>>(
write: W,
spawner: &mut T,
) -> io::Result<Self> {
let wp = Arc::new(WritePipeInner::default());
let wpc = Arc::clone(&wp);
spawner(Box::new(move || {
wpc.handle(write);
}))?;
Ok(Self(wp))
}
pub fn dup_queue(&self) -> Arc<Queue> {
Arc::clone(&self.0.queue)
}
fn fetch_err(&self) -> io::Error {
self.0.queue.kill();
if let Some(err) = self.0.error.get().copied() {
return io::Error::from(err);
}
io::Error::from(ErrorKind::BrokenPipe)
}
}
impl Write for WritePipe {
fn write(&mut self, buf: &[u8]) -> io::Result<usize> {
match self.0.queue.push(buf.to_vec()) {
Ok(()) => Ok(buf.len()),
Err(err) => {
_ = self.0.error.set(err.kind());
Err(self.fetch_err())
}
}
}
fn flush(&mut self) -> io::Result<()> {
Ok(())
}
}