Compare commits

..
Author SHA1 Message Date
62beb80fa1 Update Rust crate aes-gcm to 0.11
Some checks failed
renovate/stability-days Updates have met minimum release age requirement
CI / checks (pull_request) Failing after 8s
2026-08-20 22:01:10 +03:00
687d56f9f1
Merge remote-tracking branch 'refs/remotes/origin/master'
Some checks failed
CI / checks (push) Failing after 2s
2026-08-20 20:55:52 +02:00
4c10b56a6c
Update workflows 2026-08-20 20:55:43 +02:00
Alex Emmet
bd660b2afb
[Debug]
Some checks failed
CI / checks (push) Failing after 2m17s
2026-08-20 20:37:25 +02:00
Alex Emmet
2b0bdc3257
[Fix] Connections
Some checks failed
CI / checks (push) Failing after 3m14s
2026-08-20 20:11:08 +02:00
8 changed files with 193 additions and 55 deletions

View file

@ -7,16 +7,12 @@ 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

View file

@ -14,16 +14,10 @@ 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:
@ -32,9 +26,6 @@ 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

View file

@ -217,6 +217,12 @@ impl HandshakeEngine {
_authentication_context: &AuthenticationContext,
) -> Result<HandshakeResult, AcceptError> {
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<S: HandshakeSender>(
.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<S: HandshakeSender>(
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
}

View file

@ -143,6 +143,11 @@ 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(),

View file

@ -32,6 +32,8 @@ pub struct H3TransportSender {
pub struct H3TransportReceiver {
stream: H3RecvStream,
quinn: quinn::Connection,
read_exact_calls: u64,
}
impl H3TransportConnection {
@ -76,18 +78,43 @@ 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(|_| ())
.map(|_| {
if first_read {
tracing::debug!(
remote = %self.quinn.remote_address(),
bytes = buf.len(),
header = ?buf,
"received first bytes from WebTransport MTP stream"
);
}
})
.map_err(|error| {
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.
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.
*/
return CommunicationError::StreamClosed;
}
error!("[mtp-webserver] receive stream read_exact failed ({} bytes): {error}", buf.len());
tracing::warn!(len = buf.len(), %error, "WebTransport receive stream read_exact failed");
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"
);
CommunicationError::StreamError
})
}
@ -101,6 +128,9 @@ 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
@ -167,10 +197,27 @@ impl TransportConnection for H3TransportConnection {
loop {
match self.session.accept_uni().await {
Ok(Some((id, stream))) if id == self.session.session_id() => {
return Ok(H3TransportReceiver { stream });
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,
});
}
Ok(Some(_)) => {
Ok(Some((stream_session_id, _stream))) => {
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),
@ -306,7 +353,18 @@ async fn accept_web_connection_inner(
connection_id,
},
)
.await?;
.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?;
#[cfg(not(feature = "crypto"))]
let result = engine.accept(&sender, &receiver).await?;

View file

@ -333,6 +333,18 @@ impl<C: TransportConnection> GenericReceiver<C> {
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;
}
}
@ -376,6 +388,9 @@ impl<C: TransportConnection> GenericReceiver<C> {
)
.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,
@ -408,6 +423,12 @@ impl<C: TransportConnection> GenericReceiver<C> {
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);

View file

@ -8,6 +8,14 @@ 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 {
@ -79,6 +87,18 @@ 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"))?,
@ -628,3 +648,24 @@ 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")
);
}
}

View file

@ -128,8 +128,6 @@ pub struct WasmTransport {
buffer: Rc<RefCell<Vec<u8>>>,
/// Set to `true` when `open_next_stream` succeeds; cleared after the first frame is parsed.
new_stream_frame: Rc<Cell<bool>>,
/// A single ordered browser send stream shared by all cloned transports.
outgoing_writer: Rc<RefCell<Option<JsValue>>>,
/// Serializes stream creation and writes across concurrent callers.
send_lock: Rc<AsyncMutex<()>>,
type_map: Rc<RefCell<TypeMap>>,
@ -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::<js_sys::Function>()
.map_err(|_| js_error("createUnidirectionalStream not a function"))?;
let stream_promise = create_stream
.call0(&self.inner)?
.dyn_into::<js_sys::Promise>()
.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::<js_sys::Function>()
.map_err(|_| js_error("createUnidirectionalStream not a function"))?;
let stream_promise = create_stream
.call0(&self.inner)?
.dyn_into::<js_sys::Promise>()
.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::<js_sys::Function>()
.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::<js_sys::Promise>()).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::<js_sys::Function>()
.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::<js_sys::Promise>()).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);