From bd660b2afbf08bb1aa3fef9cf93a11f3bdafa1d4 Mon Sep 17 00:00:00 2001 From: Alex Emmet <111742636+Alex-Emmet@users.noreply.github.com> Date: Thu, 20 Aug 2026 20:37:25 +0200 Subject: [PATCH] [Debug] --- host/src/engine.rs | 22 +++++++++ mtp-webserver/src/transport.rs | 1 + transport/src/generic_connection.rs | 6 +++ wasm/src/transport.rs | 70 +++++++++++++++-------------- 4 files changed, 66 insertions(+), 33 deletions(-) diff --git a/host/src/engine.rs b/host/src/engine.rs index 38e8422..d1ebd19 100644 --- a/host/src/engine.rs +++ b/host/src/engine.rs @@ -217,6 +217,12 @@ impl HandshakeEngine { _authentication_context: &AuthenticationContext, ) -> Result { let mut first_msg = receiver.receive().await.map_err(AcceptError::Receive)?; + tracing::debug!( + message_type = ?first_msg.get_type(), + version = ?first_msg.get_str(DataType::Version), + client_id = ?first_msg.get_data(DataType::Id), + "received MTP opening message" + ); let version_str = match first_msg.get_data(DataType::Version) { Some(DataValue::Str(s)) => s.clone(), @@ -271,6 +277,11 @@ impl HandshakeEngine { return Err(AcceptError::UnsupportedVersion(client_version)); } }; + tracing::debug!( + client_version = %client_version, + negotiated_version = %negotiated, + "MTP protocol version negotiated" + ); let codec = VersionedCodec::for_version(self.registry.clone(), negotiated.clone()) .ok_or_else(|| AcceptError::UnsupportedVersion(negotiated.clone()))?; @@ -1105,6 +1116,12 @@ async fn send_rejection_generic( .add_typed_default(DataType::Connected, DataValue::BoolFalse) .add_typed_default(DataType::ErrorMessage, DataValue::Str(reason.to_string())), }; + tracing::debug!( + reason = %reason, + response_type = ?response.get_type(), + has_version = response.get_data(DataType::Version).is_some(), + "sending MTP handshake rejection" + ); let _ = sender.send(&response).await; } @@ -1138,6 +1155,11 @@ async fn send_accepted_generic( if let Some(id) = assigned_id { response = response.add_typed_default(DataType::Id, DataValue::UnsignedNumber(id as u128)); } + tracing::debug!( + version = %version, + assigned_id = ?assigned_id, + "sending accepted MTP handshake response" + ); sender.send(&response).await?; sender.finish_stream().await } diff --git a/mtp-webserver/src/transport.rs b/mtp-webserver/src/transport.rs index 31d8376..9b7de76 100644 --- a/mtp-webserver/src/transport.rs +++ b/mtp-webserver/src/transport.rs @@ -88,6 +88,7 @@ impl TransportRecvStream for H3TransportReceiver { tracing::debug!( remote = %self.quinn.remote_address(), bytes = buf.len(), + header = ?buf, "received first bytes from WebTransport MTP stream" ); } diff --git a/transport/src/generic_connection.rs b/transport/src/generic_connection.rs index 4fa6433..830fb40 100644 --- a/transport/src/generic_connection.rs +++ b/transport/src/generic_connection.rs @@ -423,6 +423,12 @@ impl GenericReceiver { break; } }; + tracing::debug!( + frames, + frame_len, + message_type = ?message.get_type(), + "decoded MTP receive frame" + ); let negotiated_type_map = type_map.read().await.clone(); message.set_type_map(&negotiated_type_map); diff --git a/wasm/src/transport.rs b/wasm/src/transport.rs index a4deea6..dc0091e 100644 --- a/wasm/src/transport.rs +++ b/wasm/src/transport.rs @@ -128,8 +128,6 @@ pub struct WasmTransport { buffer: Rc>>, /// Set to `true` when `open_next_stream` succeeds; cleared after the first frame is parsed. new_stream_frame: Rc>, - /// A single ordered browser send stream shared by all cloned transports. - outgoing_writer: Rc>>, /// Serializes stream creation and writes across concurrent callers. send_lock: Rc>, type_map: Rc>, @@ -209,7 +207,6 @@ impl WasmTransport { stream_reader: Rc::new(RefCell::new(None)), buffer: Rc::new(RefCell::new(Vec::new())), new_stream_frame: Rc::new(Cell::new(false)), - outgoing_writer: Rc::new(RefCell::new(None)), send_lock: Rc::new(AsyncMutex::new(())), type_map: Rc::new(RefCell::new(TypeMap::latest())), decode_limits: Rc::new(RefCell::new(decode_limits)), @@ -247,30 +244,30 @@ impl WasmTransport { return Err(js_error("message too large")); } - let writer_val = if let Some(writer) = self.outgoing_writer.borrow().clone() { - writer - } else { - let create_stream = js_sys::Reflect::get( - &self.inner, - &JsValue::from_str("createUnidirectionalStream"), - )? + // Use one WebTransport uni-stream per MTP frame. Chromium reliably + // publishes a browser-created uni-stream to the peer when it is + // closed; leaving a shared stream open can leave the server waiting + // in accept_uni() until the authentication deadline. The bytes are + // already the canonical MTP self-framed value, so no extra stream + // length prefix is added here. + let create_stream = js_sys::Reflect::get( + &self.inner, + &JsValue::from_str("createUnidirectionalStream"), + )? + .dyn_into::() + .map_err(|_| js_error("createUnidirectionalStream not a function"))?; + let stream_promise = create_stream + .call0(&self.inner)? + .dyn_into::() + .map_err(|_| js_error("createUnidirectionalStream did not return a Promise"))?; + let stream = JsFuture::from(stream_promise).await?; + let writable_or_stream = resolve_stream_writable(&stream)?; + let writer_val = js_sys::Reflect::get(&writable_or_stream, &JsValue::from_str("getWriter")) + .map_err(|_| js_error("missing getWriter"))? .dyn_into::() - .map_err(|_| js_error("createUnidirectionalStream not a function"))?; - let stream_promise = create_stream - .call0(&self.inner)? - .dyn_into::() - .map_err(|_| js_error("createUnidirectionalStream did not return a Promise"))?; - let stream = JsFuture::from(stream_promise).await?; - let writable_or_stream = resolve_stream_writable(&stream)?; - let writer = js_sys::Reflect::get(&writable_or_stream, &JsValue::from_str("getWriter")) - .map_err(|_| js_error("missing getWriter"))? - .dyn_into::() - .map_err(|_| js_error("getWriter not a function"))? - .call0(&writable_or_stream) - .map_err(|_| js_error("getWriter call failed"))?; - *self.outgoing_writer.borrow_mut() = Some(writer.clone()); - writer - }; + .map_err(|_| js_error("getWriter not a function"))? + .call0(&writable_or_stream) + .map_err(|_| js_error("getWriter call failed"))?; let chunk = js_sys::Uint8Array::from(frame); @@ -283,11 +280,24 @@ impl WasmTransport { .map_err(|e| js_error(format!("write failed: {:?}", e)))?; if let Err(e) = JsFuture::from(write_promise.unchecked_into::()).await { log_stream_error_code(&e, "send_frame write"); - self.outgoing_writer.borrow_mut().take(); release_writer_lock(&writer_val); return Err(e); } + let close_fn = js_sys::Reflect::get(&writer_val, &JsValue::from_str("close")) + .map_err(|_| js_error("missing close"))? + .dyn_into::() + .map_err(|_| js_error("close not a function"))?; + let close_promise = close_fn + .call0(&writer_val) + .map_err(|e| js_error(format!("close failed: {:?}", e)))?; + if let Err(e) = JsFuture::from(close_promise.unchecked_into::()).await { + // The frame was already written; do not retry it merely because + // FIN failed, as that would duplicate the MTP frame. + log_stream_error_code(&e, "send_frame close"); + } + release_writer_lock(&writer_val); + Ok(()) } @@ -677,12 +687,6 @@ impl WasmTransport { } pub fn close(&self) { - if let Some(writer) = self.outgoing_writer.borrow_mut().take() { - // The WebTransport session close below terminates the stream. The - // lock must be released first so dropping it is not interpreted as - // an application abort. - release_writer_lock(&writer); - } // Release reader locks before closing so they aren't treated as cancels. if let Some(reader) = self.stream_reader.borrow_mut().take() { release_reader_lock(&reader);