use crossbeam::channel::Receiver;
use std::fs::File;
use std::io::{self, BufWriter, ErrorKind, Result, Write};
pub fn write_loop(outfile: Option<&str>, write_rx: Receiver<Vec<u8>>) -> Result<()> {
let mut writer: Box<dyn Write> = match outfile {
Some(path) => Box::new(BufWriter::new(File::create(path)?)),
None => Box::new(BufWriter::new(io::stdout())),
};
loop {
let buffer = write_rx.recv().unwrap();
if buffer.is_empty() {
break;
}
if let Err(e) = writer.write_all(&buffer) {
if e.kind() == ErrorKind::BrokenPipe {
return Ok(());
}
return Err(e);
}
}
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
use crossbeam::channel::unbounded;
use std::fs::File;
use std::io::Read;
#[test]
fn test_write_loop() -> Result<()> {
let (write_tx, write_rx) = unbounded();
let test_data = b"Pipe Progress...";
write_tx.send(test_data.to_vec()).unwrap();
write_tx.send(Vec::new()).unwrap();
let test_output = "test_output.txt";
write_loop(Some(test_output), write_rx)?;
let mut file = File::open(test_output)?;
let mut received_data = Vec::new();
file.read_to_end(&mut received_data)?;
assert_eq!(test_data, received_data.as_slice());
std::fs::remove_file(test_output)?;
Ok(())
}
}