use std::io::Write;
use pdfboss_core::source::BoxFuture;
use crate::error::{Error, Result};
pub trait AsyncByteSink {
fn write_all<'a>(&'a mut self, buf: &'a [u8]) -> BoxFuture<'a, Result<()>>;
}
impl<S: AsyncByteSink + ?Sized> AsyncByteSink for &mut S {
fn write_all<'a>(&'a mut self, buf: &'a [u8]) -> BoxFuture<'a, Result<()>> {
(**self).write_all(buf)
}
}
impl AsyncByteSink for Vec<u8> {
fn write_all<'a>(&'a mut self, buf: &'a [u8]) -> BoxFuture<'a, Result<()>> {
self.extend_from_slice(buf);
Box::pin(std::future::ready(Ok(())))
}
}
#[derive(Debug, Clone)]
pub struct Immediate<W>(pub W);
impl<W: Write> AsyncByteSink for Immediate<W> {
fn write_all<'a>(&'a mut self, buf: &'a [u8]) -> BoxFuture<'a, Result<()>> {
let outcome = self.0.write_all(buf).map_err(Error::from);
Box::pin(std::future::ready(outcome))
}
}
#[cfg(test)]
mod tests {
use std::cell::RefCell;
use std::io::Write;
use std::rc::Rc;
use pdfboss_core::block_on;
use super::{AsyncByteSink, Immediate};
#[test]
fn a_vec_is_a_sink() {
let mut sink = Vec::new();
block_on(AsyncByteSink::write_all(&mut sink, b"abc")).expect("a Vec accepts everything");
block_on(AsyncByteSink::write_all(&mut sink, b"def")).expect("a Vec accepts everything");
assert_eq!(sink, b"abcdef");
}
#[test]
fn a_reference_to_a_sink_is_a_sink() {
fn feed<S: AsyncByteSink>(mut sink: S) -> S {
block_on(sink.write_all(b"xy")).expect("the test sinks accept everything");
sink
}
let mut sink = Vec::new();
feed(&mut sink);
let sink = feed(sink);
assert_eq!(sink, b"xyxy");
}
#[test]
fn immediate_writes_through_to_the_writer() {
let mut sink = Immediate(Vec::new());
block_on(sink.write_all(b"hello")).expect("a Vec accepts everything");
assert_eq!(sink.0, b"hello");
}
struct SharedWriter(Rc<RefCell<Vec<u8>>>);
impl Write for SharedWriter {
fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
self.0.borrow_mut().extend_from_slice(buf);
Ok(buf.len())
}
fn flush(&mut self) -> std::io::Result<()> {
Ok(())
}
}
#[test]
fn immediate_futures_are_send_even_over_a_non_send_writer() {
fn assert_send<T: Send>(_: &T) {}
let shared = Rc::new(RefCell::new(Vec::new()));
let mut sink = Immediate(SharedWriter(Rc::clone(&shared)));
let future = sink.write_all(b"eager");
assert_send(&future);
block_on(future).expect("the shared writer accepts everything");
assert_eq!(*shared.borrow(), b"eager");
}
#[test]
fn immediate_writes_eagerly() {
let shared = Rc::new(RefCell::new(Vec::new()));
let mut sink = Immediate(SharedWriter(Rc::clone(&shared)));
let unpolled = sink.write_all(b"already there");
assert_eq!(*shared.borrow(), b"already there");
drop(unpolled);
}
struct Refusing;
impl Write for Refusing {
fn write(&mut self, _: &[u8]) -> std::io::Result<usize> {
Err(std::io::Error::other("refused"))
}
fn flush(&mut self) -> std::io::Result<()> {
Ok(())
}
}
#[test]
fn immediate_surfaces_write_errors() {
let mut sink = Immediate(Refusing);
let err = block_on(sink.write_all(b"x")).unwrap_err();
assert!(matches!(err, crate::error::Error::Io(_)));
}
}