From aefcae977f549c58e6950e277a110f4afc454cfd Mon Sep 17 00:00:00 2001 From: Alex Emmet <111742636+Alex-Emmet@users.noreply.github.com> Date: Fri, 26 Jun 2026 23:10:54 +0200 Subject: [PATCH] WASM & Flake --- example-usage/host_keys.json.stale-no-kem | 4 - flake.lock | 12 +- wasm/src/client.rs | 18 +- wasm/src/transport.rs | 377 ++++++++++------------ 4 files changed, 182 insertions(+), 229 deletions(-) delete mode 100644 example-usage/host_keys.json.stale-no-kem diff --git a/example-usage/host_keys.json.stale-no-kem b/example-usage/host_keys.json.stale-no-kem deleted file mode 100644 index f38d29b..0000000 --- a/example-usage/host_keys.json.stale-no-kem +++ /dev/null @@ -1,4 +0,0 @@ -{ - "host_id": 1, - "keyring": "0000000007a0cc83597da99032dfe763d1e27eb368420b023be3518323638a09c10446e5a1e182beca1c0e8dda17a7b4b69cc0bbf1ec45f0b03c6b3c1052a540f69e4140f7d1d11039398f77a45630b7fe3254b001ce08d5dd92c6cba881ed20f4eefa6153608d7fb766bdb038939348ba091806c8bc5608650ae8c187be3bf90384b3cd227041462ba409d521c3970c7cc3fce096058f2648c6e2e242273e0af2c47fbf4d899be75b910e7d357618248d51c6e608e9c7a247e88f5961cad53bc163cea869a54c4abbc3d6ebbffa918296f4a21694f83941335ac061068ae67edce3783b8a6afb2676171906ba862a97df2b11dc9911d7e654718d2ecfb9e4b69217f50065981f15ce1d7b0b6d249f68a357c9c56b6c3edfbcf876713e7f6376baa181a4535d0ab5811c2b7e7ffef23cfb6a4af0ea4e6e2f55685d594e97a20c0fac2246618120117879f521f7bb2a26a4f697e9451e3f3653695235dac13feeb0cb8999ff41fad6bf55426877bbae8fd1cd8622c2cf90a57701a7136f01a3d03202b6b1488f70c5840ced7509425b1c3ca019bd55060509a8f2bb42977c78187d7d2801309af9ae2f609cb918c1956b0f282fa4eb591b530c7f5f83c1d54bacabe4abe34780e6c6802509082ad38ec1bcc6105eae77e67e3d414c7eb248d3f5582681e9fbbd5183fa978604f094b9013804055745deba587086f2eeb95a911d1cb4d0dd960c64ecfaad2ae0981e44d38c2441b5df03c9ea23338be78dcd0d4e8712ff6e9ab5d751870e7243c06835e2e7a14b7d3005251ab66fbff4820c0e03317035ead18e213846af6325d3e871524708eaa3093f883e223281a8f1be55cd3dc084dbdd461566c35f6a0d418dcc66ab1697242bf1016283f4e50e76e8f42feae937d4b908d6132c7d80755a58120d75489931f58de641f2f8676fca6b60da18e977b38108d8ffe98f96165222c07520c06974b39aa8b5139458ea101727f95818e69a6afcf0fa028224e8c0af5cd4f32db0a0583eb42a11fe7fa08c424563dde4e2977c898f2f68b2f183016c44dba013851d909844981d0d2635e7fe4af8f80ca3acbea30d11a2b92cf4c4866765493217fd078aa8903ea50aeace282d90a156cc37714f95d109c5421f947209656305beffa891717c7bc8802bb3eab87fc753b0450155e5f0b0c7f14bd66db4e69d330b64e8fc6cbc946e515536b06f65966a982435af15c4b612d25434ef3bc98c33c47abc3b98159733c5f03f541342e6d5e1c0ca6caa00e5616cee151771dbf3a73d571ee0b37ae54d4672e411feb1364f5b820dbc0912e8c43a8213d41576d8b207ab8564537a23aa99c08f6079b83a63c376d10b919023cb57363e7562527e634cc45f5131ae3204615ac078bb216677e65cafb2a95246635f76e1a1a7659d8be263ef403395e16ca388ec1b58b51376b221ab64309bb493687e6cfa7587b603a53e00b1eb4d413f299d7741e1de70df3dc6d3190a8c84fffb7a7901d524f4b5d9a7229c5c99f61faf3d5295760f35d3d54bbea944f044c35a2202e6d2865369809d0df2f64151fb83b3a74d8ad1e7b6cc6d0088a535f7ca6dc384bd93cef50132a0d882caac96acecd15d6eeb600ca6c25f7ff19795cf8595d62bdf5022eadc99a471a36a2b46cd1c4f4d999560eadf4994f77c50083eaad96f574100b594a5272b661c80ffed36a364e7aa69e5a0c69f1af9ba44eeb50e892309ca442ec09f0c716d1483f4d74efbf12ff845f5252e83e64d28a83541b5fbef87296e8a19c84c7c6fa46f61299f7375d74c8441e432613a5f721548a010406ab56f18c997a48d74970f6e651f5d1223e233b5b2ff952e3824b72f92755f491110517d71878e7469beaab6c0835c6467598f7969f1b85157dabebb1003b1f9081bab2aa2296a269a35c34f494e706291db5923690c19e4eaafa5bc3e838e3c6434d0ab91c0b95ce4221ed9cfbe011fa24168d821bed5af6f2fbcb75db1e1ea2a71e5a5f38ad4fab3ebf61983745f4464a0fcf5917415b01846d037fdb628637241070287cbcdcdf638353f01c03030bb36df7f48b5b1eba6152d4f569d034b5f2c9b61b1cf0402df910947d100f54bd0b0261445e60c0f4b1e1df292dbf26bb85ab327a19a33bca531d497b4e7181ba54881cfa1b43d35da64436850d298077877bd41cfd172d91413fee1888ad603a2368fedd611690878ca1f789a6fd21b6932b868bb3084e689c8995113478bbf6ba8ad04e937be34dce58b1117ebbd2eddde489fb34028e117b59fac89bf9a6b5970a989840447d3649f6e955fd0561a56d4038fac4fa4f67515400ecbf256a9339a6b61a62ffc101fc7241422e991b54215f2f6ae4d719d99e660cc6800bf33ce2a00a2b8d101908eb32a31ecb311f8d2da5509a4f0800515193d336a716a52868813a8ea5b878aacd5dca3753c66ea38c3f1f48529746581cbc375cdef323d885ee119f4b1397f5828dec90923caae9845c1f189c482d68342a70d979b95943287127996ce137cfdc2ce3f6ecbfe5f5f5b90b8260fc662c897d4cc1414ad77314f4c15b84308f5836d33d8183388dbbc0d15aa2a9c4c14f447450861e373544c25a61b874ba6716af6bce7bec2a42d6408e73fbf0a5b27fbd9aa3055bfd7b0a79cac66c6f51f00ab90bcb0dc28254a8a4fdbfb7c5bc63a6aa7227d0e335eb0eb01124582db8c7b92f8e84515acddc7bb24e2b3059329228c197e3432bead261b2dbf1798491ebb11658700201d4224bc1f165129ac6fac5909a664ef107f70035019ff7051bba3dc980ac3da0020faa51b18eb51dd316605d3842f71c23a1a15ec71c1cc5d650dabe60da9dc7c5600209300a95b9dcf41d3ee469fc9937737d572915f10e5ea2da286e13bf2dbb54ef6" -} \ No newline at end of file diff --git a/flake.lock b/flake.lock index 8c32242..0244fe2 100644 --- a/flake.lock +++ b/flake.lock @@ -20,11 +20,11 @@ }, "nixpkgs": { "locked": { - "lastModified": 1781577229, - "narHash": "sha256-lrp67w8AulE9Ks53n27I45ADSzbOCn4H+CNW1Ck8B+8=", + "lastModified": 1782467914, + "narHash": "sha256-pGvFkM8N0xEkIIXDe5YYfbEAvHrk4IxBrjB/x8OomhE=", "owner": "NixOS", "repo": "nixpkgs", - "rev": "567a49d1913ce81ac6e9582e3553dd90a955875f", + "rev": "e73de5be04e0eff4190a1432b946d469c794e7b4", "type": "github" }, "original": { @@ -62,11 +62,11 @@ "nixpkgs": "nixpkgs_2" }, "locked": { - "lastModified": 1782357464, - "narHash": "sha256-mXgoT1qDHCdSfF9IvhMtEEFNy9dxrmUfSViwP7RpzOQ=", + "lastModified": 1782443907, + "narHash": "sha256-P+pADLtK7qC1mz0/5Xq9uF77oahUR4zYLTaitiHsUHg=", "owner": "oxalica", "repo": "rust-overlay", - "rev": "77a8263847fb02dc49dbe377278ef6b952f1c6bb", + "rev": "4b06ff4acf3491ff69721df852507fcc51d0a13d", "type": "github" }, "original": { diff --git a/wasm/src/client.rs b/wasm/src/client.rs index f9ea26f..36604d6 100644 --- a/wasm/src/client.rs +++ b/wasm/src/client.rs @@ -193,7 +193,6 @@ impl WasmClient { self.set_state(ConnectionState::Connecting); let transport = WasmTransport::connect(&config.url, config.server_certificate_hashes.clone()).await?; - let inner = transport.inner().clone(); let version_str = format!("{}", PROTOCOL_VERSION); let ident = CommunicationValue::new(CommunicationType::Identification) @@ -207,6 +206,7 @@ impl WasmClient { .map_err(|e| js_error(&format!("encode failed: {}", e)))?; transport.send_frame(&ident_bytes).await?; + let loop_transport = transport.clone(); self.transport = Some(transport); self.set_state(ConnectionState::Connected); @@ -214,9 +214,7 @@ impl WasmClient { let on_msg = self.on_message.clone(); let on_err = self.on_error.clone(); wasm_bindgen_futures::spawn_local(async move { - WasmTransport::from_inner(inner) - .receive_loop(on_msg, on_err) - .await; + loop_transport.receive_loop(on_msg, on_err).await; state.set(ConnectionState::Disconnected); }); Ok(()) @@ -250,7 +248,6 @@ impl WasmClient { let transport = WasmTransport::connect(&config.url, config.server_certificate_hashes.clone()).await?; - let inner = transport.inner().clone(); // 1. Send the unsigned Identification hello. let hello = CommunicationValue::new(CommunicationType::Identification) @@ -360,6 +357,7 @@ impl WasmClient { } }; + let loop_transport = transport.clone(); self.transport = Some(transport); self.set_state(ConnectionState::Connected); @@ -367,9 +365,7 @@ impl WasmClient { let on_msg = self.on_message.clone(); let on_err = self.on_error.clone(); wasm_bindgen_futures::spawn_local(async move { - WasmTransport::from_inner(inner) - .receive_loop(on_msg, on_err) - .await; + loop_transport.receive_loop(on_msg, on_err).await; state.set(ConnectionState::Disconnected); }); @@ -403,7 +399,6 @@ impl WasmClient { let transport = WasmTransport::connect(&config.url, config.server_certificate_hashes.clone()).await?; - let inner = transport.inner().clone(); // 1. Send the unsigned Register hello (version + public-key bundle). let hello = CommunicationValue::new(CommunicationType::Register) @@ -510,6 +505,7 @@ impl WasmClient { return Err(e); } + let loop_transport = transport.clone(); self.transport = Some(transport); self.set_state(ConnectionState::Connected); @@ -517,9 +513,7 @@ impl WasmClient { let on_msg = self.on_message.clone(); let on_err = self.on_error.clone(); wasm_bindgen_futures::spawn_local(async move { - WasmTransport::from_inner(inner) - .receive_loop(on_msg, on_err) - .await; + loop_transport.receive_loop(on_msg, on_err).await; state.set(ConnectionState::Disconnected); }); diff --git a/wasm/src/transport.rs b/wasm/src/transport.rs index 06713b7..828473a 100644 --- a/wasm/src/transport.rs +++ b/wasm/src/transport.rs @@ -1,3 +1,6 @@ +use std::cell::RefCell; +use std::rc::Rc; + use wasm_bindgen::JsCast; use wasm_bindgen::prelude::*; use wasm_bindgen_futures::JsFuture; @@ -27,9 +30,37 @@ fn resolve_stream_readable(recv_stream: &JsValue) -> Result { } } +/// Outcome of reading the next framed message from the incoming stream(s). +enum FrameOutcome { + /// A complete application frame. + Frame(Vec), + /// The peer sent an explicit close frame (length == `u32::MAX`). + Closed, + /// The incoming-streams readable ended (transport gone), no more frames. + Ended, +} + +/* + * WebTransport client transport. + * + * The native host sends with a *persistent* uni-directional stream: the auth + * `Challenge` and the final `IdentificationResponse`/`RegisterResponse` arrive + * as two length-prefixed frames on the *same* QUIC stream, and later + * application messages may arrive on subsequent streams. The reader state + * (`streams_reader`, `stream_reader`, `buffer`) is therefore shared via `Rc` + * between the handshake (`read_one_frame`) and the background `receive_loop`, + * so frames are never lost across the boundary and multiple frames can be read + * from one stream. + */ #[derive(Clone)] pub struct WasmTransport { inner: WebTransport, + /// Reader over `incoming_unidirectional_streams()` (a singleton stream of streams). + streams_reader: Rc>>, + /// Reader over the host's current uni-directional stream, if one is open. + stream_reader: Rc>>, + /// Bytes already read from the current stream but not yet consumed as a frame. + buffer: Rc>>, } impl WasmTransport { @@ -58,17 +89,18 @@ impl WasmTransport { JsFuture::from(transport.ready()) .await .map_err(|e| js_error(&format!("WebTransport ready failed: {:?}", e)))?; - Ok(Self { inner: transport }) + Ok(Self { + inner: transport, + streams_reader: Rc::new(RefCell::new(None)), + stream_reader: Rc::new(RefCell::new(None)), + buffer: Rc::new(RefCell::new(Vec::new())), + }) } pub fn inner(&self) -> &WebTransport { &self.inner } - pub fn from_inner(inner: WebTransport) -> Self { - Self { inner } - } - pub async fn send_frame(&self, frame: &[u8]) -> Result<(), JsValue> { let stream_promise = self.inner.create_unidirectional_stream(); let stream = JsFuture::from(stream_promise).await?; @@ -110,247 +142,178 @@ impl WasmTransport { Ok(()) } - /// Read exactly one frame from incoming uni streams, then release the reader - /// so `receive_loop` can pick up from where we left off. - pub async fn read_one_frame(&self) -> Result, JsValue> { + /// Get (creating once) the reader over `incoming_unidirectional_streams()`. + fn ensure_streams_reader(&self) -> Result { + if let Some(reader) = self.streams_reader.borrow().clone() { + return Ok(reader); + } let incoming = self.inner.incoming_unidirectional_streams(); - - let reader_fn = js_sys::Reflect::get(&incoming, &JsValue::from_str("getReader")) + let reader = js_sys::Reflect::get(&incoming, &JsValue::from_str("getReader")) .map_err(|_| js_error("missing getReader"))? .dyn_into::() - .map_err(|_| js_error("getReader not a function"))?; - let reader_val = reader_fn + .map_err(|_| js_error("getReader not a function"))? .call0(&incoming) .map_err(|_| js_error("getReader call failed"))?; + *self.streams_reader.borrow_mut() = Some(reader.clone()); + Ok(reader) + } - let read_fn = js_sys::Reflect::get(&reader_val, &JsValue::from_str("read")) + /// Accept the next incoming uni-directional stream and make it current. + /// Returns `false` if the incoming-streams readable has ended. + async fn open_next_stream(&self) -> Result { + let streams_reader = self.ensure_streams_reader()?; + + let read_fn = js_sys::Reflect::get(&streams_reader, &JsValue::from_str("read")) .map_err(|_| js_error("missing read"))? .dyn_into::() .map_err(|_| js_error("read not a function"))?; - let result_promise = read_fn - .call0(&reader_val) - .map_err(|_| js_error("read call failed"))?; - let result = JsFuture::from(result_promise.unchecked_into::()) - .await - .map_err(|e| js_error(&format!("read failed: {:?}", e)))?; - - // Release the reader lock so receive_loop can create its own reader - if let Some(release_fn) = - js_sys::Reflect::get(&reader_val, &JsValue::from_str("releaseLock")) - .ok() - .and_then(|f| f.dyn_into::().ok()) - { - let _ = release_fn.call0(&reader_val); - } + let result = JsFuture::from( + read_fn + .call0(&streams_reader) + .map_err(|_| js_error("read call failed"))? + .unchecked_into::(), + ) + .await + .map_err(|e| js_error(&format!("accept stream failed: {:?}", e)))?; let done = js_sys::Reflect::get(&result, &JsValue::from_str("done")) .ok() .and_then(|v| v.as_bool()) .unwrap_or(false); if done { - return Err(js_error("stream ended before frame")); + return Ok(false); } let recv_stream = js_sys::Reflect::get(&result, &JsValue::from_str("value")) .map_err(|_| js_error("missing value"))?; - - let readable_or_stream = resolve_stream_readable(&recv_stream)?; - let stream_reader_fn = - js_sys::Reflect::get(&readable_or_stream, &JsValue::from_str("getReader")) - .map_err(|_| js_error("missing stream getReader"))? - .dyn_into::() - .map_err(|_| js_error("stream getReader not a function"))?; - let stream_reader = stream_reader_fn - .call0(&readable_or_stream) + let readable = resolve_stream_readable(&recv_stream)?; + let reader = js_sys::Reflect::get(&readable, &JsValue::from_str("getReader")) + .map_err(|_| js_error("missing stream getReader"))? + .dyn_into::() + .map_err(|_| js_error("stream getReader not a function"))? + .call0(&readable) .map_err(|_| js_error("stream getReader call failed"))?; - let mut buffer: Vec = Vec::new(); + *self.stream_reader.borrow_mut() = Some(reader); + Ok(true) + } + + /// Read one chunk from the current stream. `Ok(None)` means the stream ended. + async fn read_chunk(&self) -> Result>, JsValue> { + let reader = match self.stream_reader.borrow().clone() { + Some(r) => r, + None => return Ok(None), + }; + + let read_fn = js_sys::Reflect::get(&reader, &JsValue::from_str("read")) + .map_err(|_| js_error("missing read"))? + .dyn_into::() + .map_err(|_| js_error("read not a function"))?; + let result = JsFuture::from( + read_fn + .call0(&reader) + .map_err(|_| js_error("read call failed"))? + .unchecked_into::(), + ) + .await + .map_err(|e| js_error(&format!("read failed: {:?}", e)))?; + + let done = js_sys::Reflect::get(&result, &JsValue::from_str("done")) + .ok() + .and_then(|v| v.as_bool()) + .unwrap_or(true); + if done { + return Ok(None); + } + + let value = js_sys::Reflect::get(&result, &JsValue::from_str("value")) + .map_err(|_| js_error("missing value"))?; + Ok(Some(js_sys::Uint8Array::new(&value).to_vec())) + } + + /// Try to pull one complete frame out of the buffer without reading more. + fn parse_buffer(&self) -> Result, JsValue> { + let buf = self.buffer.borrow(); + if buf.len() < 4 { + return Ok(None); + } + let frame_len = u32::from_be_bytes([buf[0], buf[1], buf[2], buf[3]]); + if frame_len == CLOSE_FRAME_LEN { + return Ok(Some(FrameOutcome::Closed)); + } + let frame_len = frame_len as usize; + let Some(frame_end) = 4usize.checked_add(frame_len) else { + return Err(js_error("invalid frame length")); + }; + if frame_end > buf.len() { + return Ok(None); + } + let frame = buf[4..frame_end].to_vec(); + drop(buf); + self.buffer.borrow_mut().drain(..frame_end); + Ok(Some(FrameOutcome::Frame(frame))) + } + + /* + * Read the next framed message from the host. Frames are length-prefixed + * (u32 big-endian) and may be packed several-per-stream (the host reuses a + * persistent uni stream) or one-per-stream; both are handled by buffering + * across reads and advancing to the next stream when the current one ends. + */ + async fn next_frame(&self) -> Result { loop { - let stream_read_fn = - match js_sys::Reflect::get(&stream_reader, &JsValue::from_str("read")) - .ok() - .and_then(|f| f.dyn_into::().ok()) - { - Some(f) => f, - None => break, - }; - let chunk_promise = match stream_read_fn.call0(&stream_reader) { - Ok(p) => p, - Err(_) => break, - }; - let chunk_result = - match JsFuture::from(chunk_promise.unchecked_into::()).await { - Ok(v) => v, - Err(_) => break, - }; - - let chunk_done = js_sys::Reflect::get(&chunk_result, &JsValue::from_str("done")) - .ok() - .and_then(|v| v.as_bool()) - .unwrap_or(true); - if chunk_done { - return Err(js_error("stream closed before complete frame")); + if let Some(outcome) = self.parse_buffer()? { + return Ok(outcome); } - if let Ok(chunk_val) = js_sys::Reflect::get(&chunk_result, &JsValue::from_str("value")) - { - let arr = js_sys::Uint8Array::new(&chunk_val).to_vec(); - if !arr.is_empty() { - buffer.extend_from_slice(&arr); - } + let have_stream = self.stream_reader.borrow().is_some(); + if !have_stream && !self.open_next_stream().await? { + return Ok(FrameOutcome::Ended); } - if buffer.len() >= 4 { - let frame_len = u32::from_be_bytes([buffer[0], buffer[1], buffer[2], buffer[3]]); - if frame_len == CLOSE_FRAME_LEN { - return Err(js_error("connection closed before frame")); + match self.read_chunk().await? { + Some(chunk) => { + if !chunk.is_empty() { + self.buffer.borrow_mut().extend_from_slice(&chunk); + } } - - let frame_len = frame_len as usize; - let Some(frame_end) = 4usize.checked_add(frame_len) else { - return Err(js_error("invalid frame length")); - }; - if frame_end <= buffer.len() { - return Ok(buffer[4..frame_end].to_vec()); + None => { + // Current stream finished; the next frame (if any) is on a + // subsequent stream. Any trailing partial bytes are dropped + // since the host never splits a frame across streams. + *self.stream_reader.borrow_mut() = None; + self.buffer.borrow_mut().clear(); } } } - - Err(js_error("stream ended before frame complete")) } + /// Read exactly one application frame (used during the auth handshake). + pub async fn read_one_frame(&self) -> Result, JsValue> { + match self.next_frame().await? { + FrameOutcome::Frame(frame) => Ok(frame), + FrameOutcome::Closed => Err(js_error("connection closed before frame")), + FrameOutcome::Ended => Err(js_error("stream ended before frame")), + } + } + + /// Background loop: deliver every incoming frame to `on_message` until the + /// connection closes. Shares reader state with `read_one_frame`, so frames + /// buffered during the handshake are not lost. pub async fn receive_loop(&self, on_message: js_sys::Function, on_error: js_sys::Function) { - let incoming = self.inner.incoming_unidirectional_streams(); - - let reader_fn = match js_sys::Reflect::get(&incoming, &JsValue::from_str("getReader")) { - Ok(f) => match f.dyn_into::() { - Ok(f) => f, - Err(_) => return, - }, - Err(_) => return, - }; - let reader_val = match reader_fn.call0(&incoming) { - Ok(v) => v, - Err(_) => return, - }; - loop { - let read_fn = match js_sys::Reflect::get(&reader_val, &JsValue::from_str("read")) { - Ok(f) => match f.dyn_into::() { - Ok(f) => f, - Err(_) => break, - }, - Err(_) => break, - }; - let result = match read_fn.call0(&reader_val) { - Ok(p) => match JsFuture::from(p.unchecked_into::()).await { - Ok(v) => v, - Err(_) => break, - }, - Err(_) => break, - }; - - let done = js_sys::Reflect::get(&result, &JsValue::from_str("done")) - .ok() - .and_then(|v| v.as_bool()) - .unwrap_or(false); - if done { - break; - } - - let recv_stream = match js_sys::Reflect::get(&result, &JsValue::from_str("value")) { - Ok(v) => v, - Err(_) => continue, - }; - - let this = self.clone(); - let on_msg = on_message.clone(); - let on_err = on_error.clone(); - wasm_bindgen_futures::spawn_local(async move { - let _ = this.handle_stream(recv_stream, on_msg, on_err).await; - }); - } - } - - async fn handle_stream( - &self, - recv_stream: JsValue, - on_message: js_sys::Function, - _on_error: js_sys::Function, - ) -> Result<(), JsValue> { - let readable_or_stream = resolve_stream_readable(&recv_stream)?; - let stream_reader_fn = - js_sys::Reflect::get(&readable_or_stream, &JsValue::from_str("getReader"))? - .dyn_into::()?; - let stream_reader = stream_reader_fn.call0(&readable_or_stream)?; - - let mut buffer: Vec = Vec::new(); - - loop { - let stream_read_fn = js_sys::Reflect::get(&stream_reader, &JsValue::from_str("read"))? - .dyn_into::()?; - let chunk_promise = stream_read_fn.call0(&stream_reader)?; - let chunk_result = - JsFuture::from(chunk_promise.unchecked_into::()).await?; - - let chunk_done = js_sys::Reflect::get(&chunk_result, &JsValue::from_str("done")) - .ok() - .and_then(|v| v.as_bool()) - .unwrap_or(true); - if chunk_done { - break; - } - - if let Ok(chunk_val) = js_sys::Reflect::get(&chunk_result, &JsValue::from_str("value")) - { - let arr = js_sys::Uint8Array::new(&chunk_val).to_vec(); - if !arr.is_empty() { - buffer.extend_from_slice(&arr); + match self.next_frame().await { + Ok(FrameOutcome::Frame(frame)) => { + let arr = js_sys::Uint8Array::from(&frame[..]); + let _ = on_message.call1(&JsValue::NULL, &arr); } - } - - // Extract all complete frames from the buffer - while buffer.len() >= 4 { - let frame_len = u32::from_be_bytes([buffer[0], buffer[1], buffer[2], buffer[3]]); - if frame_len == CLOSE_FRAME_LEN { - return Ok(()); - } - - let frame_len = frame_len as usize; - let Some(frame_end) = 4usize.checked_add(frame_len) else { - return Err(js_error("invalid frame length")); - }; - if frame_end > buffer.len() { + Ok(FrameOutcome::Closed) | Ok(FrameOutcome::Ended) => break, + Err(e) => { + let _ = on_error.call1(&JsValue::NULL, &e); break; } - let frame = buffer[4..frame_end].to_vec(); - let arr = js_sys::Uint8Array::from(&frame[..]); - let _ = on_message.call1(&JsValue::NULL, &arr); - buffer.drain(..frame_end); } } - - // Process any remaining complete frames after stream closes - while buffer.len() >= 4 { - let frame_len = u32::from_be_bytes([buffer[0], buffer[1], buffer[2], buffer[3]]); - if frame_len == CLOSE_FRAME_LEN { - return Ok(()); - } - - let frame_len = frame_len as usize; - let Some(frame_end) = 4usize.checked_add(frame_len) else { - return Err(js_error("invalid frame length")); - }; - if frame_end > buffer.len() { - break; - } - let frame = buffer[4..frame_end].to_vec(); - let arr = js_sys::Uint8Array::from(&frame[..]); - let _ = on_message.call1(&JsValue::NULL, &arr); - buffer.drain(..frame_end); - } - - Ok(()) } pub fn close(&self) {