From a31dfda0e9162d43429aa87a8f7f99cc2273a2cc Mon Sep 17 00:00:00 2001 From: clover caruso Date: Mon, 5 Oct 2026 21:44:54 -0700 Subject: [PATCH] fix: keep large Live Share section updates within relay limits Bound differential section data by the encoded chunk budget; larger updates retain the existing notification and chunked read path. The regression reproduces a relay disconnect on the previous writer. fixes #94 Assisted-by: gpt-6.1-sol --- crates/notebook/src/live/share.rs | 20 ++++++- crates/notebook/tests/live_share.rs | 90 +++++++++++++++++++++++++++++ 2 files changed, 108 insertions(+), 2 deletions(-) diff --git a/crates/notebook/src/live/share.rs b/crates/notebook/src/live/share.rs index ca8a440d2aacb6f042c434d90c0384914ce817d0..7b966c3e5e9f0fef0a429111ab28d411a0dba72c 100644 --- a/crates/notebook/src/live/share.rs +++ b/crates/notebook/src/live/share.rs @@ -822,7 +822,7 @@ fn digest(stamp: &Stamp) -> u64 { } /// The writes that make `after` of `before`, a commit's appended bytes, patches and header, -/// where they are much less than `after` itself. +/// where their encoding fits a chunk and they are much less than `after` itself. fn delta(before: &[u8], after: &[u8]) -> Option> { const BLOCK: usize = 4096; if before.len() < 1024 || after.len() < before.len() { @@ -858,7 +858,8 @@ fn delta(before: &[u8], after: &[u8]) -> Option> { }); } let sent: usize = writes.iter().map(|write| write.bytes.len()).sum(); - (sent <= after.len() / 2).then_some(writes) + (sent <= after.len() / 2 && sent <= CHUNK && minicbor::to_vec(&writes).ok()?.len() <= CHUNK) + .then_some(writes) } /// `image` with `writes`, `length` long: none where a write falls outside it. @@ -2060,4 +2061,19 @@ mod tests { }]; assert!(written(&before, after.len() as u64, &outside).is_none()); } + + #[test] + fn differential_reads_include_metadata_in_the_chunk_budget() { + let before = vec![0; CHUNK * 4]; + let mut after = before.clone(); + after.resize(before.len() + CHUNK - 1024, 1); + assert!(delta(&before, &after).is_none()); + after.truncate(before.len() + CHUNK / 2); + let writes = delta(&before, &after).unwrap(); + assert!(minicbor::to_vec(&writes).unwrap().len() <= CHUNK); + assert_eq!( + written(&before, after.len() as u64, &writes).unwrap(), + after + ); + } } diff --git a/crates/notebook/tests/live_share.rs b/crates/notebook/tests/live_share.rs index 6488732a92e816e0cb10ab96f4eac96a7d944dfd..885534822a2b6c9ed79d1bbd1c32915450098769 100644 --- a/crates/notebook/tests/live_share.rs +++ b/crates/notebook/tests/live_share.rs @@ -305,6 +305,96 @@ fn large_files_travel_in_chunks() { ); } +#[test] +fn large_section_appends_keep_the_room_connected_and_guests_converge() { + use notebook::{EditStatus, Replica, live::share::HostedRemote}; + use onestore::op::{Edit, Op, PageOp}; + use std::time::{Duration, Instant}; + + let directory = tempfile::tempdir().unwrap(); + let folder = notebook(directory.path()); + let file = folder.join("Garden.one"); + let mut original = "A".repeat(512 << 10); + let image = onestore::create_section("Garden.one", &original, "Fixture").unwrap(); + std::fs::write(&file, &image).unwrap(); + let (space, text, _) = server::text(&image); + let arena = Arena::default(); + let mut section = onestore::Section::open(&arena, image).unwrap(); + section + .apply( + "Fixture", + &Edit { + at: 134_000_000_000_000_000, + ops: vec![Op::Page { + space, + op: PageOp::Text { + text, + range: 0..1, + with: "a".into(), + }, + }], + }, + ) + .unwrap(); + section.seal().unwrap().unwrap().commit_file(&file).unwrap(); + original.replace_range(0..1, "a"); + let url = relay(Default::default()); + let sharing = Sharing::new("").unwrap(); + let host = host(&folder, &directory.path().join("host"), &sharing, &url); + let code = code(&host); + let (alice, alice_notebook) = guest("Alice", &code, &url, &directory.path().join("alice")); + let (bob, bob_notebook) = guest("Bob", &code, &url, &directory.path().join("bob")); + until("the guests never joined the room", || { + host.guests().len() == 2 + }); + let image = bob_notebook.read_section("Garden.one").unwrap(); + assert_eq!(alice_notebook.read_section("Garden.one").unwrap(), image); + let (space, text, _) = server::text(&image); + let writer = Replica::create(directory.path().join("writer.sqlite"), &image).unwrap(); + let reader = Replica::create(directory.path().join("reader.sqlite"), &image).unwrap(); + let append = "B".repeat(192 << 10); + let expected = format!("{original}{append}"); + let id = writer + .apply( + "Alice", + Edit { + at: 134_000_000_000_000_000, + ops: vec![Op::Page { + space, + op: PageOp::Text { + text, + range: original.len() as u32..original.len() as u32, + with: append, + }, + }], + }, + ) + .unwrap(); + writer + .sync_once(&mut HostedRemote::new(&alice, "Garden.one")) + .unwrap(); + let stored = std::fs::read(&file).unwrap(); + let added = stored.len() - image.len(); + assert!(added > relay::server::Config::default().max_message); + assert!(added + 1024 < stored.len() / 2); + assert_eq!(server::text(&stored).2, expected); + let Some(EditStatus::Published { revision }) = writer.status(id).unwrap() else { + panic!("the append has no durable receipt"); + }; + let store = onestore::Store::parse(&stored).unwrap(); + let index = onestore::RevisionIndex::parse(&store).unwrap(); + assert!(index.spaces[&space].revisions.contains_key(&revision)); + assert!(writer.pending().unwrap().is_empty()); + let mut remote = HostedRemote::new(&bob, "Garden.one"); + reader.sync_once(&mut remote).unwrap(); + assert_eq!(reader.snapshot().unwrap(), stored); + let deadline = Instant::now() + Duration::from_millis(1200); + while Instant::now() < deadline { + assert_eq!(host.guests().len(), 2, "the section append broke the room"); + std::thread::sleep(Duration::from_millis(10)); + } +} + /// A guest that floods its host with requests is hung up on once too many wait, having had /// answers to few of them, and the host goes on serving the others. #[test] -- 2.54.0