From eadaca2d408a9c1d547aa069949aeb8eb6fde3bd Mon Sep 17 00:00:00 2001 From: Luke Curley Date: Sat, 26 Sep 2026 10:03:12 -0700 Subject: [PATCH 1/5] quest: claim close-codes From c50ec6f1338b732ce69d4879421d8765fc894fc0 Mon Sep 17 00:00:00 2001 From: Luke Curley Date: Sat, 26 Sep 2026 10:23:01 -0700 Subject: [PATCH 2/5] test(tokio): the peer's close code reaches the client on every transport Runs the #4249 cases (abort after accept, abort then drop, reject during the handshake) over https://, moqt://, and ws://. Fails until the qmux and web-transport-moq fixes are released and pinned. Co-Authored-By: Claude Opus 5.5 --- quest/m1/README.md | 1 - quest/m1/close-codes.md | 39 --------- quest/m1/quic/qmux.md | 4 +- rs/moq-tokio/tests/close_code.rs | 131 +++++++++++++++++++++++++++++++ 4 files changed, 133 insertions(+), 42 deletions(-) delete mode 100644 quest/m1/close-codes.md create mode 100644 rs/moq-tokio/tests/close_code.rs diff --git a/quest/m1/README.md b/quest/m1/README.md index 5f91306218..126e3c3466 100644 --- a/quest/m1/README.md +++ b/quest/m1/README.md @@ -30,7 +30,6 @@ transport, benchmark tooling); worktrees isolate commits, not semantics. - [Kotlin end](/quest/m1/kotlin-end.md) - Kotlin exposes `close()` as `end()`, since `AutoCloseable.close()` takes the name - [Remove finish](/quest/m1/broadcast-remove.md) - on dev, the deprecated broadcast end APIs are gone and `closed()` carries no cause - [Session close](/quest/m1/session-close.md) - a graceful session end withdraws announces and waits one second for the ack -- [Close codes](/quest/m1/close-codes.md) - a client sees the peer's application close code over WebSocket and raw QUIC, like WebTransport - [JS caught up](/quest/m1/js-announce-caught-up.md) - @moq/net's announce consumer says when the initial set has landed, like Rust - [Bindings caught up](/quest/m1/announce-live-bindings.md) - moq-ffi, libmoq, and every wrapper yield the same flat announce event, `Live` included - [IETF announce count](/quest/m1/ietf-announce-count.md) - an opt-in moq-transport extension carries the replay count, so IETF announce consumers go live without a timer diff --git a/quest/m1/close-codes.md b/quest/m1/close-codes.md deleted file mode 100644 index f1e3e9a34f..0000000000 --- a/quest/m1/close-codes.md +++ /dev/null @@ -1,39 +0,0 @@ -# [M] Close codes on every transport - -## Goal - -A client sees the application close code its peer sent, for example -`SessionError::App(4011)` from `Session::abort` or `Request::reject`, over -WebSocket (qmux) and raw QUIC (`moqt://`), as it already does over -WebTransport. Never `Transport("connection closed")` or its own `Internal`. - -## Plan - -Both bugs are upstream; fix them at the source, release, and bump the pins. - -- qmux 0.5.1 (`moq-dev/web-transport`) lets later writes overwrite the - recorded close in `session.rs`: the WS Close frame read after - APPLICATION_CLOSE (reader loop, backend `send_replace`), a local `close()` - after the peer closed, and a second peer APPLICATION_CLOSE. `accept_uni` and - `accept_bi` also return a bare `Closed`. Make the first close win, make - `close()` a no-op once closed, and have `accept_*` return the recorded - reason. Add a qmux test where APPLICATION_CLOSE and EOF arrive together. -- `web-transport-moq` 1.3.1 (`moq-dev/noq`) maps `ApplicationClosed` only - through the HTTP/3 code space in `error.rs`, so a raw `moqt://` code yields - no `session_error()`. Map raw QUIC codes directly. -- moq-net's `close(Internal)` after a transport error is correct: closing a - closed connection does nothing. Do not work around it here. -- One moq-tokio regression runs the issue's three cases (abort after accept, - abort then drop, reject during handshake) over `https://`, `ws://`, and - `moqt://`, and fails on the current pins. -- Check whether `@moq/net`'s qmux peer keeps the first close too; fix it in - the same PR if not. - -## Closes - -- [#4249](https://github.com/moq-dev/moq/issues/4249) - application close code is lost over the WebSocket (qmux) transport - -## Related - -- [qmux on noq](/quest/m1/quic/qmux.md) - the rewrite must keep first-close-wins -- [io_uring close](/quest/m1/quic/uring-close.md) - the same symptom class on the io_uring backend diff --git a/quest/m1/quic/qmux.md b/quest/m1/quic/qmux.md index c67682ead8..aaf405fdb3 100644 --- a/quest/m1/quic/qmux.md +++ b/quest/m1/quic/qmux.md @@ -38,8 +38,8 @@ Add the missing wire evidence before release: golden draft-02 vectors, bidirectional interoperability against the published `qmux` 0.5.x crate, and the TypeScript qmux/WebSocket peer used by `js/net`. Preserve rejection of prohibited QUIC frames, params-first setup, record-size validation, close and -reset semantics (the first recorded close wins, per -[close codes](/quest/m1/close-codes.md)), keep-alive behavior, and bounded +reset semantics (the first recorded close wins, which moq-tokio's +`close_code` test checks end to end), keep-alive behavior, and bounded flow-control tests. The crate lives in the fork's workspace as `moq-noq-qmux`, so the stream diff --git a/rs/moq-tokio/tests/close_code.rs b/rs/moq-tokio/tests/close_code.rs new file mode 100644 index 0000000000..3dd3346636 --- /dev/null +++ b/rs/moq-tokio/tests/close_code.rs @@ -0,0 +1,131 @@ +//! A client sees the application close code its server sent, on every transport. +//! +//! Each case runs several fresh connections, since the losses this guards against were +//! races between the peer's close and whatever the transport reported next. + +#![cfg(all(feature = "noq", feature = "websocket"))] + +use std::time::Duration; + +use moq_tokio::moq_net; + +const CODE: u16 = 4011; +const RUNS: usize = 5; + +#[derive(Clone, Copy, Debug)] +enum Case { + /// `Session::abort` after accepting, keeping the session handle. + Abort, + /// `Session::abort` after accepting, then dropping the session handle at once. + AbortThenDrop, + /// `Request::reject` during the handshake. + Reject, +} + +struct Server { + quic: u16, + websocket: u16, +} + +async fn serve(case: Case) -> Server { + let mut config = moq_tokio::server::Config::default(); + config.listen.bind = Some("127.0.0.1:0".parse().unwrap()); + config.listen.tls.generate = vec!["localhost".into()]; + config.websocket = Some( + moq_tokio::websocket::Listener::bind("127.0.0.1:0".parse().unwrap()) + .await + .expect("failed to bind websocket"), + ); + let mut listener = config + .init() + .expect("failed to init server") + .listen() + .await + .expect("failed to listen"); + let server = Server { + quic: listener.local_addr().expect("no quic addr").port(), + websocket: listener.websocket_local_addr().expect("no websocket addr").port(), + }; + + tokio::spawn(async move { + while let Some(request) = listener.accept().await { + tokio::spawn(async move { + match case { + Case::Reject => { + request.reject(moq_tokio::server::Reject::App(CODE)).await.ok(); + } + Case::Abort | Case::AbortThenDrop => { + let Ok(session) = request.ok().await else { return }; + session.abort(moq_net::Error::App(CODE)); + if let Case::Abort = case { + tokio::time::sleep(Duration::from_secs(5)).await; + } + drop(session); + } + } + }); + } + }); + + server +} + +/// The client's terminal error, whether the connect or the established session failed. +async fn client_error(url: &str, websocket: bool) -> moq_tokio::Error { + let mut config = moq_tokio::connect::Config::default(); + config.tls.insecure = Some(true); + config.websocket.enabled = Some(websocket); + let client = config.init(Default::default()).expect("failed to init client"); + + let connection = client.with_reconnect(false).connect(url.parse::().unwrap()); + tokio::time::timeout(Duration::from_secs(10), async { + match connection.established().await { + Ok(connection) => connection.closed().await.expect_err("closed without an error"), + Err(err) => err, + } + }) + .await + .expect("the connection did not close") +} + +async fn check(case: Case) { + let server = serve(case).await; + let urls = [ + (format!("https://localhost:{}/", server.quic), false), + (format!("moqt://localhost:{}/", server.quic), false), + (format!("ws://127.0.0.1:{}/", server.websocket), true), + ]; + + let mut failures = Vec::new(); + for (url, websocket) in &urls { + for _ in 0..RUNS { + let err = client_error(url, *websocket).await; + if !matches!( + err, + moq_tokio::Error::MoqNet(moq_net::Error::Session(moq_net::SessionError::App(CODE))) + ) { + failures.push(format!("{url}: {err:?}")); + } + } + } + assert!( + failures.is_empty(), + "{case:?} lost the close code:\n{}", + failures.join("\n") + ); +} + +#[tokio::test] +async fn abort() { + check(Case::Abort).await; +} + +#[tokio::test] +async fn abort_then_drop() { + check(Case::AbortThenDrop).await; +} + +#[tokio::test] +async fn reject() { + check(Case::Reject).await; +} From c12415dc464f38a94b98004d103536b66ef756ec Mon Sep 17 00:00:00 2001 From: Luke Curley Date: Sat, 26 Sep 2026 11:55:58 -0700 Subject: [PATCH 3/5] feat(net)!: move to web-transport-trait 0.5 SendStream::set_priority takes an i32 send order in trait 0.5. moq-net keeps its u8 send orders and widens them at the trait boundary. Pins qmux 0.6, web-transport-iroh 0.8, and web-transport-wasm 0.7, which implement the new trait. Co-Authored-By: Claude Opus 5.5 --- Cargo.toml | 8 ++++---- rs/moq-e2ee/tests/support/mock.rs | 2 +- rs/moq-net/src/client.rs | 4 ++-- rs/moq-net/src/coding/writer.rs | 4 ++-- rs/moq-net/src/ietf/adapter.rs | 4 ++-- rs/moq-net/src/ietf/publisher.rs | 4 ++-- rs/moq-net/src/lite/test_transport.rs | 3 ++- rs/moq-net/src/server.rs | 2 +- rs/moq-net/tests/support/mock.rs | 2 +- rs/moq-tokio/src/transport.rs | 8 ++++---- rs/moq-uring/src/quic/noq/stream.rs | 9 ++------- rs/moq-uring/src/quic/web.rs | 2 +- 12 files changed, 24 insertions(+), 28 deletions(-) diff --git a/Cargo.toml b/Cargo.toml index 5f9dcd4431..98abc69152 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -203,7 +203,7 @@ object_store = { version = "0.14", default-features = false } percent-encoding = "2" pollster = "1.0" pulldown-cmark = { version = "0.13", default-features = false } -qmux = { version = "0.5.1", default-features = false } +qmux = { version = "0.6", default-features = false } rand = "0.10.1" # default-features off so each consumer picks its own TLS backend and body # features; nothing here wants the default native-tls stack. @@ -233,14 +233,14 @@ url = "2" usage = { package = "usage-rs", version = "6.11.1" } uuid = { version = "1.26", features = ["v7"] } web-async = { version = "0.1.5", features = ["tracing"] } -web-transport-iroh = "0.7" +web-transport-iroh = "0.8" # default-features off so the QUIC crypto provider is chosen by moq-tokio's # aws-lc-rs / ring features. # The WebTransport adapter released with moq-noq, from the same repository. web-transport-moq = { version = "1.3.1", default-features = false } web-transport-proto = "0.6" -web-transport-trait = "0.4" -web-transport-wasm = "0.6" +web-transport-trait = "0.5" +web-transport-wasm = "0.7" x11rb = { version = "0.14.0", features = ["randr", "xfixes"] } zeroize = { version = "1", features = ["derive"] } diff --git a/rs/moq-e2ee/tests/support/mock.rs b/rs/moq-e2ee/tests/support/mock.rs index 35caf23c7a..7757ceb7a6 100644 --- a/rs/moq-e2ee/tests/support/mock.rs +++ b/rs/moq-e2ee/tests/support/mock.rs @@ -118,7 +118,7 @@ impl poll::SendStream for MockSendStream { } } - fn set_priority(&mut self, _order: u8) {} + fn set_priority(&mut self, _order: i32) {} fn finish(&mut self) -> Result<(), Self::Error> { if let Some(tx) = self.tx.take() { diff --git a/rs/moq-net/src/client.rs b/rs/moq-net/src/client.rs index e7292ed723..cc03dc73b1 100644 --- a/rs/moq-net/src/client.rs +++ b/rs/moq-net/src/client.rs @@ -604,7 +604,7 @@ mod tests { Poll::Ready(Ok(buf.len())) } - fn set_priority(&mut self, _order: u8) {} + fn set_priority(&mut self, _order: i32) {} fn finish(&mut self) -> Result<(), Self::Error> { Ok(()) @@ -955,7 +955,7 @@ mod tests { self.inner.poll_write(cx, buf) } - fn set_priority(&mut self, order: u8) { + fn set_priority(&mut self, order: i32) { self.inner.set_priority(order); } diff --git a/rs/moq-net/src/coding/writer.rs b/rs/moq-net/src/coding/writer.rs index 21dd32863b..33ec964ae3 100644 --- a/rs/moq-net/src/coding/writer.rs +++ b/rs/moq-net/src/coding/writer.rs @@ -174,7 +174,7 @@ impl Writer { /// where higher values preempt lower ones. The lite priority queue's rank is the /// opposite (0 = most urgent); rank holders convert via `PriorityHandle::send_order`. pub fn set_priority(&mut self, send_order: u8) { - self.stream.as_mut().unwrap().set_priority(send_order); + self.stream.as_mut().unwrap().set_priority(send_order.into()); } /// Cast the writer to a different version, used during version negotiation. @@ -235,7 +235,7 @@ mod tests { Poll::Ready(Err(self.clone())) } - fn set_priority(&mut self, _: u8) {} + fn set_priority(&mut self, _: i32) {} fn finish(&mut self) -> Result<(), Self::Error> { Err(self.clone()) diff --git a/rs/moq-net/src/ietf/adapter.rs b/rs/moq-net/src/ietf/adapter.rs index e2ab2f2756..2ec6597457 100644 --- a/rs/moq-net/src/ietf/adapter.rs +++ b/rs/moq-net/src/ietf/adapter.rs @@ -376,7 +376,7 @@ impl web_transport_trait::poll::SendStream for VirtualSendStream { Poll::Ready(self.push(chunk).map(|()| len)) } - fn set_priority(&mut self, _order: u8) {} + fn set_priority(&mut self, _order: i32) {} fn finish(&mut self) -> Result<(), Self::Error> { // Flush any remaining buffered data (e.g. if registration never completed). @@ -419,7 +419,7 @@ impl web_transport_trait::poll::SendStream f } } - fn set_priority(&mut self, order: u8) { + fn set_priority(&mut self, order: i32) { match self { Self::Real(s) => s.set_priority(order), Self::Virtual(s) => s.set_priority(order), diff --git a/rs/moq-net/src/ietf/publisher.rs b/rs/moq-net/src/ietf/publisher.rs index 8fa8c7ad11..dcb09ce966 100644 --- a/rs/moq-net/src/ietf/publisher.rs +++ b/rs/moq-net/src/ietf/publisher.rs @@ -2031,7 +2031,7 @@ impl TrackServe { let mut stream = std::future::poll_fn(|cx| self.session.poll_open_uni(cx)) .await .map_err(Error::from_transport)?; - stream.set_priority(priority); + stream.set_priority(priority.into()); let mut writer = Writer::new(stream, self.version); writer.buffer(&ietf::GroupHeader { @@ -2222,7 +2222,7 @@ impl GroupServe { }; self.opened.fetch_add(1, Ordering::Relaxed); let mut stream = stream; - stream.set_priority(self.priority); + stream.set_priority(self.priority.into()); let mut writer = Writer::new(stream, self.version); if let Err(err) = writer.buffer(&self.msg) { diff --git a/rs/moq-net/src/lite/test_transport.rs b/rs/moq-net/src/lite/test_transport.rs index cfac22a22e..0f3ac4a6c4 100644 --- a/rs/moq-net/src/lite/test_transport.rs +++ b/rs/moq-net/src/lite/test_transport.rs @@ -133,7 +133,8 @@ impl poll::SendStream for SinkSend { Poll::Ready(Ok(buf.len())) } - fn set_priority(&mut self, order: u8) { + fn set_priority(&mut self, order: i32) { + let order = u8::try_from(order).expect("moq-net sends u8 send orders"); self.log.priorities.lock().unwrap().push(order); } diff --git a/rs/moq-net/src/server.rs b/rs/moq-net/src/server.rs index e26f5e632d..a716f27bdc 100644 --- a/rs/moq-net/src/server.rs +++ b/rs/moq-net/src/server.rs @@ -877,7 +877,7 @@ mod tests { ) -> std::task::Poll> { std::task::Poll::Ready(Ok(buf.len())) } - fn set_priority(&mut self, _order: u8) {} + fn set_priority(&mut self, _order: i32) {} fn finish(&mut self) -> Result<(), Self::Error> { Ok(()) } diff --git a/rs/moq-net/tests/support/mock.rs b/rs/moq-net/tests/support/mock.rs index ed5ad7fa3c..3bc1cc63ea 100644 --- a/rs/moq-net/tests/support/mock.rs +++ b/rs/moq-net/tests/support/mock.rs @@ -144,7 +144,7 @@ impl poll::SendStream for MockSendStream { ) } - fn set_priority(&mut self, _order: u8) {} + fn set_priority(&mut self, _order: i32) {} fn finish(&mut self) -> Result<(), Self::Error> { if self.tx.is_some() { diff --git a/rs/moq-tokio/src/transport.rs b/rs/moq-tokio/src/transport.rs index ffc046c658..af2b4e0a9f 100644 --- a/rs/moq-tokio/src/transport.rs +++ b/rs/moq-tokio/src/transport.rs @@ -211,7 +211,7 @@ pub struct SendStream { // `Sync` off the containing session types. state: std::sync::Mutex>>, // Deferred actions, applied when the in-flight operation settles. - priority: Option, + priority: Option, finish: bool, reset: Option, /// Whether `finish` or `reset` has been observed, so no further writes can @@ -370,7 +370,7 @@ impl wt_poll::SendStream for SendS } } - fn set_priority(&mut self, order: u8) { + fn set_priority(&mut self, order: i32) { match self.state.get_mut().unwrap().as_mut() { Some(SendState::Idle(stream)) => stream.set_priority(order), _ => self.priority = Some(order), @@ -807,7 +807,7 @@ mod tests { struct FakeSend { writes: Arc>>, blocked: Arc, - priorities: Arc>>, + priorities: Arc>>, finished: Arc, resets: Arc>>, /// The peer never acknowledges the FIN: closed() stays pending forever. @@ -830,7 +830,7 @@ mod tests { Ok(buf.len()) } - fn set_priority(&mut self, order: u8) { + fn set_priority(&mut self, order: i32) { self.priorities.lock().unwrap().push(order); } diff --git a/rs/moq-uring/src/quic/noq/stream.rs b/rs/moq-uring/src/quic/noq/stream.rs index 0c6802fa2c..bb9d90110e 100644 --- a/rs/moq-uring/src/quic/noq/stream.rs +++ b/rs/moq-uring/src/quic/noq/stream.rs @@ -118,15 +118,10 @@ impl web_transport_trait::poll::SendStream for SendStream { } } - fn set_priority(&mut self, order: u8) { + fn set_priority(&mut self, order: i32) { // The trait (like W3C sendOrder) sends HIGHER values first, and so // does noq. - let _ = self - .shared - .conn - .borrow_mut() - .send_stream(self.id) - .set_priority(i32::from(order)); + let _ = self.shared.conn.borrow_mut().send_stream(self.id).set_priority(order); } fn finish(&mut self) -> Result<(), Self::Error> { diff --git a/rs/moq-uring/src/quic/web.rs b/rs/moq-uring/src/quic/web.rs index 1100165c20..957490ab4b 100644 --- a/rs/moq-uring/src/quic/web.rs +++ b/rs/moq-uring/src/quic/web.rs @@ -859,7 +859,7 @@ impl web_transport_trait::poll::SendStream for SendStream { } } - fn set_priority(&mut self, order: u8) { + fn set_priority(&mut self, order: i32) { web_transport_trait::poll::SendStream::set_priority(&mut self.inner, order); } From d43956dfcff0968b4609eac551b7c99427178734 Mon Sep 17 00:00:00 2001 From: Luke Curley Date: Sat, 26 Sep 2026 12:33:53 -0700 Subject: [PATCH 4/5] chore: pin web-transport-moq 1.3.2 for the raw QUIC close code Co-Authored-By: Claude Opus 5.5 --- Cargo.toml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/Cargo.toml b/Cargo.toml index 98abc69152..d170a5a385 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -237,7 +237,7 @@ web-transport-iroh = "0.8" # default-features off so the QUIC crypto provider is chosen by moq-tokio's # aws-lc-rs / ring features. # The WebTransport adapter released with moq-noq, from the same repository. -web-transport-moq = { version = "1.3.1", default-features = false } +web-transport-moq = { version = "1.3.2", default-features = false } web-transport-proto = "0.6" web-transport-trait = "0.5" web-transport-wasm = "0.7" From f97d647c7c5fb128fbf75cfee722e9fdb008fdda Mon Sep 17 00:00:00 2001 From: Luke Curley Date: Mon, 28 Sep 2026 12:13:13 -0700 Subject: [PATCH 5/5] test(tokio): hold the aborted session until the runtime shuts down The Abort case slept 5s before dropping the handle, so a stalled worker could collapse it into AbortThenDrop. Co-Authored-By: Claude Opus 5.5 --- rs/moq-tokio/tests/close_code.rs | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/rs/moq-tokio/tests/close_code.rs b/rs/moq-tokio/tests/close_code.rs index 3dd3346636..9d5f287e97 100644 --- a/rs/moq-tokio/tests/close_code.rs +++ b/rs/moq-tokio/tests/close_code.rs @@ -58,7 +58,8 @@ async fn serve(case: Case) -> Server { let Ok(session) = request.ok().await else { return }; session.abort(moq_net::Error::App(CODE)); if let Case::Abort = case { - tokio::time::sleep(Duration::from_secs(5)).await; + // Hold the handle until the test's runtime shuts down. + std::future::pending::<()>().await; } drop(session); }