Skip to main content

xwt_web/
lib.rs

1#![cfg_attr(
2    target_family = "wasm",
3    doc = "The [`web_wt_sys`]-powered implementation of [`xwt_core`]."
4)]
5#![cfg_attr(
6    not(target_family = "wasm"),
7    doc = "The `web_wt_sys`-powered implementation of `xwt_core`."
8)]
9#![cfg(target_family = "wasm")]
10
11use std::{num::NonZeroUsize, rc::Rc};
12
13use wasm_bindgen::prelude::*;
14
15mod error;
16mod error_as_error_code;
17mod options;
18
19pub use web_sys;
20pub use web_wt_sys;
21pub use xwt_core as core;
22
23pub use {error::*, options::*};
24
25/// An endpoint for the xwt.
26///
27/// Internally holds the connection options and can create
28/// a new `WebTransport` object on the "web" side on a connection request.
29#[derive(Debug, Clone, Default)]
30pub struct Endpoint {
31    /// The options to use to create the `WebTransport`s.
32    pub options: web_wt_sys::WebTransportOptions,
33}
34
35impl xwt_core::endpoint::Connect for Endpoint {
36    type Error = Error;
37    type Connecting = Connecting;
38
39    async fn connect(&self, url: &str) -> Result<Self::Connecting, Self::Error> {
40        let transport = web_wt_sys::WebTransport::new_with_options(url, &self.options)?;
41        Ok(Connecting { transport })
42    }
43}
44
45/// Connecting represents the transient connection state when
46/// the [`web_wt_sys::WebTransport`] has been created but is not ready yet.
47#[derive(Debug)]
48pub struct Connecting {
49    /// The WebTransport instance.
50    pub transport: web_wt_sys::WebTransport,
51}
52
53impl xwt_core::endpoint::connect::Connecting for Connecting {
54    type Session = Session;
55    type Error = Error;
56
57    async fn wait_connect(self) -> Result<Self::Session, Self::Error> {
58        let Connecting { transport } = self;
59
60        // Wait for whichever of the two settles first.
61        //
62        // Waiting on `ready` alone is not enough: Safari does not reliably
63        // settle it when the session fails to establish, while `closed` does
64        // settle in that case - so watching both keeps a refused session from
65        // hanging forever.
66        let ready = transport.ready();
67        let closed = transport.closed();
68        let candidates = js_sys::Array::of2(&ready, &closed);
69        let outcome =
70            wasm_bindgen_futures::JsFuture::from(js_sys::Promise::race(&candidates)).await?;
71
72        // `ready` fulfills with `undefined`, while `closed` fulfills with
73        // a close info, so any other value means the session was closed before
74        // it ever became ready.
75        if !outcome.is_undefined() {
76            return Err(Error(
77                JsError::new("xwt: the session was closed before it became ready").into(),
78            ));
79        }
80
81        Ok(Session::new(transport))
82    }
83}
84
85/// Session holds the [`web_wt_sys::WebTransport`] and is responsible for
86/// providing access to the Web API of WebTransport in a way that is portable.
87/// It also holds handles to the datagram reader and writer, as well as
88/// the datagram reader state.
89#[derive(Debug)]
90pub struct Session {
91    /// The WebTransport instance.
92    transport: Option<Rc<web_wt_sys::WebTransport>>,
93
94    /// The datagrams state for this session.
95    pub datagrams: Datagrams,
96
97    /// Whether to close the session on drop.
98    pub close_on_drop: bool,
99}
100
101impl Session {
102    /// Construct a new session from a [`web_wt_sys::WebTransport`].
103    pub fn new(transport: web_wt_sys::WebTransport) -> Self {
104        let datagrams = Datagrams::from_transport(&transport);
105        Self {
106            transport: Some(Rc::new(transport)),
107            datagrams,
108            close_on_drop: true,
109        }
110    }
111
112    /// If possible, relieves the underlying [`web_wt_sys::WebTransport`] of
113    /// any `xwt-web`-held locks and dependencies and exposes it.
114    pub fn try_unwrap(mut self) -> Result<web_wt_sys::WebTransport, Self> {
115        // Take the transport out of the session; this state is not valid
116        // "publicly", only while we're in this fn.
117        let transport = self.transport.take().unwrap();
118
119        // We want to ensure there are no other references to
120        // the transport's [`Rc`], otherwise something must still be using it.
121        // If we can't unwrap it successfully - we reinsert the transport back
122        // into `self` and return it as an `Err` - it will thus remain
123        // operational to permit doing whatever work may have to be done
124        // further.
125        let unwrapped = match Rc::try_unwrap(transport) {
126            Ok(unwrapped) => unwrapped,
127            Err(transport) => {
128                let _ = self.transport.insert(transport);
129                return Err(self);
130            }
131        };
132
133        // Do not close the transport (we have taken it out anyway).
134        self.close_on_drop = false;
135
136        // Drop the session to release the datagram readers/writers.
137        drop(self);
138
139        // Return the unwrapped transport.
140        Ok(unwrapped)
141    }
142
143    /// Obtain a transport ref.
144    pub const fn transport_ref(&self) -> &Rc<web_wt_sys::WebTransport> {
145        // Trnasport should never be gone generally, only inside of
146        // the `try_unwrap`.
147        self.transport.as_ref().unwrap()
148    }
149}
150
151impl Drop for Session {
152    fn drop(&mut self) {
153        if self.close_on_drop {
154            self.transport_ref().close();
155        }
156    }
157}
158
159/// The reader for the datagrams readable stream.
160///
161/// Per the WebTransport spec, the datagrams readable is a regular (non-byte)
162/// [`web_sys::ReadableStream`], which only supports default readers
163/// (this is what Firefox implements); Chrome, however, exposes it as
164/// a byte stream, allowing BYOB reads.
165/// We feature-detect BYOB support and fall back to a default reader.
166#[derive(Debug)]
167pub enum DatagramsReader {
168    /// A BYOB reader, used when the datagrams readable is a byte stream.
169    Byob(web_sys::ReadableStreamByobReader),
170
171    /// A default reader, used when the datagrams readable is not
172    /// a byte stream.
173    Default(web_sys::ReadableStreamDefaultReader),
174}
175
176impl DatagramsReader {
177    /// Acquire a reader for the given datagrams readable stream, preferring
178    /// a BYOB reader when the stream supports it.
179    pub fn for_stream(readable_stream: web_sys::ReadableStream) -> Self {
180        match web_sys_stream_utils::try_get_reader_byob(readable_stream.clone()) {
181            Ok(reader) => Self::Byob(reader),
182            Err(_) => Self::Default(web_sys_stream_utils::get_reader(readable_stream)),
183        }
184    }
185
186    /// Release the stream lock held by the reader.
187    pub fn release_lock(&self) {
188        match self {
189            Self::Byob(reader) => reader.release_lock(),
190            Self::Default(reader) => reader.release_lock(),
191        }
192    }
193}
194
195/// Datagrams hold the portions of the session that are responsible for working
196/// with the datagrams.
197#[derive(Debug)]
198pub struct Datagrams {
199    /// The datagram reader.
200    pub readable_stream_reader: DatagramsReader,
201
202    /// The datagram writer.
203    pub writable_stream_writer: web_sys::WritableStreamDefaultWriter,
204
205    /// The desired size of the datagram read buffer.
206    /// Used to allocate the datagram read buffer in case it gets lost.
207    pub read_buffer_size: u32,
208
209    /// The datagram read internal buffer.
210    pub read_buffer: tokio::sync::Mutex<Option<js_sys::ArrayBuffer>>,
211
212    /// Unlock the streams on drop.
213    pub unlock_streams_on_drop: bool,
214}
215
216impl Datagrams {
217    /// Create a datagrams state from the transport.
218    pub fn from_transport(transport: &web_wt_sys::WebTransport) -> Self {
219        Self::from_transport_datagrams(&transport.datagrams())
220    }
221
222    /// Create a datagrams state from the transport datagrams.
223    pub fn from_transport_datagrams(
224        datagrams: &web_wt_sys::WebTransportDatagramDuplexStream,
225    ) -> Self {
226        let read_buffer_size = 65536; // 65k buffers as per spec recommendation
227
228        let readable_stream_reader = DatagramsReader::for_stream(datagrams.readable());
229        // Feature-detect `createWritable`; fall back to the legacy `writable`
230        // attribute on browsers that do not implement it yet.
231        let writable: web_sys::WritableStream = if datagrams.has_create_writable() {
232            datagrams.create_writable().unwrap().into()
233        } else {
234            #[expect(deprecated)]
235            let writable = datagrams.writable();
236            writable
237        };
238        let writable_stream_writer = web_sys_stream_utils::get_writer(writable);
239
240        let read_buffer = js_sys::ArrayBuffer::new(read_buffer_size);
241        let read_buffer = tokio::sync::Mutex::new(Some(read_buffer));
242
243        Self {
244            readable_stream_reader,
245            writable_stream_writer,
246            read_buffer_size,
247            read_buffer,
248            unlock_streams_on_drop: true,
249        }
250    }
251}
252
253impl Drop for Datagrams {
254    fn drop(&mut self) {
255        if self.unlock_streams_on_drop {
256            self.readable_stream_reader.release_lock();
257            self.writable_stream_writer.release_lock();
258        }
259    }
260}
261
262impl xwt_core::session::stream::SendSpec for Session {
263    type SendStream = SendStream;
264}
265
266impl xwt_core::session::stream::RecvSpec for Session {
267    type RecvStream = RecvStream;
268}
269
270/// Send the data into a WebTransport stream.
271pub struct SendStream {
272    /// The WebTransport instance.
273    pub transport: Rc<web_wt_sys::WebTransport>,
274
275    /// The handle to the stream to write to.
276    pub stream: web_wt_sys::WebTransportSendStream,
277
278    /// A writer to conduct the operation.
279    pub writer: web_sys_async_io::Writer,
280
281    /// Unlock the writer on drop.
282    pub unlock_writer_on_drop: bool,
283}
284
285impl Drop for SendStream {
286    fn drop(&mut self) {
287        if self.unlock_writer_on_drop {
288            self.writer.inner.release_lock();
289        }
290    }
291}
292
293/// Recv the data from a WebTransport stream.
294pub struct RecvStream {
295    /// The WebTransport instance.
296    pub transport: Rc<web_wt_sys::WebTransport>,
297
298    /// The handle to the stream to read from.
299    pub stream: web_wt_sys::WebTransportReceiveStream,
300
301    /// A reader to conduct the operation.
302    pub reader: web_sys_async_io::Reader,
303
304    /// Unlock the reader on drop.
305    pub unlock_reader_on_drop: bool,
306}
307
308impl Drop for RecvStream {
309    fn drop(&mut self) {
310        if self.unlock_reader_on_drop {
311            self.reader.inner.release_lock();
312        }
313    }
314}
315
316/// Open a reader for the given stream and create a [`RecvStream`].
317fn wrap_recv_stream(
318    transport: &Rc<web_wt_sys::WebTransport>,
319    stream: web_wt_sys::WebTransportReceiveStream,
320) -> RecvStream {
321    // Per the WebTransport spec, the receive stream is a byte stream
322    // supporting BYOB reads (this is what Chrome and Firefox implement);
323    // Safari, however, does not implement it as a byte stream.
324    // We feature-detect BYOB support and fall back to a default reader.
325    let reader = match web_sys_stream_utils::try_get_reader_byob(stream.clone()) {
326        Ok(reader) => web_sys_async_io::reader::Mode::Byob {
327            reader,
328            internal_buf: None,
329        },
330        Err(_) => web_sys_async_io::reader::Mode::Default {
331            reader: web_sys_stream_utils::get_reader(stream.clone()),
332        },
333    };
334    let reader = web_sys_async_io::Reader::new(reader);
335
336    RecvStream {
337        transport: Rc::clone(transport),
338        stream,
339        reader,
340        unlock_reader_on_drop: true,
341    }
342}
343
344/// Open a writer for the given stream and create a [`SendStream`].
345fn wrap_send_stream(
346    transport: &Rc<web_wt_sys::WebTransport>,
347    stream: web_wt_sys::WebTransportSendStream,
348) -> SendStream {
349    let writer = stream.get_writer().unwrap();
350    let writer = web_sys_async_io::Writer::new(writer.into());
351    SendStream {
352        transport: Rc::clone(transport),
353        stream,
354        writer,
355        unlock_writer_on_drop: true,
356    }
357}
358
359/// Take a bidi stream and wrap a reader and writer for it.
360fn wrap_bi_stream(
361    transport: &Rc<web_wt_sys::WebTransport>,
362    stream: web_wt_sys::WebTransportBidirectionalStream,
363) -> (SendStream, RecvStream) {
364    let writable = stream.writable();
365    let readable = stream.readable();
366
367    let send_stream = wrap_send_stream(transport, writable);
368    let recv_stream = wrap_recv_stream(transport, readable);
369
370    (send_stream, recv_stream)
371}
372
373impl xwt_core::session::stream::OpenBi for Session {
374    type Opening = xwt_core::utils::dummy::OpeningBiStream<Session>;
375
376    type Error = Error;
377
378    async fn open_bi(&self) -> Result<Self::Opening, Self::Error> {
379        let transport = self.transport_ref();
380        let value =
381            wasm_bindgen_futures::JsFuture::from(transport.create_bidirectional_stream()).await?;
382        let value = wrap_bi_stream(transport, value);
383        Ok(xwt_core::utils::dummy::OpeningBiStream(value))
384    }
385}
386
387impl xwt_core::session::stream::AcceptBi for Session {
388    type Error = Error;
389
390    async fn accept_bi(&self) -> Result<(Self::SendStream, Self::RecvStream), Self::Error> {
391        let transport = self.transport_ref();
392        let incoming: web_sys::ReadableStream = transport.incoming_bidirectional_streams();
393        let reader: JsValue = incoming.get_reader().into();
394        let reader: web_sys::ReadableStreamDefaultReader = reader.into();
395        let read_result = wasm_bindgen_futures::JsFuture::from(reader.read()).await?;
396        let read_result: web_wt_sys::ReadableStreamReadResult<
397            web_wt_sys::WebTransportBidirectionalStream,
398        > = read_result.unchecked_into();
399        if read_result.is_done() {
400            return Err(Error(JsError::new("xwt: accept bi reader is done").into()));
401        }
402        let Some(value) = read_result.get_value() else {
403            return Err(Error(
404                JsError::new("xwt: accept bi read result has no value").into(),
405            ));
406        };
407        let value = wrap_bi_stream(transport, value);
408        Ok(value)
409    }
410}
411
412impl xwt_core::session::stream::OpenUni for Session {
413    type Opening = xwt_core::utils::dummy::OpeningUniStream<Session>;
414    type Error = Error;
415
416    async fn open_uni(&self) -> Result<Self::Opening, Self::Error> {
417        let transport = self.transport_ref();
418        let value =
419            wasm_bindgen_futures::JsFuture::from(transport.create_unidirectional_stream()).await?;
420        let send_stream = wrap_send_stream(transport, value);
421        Ok(xwt_core::utils::dummy::OpeningUniStream(send_stream))
422    }
423}
424
425impl xwt_core::session::stream::AcceptUni for Session {
426    type Error = Error;
427
428    async fn accept_uni(&self) -> Result<Self::RecvStream, Self::Error> {
429        let transport = self.transport_ref();
430        let incoming: web_sys::ReadableStream = transport.incoming_unidirectional_streams();
431        let reader: JsValue = incoming.get_reader().into();
432        let reader: web_sys::ReadableStreamDefaultReader = reader.into();
433        let read_result = wasm_bindgen_futures::JsFuture::from(reader.read()).await?;
434        let read_result: web_wt_sys::ReadableStreamReadResult<
435            web_wt_sys::WebTransportReceiveStream,
436        > = read_result.unchecked_into();
437        if read_result.is_done() {
438            return Err(Error(JsError::new("xwt: accept uni reader is done").into()));
439        }
440        let Some(value) = read_result.get_value() else {
441            return Err(Error(
442                JsError::new("xwt: accept uni read result has no value").into(),
443            ));
444        };
445        let recv_stream = wrap_recv_stream(transport, value);
446        Ok(recv_stream)
447    }
448}
449
450impl tokio::io::AsyncWrite for SendStream {
451    fn poll_write(
452        mut self: std::pin::Pin<&mut Self>,
453        cx: &mut std::task::Context<'_>,
454        buf: &[u8],
455    ) -> std::task::Poll<Result<usize, std::io::Error>> {
456        std::pin::Pin::new(&mut self.writer).poll_write(cx, buf)
457    }
458
459    fn poll_flush(
460        mut self: std::pin::Pin<&mut Self>,
461        cx: &mut std::task::Context<'_>,
462    ) -> std::task::Poll<Result<(), std::io::Error>> {
463        std::pin::Pin::new(&mut self.writer).poll_flush(cx)
464    }
465
466    fn poll_shutdown(
467        mut self: std::pin::Pin<&mut Self>,
468        cx: &mut std::task::Context<'_>,
469    ) -> std::task::Poll<Result<(), std::io::Error>> {
470        std::pin::Pin::new(&mut self.writer).poll_shutdown(cx)
471    }
472}
473
474impl tokio::io::AsyncRead for RecvStream {
475    fn poll_read(
476        mut self: std::pin::Pin<&mut Self>,
477        cx: &mut std::task::Context<'_>,
478        buf: &mut tokio::io::ReadBuf<'_>,
479    ) -> std::task::Poll<std::io::Result<()>> {
480        std::pin::Pin::new(&mut self.reader).poll_read(cx, buf)
481    }
482}
483
484/// An error that can occur during the stream writes.
485#[derive(Debug, thiserror::Error)]
486pub enum StreamWriteError {
487    /// The write was called with a zero-size write buffer.
488    #[error("zero size write buffer")]
489    ZeroSizeWriteBuffer,
490
491    /// The write call thrown an exception.
492    #[error("write error: {0}")]
493    Write(Error),
494}
495
496impl xwt_core::stream::Write for SendStream {
497    type Error = StreamWriteError;
498
499    async fn write(&mut self, buf: &[u8]) -> Result<NonZeroUsize, Self::Error> {
500        let Some(buf_len) = NonZeroUsize::new(buf.len()) else {
501            return Err(StreamWriteError::ZeroSizeWriteBuffer);
502        };
503
504        web_sys_stream_utils::write(&self.writer.inner, buf)
505            .await
506            .map_err(|err| StreamWriteError::Write(err.into()))?;
507
508        Ok(buf_len)
509    }
510}
511
512/// Build the abort reason that carries the given stream error code.
513///
514/// The code only reaches the peer when the reason is a
515/// [`web_wt_sys::WebTransportError`] with the `streamErrorCode` set; for any
516/// other reason the code `0` is sent instead.
517fn stream_abort_reason(error_code: xwt_core::stream::ErrorCode) -> JsValue {
518    let options = web_wt_sys::WebTransportErrorOptions::new();
519    options.set_source(web_wt_sys::WebTransportErrorSource::Stream);
520    options.set_stream_error_code(error_code);
521    web_wt_sys::WebTransportError::new_with_init(&options).into()
522}
523
524impl xwt_core::stream::WriteAbort for SendStream {
525    type Error = Error;
526
527    async fn abort(self, error_code: xwt_core::stream::ErrorCode) -> Result<(), Self::Error> {
528        wasm_bindgen_futures::JsFuture::from(
529            self.writer
530                .inner
531                .abort_with_reason(&stream_abort_reason(error_code)),
532        )
533        .await
534        .map(|val| {
535            debug_assert!(val.is_undefined());
536        })
537        .map_err(Error::from)
538    }
539}
540
541impl xwt_core::stream::WriteAborted for SendStream {
542    type Error = Error;
543
544    async fn aborted(self) -> Result<xwt_core::stream::ErrorCode, Self::Error> {
545        // Hack our way through...
546        let result = wasm_bindgen_futures::JsFuture::from(self.writer.inner.closed()).await;
547        match result {
548            Ok(value) => {
549                debug_assert!(value.is_undefined());
550                Ok(0)
551            }
552            Err(value) => {
553                let error: web_wt_sys::WebTransportError = value.dyn_into().unwrap();
554                if error.source() != web_wt_sys::WebTransportErrorSource::Stream {
555                    return Err(Error(error.into()));
556                }
557                let Some(code) = error.stream_error_code() else {
558                    return Err(Error(error.into()));
559                };
560                Ok(code)
561            }
562        }
563    }
564}
565
566impl xwt_core::stream::Finish for SendStream {
567    type Error = Error;
568
569    async fn finish(self) -> Result<(), Self::Error> {
570        wasm_bindgen_futures::JsFuture::from(self.writer.inner.close())
571            .await
572            .map(|val| {
573                debug_assert!(val.is_undefined());
574            })
575            .map_err(Error::from)
576    }
577}
578
579impl xwt_core::stream::Finished for RecvStream {
580    type Error = Error;
581
582    async fn finished(self) -> Result<(), Self::Error> {
583        wasm_bindgen_futures::JsFuture::from(self.reader.inner.closed())
584            .await
585            .map(|val| {
586                debug_assert!(val.is_undefined());
587            })
588            .map_err(Error::from)
589    }
590}
591
592/// An error that can occur while reading stream data.
593#[derive(Debug, thiserror::Error)]
594pub enum StreamReadError {
595    /// This is an odd case, which is still tbd what to do with.
596    #[error("byob read consumed the buffer and didn't provide a new one")]
597    ByobReadConsumedBuffer,
598
599    /// The `read_byob` call thrown an exception.
600    #[error("read error: {0}")]
601    Read(Error),
602
603    /// The stream was closed, and there is no more data to expect there.
604    #[error("stream closed")]
605    Closed,
606}
607
608impl From<web_sys_async_io::ReadError> for StreamReadError {
609    fn from(err: web_sys_async_io::ReadError) -> Self {
610        match err {
611            web_sys_async_io::ReadError::Read(err) => Self::Read(err.into()),
612            web_sys_async_io::ReadError::ByobReadConsumedBuffer => Self::ByobReadConsumedBuffer,
613        }
614    }
615}
616
617impl xwt_core::stream::Read for RecvStream {
618    type Error = StreamReadError;
619
620    async fn read(&mut self, buf: &mut [u8]) -> Result<NonZeroUsize, Self::Error> {
621        let len = self.reader.read_into(buf).await?;
622
623        // Detect when the read is aborted because the stream was closed
624        // without an error.
625        NonZeroUsize::new(len).ok_or(StreamReadError::Closed)
626    }
627}
628
629impl xwt_core::stream::ReadAbort for RecvStream {
630    type Error = Error;
631
632    async fn abort(self, error_code: xwt_core::stream::ErrorCode) -> Result<(), Self::Error> {
633        wasm_bindgen_futures::JsFuture::from(
634            self.reader
635                .inner
636                .cancel_with_reason(&stream_abort_reason(error_code)),
637        )
638        .await
639        .map(|_| ())
640        .map_err(Error::from)
641    }
642}
643
644impl xwt_core::stream::ReadAborted for RecvStream {
645    type Error = Error;
646
647    async fn aborted(self) -> Result<xwt_core::stream::ErrorCode, Self::Error> {
648        // Hack our way through...
649        let result = wasm_bindgen_futures::JsFuture::from(self.reader.inner.closed()).await;
650        match result {
651            Ok(value) => {
652                debug_assert!(value.is_undefined());
653                Ok(0)
654            }
655            Err(value) => {
656                let error: web_wt_sys::WebTransportError = value.dyn_into().unwrap();
657                if error.source() != web_wt_sys::WebTransportErrorSource::Stream {
658                    return Err(Error(error.into()));
659                }
660                let Some(code) = error.stream_error_code() else {
661                    return Err(Error(error.into()));
662                };
663                Ok(code)
664            }
665        }
666    }
667}
668
669impl Datagrams {
670    /// Receive the datagram and handle the buffer with the given function.
671    ///
672    /// Cloning the buffer in the `f` will result in the undefined behaviour,
673    /// because it will create a second reference to an object that is intended
674    /// to be under a `mut ref`.
675    /// Although is would not teachnically be unsafe, it would violate
676    /// the borrow checker rules.
677    pub async fn receive_with<R>(
678        &self,
679        max_read_size: Option<u32>,
680        f: impl FnOnce(&mut js_sys::Uint8Array) -> R,
681    ) -> Result<R, Error> {
682        let mut buffer_guard = self.read_buffer.lock().await;
683
684        match &self.readable_stream_reader {
685            DatagramsReader::Byob(reader) => {
686                let buffer = buffer_guard
687                    .take()
688                    .unwrap_or_else(|| js_sys::ArrayBuffer::new(self.read_buffer_size));
689                let view = if let Some(max_read_size) = max_read_size {
690                    let desired_buffer_length = buffer.byte_length().min(max_read_size);
691                    js_sys::Uint8Array::new_with_byte_offset_and_length(
692                        &buffer,
693                        0,
694                        desired_buffer_length,
695                    )
696                } else {
697                    js_sys::Uint8Array::new(&buffer)
698                };
699
700                let maybe_view = web_sys_stream_utils::read_byob(reader, view).await?;
701                let Some(mut view) = maybe_view else {
702                    return Err(wasm_bindgen::JsError::new("unexpected stream termination").into());
703                };
704
705                let result = f(&mut view);
706
707                *buffer_guard = Some(view.buffer());
708                Ok(result)
709            }
710            DatagramsReader::Default(reader) => {
711                let maybe_view = web_sys_stream_utils::read_uint8array(reader).await?;
712                let Some(view) = maybe_view else {
713                    return Err(wasm_bindgen::JsError::new("unexpected stream termination").into());
714                };
715
716                // A default reader always yields the whole datagram; honor
717                // the requested read size limit by truncating the view.
718                let mut view = match max_read_size {
719                    Some(max_read_size) if view.length() > max_read_size => {
720                        view.subarray(0, max_read_size)
721                    }
722                    _ => view,
723                };
724
725                let result = f(&mut view);
726
727                Ok(result)
728            }
729        }
730    }
731}
732
733impl xwt_core::session::datagram::MaxSize for Session {
734    fn max_datagram_size(&self) -> Option<usize> {
735        let transport = self.transport_ref();
736        let max_datagram_size = transport.datagrams().max_datagram_size();
737        Some(usize::try_from(max_datagram_size).unwrap()) // u32 should fit in a usize on WASM
738    }
739}
740
741impl xwt_core::session::datagram::Receive for Session {
742    type Datagram = Vec<u8>;
743    type Error = Error;
744
745    async fn receive_datagram(&self) -> Result<Self::Datagram, Self::Error> {
746        self.datagrams
747            .receive_with(None, |buffer| buffer.to_vec())
748            .await
749    }
750}
751
752impl xwt_core::session::datagram::ReceiveInto for Session {
753    type Error = Error;
754
755    async fn receive_datagram_into(&self, buf: &mut [u8]) -> Result<usize, Self::Error> {
756        let max_read_size = buf.len().try_into().unwrap();
757        self.datagrams
758            .receive_with(Some(max_read_size), |buffer| {
759                let len = buffer.length() as usize;
760                buffer.copy_to(&mut buf[..len]);
761                len
762            })
763            .await
764    }
765}
766
767impl xwt_core::session::datagram::Send for Session {
768    type Error = Error;
769
770    async fn send_datagram<D>(&self, payload: D) -> Result<(), Self::Error>
771    where
772        D: AsRef<[u8]>,
773    {
774        web_sys_stream_utils::write(&self.datagrams.writable_stream_writer, payload.as_ref())
775            .await?;
776        Ok(())
777    }
778}