diff --git a/.forgejo/workflows/ci.yml b/.forgejo/workflows/ci.yml index 2d870bc..70e6c39 100644 --- a/.forgejo/workflows/ci.yml +++ b/.forgejo/workflows/ci.yml @@ -7,12 +7,16 @@ on: env: CARGO_TERM_COLOR: always + NIX_CONFIG: experimental-features = nix-command flakes jobs: checks: name: checks runs-on: nixos steps: + - name: Install node + run: nix profile add nixpkgs#nodejs_24 + - name: Checkout uses: https://data.forgejo.org/actions/checkout@v7 diff --git a/.forgejo/workflows/release.yml b/.forgejo/workflows/release.yml index 5996da0..05bc6e2 100644 --- a/.forgejo/workflows/release.yml +++ b/.forgejo/workflows/release.yml @@ -14,10 +14,16 @@ on: required: true type: string +env: + NIX_CONFIG: experimental-features = nix-command flakes + jobs: release: runs-on: nixos steps: + - name: Install node & bun + run: nix profile add nixpkgs#nodejs_24 nixpkgs#bun + - name: Check out repo uses: https://data.forgejo.org/actions/checkout@v7 with: @@ -26,6 +32,9 @@ jobs: - name: Install dependencies run: bun install + - name: Install cc linker, sed & jq + run: nix profile add nixpkgs#stdenv.cc nixpkgs#gnused nixpkgs#jq + - name: Build all run: bun build:all diff --git a/host/src/engine.rs b/host/src/engine.rs index d1ebd19..38e8422 100644 --- a/host/src/engine.rs +++ b/host/src/engine.rs @@ -217,12 +217,6 @@ 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(), @@ -277,11 +271,6 @@ 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()))?; @@ -1116,12 +1105,6 @@ 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; } @@ -1155,11 +1138,6 @@ 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/h3.rs b/mtp-webserver/src/h3.rs index abb17aa..c09379c 100644 --- a/mtp-webserver/src/h3.rs +++ b/mtp-webserver/src/h3.rs @@ -143,11 +143,6 @@ pub(crate) async fn run_driver( return; } }; - tracing::debug!( - remote = %remote_addr, - session_id = ?session.session_id(), - "accepted WebTransport MTP session" - ); tokio::spawn(run_session_requests( session.clone(), router.clone(), diff --git a/mtp-webserver/src/transport.rs b/mtp-webserver/src/transport.rs index 9b7de76..a484f26 100644 --- a/mtp-webserver/src/transport.rs +++ b/mtp-webserver/src/transport.rs @@ -32,8 +32,6 @@ pub struct H3TransportSender { pub struct H3TransportReceiver { stream: H3RecvStream, - quinn: quinn::Connection, - read_exact_calls: u64, } impl H3TransportConnection { @@ -78,43 +76,18 @@ impl TransportSendStream for H3TransportSender { #[async_trait::async_trait] impl TransportRecvStream for H3TransportReceiver { async fn read_exact(&mut self, buf: &mut [u8]) -> Result<(), CommunicationError> { - let first_read = self.read_exact_calls == 0; - self.read_exact_calls += 1; self.stream .read_exact(buf) .await - .map(|_| { - if first_read { - tracing::debug!( - remote = %self.quinn.remote_address(), - bytes = buf.len(), - header = ?buf, - "received first bytes from WebTransport MTP stream" - ); - } - }) + .map(|_| ()) .map_err(|error| { - if error.kind() == std::io::ErrorKind::UnexpectedEof - || self.quinn.close_reason().is_some() - { - /* - * Reaching FIN, or losing the enclosing QUIC connection, - * is a normal stream-closure path. Do not turn it into a - * frame-header failure and close the connection again. - */ + if error.kind() == std::io::ErrorKind::UnexpectedEof { + // Browser control frames are sent on one-frame uni streams. + // Reaching FIN while looking for another frame is normal. return CommunicationError::StreamClosed; } - error!( - "[mtp-webserver] receive stream read_exact failed ({} bytes): {error}", - buf.len() - ); - tracing::warn!( - remote = %self.quinn.remote_address(), - first_read, - len = buf.len(), - %error, - "WebTransport receive stream read_exact failed" - ); + error!("[mtp-webserver] receive stream read_exact failed ({} bytes): {error}", buf.len()); + tracing::warn!(len = buf.len(), %error, "WebTransport receive stream read_exact failed"); CommunicationError::StreamError }) } @@ -128,9 +101,6 @@ impl TransportRecvStream for H3TransportReceiver { Ok(Some(buf)) } Err(error) => { - if self.quinn.close_reason().is_some() { - return Err(CommunicationError::StreamClosed); - } error!( "[mtp-webserver] receive stream read failed (max {} bytes): {error}", max @@ -197,27 +167,10 @@ impl TransportConnection for H3TransportConnection { loop { match self.session.accept_uni().await { Ok(Some((id, stream))) if id == self.session.session_id() => { - let stream_id = h3::quic::RecvStream::recv_id(&stream); - tracing::debug!( - remote = %self.quinn.remote_address(), - session_id = ?self.session.session_id(), - stream_id = ?stream_id, - "accepted WebTransport MTP receive stream" - ); - return Ok(H3TransportReceiver { - stream, - quinn: self.quinn.clone(), - read_exact_calls: 0, - }); + return Ok(H3TransportReceiver { stream }); } - Ok(Some((stream_session_id, _stream))) => { + Ok(Some(_)) => { consecutive_errors = 0; - tracing::debug!( - remote = %self.quinn.remote_address(), - session_id = ?self.session.session_id(), - stream_session_id = ?stream_session_id, - "ignored WebTransport receive stream belonging to another session" - ); continue; } Ok(None) => return Err(CommunicationError::StreamClosed), @@ -353,18 +306,7 @@ async fn accept_web_connection_inner( connection_id, }, ) - .await; - #[cfg(feature = "crypto")] - if let Err(error) = &result { - tracing::warn!( - remote = %remote_addr, - connection_id, - %error, - "WebTransport MTP handshake failed" - ); - } - #[cfg(feature = "crypto")] - let result = result?; + .await?; #[cfg(not(feature = "crypto"))] let result = engine.accept(&sender, &receiver).await?; diff --git a/transport/src/generic_connection.rs b/transport/src/generic_connection.rs index 830fb40..23f3fe9 100644 --- a/transport/src/generic_connection.rs +++ b/transport/src/generic_connection.rs @@ -333,18 +333,6 @@ impl GenericReceiver { break 'stream; } Err(_) => { - if frames == 0 { - tracing::warn!( - timeout = ?policy.read_timeout, - "MTP receive stream timed out before its first complete frame" - ); - } else { - tracing::debug!( - frames, - timeout = ?policy.read_timeout, - "MTP receive stream idle timeout" - ); - } break; } } @@ -388,9 +376,6 @@ impl GenericReceiver { ) .await; if !matches!(&body_read, Ok(Ok(()))) { - if matches!(&body_read, Ok(Err(CommunicationError::StreamClosed))) { - break 'stream; - } tracing::warn!( pipe_chunk_len = chunk_len, ?body_read, @@ -423,12 +408,6 @@ 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/client/authentication.rs b/wasm/src/client/authentication.rs index df7989e..45cec8d 100644 --- a/wasm/src/client/authentication.rs +++ b/wasm/src/client/authentication.rs @@ -8,14 +8,6 @@ use crate::config::ConnectionConfig; use crate::error::js_error; use crate::transport::WasmTransport; -fn server_rejection_message(outcome: &CommunicationValue) -> Option<&str> { - (outcome.get_data(DataType::Connected) == Some(&DataValue::BoolFalse)).then(|| { - outcome - .get_str(DataType::ErrorMessage) - .unwrap_or("host rejected the connection") - }) -} - #[wasm_bindgen] #[allow(deprecated)] impl WasmClient { @@ -87,18 +79,6 @@ impl WasmClient { .unwrap_or("host does not support this protocol version"), )); } - - // Generic host rejections are IdentificationResponse frames with - // Connected=false. They intentionally do not carry a negotiated - // Version because negotiation never completed. Check this before - // reading Version, otherwise a useful server error such as an - // authentication timeout is reported as the misleading - // "host omitted a valid negotiated protocol version". - if let Some(message) = server_rejection_message(&outcome) { - self.set_state_if_current(generation, ConnectionState::Disconnected); - return Err(js_error(message)); - } - let negotiated_version = match outcome.get_data(DataType::Version) { Some(DataValue::Str(version)) => mtp_codec::Version::parse(version) .ok_or_else(|| js_error("host omitted a valid negotiated protocol version"))?, @@ -648,24 +628,3 @@ impl WasmClient { Ok(server_challenge) } } - -#[cfg(test)] -mod tests { - use super::server_rejection_message; - use mtp_codec::{CommunicationType, CommunicationValue, DataType, DataValue}; - - #[test] - fn reports_rejection_reason_without_a_negotiated_version() { - let response = CommunicationValue::new(CommunicationType::IdentificationResponse) - .add_typed_default(DataType::Connected, DataValue::BoolFalse) - .add_typed_default( - DataType::ErrorMessage, - DataValue::Str("authentication handshake timed out".into()), - ); - - assert_eq!( - server_rejection_message(&response), - Some("authentication handshake timed out") - ); - } -} diff --git a/wasm/src/transport.rs b/wasm/src/transport.rs index dc0091e..a4deea6 100644 --- a/wasm/src/transport.rs +++ b/wasm/src/transport.rs @@ -128,6 +128,8 @@ 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>, @@ -207,6 +209,7 @@ 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)), @@ -244,30 +247,30 @@ impl WasmTransport { return Err(js_error("message too large")); } - // 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"))? + 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"), + )? .dyn_into::() - .map_err(|_| js_error("getWriter not a function"))? - .call0(&writable_or_stream) - .map_err(|_| js_error("getWriter call failed"))?; + .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 + }; let chunk = js_sys::Uint8Array::from(frame); @@ -280,24 +283,11 @@ 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(()) } @@ -687,6 +677,12 @@ 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);