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#[derive(Debug, Clone, Default)]
30pub struct Endpoint {
31 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#[derive(Debug)]
48pub struct Connecting {
49 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 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 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#[derive(Debug)]
90pub struct Session {
91 transport: Option<Rc<web_wt_sys::WebTransport>>,
93
94 pub datagrams: Datagrams,
96
97 pub close_on_drop: bool,
99}
100
101impl Session {
102 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 pub fn try_unwrap(mut self) -> Result<web_wt_sys::WebTransport, Self> {
115 let transport = self.transport.take().unwrap();
118
119 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 self.close_on_drop = false;
135
136 drop(self);
138
139 Ok(unwrapped)
141 }
142
143 pub const fn transport_ref(&self) -> &Rc<web_wt_sys::WebTransport> {
145 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#[derive(Debug)]
167pub enum DatagramsReader {
168 Byob(web_sys::ReadableStreamByobReader),
170
171 Default(web_sys::ReadableStreamDefaultReader),
174}
175
176impl DatagramsReader {
177 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 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#[derive(Debug)]
198pub struct Datagrams {
199 pub readable_stream_reader: DatagramsReader,
201
202 pub writable_stream_writer: web_sys::WritableStreamDefaultWriter,
204
205 pub read_buffer_size: u32,
208
209 pub read_buffer: tokio::sync::Mutex<Option<js_sys::ArrayBuffer>>,
211
212 pub unlock_streams_on_drop: bool,
214}
215
216impl Datagrams {
217 pub fn from_transport(transport: &web_wt_sys::WebTransport) -> Self {
219 Self::from_transport_datagrams(&transport.datagrams())
220 }
221
222 pub fn from_transport_datagrams(
224 datagrams: &web_wt_sys::WebTransportDatagramDuplexStream,
225 ) -> Self {
226 let read_buffer_size = 65536; let readable_stream_reader = DatagramsReader::for_stream(datagrams.readable());
229 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
270pub struct SendStream {
272 pub transport: Rc<web_wt_sys::WebTransport>,
274
275 pub stream: web_wt_sys::WebTransportSendStream,
277
278 pub writer: web_sys_async_io::Writer,
280
281 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
293pub struct RecvStream {
295 pub transport: Rc<web_wt_sys::WebTransport>,
297
298 pub stream: web_wt_sys::WebTransportReceiveStream,
300
301 pub reader: web_sys_async_io::Reader,
303
304 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
316fn wrap_recv_stream(
318 transport: &Rc<web_wt_sys::WebTransport>,
319 stream: web_wt_sys::WebTransportReceiveStream,
320) -> RecvStream {
321 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
344fn 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
359fn 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#[derive(Debug, thiserror::Error)]
486pub enum StreamWriteError {
487 #[error("zero size write buffer")]
489 ZeroSizeWriteBuffer,
490
491 #[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
512fn 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 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#[derive(Debug, thiserror::Error)]
594pub enum StreamReadError {
595 #[error("byob read consumed the buffer and didn't provide a new one")]
597 ByobReadConsumedBuffer,
598
599 #[error("read error: {0}")]
601 Read(Error),
602
603 #[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 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 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 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 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()) }
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}