Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 516e9639c5 |
8 changed files with 55 additions and 193 deletions
|
|
@ -7,12 +7,16 @@ on:
|
||||||
|
|
||||||
env:
|
env:
|
||||||
CARGO_TERM_COLOR: always
|
CARGO_TERM_COLOR: always
|
||||||
|
NIX_CONFIG: experimental-features = nix-command flakes
|
||||||
|
|
||||||
jobs:
|
jobs:
|
||||||
checks:
|
checks:
|
||||||
name: checks
|
name: checks
|
||||||
runs-on: nixos
|
runs-on: nixos
|
||||||
steps:
|
steps:
|
||||||
|
- name: Install node
|
||||||
|
run: nix profile add nixpkgs#nodejs_24
|
||||||
|
|
||||||
- name: Checkout
|
- name: Checkout
|
||||||
uses: https://data.forgejo.org/actions/checkout@v7
|
uses: https://data.forgejo.org/actions/checkout@v7
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -14,10 +14,16 @@ on:
|
||||||
required: true
|
required: true
|
||||||
type: string
|
type: string
|
||||||
|
|
||||||
|
env:
|
||||||
|
NIX_CONFIG: experimental-features = nix-command flakes
|
||||||
|
|
||||||
jobs:
|
jobs:
|
||||||
release:
|
release:
|
||||||
runs-on: nixos
|
runs-on: nixos
|
||||||
steps:
|
steps:
|
||||||
|
- name: Install node & bun
|
||||||
|
run: nix profile add nixpkgs#nodejs_24 nixpkgs#bun
|
||||||
|
|
||||||
- name: Check out repo
|
- name: Check out repo
|
||||||
uses: https://data.forgejo.org/actions/checkout@v7
|
uses: https://data.forgejo.org/actions/checkout@v7
|
||||||
with:
|
with:
|
||||||
|
|
@ -26,6 +32,9 @@ jobs:
|
||||||
- name: Install dependencies
|
- name: Install dependencies
|
||||||
run: bun install
|
run: bun install
|
||||||
|
|
||||||
|
- name: Install cc linker, sed & jq
|
||||||
|
run: nix profile add nixpkgs#stdenv.cc nixpkgs#gnused nixpkgs#jq
|
||||||
|
|
||||||
- name: Build all
|
- name: Build all
|
||||||
run: bun build:all
|
run: bun build:all
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -217,12 +217,6 @@ impl HandshakeEngine {
|
||||||
_authentication_context: &AuthenticationContext,
|
_authentication_context: &AuthenticationContext,
|
||||||
) -> Result<HandshakeResult, AcceptError> {
|
) -> Result<HandshakeResult, AcceptError> {
|
||||||
let mut first_msg = receiver.receive().await.map_err(AcceptError::Receive)?;
|
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) {
|
let version_str = match first_msg.get_data(DataType::Version) {
|
||||||
Some(DataValue::Str(s)) => s.clone(),
|
Some(DataValue::Str(s)) => s.clone(),
|
||||||
|
|
@ -277,11 +271,6 @@ impl HandshakeEngine {
|
||||||
return Err(AcceptError::UnsupportedVersion(client_version));
|
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())
|
let codec = VersionedCodec::for_version(self.registry.clone(), negotiated.clone())
|
||||||
.ok_or_else(|| AcceptError::UnsupportedVersion(negotiated.clone()))?;
|
.ok_or_else(|| AcceptError::UnsupportedVersion(negotiated.clone()))?;
|
||||||
|
|
@ -1116,12 +1105,6 @@ async fn send_rejection_generic<S: HandshakeSender>(
|
||||||
.add_typed_default(DataType::Connected, DataValue::BoolFalse)
|
.add_typed_default(DataType::Connected, DataValue::BoolFalse)
|
||||||
.add_typed_default(DataType::ErrorMessage, DataValue::Str(reason.to_string())),
|
.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;
|
let _ = sender.send(&response).await;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -1155,11 +1138,6 @@ async fn send_accepted_generic<S: HandshakeSender>(
|
||||||
if let Some(id) = assigned_id {
|
if let Some(id) = assigned_id {
|
||||||
response = response.add_typed_default(DataType::Id, DataValue::UnsignedNumber(id as u128));
|
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.send(&response).await?;
|
||||||
sender.finish_stream().await
|
sender.finish_stream().await
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -143,11 +143,6 @@ pub(crate) async fn run_driver(
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
};
|
};
|
||||||
tracing::debug!(
|
|
||||||
remote = %remote_addr,
|
|
||||||
session_id = ?session.session_id(),
|
|
||||||
"accepted WebTransport MTP session"
|
|
||||||
);
|
|
||||||
tokio::spawn(run_session_requests(
|
tokio::spawn(run_session_requests(
|
||||||
session.clone(),
|
session.clone(),
|
||||||
router.clone(),
|
router.clone(),
|
||||||
|
|
|
||||||
|
|
@ -32,8 +32,6 @@ pub struct H3TransportSender {
|
||||||
|
|
||||||
pub struct H3TransportReceiver {
|
pub struct H3TransportReceiver {
|
||||||
stream: H3RecvStream,
|
stream: H3RecvStream,
|
||||||
quinn: quinn::Connection,
|
|
||||||
read_exact_calls: u64,
|
|
||||||
}
|
}
|
||||||
|
|
||||||
impl H3TransportConnection {
|
impl H3TransportConnection {
|
||||||
|
|
@ -78,43 +76,18 @@ impl TransportSendStream for H3TransportSender {
|
||||||
#[async_trait::async_trait]
|
#[async_trait::async_trait]
|
||||||
impl TransportRecvStream for H3TransportReceiver {
|
impl TransportRecvStream for H3TransportReceiver {
|
||||||
async fn read_exact(&mut self, buf: &mut [u8]) -> Result<(), CommunicationError> {
|
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
|
self.stream
|
||||||
.read_exact(buf)
|
.read_exact(buf)
|
||||||
.await
|
.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| {
|
.map_err(|error| {
|
||||||
if error.kind() == std::io::ErrorKind::UnexpectedEof
|
if error.kind() == std::io::ErrorKind::UnexpectedEof {
|
||||||
|| self.quinn.close_reason().is_some()
|
// Browser control frames are sent on one-frame uni streams.
|
||||||
{
|
// Reaching FIN while looking for another frame is normal.
|
||||||
/*
|
|
||||||
* 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;
|
return CommunicationError::StreamClosed;
|
||||||
}
|
}
|
||||||
error!(
|
error!("[mtp-webserver] receive stream read_exact failed ({} bytes): {error}", buf.len());
|
||||||
"[mtp-webserver] receive stream read_exact failed ({} bytes): {error}",
|
tracing::warn!(len = buf.len(), %error, "WebTransport receive stream read_exact failed");
|
||||||
buf.len()
|
|
||||||
);
|
|
||||||
tracing::warn!(
|
|
||||||
remote = %self.quinn.remote_address(),
|
|
||||||
first_read,
|
|
||||||
len = buf.len(),
|
|
||||||
%error,
|
|
||||||
"WebTransport receive stream read_exact failed"
|
|
||||||
);
|
|
||||||
CommunicationError::StreamError
|
CommunicationError::StreamError
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
@ -128,9 +101,6 @@ impl TransportRecvStream for H3TransportReceiver {
|
||||||
Ok(Some(buf))
|
Ok(Some(buf))
|
||||||
}
|
}
|
||||||
Err(error) => {
|
Err(error) => {
|
||||||
if self.quinn.close_reason().is_some() {
|
|
||||||
return Err(CommunicationError::StreamClosed);
|
|
||||||
}
|
|
||||||
error!(
|
error!(
|
||||||
"[mtp-webserver] receive stream read failed (max {} bytes): {error}",
|
"[mtp-webserver] receive stream read failed (max {} bytes): {error}",
|
||||||
max
|
max
|
||||||
|
|
@ -197,27 +167,10 @@ impl TransportConnection for H3TransportConnection {
|
||||||
loop {
|
loop {
|
||||||
match self.session.accept_uni().await {
|
match self.session.accept_uni().await {
|
||||||
Ok(Some((id, stream))) if id == self.session.session_id() => {
|
Ok(Some((id, stream))) if id == self.session.session_id() => {
|
||||||
let stream_id = h3::quic::RecvStream::recv_id(&stream);
|
return Ok(H3TransportReceiver { 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((stream_session_id, _stream))) => {
|
Ok(Some(_)) => {
|
||||||
consecutive_errors = 0;
|
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;
|
continue;
|
||||||
}
|
}
|
||||||
Ok(None) => return Err(CommunicationError::StreamClosed),
|
Ok(None) => return Err(CommunicationError::StreamClosed),
|
||||||
|
|
@ -353,18 +306,7 @@ async fn accept_web_connection_inner(
|
||||||
connection_id,
|
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"))]
|
#[cfg(not(feature = "crypto"))]
|
||||||
let result = engine.accept(&sender, &receiver).await?;
|
let result = engine.accept(&sender, &receiver).await?;
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -333,18 +333,6 @@ impl<C: TransportConnection> GenericReceiver<C> {
|
||||||
break 'stream;
|
break 'stream;
|
||||||
}
|
}
|
||||||
Err(_) => {
|
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;
|
break;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
@ -388,9 +376,6 @@ impl<C: TransportConnection> GenericReceiver<C> {
|
||||||
)
|
)
|
||||||
.await;
|
.await;
|
||||||
if !matches!(&body_read, Ok(Ok(()))) {
|
if !matches!(&body_read, Ok(Ok(()))) {
|
||||||
if matches!(&body_read, Ok(Err(CommunicationError::StreamClosed))) {
|
|
||||||
break 'stream;
|
|
||||||
}
|
|
||||||
tracing::warn!(
|
tracing::warn!(
|
||||||
pipe_chunk_len = chunk_len,
|
pipe_chunk_len = chunk_len,
|
||||||
?body_read,
|
?body_read,
|
||||||
|
|
@ -423,12 +408,6 @@ impl<C: TransportConnection> GenericReceiver<C> {
|
||||||
break;
|
break;
|
||||||
}
|
}
|
||||||
};
|
};
|
||||||
tracing::debug!(
|
|
||||||
frames,
|
|
||||||
frame_len,
|
|
||||||
message_type = ?message.get_type(),
|
|
||||||
"decoded MTP receive frame"
|
|
||||||
);
|
|
||||||
let negotiated_type_map = type_map.read().await.clone();
|
let negotiated_type_map = type_map.read().await.clone();
|
||||||
message.set_type_map(&negotiated_type_map);
|
message.set_type_map(&negotiated_type_map);
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -8,14 +8,6 @@ use crate::config::ConnectionConfig;
|
||||||
use crate::error::js_error;
|
use crate::error::js_error;
|
||||||
use crate::transport::WasmTransport;
|
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]
|
#[wasm_bindgen]
|
||||||
#[allow(deprecated)]
|
#[allow(deprecated)]
|
||||||
impl WasmClient {
|
impl WasmClient {
|
||||||
|
|
@ -87,18 +79,6 @@ impl WasmClient {
|
||||||
.unwrap_or("host does not support this protocol version"),
|
.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) {
|
let negotiated_version = match outcome.get_data(DataType::Version) {
|
||||||
Some(DataValue::Str(version)) => mtp_codec::Version::parse(version)
|
Some(DataValue::Str(version)) => mtp_codec::Version::parse(version)
|
||||||
.ok_or_else(|| js_error("host omitted a valid negotiated protocol version"))?,
|
.ok_or_else(|| js_error("host omitted a valid negotiated protocol version"))?,
|
||||||
|
|
@ -648,24 +628,3 @@ impl WasmClient {
|
||||||
Ok(server_challenge)
|
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")
|
|
||||||
);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
|
||||||
|
|
@ -128,6 +128,8 @@ pub struct WasmTransport {
|
||||||
buffer: Rc<RefCell<Vec<u8>>>,
|
buffer: Rc<RefCell<Vec<u8>>>,
|
||||||
/// Set to `true` when `open_next_stream` succeeds; cleared after the first frame is parsed.
|
/// Set to `true` when `open_next_stream` succeeds; cleared after the first frame is parsed.
|
||||||
new_stream_frame: Rc<Cell<bool>>,
|
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.
|
/// Serializes stream creation and writes across concurrent callers.
|
||||||
send_lock: Rc<AsyncMutex<()>>,
|
send_lock: Rc<AsyncMutex<()>>,
|
||||||
type_map: Rc<RefCell<TypeMap>>,
|
type_map: Rc<RefCell<TypeMap>>,
|
||||||
|
|
@ -207,6 +209,7 @@ impl WasmTransport {
|
||||||
stream_reader: Rc::new(RefCell::new(None)),
|
stream_reader: Rc::new(RefCell::new(None)),
|
||||||
buffer: Rc::new(RefCell::new(Vec::new())),
|
buffer: Rc::new(RefCell::new(Vec::new())),
|
||||||
new_stream_frame: Rc::new(Cell::new(false)),
|
new_stream_frame: Rc::new(Cell::new(false)),
|
||||||
|
outgoing_writer: Rc::new(RefCell::new(None)),
|
||||||
send_lock: Rc::new(AsyncMutex::new(())),
|
send_lock: Rc::new(AsyncMutex::new(())),
|
||||||
type_map: Rc::new(RefCell::new(TypeMap::latest())),
|
type_map: Rc::new(RefCell::new(TypeMap::latest())),
|
||||||
decode_limits: Rc::new(RefCell::new(decode_limits)),
|
decode_limits: Rc::new(RefCell::new(decode_limits)),
|
||||||
|
|
@ -244,12 +247,9 @@ impl WasmTransport {
|
||||||
return Err(js_error("message too large"));
|
return Err(js_error("message too large"));
|
||||||
}
|
}
|
||||||
|
|
||||||
// Use one WebTransport uni-stream per MTP frame. Chromium reliably
|
let writer_val = if let Some(writer) = self.outgoing_writer.borrow().clone() {
|
||||||
// publishes a browser-created uni-stream to the peer when it is
|
writer
|
||||||
// closed; leaving a shared stream open can leave the server waiting
|
} else {
|
||||||
// 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(
|
let create_stream = js_sys::Reflect::get(
|
||||||
&self.inner,
|
&self.inner,
|
||||||
&JsValue::from_str("createUnidirectionalStream"),
|
&JsValue::from_str("createUnidirectionalStream"),
|
||||||
|
|
@ -262,12 +262,15 @@ impl WasmTransport {
|
||||||
.map_err(|_| js_error("createUnidirectionalStream did not return a Promise"))?;
|
.map_err(|_| js_error("createUnidirectionalStream did not return a Promise"))?;
|
||||||
let stream = JsFuture::from(stream_promise).await?;
|
let stream = JsFuture::from(stream_promise).await?;
|
||||||
let writable_or_stream = resolve_stream_writable(&stream)?;
|
let writable_or_stream = resolve_stream_writable(&stream)?;
|
||||||
let writer_val = js_sys::Reflect::get(&writable_or_stream, &JsValue::from_str("getWriter"))
|
let writer = js_sys::Reflect::get(&writable_or_stream, &JsValue::from_str("getWriter"))
|
||||||
.map_err(|_| js_error("missing getWriter"))?
|
.map_err(|_| js_error("missing getWriter"))?
|
||||||
.dyn_into::<js_sys::Function>()
|
.dyn_into::<js_sys::Function>()
|
||||||
.map_err(|_| js_error("getWriter not a function"))?
|
.map_err(|_| js_error("getWriter not a function"))?
|
||||||
.call0(&writable_or_stream)
|
.call0(&writable_or_stream)
|
||||||
.map_err(|_| js_error("getWriter call failed"))?;
|
.map_err(|_| js_error("getWriter call failed"))?;
|
||||||
|
*self.outgoing_writer.borrow_mut() = Some(writer.clone());
|
||||||
|
writer
|
||||||
|
};
|
||||||
|
|
||||||
let chunk = js_sys::Uint8Array::from(frame);
|
let chunk = js_sys::Uint8Array::from(frame);
|
||||||
|
|
||||||
|
|
@ -280,24 +283,11 @@ impl WasmTransport {
|
||||||
.map_err(|e| js_error(format!("write failed: {:?}", e)))?;
|
.map_err(|e| js_error(format!("write failed: {:?}", e)))?;
|
||||||
if let Err(e) = JsFuture::from(write_promise.unchecked_into::<js_sys::Promise>()).await {
|
if let Err(e) = JsFuture::from(write_promise.unchecked_into::<js_sys::Promise>()).await {
|
||||||
log_stream_error_code(&e, "send_frame write");
|
log_stream_error_code(&e, "send_frame write");
|
||||||
|
self.outgoing_writer.borrow_mut().take();
|
||||||
release_writer_lock(&writer_val);
|
release_writer_lock(&writer_val);
|
||||||
return Err(e);
|
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(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -687,6 +677,12 @@ impl WasmTransport {
|
||||||
}
|
}
|
||||||
|
|
||||||
pub fn close(&self) {
|
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.
|
// Release reader locks before closing so they aren't treated as cancels.
|
||||||
if let Some(reader) = self.stream_reader.borrow_mut().take() {
|
if let Some(reader) = self.stream_reader.borrow_mut().take() {
|
||||||
release_reader_lock(&reader);
|
release_reader_lock(&reader);
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue