| author | |
| committer | |
| log | 8b4ea55788421f43ea8cdaa7fe335fc8b444c340 |
| tree | 953e0e7e0c951603a9fb6e2120eaad4e9138ecbc |
| parent | 9b214d46d7ee55e79d02af499c16cfdba098266f |
| signature | Signed by SSH key SHA256:52mNGHRsVFBDED9IAX5pe+LRWUefqTbxEReunq21QvU |
Retire rejected or partial sends before any later encrypted frame can reuse the socket, invalidate retained peer lines, and exercise the JavaScript boundary in the local gate.
fixes #93
Assisted-by: gpt-6.1-sol5 files changed, 138 insertions(+), 13 deletions(-)
crates/notebook/src/live/web.js+9-3| ... | ... | @@ -7,6 +7,7 @@ export function liveConnect(url, message, closed) { |
| 7 | 7 | socket.binaryType = 'arraybuffer'; |
| 8 | 8 | socket.onmessage = event => message(event.data); |
| 9 | 9 | socket.onclose = () => { |
| 10 | if (!sockets.has(id)) return; | |
| 10 | 11 | liveClose(id); |
| 11 | 12 | closed(); |
| 12 | 13 | }; |
| ... | ... | @@ -16,9 +17,14 @@ export function liveConnect(url, message, closed) { |
| 16 | 17 | |
| 17 | 18 | export function liveSend(id, bytes) { |
| 18 | 19 | const socket = sockets.get(id); |
| 19 | if (!socket || socket.readyState !== WebSocket.OPEN) throw new Error('Relay disconnected'); | |
| 20 | if (socket.bufferedAmount > 1048576) throw new Error('Relay cannot keep up'); | |
| 21 | socket.send(bytes); | |
| 20 | try { | |
| 21 | if (!socket || socket.readyState !== WebSocket.OPEN) throw new Error('Relay disconnected'); | |
| 22 | if (socket.bufferedAmount > 1048576) throw new Error('Relay cannot keep up'); | |
| 23 | socket.send(bytes); | |
| 24 | } catch (error) { | |
| 25 | socket?.onclose?.(); | |
| 26 | throw error; | |
| 27 | } | |
| 22 | 28 | } |
| 23 | 29 | |
| 24 | 30 | export function liveClose(id) { |
crates/notebook/src/live/web.rs+34-10| ... | ... | @@ -80,12 +80,17 @@ impl Line { |
| 80 | 80 | pub fn send(&self, kind: u16, body: &impl Encode<()>) -> io::Result<()> { |
| 81 | 81 | let shared = self.shared.upgrade().ok_or(io::ErrorKind::NotConnected)?; |
| 82 | 82 | let mut bytes = Vec::new(); |
| 83 | self.sealer | |
| 83 | let sealed = self | |
| 84 | .sealer | |
| 84 | 85 | .lock() |
| 85 | 86 | .unwrap() |
| 86 | 87 | .as_mut() |
| 87 | 88 | .ok_or(io::ErrorKind::NotConnected)? |
| 88 | .send(&mut bytes, kind, body)?; | |
| 89 | .send(&mut bytes, kind, body); | |
| 90 | if let Err(error) = sealed { | |
| 91 | shared.disconnected(); | |
| 92 | return Err(error); | |
| 93 | } | |
| 89 | 94 | shared.stream(self.slot, &bytes) |
| 90 | 95 | } |
| 91 | 96 | |
| ... | ... | @@ -321,11 +326,18 @@ impl Shared { |
| 321 | 326 | fn disconnected(self: &Arc<Self>) { |
| 322 | 327 | let peers = { |
| 323 | 328 | let mut state = self.state.lock().unwrap(); |
| 324 | close(state.socket); | |
| 329 | let socket = std::mem::replace(&mut state.socket, 0); | |
| 330 | if socket == 0 { | |
| 331 | return; | |
| 332 | } | |
| 333 | close(socket); | |
| 325 | 334 | let peers: Vec<_> = state |
| 326 | 335 | .streams |
| 327 | 336 | .values() |
| 328 | .filter_map(|s| s.peer.as_ref().map(|p| p.hello.clone())) | |
| 337 | .filter_map(|s| { | |
| 338 | *s.line.sealer.lock().unwrap() = None; | |
| 339 | s.peer.as_ref().map(|p| p.hello.clone()) | |
| 340 | }) | |
| 329 | 341 | .collect(); |
| 330 | 342 | state.streams.clear(); |
| 331 | 343 | state.members.clear(); |
| ... | ... | @@ -358,7 +370,12 @@ impl Shared { |
| 358 | 370 | Ok(()) |
| 359 | 371 | } |
| 360 | 372 | |
| 361 | fn group(&self, kind: u16, body: &impl Encode<()>, to: Option<u32>) -> io::Result<()> { | |
| 373 | fn group( | |
| 374 | self: &Arc<Self>, | |
| 375 | kind: u16, | |
| 376 | body: &impl Encode<()>, | |
| 377 | to: Option<u32>, | |
| 378 | ) -> io::Result<()> { | |
| 362 | 379 | let body = minicbor::to_vec(body).map_err(io::Error::other)?; |
| 363 | 380 | let (socket, sealed) = { |
| 364 | 381 | let mut state = self.state.lock().unwrap(); |
| ... | ... | @@ -366,9 +383,16 @@ impl Shared { |
| 366 | 383 | .sealer |
| 367 | 384 | .as_mut() |
| 368 | 385 | .ok_or(io::ErrorKind::NotConnected)? |
| 369 | .seal(kind, &body, to.is_none())?; | |
| 386 | .seal(kind, &body, to.is_none()); | |
| 370 | 387 | (state.socket, sealed) |
| 371 | 388 | }; |
| 389 | let sealed = match sealed { | |
| 390 | Ok(sealed) => sealed, | |
| 391 | Err(error) => { | |
| 392 | self.disconnected(); | |
| 393 | return Err(error); | |
| 394 | } | |
| 395 | }; | |
| 372 | 396 | let mut bytes = match to { |
| 373 | 397 | Some(slot) => [&(::relay::GROUP | 1).to_be_bytes()[..], &slot.to_be_bytes()].concat(), |
| 374 | 398 | None => ::relay::BROADCAST.to_be_bytes().to_vec(), |
| ... | ... | @@ -434,10 +458,10 @@ impl Shared { |
| 434 | 458 | let hello = { |
| 435 | 459 | let mut state = self.state.lock().unwrap(); |
| 436 | 460 | state.members.remove(&slot); |
| 437 | state | |
| 438 | .streams | |
| 439 | .remove(&slot) | |
| 440 | .and_then(|s| s.peer.map(|p| p.hello)) | |
| 461 | state.streams.remove(&slot).and_then(|s| { | |
| 462 | *s.line.sealer.lock().unwrap() = None; | |
| 463 | s.peer.map(|p| p.hello) | |
| 464 | }) | |
| 441 | 465 | }; |
| 442 | 466 | if let Some(hello) = hello { |
| 443 | 467 | (self.events)(Event::Left(&hello)); |
tools/TESTING.md+1| ... | ... | @@ -20,6 +20,7 @@ another, and the exit status is the result. |
| 20 | 20 | | `windows-aarch64`, `linux-*` | Clippy `-D warnings` on what ships (libraries and binaries but `mobile`), then the `snowbound` build; Linux through `platform/linux/cargo.sh`, which links with zig against glibc 2.17 | |
| 21 | 21 | | `ios` | `xcodebuild` of the simulator app, unsigned | |
| 22 | 22 | | `web` | Clippy `-D warnings` on `snowbound` for `wasm32-unknown-unknown`, SQLite built by nixpkgs' clang; `release_web.py` links, optimizes and deploys the static folder | |
| 23 | | `web-js` | Node boundary tests for Live Share browser sockets | | |
| 23 | 24 | | `macos-10.6` | `platform/snow-leopard/cargo.sh` build of `snowbound`; skipped, saying why, without the SDK or nightly `rust-src` | |
| 24 | 25 | |
| 25 | 26 | ```sh |
tools/ci.py+4| ... | ... | @@ -78,6 +78,10 @@ def lanes(): |
| 78 | 78 | commands=[cross(linux, arch, *cross_clippy), cross(linux, arch, 'build', '-p', 'snowbound')], |
| 79 | 79 | environment={'CARGO_TARGET_DIR': str(TARGET / 'linux')}, |
| 80 | 80 | missing=None if shutil.which('zig') else 'needs zig')) |
| 81 | result.append(dict( | |
| 82 | name='web-js', minutes=5, packages=['notebook'], paths=('tools/web/',), | |
| 83 | commands=[['node', '--test', 'tools/web/test_live.mjs']], | |
| 84 | missing=None if shutil.which('node') else 'needs node')) | |
| 81 | 85 | targets = subprocess.run(['rustup', 'target', 'list', '--installed'], capture_output=True, text=True).stdout |
| 82 | 86 | wasm = web_environment() |
| 83 | 87 | result.append(dict( |
tools/web/test_live.mjs created+90| ... | ... | @@ -0,0 +1,90 @@ |
| 1 | import assert from 'node:assert/strict'; | |
| 2 | import { readFile } from 'node:fs/promises'; | |
| 3 | import test from 'node:test'; | |
| 4 | ||
| 5 | const opened = []; | |
| 6 | class WebSocket { | |
| 7 | static OPEN = 1; | |
| 8 | constructor() { | |
| 9 | this.readyState = WebSocket.OPEN; | |
| 10 | this.bufferedAmount = 0; | |
| 11 | this.sent = []; | |
| 12 | this.closes = 0; | |
| 13 | opened.push(this); | |
| 14 | } | |
| 15 | send(bytes) { | |
| 16 | if (this.failure) throw this.failure; | |
| 17 | this.sent.push(bytes.slice()); | |
| 18 | this.bufferedAmount += bytes.length; | |
| 19 | } | |
| 20 | close() { | |
| 21 | this.closes++; | |
| 22 | this.readyState = 3; | |
| 23 | } | |
| 24 | } | |
| 25 | globalThis.WebSocket = WebSocket; | |
| 26 | const source = await readFile(new URL('../../crates/notebook/src/live/web.js', import.meta.url), 'utf8'); | |
| 27 | const { liveConnect, liveSend, liveClose } = await import( | |
| 28 | `data:text/javascript;base64,${Buffer.from(source).toString('base64')}` | |
| 29 | ); | |
| 30 | ||
| 31 | for (const accepted of [0, 1]) { | |
| 32 | test(`backpressure after ${accepted} chunks retires the socket before notifying`, () => { | |
| 33 | let disconnected = 0; | |
| 34 | const id = liveConnect('ws://disposable.invalid', () => {}, () => { | |
| 35 | disconnected++; | |
| 36 | assert.throws(() => liveSend(id, new Uint8Array(1)), /Relay disconnected/); | |
| 37 | }); | |
| 38 | const socket = opened.at(-1); | |
| 39 | const onclose = socket.onclose; | |
| 40 | const chunk = new Uint8Array(64 << 10); | |
| 41 | socket.bufferedAmount = accepted ? 1048576 - chunk.length + 1 : 1048577; | |
| 42 | if (accepted) liveSend(id, chunk); | |
| 43 | assert.throws(() => liveSend(id, chunk), /Relay cannot keep up/); | |
| 44 | assert.equal(socket.sent.length, accepted); | |
| 45 | assert.equal(socket.closes, 1); | |
| 46 | assert.equal(disconnected, 1); | |
| 47 | assert.equal(socket.onmessage, null); | |
| 48 | assert.equal(socket.onclose, null); | |
| 49 | onclose(); | |
| 50 | liveClose(id); | |
| 51 | assert.equal(socket.closes, 1); | |
| 52 | assert.equal(disconnected, 1); | |
| 53 | ||
| 54 | const replacement = liveConnect('ws://disposable.invalid', () => {}, () => disconnected++); | |
| 55 | liveSend(replacement, chunk); | |
| 56 | assert.throws(() => liveSend(id, chunk), /Relay disconnected/); | |
| 57 | onclose(); | |
| 58 | assert.equal(opened.at(-1).sent.length, 1); | |
| 59 | assert.equal(opened.at(-1).closes, 0); | |
| 60 | assert.equal(disconnected, 1); | |
| 61 | liveClose(replacement); | |
| 62 | }); | |
| 63 | } | |
| 64 | ||
| 65 | test('a WebSocket send exception retires the socket and preserves the error', () => { | |
| 66 | let disconnected = 0; | |
| 67 | const id = liveConnect('ws://disposable.invalid', () => {}, () => disconnected++); | |
| 68 | const socket = opened.at(-1); | |
| 69 | socket.failure = new Error('Transport failed'); | |
| 70 | assert.throws(() => liveSend(id, new Uint8Array(1)), error => error === socket.failure); | |
| 71 | assert.equal(socket.sent.length, 0); | |
| 72 | assert.equal(socket.closes, 1); | |
| 73 | assert.equal(disconnected, 1); | |
| 74 | }); | |
| 75 | ||
| 76 | test('a non-open socket is retired, while deliberate close never reconnects', () => { | |
| 77 | let disconnected = 0; | |
| 78 | const id = liveConnect('ws://disposable.invalid', () => {}, () => disconnected++); | |
| 79 | const socket = opened.at(-1); | |
| 80 | socket.readyState = 2; | |
| 81 | assert.throws(() => liveSend(id, new Uint8Array(1)), /Relay disconnected/); | |
| 82 | assert.equal(socket.closes, 1); | |
| 83 | assert.equal(disconnected, 1); | |
| 84 | const next = liveConnect('ws://disposable.invalid', () => {}, () => disconnected++); | |
| 85 | const onclose = opened.at(-1).onclose; | |
| 86 | liveClose(next); | |
| 87 | onclose(); | |
| 88 | assert.equal(opened.at(-1).closes, 1); | |
| 89 | assert.equal(disconnected, 1); | |
| 90 | }); |