File
Blob: src/rust/kj/io.rs
| 1 | use std::pin::Pin; |
| 2 | |
| 3 | use kj_rs::KjOwn; |
| 4 | |
| 5 | use crate::OwnOrMut; |
| 6 | use crate::Result; |
| 7 | |
| 8 | #[cxx::bridge(namespace = "kj::rust")] |
| 9 | pub mod ffi { |
| 10 | unsafe extern "C++" { |
| 11 | include!("workerd/rust/kj/ffi.h"); |
| 12 | |
| 13 | type AsyncInputStream; |
| 14 | type AsyncIoStream; |
| 15 | type AsyncOutputStream; |
| 16 | |
| 17 | async fn async_output_stream_write( |
| 18 | this_: Pin<&mut AsyncOutputStream>, |
| 19 | buffer: &[u8], |
| 20 | ) -> Result<()>; |
| 21 | |
| 22 | async fn async_output_stream_when_write_disconnected( |
| 23 | this_: Pin<&mut AsyncOutputStream>, |
| 24 | ) -> Result<()>; |
| 25 | } |
| 26 | } |
| 27 | |
| 28 | pub type AsyncInputStream = ffi::AsyncInputStream; |
| 29 | pub type AsyncIoStream = ffi::AsyncIoStream; |
| 30 | |
| 31 | /// Owned-or-borrowed wrapper for `kj::AsyncOutputStream`. |
| 32 | pub struct AsyncOutputStream<'a>(OwnOrMut<'a, ffi::AsyncOutputStream>); |
| 33 | |
| 34 | impl AsyncOutputStream<'_> { |
| 35 | pub async fn write(&mut self, buffer: &[u8]) -> Result<()> { |
| 36 | let stream = self.0.as_mut(); |
| 37 | ffi::async_output_stream_write(stream, buffer).await?; |
| 38 | Ok(()) |
| 39 | } |
| 40 | |
| 41 | pub async fn when_write_disconnected(&mut self) -> Result<()> { |
| 42 | let stream = self.0.as_mut(); |
| 43 | ffi::async_output_stream_when_write_disconnected(stream).await?; |
| 44 | Ok(()) |
| 45 | } |
| 46 | } |
| 47 | |
| 48 | impl From<KjOwn<ffi::AsyncOutputStream>> for AsyncOutputStream<'_> { |
| 49 | fn from(value: KjOwn<ffi::AsyncOutputStream>) -> Self { |
| 50 | Self(value.into()) |
| 51 | } |
| 52 | } |
| 53 | |
| 54 | impl<'a> From<Pin<&'a mut ffi::AsyncOutputStream>> for AsyncOutputStream<'a> { |
| 55 | fn from(value: Pin<&'a mut ffi::AsyncOutputStream>) -> Self { |
| 56 | Self(value.into()) |
| 57 | } |
| 58 | } |