Skip to content
File

Blob: src/rust/kj/io.rs

rust59 lines
1use std::pin::Pin;
2 
3use kj_rs::KjOwn;
4 
5use crate::OwnOrMut;
6use crate::Result;
7 
8#[cxx::bridge(namespace = "kj::rust")]
9pub 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 
28pub type AsyncInputStream = ffi::AsyncInputStream;
29pub type AsyncIoStream = ffi::AsyncIoStream;
30 
31/// Owned-or-borrowed wrapper for `kj::AsyncOutputStream`.
32pub struct AsyncOutputStream<'a>(OwnOrMut<'a, ffi::AsyncOutputStream>);
33 
34impl 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 
48impl From<KjOwn<ffi::AsyncOutputStream>> for AsyncOutputStream<'_> {
49 fn from(value: KjOwn<ffi::AsyncOutputStream>) -> Self {
50 Self(value.into())
51 }
52}
53 
54impl<'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}