use crate::queue::Queue;
use defer_heavy::defer;
use std::io;
use std::io::{Cursor, ErrorKind, Read};
use std::sync::{Arc, OnceLock};
#[derive(Debug, Default)]
struct ReadPipeInner {
queue: Arc<Queue>,
error: OnceLock<ErrorKind>,
}
impl ReadPipeInner {
fn handle<T: Read + Send>(&self, mut read: T) {
defer! {
self.queue.kill();
}
let mut buffer = vec![0u8; 0x1_00_00];
loop {
let packet = match read.read(buffer.as_mut_slice()) {
Ok(count) => buffer[0..count].to_vec(),
Err(err) => {
_ = self.error.set(err.kind());
return;
}
};
if packet.is_empty() {
if let Err(err) = self.queue.push(packet) {
_ = self.error.set(err.kind());
}
return;
}
if let Err(err) = self.queue.push(packet) {
_ = self.error.set(err.kind());
}
}
}
}
#[derive(Debug)]
pub struct ReadPipe {
eof: bool,
nb: bool,
pipe: Arc<ReadPipeInner>,
cursor: Cursor<Vec<u8>>,
}
impl Drop for ReadPipe {
fn drop(&mut self) {
self.pipe.queue.kill();
}
}
impl ReadPipe {
pub fn new<R: Read + Send + 'static, T: FnMut(Box<dyn FnOnce() + Send>) -> io::Result<()>>(
write: R,
spawner: &mut T,
initial_data: Vec<u8>,
) -> io::Result<Self> {
let wp = Arc::new(ReadPipeInner::default());
let wpc = Arc::clone(&wp);
spawner(Box::new(move || {
wpc.handle(write);
}))?;
Ok(Self {
nb: false,
eof: false,
pipe: wp,
cursor: Cursor::new(initial_data),
})
}
pub fn nb(&mut self, value: bool) {
self.nb = value;
}
pub fn dup_queue(&self) -> Arc<Queue> {
Arc::clone(&self.pipe.queue)
}
fn fetch_err(&self) -> io::Error {
self.pipe.queue.kill();
if let Some(err) = self.pipe.error.get().copied() {
return io::Error::from(err);
}
io::Error::from(ErrorKind::BrokenPipe)
}
}
impl Read for ReadPipe {
fn read(&mut self, buf: &mut [u8]) -> io::Result<usize> {
if self.eof || buf.is_empty() {
return Ok(0);
}
loop {
let size = self.cursor.read(buf)?;
if size != 0 {
return Ok(size);
}
if self.nb {
match self.pipe.queue.try_pop() {
Ok(Some(data)) => {
if data.is_empty() {
self.eof = true;
return Ok(0);
}
self.cursor = Cursor::new(data);
continue;
}
Ok(None) => {
return Err(io::Error::from(ErrorKind::WouldBlock)); }
Err(err) => {
_ = self.pipe.error.set(err.kind());
return Err(self.fetch_err());
}
}
}
match self.pipe.queue.pop() {
Ok(data) => {
if data.is_empty() {
self.eof = true;
return Ok(0);
}
self.cursor = Cursor::new(data);
}
Err(err) => {
_ = self.pipe.error.set(err.kind());
return Err(self.fetch_err());
}
}
}
}
}