(fix): wasm rust runtime error (qol): update comments
This commit is contained in:
parent
92f1190b11
commit
c1761aae2b
3 changed files with 60 additions and 42 deletions
|
|
@ -19,8 +19,8 @@ export interface MTPCredentialStorage {
|
||||||
export type MTPStorage = MTPCredentialStorage;
|
export type MTPStorage = MTPCredentialStorage;
|
||||||
|
|
||||||
export type MTPLogEvent =
|
export type MTPLogEvent =
|
||||||
| { hint: "info" | "warning"; type: string; data: unknown }
|
| { hint: "info" | "warning"; type: string; data: unknown; direction?: "send" | "recv" }
|
||||||
| { hint: "error"; type: string | "error"; error: string; data?: unknown };
|
| { hint: "error"; type: string | "error"; error: string; data?: unknown; direction?: "send" | "recv" };
|
||||||
|
|
||||||
export type ParsedFrame = RawBindings.ParsedFrame;
|
export type ParsedFrame = RawBindings.ParsedFrame;
|
||||||
|
|
||||||
|
|
@ -605,13 +605,14 @@ export class MTPClient {
|
||||||
try {
|
try {
|
||||||
const frame = this.raw.bindings.parse_frame(message);
|
const frame = this.raw.bindings.parse_frame(message);
|
||||||
emit(this.#options.logger, isErrorType(frame.type)
|
emit(this.#options.logger, isErrorType(frame.type)
|
||||||
? { hint: "error", type: frame.type, error: errorMessage(frame), data: frame.data }
|
? { hint: "error", type: frame.type, error: errorMessage(frame), data: frame.data, direction: "send" }
|
||||||
: { hint: "info", type: frame.type, data: frame.data });
|
: { hint: "info", type: frame.type, data: frame.data, direction: "send" });
|
||||||
} catch (error) {
|
} catch (error) {
|
||||||
emit(this.#options.logger, {
|
emit(this.#options.logger, {
|
||||||
hint: "error",
|
hint: "error",
|
||||||
type: "error",
|
type: "error",
|
||||||
error: String(error),
|
error: String(error),
|
||||||
|
direction: "send",
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -622,6 +623,19 @@ export class MTPClient {
|
||||||
async request(type: MTPCommunicationType, data: Record<string, unknown>, options?: MTPRequestOptions): Promise<ParsedFrame>;
|
async request(type: MTPCommunicationType, data: Record<string, unknown>, options?: MTPRequestOptions): Promise<ParsedFrame>;
|
||||||
async request(typeOrFrame: Uint8Array | MTPCommunicationType, data?: Record<string, unknown>, options: MTPRequestOptions = {}): Promise<ParsedFrame> {
|
async request(typeOrFrame: Uint8Array | MTPCommunicationType, data?: Record<string, unknown>, options: MTPRequestOptions = {}): Promise<ParsedFrame> {
|
||||||
const frame = this.#buildFrame(typeOrFrame, data, options);
|
const frame = this.#buildFrame(typeOrFrame, data, options);
|
||||||
|
try {
|
||||||
|
const parsed = this.raw.bindings.parse_frame(frame);
|
||||||
|
emit(this.#options.logger, isErrorType(parsed.type)
|
||||||
|
? { hint: "error", type: parsed.type, error: errorMessage(parsed), data: parsed.data, direction: "send" }
|
||||||
|
: { hint: "info", type: parsed.type, data: parsed.data, direction: "send" });
|
||||||
|
} catch (error) {
|
||||||
|
emit(this.#options.logger, {
|
||||||
|
hint: "error",
|
||||||
|
type: "error",
|
||||||
|
error: String(error),
|
||||||
|
direction: "send",
|
||||||
|
});
|
||||||
|
}
|
||||||
return await this.raw.client.request(frame, options.responseType ?? null);
|
return await this.raw.client.request(frame, options.responseType ?? null);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -643,12 +657,14 @@ export class MTPClient {
|
||||||
type: frame.type,
|
type: frame.type,
|
||||||
error: errorMessage(frame),
|
error: errorMessage(frame),
|
||||||
data: frame.data,
|
data: frame.data,
|
||||||
|
direction: "recv",
|
||||||
});
|
});
|
||||||
} else {
|
} else {
|
||||||
emit(this.#options.logger, {
|
emit(this.#options.logger, {
|
||||||
hint: "info",
|
hint: "info",
|
||||||
type: frame.type,
|
type: frame.type,
|
||||||
data: frame.data,
|
data: frame.data,
|
||||||
|
direction: "recv",
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -687,9 +687,39 @@ impl WasmClient {
|
||||||
|
|
||||||
fn set_state(&self, new_state: ConnectionState) {
|
fn set_state(&self, new_state: ConnectionState) {
|
||||||
self.state.set(new_state);
|
self.state.set(new_state);
|
||||||
let _ = self
|
|
||||||
.on_state_change
|
// Defer the callback to a microtask so re-entrant &mut self calls don't alias.
|
||||||
.call1(&JsValue::NULL, &JsValue::from(new_state as u8));
|
let cb = self.on_state_change.clone();
|
||||||
|
let val = JsValue::from(new_state as u8);
|
||||||
|
let closure = Closure::wrap(Box::new(move || {
|
||||||
|
let _ = cb.call1(&JsValue::NULL, &val);
|
||||||
|
}) as Box<dyn FnMut()>);
|
||||||
|
|
||||||
|
let global = js_sys::global();
|
||||||
|
let mut closure_opt = Some(closure);
|
||||||
|
|
||||||
|
let qmt = js_sys::Reflect::get(&global, &JsValue::from_str("queueMicrotask"))
|
||||||
|
.and_then(|f| f.dyn_into::<js_sys::Function>().map_err(Into::into));
|
||||||
|
let scheduled = match qmt {
|
||||||
|
Ok(qmt) => {
|
||||||
|
if let Some(c) = closure_opt.take() {
|
||||||
|
let _ = qmt.call1(&global, c.as_ref());
|
||||||
|
c.forget();
|
||||||
|
}
|
||||||
|
true
|
||||||
|
}
|
||||||
|
Err(_) => false,
|
||||||
|
};
|
||||||
|
if !scheduled {
|
||||||
|
if let Ok(set_timeout) = js_sys::Reflect::get(&global, &JsValue::from_str("setTimeout"))
|
||||||
|
.and_then(|f| f.dyn_into::<js_sys::Function>().map_err(Into::into))
|
||||||
|
{
|
||||||
|
if let Some(c) = closure_opt.take() {
|
||||||
|
let _ = set_timeout.call2(&global, c.as_ref(), &JsValue::from_f64(0.0));
|
||||||
|
c.forget();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
fn start_receive_loop(&mut self, transport: WasmTransport) {
|
fn start_receive_loop(&mut self, transport: WasmTransport) {
|
||||||
|
|
|
||||||
|
|
@ -10,16 +10,7 @@ use crate::frame::parse_frame_value;
|
||||||
|
|
||||||
const CLOSE_FRAME_LEN: u32 = u32::MAX;
|
const CLOSE_FRAME_LEN: u32 = u32::MAX;
|
||||||
|
|
||||||
/// Inspect a JS error value for a WebTransport **stream-level** error and, if
|
/// Logs the `streamErrorCode` from a stream-level WebTransportError (STOP_SENDING / RESET_STREAM). Session errors are skipped.
|
||||||
/// present, log the `streamErrorCode` carried by STOP_SENDING / RESET_STREAM.
|
|
||||||
///
|
|
||||||
/// Per draft-ietf-webtrans-http3-15 §4.4, a WebTransport application MUST
|
|
||||||
/// provide an error code for those operations. The browser surfaces these as
|
|
||||||
/// `WebTransportError` with `source = "stream"` and a numeric `streamErrorCode`.
|
|
||||||
///
|
|
||||||
/// Session-level errors (`source = "session"`) are normal connection
|
|
||||||
/// closures and are **not** logged here — they propagate to `on_error`
|
|
||||||
/// in the receive loop like any other transport error.
|
|
||||||
fn log_stream_error_code(error: &JsValue, context: &str) {
|
fn log_stream_error_code(error: &JsValue, context: &str) {
|
||||||
let source = js_sys::Reflect::get(error, &JsValue::from_str("source"))
|
let source = js_sys::Reflect::get(error, &JsValue::from_str("source"))
|
||||||
.ok()
|
.ok()
|
||||||
|
|
@ -78,10 +69,7 @@ fn resolve_stream_readable(recv_stream: &JsValue) -> Result<JsValue, JsValue> {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Release a `WritableStreamDefaultWriter`'s lock on its stream. Called after
|
/// Releases a writer's lock so an abandoned writer isn't treated as an abort (which sends STOP_SENDING).
|
||||||
/// `writer.close()` (or on write failure) so the runtime does not interpret an
|
|
||||||
/// abandoned locked writer as an abort, which would surface as STOP_SENDING to
|
|
||||||
/// the peer. Errors are ignored — `releaseLock` is best-effort cleanup.
|
|
||||||
fn release_writer_lock(writer: &JsValue) {
|
fn release_writer_lock(writer: &JsValue) {
|
||||||
if let Ok(release) = js_sys::Reflect::get(writer, &JsValue::from_str("releaseLock"))
|
if let Ok(release) = js_sys::Reflect::get(writer, &JsValue::from_str("releaseLock"))
|
||||||
.and_then(|f| f.dyn_into::<js_sys::Function>().map_err(Into::into))
|
.and_then(|f| f.dyn_into::<js_sys::Function>().map_err(Into::into))
|
||||||
|
|
@ -90,11 +78,7 @@ fn release_writer_lock(writer: &JsValue) {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Release a `ReadableStreamDefaultReader`'s lock on its stream. Mirrors
|
/// Releases a reader's lock so an abandoned reader isn't treated as a cancel (which sends STOP_SENDING).
|
||||||
/// `release_writer_lock`: abandoning a locked reader can be interpreted by the
|
|
||||||
/// runtime as a `reader.cancel()` (sending STOP_SENDING to the peer) even on an
|
|
||||||
/// already-closed or errored stream. Calling `releaseLock` explicitly avoids
|
|
||||||
/// that. Errors are ignored — best-effort cleanup.
|
|
||||||
fn release_reader_lock(reader: &JsValue) {
|
fn release_reader_lock(reader: &JsValue) {
|
||||||
if let Ok(release) = js_sys::Reflect::get(reader, &JsValue::from_str("releaseLock"))
|
if let Ok(release) = js_sys::Reflect::get(reader, &JsValue::from_str("releaseLock"))
|
||||||
.and_then(|f| f.dyn_into::<js_sys::Function>().map_err(Into::into))
|
.and_then(|f| f.dyn_into::<js_sys::Function>().map_err(Into::into))
|
||||||
|
|
@ -256,15 +240,11 @@ impl WasmTransport {
|
||||||
.call0(&writer_val)
|
.call0(&writer_val)
|
||||||
.map_err(|e| js_error(&format!("close failed: {:?}", e)))?;
|
.map_err(|e| js_error(&format!("close failed: {:?}", e)))?;
|
||||||
if let Err(e) = JsFuture::from(close_promise.unchecked_into::<js_sys::Promise>()).await {
|
if let Err(e) = JsFuture::from(close_promise.unchecked_into::<js_sys::Promise>()).await {
|
||||||
// The write already succeeded; a STOP_SENDING on close just means
|
// Write succeeded; STOP_SENDING on close just means peer stopped reading before FIN.
|
||||||
// the peer stopped reading before we could send FIN. The data is in
|
|
||||||
// flight, so this is not a send failure — log and return success.
|
|
||||||
log_stream_error_code(&e, "send_frame close");
|
log_stream_error_code(&e, "send_frame close");
|
||||||
}
|
}
|
||||||
|
|
||||||
// Always release the writer's lock on the WritableStream. Abandoning a
|
// Release the lock so the writer isn't treated as an abort.
|
||||||
// locked writer (e.g. via drop) can be interpreted by the runtime as an
|
|
||||||
// abort, which may surface as STOP_SENDING to the peer.
|
|
||||||
release_writer_lock(&writer_val);
|
release_writer_lock(&writer_val);
|
||||||
|
|
||||||
Ok(())
|
Ok(())
|
||||||
|
|
@ -418,12 +398,7 @@ impl WasmTransport {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
None => {
|
None => {
|
||||||
// Current stream finished; the next frame (if any) is on a
|
// Stream finished; release the reader's lock to avoid a spurious cancel.
|
||||||
// subsequent stream. Any trailing partial bytes are dropped
|
|
||||||
// since the host never splits a frame across streams.
|
|
||||||
// Release the reader's lock explicitly so the runtime does
|
|
||||||
// not treat the abandoned lock as a cancel (which would send
|
|
||||||
// STOP_SENDING to the peer on an already-closed stream).
|
|
||||||
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);
|
||||||
}
|
}
|
||||||
|
|
@ -470,10 +445,7 @@ impl WasmTransport {
|
||||||
}
|
}
|
||||||
|
|
||||||
pub fn close(&self) {
|
pub fn close(&self) {
|
||||||
// Release any held reader locks before tearing down the session, so the
|
// Release reader locks before closing so they aren't treated as cancels.
|
||||||
// runtime does not interpret an abandoned locked reader as a cancel
|
|
||||||
// (which would send STOP_SENDING to the peer). Once the locks are
|
|
||||||
// released the underlying streams can be torn down cleanly.
|
|
||||||
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