diff --git a/proxmox-async/src/blocking/mod.rs b/proxmox-async/src/blocking/mod.rs index 28247b34..06f821a9 100644 --- a/proxmox-async/src/blocking/mod.rs +++ b/proxmox-async/src/blocking/mod.rs @@ -9,3 +9,6 @@ pub use tokio_writer_adapter::TokioWriterAdapter; mod wrapped_reader_stream; pub use wrapped_reader_stream::WrappedReaderStream; + +mod sender_writer; +pub use sender_writer::SenderWriter; diff --git a/proxmox-async/src/blocking/sender_writer.rs b/proxmox-async/src/blocking/sender_writer.rs new file mode 100644 index 00000000..62682e5f --- /dev/null +++ b/proxmox-async/src/blocking/sender_writer.rs @@ -0,0 +1,47 @@ +use std::io; + +use anyhow::Error; +use tokio::sync::mpsc::Sender; + +/// Wrapper struct around [`tokio::sync::mpsc::Sender`] for `Result, Error>` that implements [`std::io::Write`] +pub struct SenderWriter { + sender: Sender, Error>>, +} + +impl SenderWriter { + pub fn from_sender(sender: tokio::sync::mpsc::Sender, Error>>) -> Self { + Self { sender } + } + + fn write_impl(&mut self, buf: &[u8]) -> io::Result { + if let Err(err) = self.sender.blocking_send(Ok(buf.to_vec())) { + return Err(io::Error::new( + io::ErrorKind::UnexpectedEof, + format!("could not send: {}", err), + )); + } + + Ok(buf.len()) + } + + fn flush_impl(&mut self) -> io::Result<()> { + Ok(()) + } +} + +impl io::Write for SenderWriter { + fn write(&mut self, buf: &[u8]) -> io::Result { + self.write_impl(buf) + } + + fn flush(&mut self) -> io::Result<()> { + self.flush_impl() + } +} + +impl Drop for SenderWriter { + fn drop(&mut self) { + // ignore errors + let _ = self.flush_impl(); + } +}