From 02acde3ba6bf8683a349deb02c298332a7c4f4ed Mon Sep 17 00:00:00 2001 From: clover caruso Date: Fri, 2 Oct 2026 21:55:05 -0700 Subject: [PATCH] feat: Live Share speaks protocol 2, and each keystroke reaches the others in tens of milliseconds A keystroke took a second or more to reach another peer, and in the app up to twenty seconds to reach the host's own window, which waited on its folder's watch. Now: - the host tells its own sections what guests change at once, and tells guests when its own open section publishes, rather than waiting on the watch; - a guest's background settles a report for 20 ms, not a second, since its host sends one a commit; - the host sends what each commit changed (Delta: the header, patches and appended bytes), and a guest holding the section takes it on without reading the section again; its own commits apply to what it holds too. Measured with examples/live_latency (20 keys, keystroke stored in one guest's replica to the change in another's, and host's file to a guest), median: before after loopback LAN 1076 ms 57 ms (host to guest 1026 -> 34 ms) relay on loopback 1071 ms 57 ms (host to guest 1027 -> 33 ms) relay.snowbound... - 210 ms (host to guest 137 ms) Protocol 2 opens with version 2 and answers a peer of another version before failing, so each end names the other's version: joining says which side must update Snowbound. Assisted-by: claude-opus-5.5 --- crates/notebook/Cargo.toml | 4 + crates/notebook/examples/live_latency.rs | 206 ++++++++++++++++++ crates/notebook/src/background.rs | 24 ++- crates/notebook/src/live.rs | 22 +- crates/notebook/src/live/share.rs | 258 +++++++++++++++++++++-- crates/notebook/src/live/tests.rs | 44 ++++ crates/notebook/src/live/wire.rs | 60 +++++- crates/snowbound/src/live.rs | 9 +- crates/snowbound/src/main.rs | 10 + crates/snowbound/src/share.rs | 8 + 10 files changed, 617 insertions(+), 28 deletions(-) create mode 100644 crates/notebook/examples/live_latency.rs diff --git a/crates/notebook/Cargo.toml b/crates/notebook/Cargo.toml index 1c12191aad2634504ded7d482ed21a8358cf3fe8..63ef60a46c7b52d4346aa3058b0923db2c9cb6c7 100644 --- a/crates/notebook/Cargo.toml +++ b/crates/notebook/Cargo.toml @@ -59,6 +59,10 @@ libc = "0.2" serde_json = "1" tempfile = "3" +[[example]] +name = "live_latency" +required-features = ["live"] + [[example]] name = "smb_offline_client" required-features = ["smb"] diff --git a/crates/notebook/examples/live_latency.rs b/crates/notebook/examples/live_latency.rs new file mode 100644 index 0000000000000000000000000000000000000000..74bb8304fc2715ea34fea9c83acc8a984196dfb1 --- /dev/null +++ b/crates/notebook/examples/live_latency.rs @@ -0,0 +1,206 @@ +//! Measures Live Share's typing latency: how long after a keystroke lands in one guest's +//! replica the other guest's section shows it, and the same from the host to a guest. +//! +//! `cargo run -p notebook --features live --example live_latency -- lan|relay [URL]`: `lan` +//! meets by mDNS on this computer's loopback; `relay` through `URL`, or a relay this run +//! starts on loopback. + +#[path = "../tests/support/server.rs"] +mod server; + +use notebook::{ + Replica, + live::{ + Hello, Reach, + share::{self, Guest, Host, Sharing}, + }, + session::{Background, Notebook, Section}, +}; +use onestore::op::{Edit, Op, PageOp}; +use std::{ + path::Path, + sync::Arc, + thread, + time::{Duration, Instant}, +}; + +const KEYS: usize = 20; + +fn until(what: &str, done: impl Fn() -> bool) { + let deadline = Instant::now() + Duration::from_secs(60); + while !done() { + assert!(Instant::now() < deadline, "{what}"); + thread::sleep(Duration::from_millis(1)); + } +} + +fn hello(name: &str) -> Hello { + Hello::new(name.into(), None).unwrap() +} + +/// A guest as the app holds one: the notebook, its background, and the section open. +struct Open { + _background: Background, + section: Section, +} + +fn guest( + name: &str, + code: &str, + reach: Option, + relay: Option<&str>, + cache: &Path, +) -> (Arc, Open) { + let welcome = share::join(hello(name), code, "", reach, relay).unwrap(); + let guest = Guest::start( + hello(name), + welcome.share, + welcome.secret, + reach, + relay, + || {}, + ) + .unwrap(); + until("no host", || guest.host().is_some()); + let mut notebook = Notebook::open_hosted(Arc::clone(&guest), cache).unwrap(); + let background = Background::hosted(Arc::clone(&guest), || {}).unwrap(); + background.watch(notebook.replicas()); + let replica = notebook.replica_path("Garden.one").unwrap(); + std::fs::create_dir_all(replica.parent().unwrap()).unwrap(); + let replica = + Replica::open_or_create(&replica, None, || notebook.read_section("Garden.one")).unwrap(); + let section = + Section::resume_hosted("Garden.one".into(), replica, Arc::clone(&guest), || {}).unwrap(); + background.hold("Garden.one", §ion); + ( + guest, + Open { + _background: background, + section, + }, + ) +} + +/// Types `key` at the end of the page's text in `section`; returns the text it then reads. +fn type_key(section: &Section, key: char) -> String { + let (space, text, before) = text_of(section); + let at = before.encode_utf16().count() as u32; + let op = PageOp::Text { + text, + range: at..at, + with: key.to_string(), + }; + let edit = Edit { + at: 134_000_000_000_000_000, + ops: vec![Op::Page { space, op }], + }; + section.replica().apply("Typist", edit).unwrap(); + format!("{before}{key}") +} + +fn text_of(section: &Section) -> (onestore::ExGuid, onestore::ExGuid, String) { + let (space, ..) = section.pages().unwrap()[0]; + let page = section.page(space).unwrap(); + page.objects + .iter() + .find_map(|object| match object { + onestore::page::PageObject::Outline(outline) => outline + .paragraphs + .iter() + .find_map(|p| p.text().map(|t| (space, t.id, t.text.text().to_owned()))), + _ => None, + }) + .unwrap() +} + +fn report(what: &str, mut times: Vec) { + times.sort(); + let ms = |at: usize| times[at].as_secs_f64() * 1000.0; + println!( + "{what}: median {:.0} ms, p90 {:.0} ms, worst {:.0} ms over {} keys", + ms(times.len() / 2), + ms(times.len() * 9 / 10), + ms(times.len() - 1), + times.len() + ); +} + +fn main() { + let args: Vec = std::env::args().skip(1).collect(); + let (reach, url) = match args.first().map(String::as_str) { + Some("lan") => (Some(Reach::Loopback), None), + Some("relay") => ( + None, + Some(args.get(1).cloned().unwrap_or_else(|| { + let listener = std::net::TcpListener::bind("127.0.0.1:0").unwrap(); + let url = format!("ws://{}", listener.local_addr().unwrap()); + thread::spawn(move || relay::server::serve(listener, Default::default())); + url + })), + ), + _ => panic!("lan or relay [URL]"), + }; + let relay = url.as_deref(); + let directory = tempfile::tempdir().unwrap(); + let folder = directory.path().join("Garden"); + std::fs::create_dir_all(&folder).unwrap(); + std::fs::write( + folder.join("Garden.one"), + onestore::create_section("Garden.one", "Typed:", "Fixture").unwrap(), + ) + .unwrap(); + let storage = Notebook::open(&folder, directory.path().join("host")) + .unwrap() + .into_storage(); + let host = Host::start( + storage, + hello("Ada"), + Sharing::new("").unwrap(), + "Garden", + reach, + relay, + || {}, + ) + .unwrap(); + until("no code", || { + host.code().is_some_and(|code| share::code(&code).is_some()) + }); + let code = host.code().unwrap(); + let (_grace, grace) = guest("Grace", &code, reach, relay, &directory.path().join("g")); + let (_alan, alan) = guest("Alan", &code, reach, relay, &directory.path().join("a")); + + let mut times = Vec::new(); + for key in "abcdefghijklmnopqrstuvwxyz".chars().take(KEYS) { + let typed = Instant::now(); + let expected = type_key(&grace.section, key); + until("the other guest never saw a key", || { + alan.section.events(); + text_of(&alan.section).2 == expected + }); + times.push(typed.elapsed()); + thread::sleep(Duration::from_millis(150)); + } + report("guest to guest", times); + + // The host's own keystrokes, published to its file as its app does, then reported to + // the share as its app reports them. + let file = folder.join("Garden.one"); + let mut times = Vec::new(); + for key in "ABCDEFGHIJKLMNOPQRSTUVWXYZ".chars().take(KEYS) { + let image = std::fs::read(&file).unwrap(); + let (space, text, before) = server::text(&image); + let at = before.encode_utf16().count() as u32; + let typed = Instant::now(); + let next = server::typed(&image, space, text, at..at, &key.to_string()); + std::fs::write(&file, &next).unwrap(); + host.touched(&["Garden.one".into()]); + let expected = format!("{before}{key}"); + until("the guest never saw the host's key", || { + alan.section.events(); + text_of(&alan.section).2 == expected + }); + times.push(typed.elapsed()); + thread::sleep(Duration::from_millis(150)); + } + report("host to guest", times); +} diff --git a/crates/notebook/src/background.rs b/crates/notebook/src/background.rs index 5e1d130905e8120176e24a4b037986e5bbc94210..3f5a6a0cc1e94d296931553f5cc2fa392466cce0 100644 --- a/crates/notebook/src/background.rs +++ b/crates/notebook/src/background.rs @@ -91,6 +91,8 @@ struct Watched { lost: bool, /// The current connection's watch reports the folder's changes. reported: bool, + /// How long a report settles, where not `SETTLE` (`Background::set_settle`). + settle: Option, } struct Watch { @@ -356,7 +358,7 @@ impl Background { guest: Arc, notify: impl Fn() + Send + 'static, ) -> Result { - Self::start( + let background = Self::start( true, move |reports| { guest.watch(reports)?; @@ -366,7 +368,17 @@ impl Background { Ok(((bind, list), true)) }, notify, - ) + )?; + background.set_settle(crate::live::share::SETTLE); + Ok(background) + } + + /// Checks what a report names after `settle` rather than a second, where reports come + /// one to a commit, as a Live Share host sends them. + pub fn set_settle(&self, settle: Duration) { + if let Ok(mut watched) = self.0.watched.lock() { + watched.settle = Some(settle); + } } /// Has `listener` hear every report of changed paths from now on, as a Live Share host @@ -638,6 +650,7 @@ impl Watched { } fn touched(&mut self, paths: &[String], now: Instant) { + let settled = now + self.settle.unwrap_or(SETTLE); for path in paths { let path = path.trim_matches('/'); let section = self @@ -647,16 +660,13 @@ impl Watched { match section { Some(index) => { let watch = &mut self.sections[index]; - watch.due = watch.due.min(now + SETTLE); + watch.due = watch.due.min(settled); watch.current = false; watch.reported = true; } None if self.sections.iter().any(|watch| within(&watch.path, path)) => { self.folders.insert(path.to_owned()); - self.relist = Some( - self.relist - .map_or(now + SETTLE, |due| due.min(now + SETTLE)), - ); + self.relist = Some(self.relist.map_or(settled, |due| due.min(settled))); } None => {} } diff --git a/crates/notebook/src/live.rs b/crates/notebook/src/live.rs index 42858132ae28436450f7ef5611b7e81fb9252709..d1bc29fe1e3360c16e32c4d6f4ae4d16b60f9ad7 100644 --- a/crates/notebook/src/live.rs +++ b/crates/notebook/src/live.rs @@ -192,6 +192,8 @@ struct State { failed: u32, /// The code admits no one new. burned: bool, + /// The Live Share version of the last peer met that speaks another. + outdated: Option, daemon: Option, } @@ -379,6 +381,12 @@ impl Live { self.shared.state.lock().unwrap().failed } + /// The Live Share version of the last peer met that speaks another, which one of the two + /// must update to meet. + pub fn other_version(&self) -> Option { + self.shared.state.lock().unwrap().outdated + } + /// Whether this end's code had too many wrong tries and admits no one new. pub fn burned(&self) -> bool { self.shared.state.lock().unwrap().burned @@ -535,12 +543,22 @@ impl Shared { stream.set_read_timeout(GONE)?; Ok((send, receive, hello)) })(); - pipe.met(met.is_ok()); + let other = (met.as_ref().err()) + .and_then(|error| error.get_ref()?.downcast_ref::()) + .copied(); + // A peer of another version guessed nothing, the keys never being agreed. + pipe.met(met.is_ok() || other.is_some()); let (send, mut receive, hello) = match met { Ok(met) => met, Err(error) => { eprintln!("Live: no meeting in {tag}: {error}"); - if error.kind() == io::ErrorKind::InvalidData { + if let Some(wire::Version(version)) = other { + self.state.lock().unwrap().outdated = Some(version); + if matches!(self.room, Room::Code { owner: false, .. }) { + self.stopped.store(true, Ordering::Release); + } + (self.events)(Event::Changed); + } else if error.kind() == io::ErrorKind::InvalidData { self.failed(); } pipe.shutdown(); diff --git a/crates/notebook/src/live/share.rs b/crates/notebook/src/live/share.rs index 00594c4e3920e36f255794877a6ab5219c2f5bb0..84d0fe1a5f612efe3a8d23383988ab25e9857bc4 100644 --- a/crates/notebook/src/live/share.rs +++ b/crates/notebook/src/live/share.rs @@ -9,7 +9,9 @@ use super::{ Event, Hello, Line, Live, Peer, Presence, Reach, Relayed, Room, - wire::{self, Failure, Reply, Request, Touched, Welcome, WireEntry, WireStamp, kind}, + wire::{ + self, Delta, Failure, Reply, Request, Touched, Welcome, WireEntry, WireStamp, Written, kind, + }, }; use crate::{Error, Result, background::Reports, discover, session::Storage}; use onestore::{CommitError, CommitState, RevisionIndex, Stamp, Store, Transaction}; @@ -37,8 +39,12 @@ const LIMIT: usize = 256 << 20; /// Snapshots of files being read, per guest, and how long one is kept unread. const SNAPSHOTS: usize = 8; const SNAPSHOT_AGE: Duration = Duration::from_secs(120); -/// Sections whose image a host keeps to check guests' commits on. +/// Sections whose image a host keeps to check guests' commits on and tell them what changed, +/// and a guest keeps to take those changes on without reading. const IMAGES: usize = 4; +/// How long a report of a change settles in a guest's background: its host reports each +/// commit once, at once. +pub const SETTLE: Duration = Duration::from_millis(20); /// What a host says leaving as it stops sharing. const STOPPED: &str = "stopped"; /// What a host says hanging up on a guest that asks too much too fast. @@ -95,6 +101,9 @@ pub fn location(share: &[u8; 16]) -> String { pub enum Refusal { /// Not a code, or one mistyped, which spends none of the relay's tries. Malformed, + /// The person sharing runs a Snowbound of another Live Share version: the newer one's + /// `true` where it is theirs, so this one should update. + Version { theirs_newer: bool }, /// The code's secret or password is wrong. Wrong, /// No one shares with the code's number now. @@ -149,6 +158,11 @@ pub fn join( if let Ok(welcome) = welcome.try_recv() { return Ok(welcome); } + if let Some(version) = live.other_version() { + return Err(Refusal::Version { + theirs_newer: version > wire::VERSION, + }); + } if live.failed() > 0 { return Err(Refusal::Wrong); } @@ -209,6 +223,7 @@ impl Host { snapshots: Mutex::default(), puts: Mutex::default(), guests: Mutex::default(), + host: Mutex::default(), }); let serving = Hello { serves: Some(sharing.share), @@ -296,8 +311,18 @@ impl Host { } } - /// Tells every guest the files at these catalog paths changed. + /// Has `listener` hear the catalog paths guests change from now on, sooner than a watch + /// on the notebook's folder would. + pub fn on_changed(&self, listener: crate::session::Listener) { + *self.served.host.lock().unwrap() = Some(listener); + } + + /// Tells every guest the files at these catalog paths changed, with what changed in the + /// sections a guest read lately. pub fn touched(&self, paths: &[String]) { + for path in paths { + self.served.changed_here(path); + } self.served.tell(paths); } @@ -355,6 +380,8 @@ struct Served { /// Bytes a later request carries, by guest and upload. puts: Mutex>>, guests: Mutex>, + /// Hears the paths guests changed, as the host's own notebook should. + host: Mutex>, } /// A guest as its host serves it: the line to it, its requests waiting for its workers, and @@ -451,6 +478,14 @@ impl Served { .retain(|(guest, _), _| guest != peer); } + /// Tells every guest, and the host, that a guest changed the files at `paths`. + fn changed(&self, paths: &[String]) { + self.tell(paths); + if let Some(listener) = &*self.host.lock().unwrap() { + listener(paths); + } + } + /// Tells every guest the files at `paths` changed. fn tell(&self, paths: &[String]) { let touched = Touched { @@ -540,7 +575,7 @@ impl Served { kind::COMMIT => { let transaction = Transaction::from_bytes(&self.carried(peer, &request)?)?; self.commit(path, &transaction)?; - self.tell(&[path.to_owned()]); + self.changed(&[path.to_owned()]); done } kind::CONFIRM => { @@ -549,12 +584,12 @@ impl Served { } kind::CREATE => { self.storage.create(path, &self.carried(peer, &request)?)?; - self.tell(&[folder(path)]); + self.changed(&[folder(path)]); done } kind::CREATE_DIRECTORY => { self.storage.create_directory(path)?; - self.tell(&[folder(path)]); + self.changed(&[folder(path)]); done } kind::HIDE => { @@ -568,12 +603,12 @@ impl Served { } else { self.storage.replace(path, to)?; } - self.tell(&[folder(path), folder(to)]); + self.changed(&[folder(path), folder(to)]); done } kind::DELETE => { self.storage.delete(path)?; - self.tell(&[folder(path)]); + self.changed(&[folder(path)]); done } kind::PLACE => { @@ -582,12 +617,12 @@ impl Served { .ok_or_else(|| refused(io::ErrorKind::InvalidInput, "No ancestor"))?; let name = request.name.as_deref().unwrap_or_default(); self.storage.place(path, ancestor, name)?; - self.tell(&[path.to_owned()]); + self.changed(&[path.to_owned()]); done } kind::SUPERSEDE => { self.storage.supersede(path, &stamp()?, to()?)?; - self.tell(&[folder(path)]); + self.changed(&[folder(path)]); done } _ => return Err(refused(io::ErrorKind::Unsupported, "An unknown request")), @@ -716,9 +751,103 @@ impl Served { ))); } self.storage.commit(path, transaction)?; + self.tell_delta(path, &image, &next); self.keep(path, Arc::new(next)); Ok(()) } + + /// Tells every guest what a commit changed in the section at `path`, from `before` to + /// `after`, where that is small enough to send. + fn tell_delta(&self, path: &str, before: &[u8], after: &[u8]) { + let (Ok(base), Some(writes)) = (Stamp::of(before), delta(before, after)) else { + return; + }; + let delta = Delta { + path: path.to_owned(), + base: (&base).into(), + length: after.len() as u64, + writes, + }; + for guest in self.guests.lock().unwrap().values() { + let _ = guest.line.send(kind::DELTA, &delta); + } + } + + /// The host's own change to the section at `path`, which guests hear as a delta where + /// the host kept the image before it. + fn changed_here(&self, path: &str) { + let kept = self + .images + .lock() + .unwrap() + .iter() + .find_map(|(held, at, image)| (held == path).then(|| (at.clone(), Arc::clone(image)))); + let Some((at, before)) = kept else { + return; + }; + if self.storage.stamp(path).is_ok_and(|now| now == at) { + return; + } + if let Ok(after) = self.storage.read(path) { + self.tell_delta(path, &before, &after); + self.keep(path, Arc::new(after)); + } + } +} + +/// The writes that make `after` of `before`, a commit's appended bytes, patches and header, +/// where 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() { + return None; + } + let mut writes = vec![Written { + offset: 0, + bytes: after[..1024].to_vec(), + }]; + let mut at = 1024; + while at < before.len() { + let end = (at + BLOCK).min(before.len()); + if before[at..end] != after[at..end] { + let first = (at..end).find(|&i| before[i] != after[i]).unwrap_or(at); + let last = (at..end).rfind(|&i| before[i] != after[i]).unwrap_or(first) + 1; + match writes.last_mut() { + Some(write) if write.offset as usize + write.bytes.len() + 64 >= first => { + let from = write.offset as usize; + write.bytes = after[from..last].to_vec(); + } + _ => writes.push(Written { + offset: first as u64, + bytes: after[first..last].to_vec(), + }), + } + } + at = end; + } + if after.len() > before.len() { + writes.push(Written { + offset: before.len() as u64, + bytes: after[before.len()..].to_vec(), + }); + } + let sent: usize = writes.iter().map(|write| write.bytes.len()).sum(); + (sent <= after.len() / 2).then_some(writes) +} + +/// `image` with `writes`, `length` long: none where a write falls outside it. +fn written(image: &[u8], length: u64, writes: &[Written]) -> Option> { + let length = usize::try_from(length) + .ok() + .filter(|length| *length >= image.len())?; + let mut next = image.to_vec(); + next.resize(length, 0); + for write in writes { + let offset = usize::try_from(write.offset).ok()?; + next.get_mut(offset..offset.checked_add(write.bytes.len())?)? + .copy_from_slice(&write.bytes); + } + Some(next) } fn within(folder: &str, name: &str) -> String { @@ -846,6 +975,8 @@ struct Inner { stopped: AtomicBool, /// Bytes of chunks asked for and not yet given. asked: (Mutex, Condvar), + /// The sections read lately, kept as the host's deltas change them, newest first. + images: Mutex>)>>, } impl Guest { @@ -867,6 +998,7 @@ impl Guest { watch: Mutex::default(), stopped: AtomicBool::new(false), asked: Default::default(), + images: Mutex::default(), }); let heard = Arc::clone(&inner); let live = Live::start(me, &Room::Notebook(secret), reach, relay, move |event| { @@ -1026,9 +1158,19 @@ impl Guest { image.extend_from_slice(&chunk); } image.truncate(length); + if kind == kind::READ { + self.inner.hold(path, Arc::new(image.clone())); + } Ok(image) } + /// The image of the section at `path` with stamp `stamp`, where this guest holds it. + fn held(&self, path: &str, stamp: &Stamp) -> Option> { + let images = self.inner.images.lock().unwrap(); + let (_, image) = images.iter().find(|(held, _)| held == path)?; + (Stamp::of(image).ok().as_ref() == Some(stamp)).then(|| image.to_vec()) + } + pub(crate) fn entries(&self, folder: &str) -> io::Result> { let reply = self.ask( kind::LIST, @@ -1101,6 +1243,46 @@ impl Guest { } impl Inner { + fn hold(&self, path: &str, image: Arc>) { + let mut images = self.images.lock().unwrap(); + images.retain(|(held, _)| held != path); + images.insert(0, (path.to_owned(), image)); + images.truncate(IMAGES); + } + + /// Takes on a delta to an image held. + fn apply(&self, delta: &Delta) { + let held = self + .images + .lock() + .unwrap() + .iter() + .find_map(|(path, image)| (*path == delta.path).then(|| Arc::clone(image))); + let base: Option = (&delta.base).try_into().ok(); + if let Some(image) = held + && Stamp::of(&image).ok() == base + && let Some(next) = written(&image, delta.length, &delta.writes) + { + self.hold(&delta.path, Arc::new(next)); + } + } + + /// Takes on this guest's own commit to an image held. + fn published(&self, path: &str, transaction: &Transaction) { + let held = self + .images + .lock() + .unwrap() + .iter() + .find_map(|(held, image)| (held == path).then(|| Arc::clone(image))); + if let Some(image) = held { + let mut next = (*image).clone(); + if transaction.apply(&mut next).is_ok() { + self.hold(path, Arc::new(next)); + } + } + } + fn heard(&self, event: Event) { let serves = |hello: &Hello| hello.serves == Some(self.share); match event { @@ -1125,6 +1307,11 @@ impl Inner { let _ = waiting.send(reply); } } + kind::DELTA if serves(from) => { + if let Ok(delta) = minicbor::decode::(body) { + self.apply(&delta); + } + } kind::TOUCHED if serves(from) => { if let Ok(touched) = minicbor::decode::(body) && let Some(reports) = &*self.watch.lock().unwrap() @@ -1150,6 +1337,8 @@ impl Inner { pub struct HostedRemote { guest: Arc, path: String, + /// The stamp last asked for, which an image the guest holds may already have. + seen: Option, } impl HostedRemote { @@ -1157,21 +1346,30 @@ impl HostedRemote { Self { guest: Arc::clone(guest), path: path.to_owned(), + seen: None, } } } impl crate::Remote for HostedRemote { fn read(&mut self) -> io::Result> { + if let Some(image) = (self.seen.as_ref()).and_then(|seen| self.guest.held(&self.path, seen)) + { + return Ok(image); + } self.guest.read(kind::READ, &self.path, LIMIT) } fn stamp(&mut self) -> io::Result { - self.guest.stamp(&self.path) + let stamp = self.guest.stamp(&self.path)?; + self.seen = Some(stamp.clone()); + Ok(stamp) } fn publish(&mut self, transaction: &Transaction) -> std::result::Result<(), CommitError> { - self.guest.commit(&self.path, transaction) + self.guest.commit(&self.path, transaction)?; + self.guest.inner.published(&self.path, transaction); + Ok(()) } fn confirm(&mut self, base: &Stamp) -> std::result::Result<(), CommitError> { @@ -1373,3 +1571,39 @@ impl Storage for Hosted { self.guest.verb(kind::SUPERSEDE, path, request) } } + +#[cfg(test)] +mod tests { + use super::*; + + /// A commit's writes, found by comparing images, rebuild the image after it from the one + /// before, and a change touching most of the file is left to be read whole. + #[test] + fn deltas_rebuild_the_image_after_a_commit() { + let before: Vec = (0..20_000u32).map(|at| (at % 251) as u8).collect(); + let mut after = before.clone(); + after[3] ^= 1; + after[5000..5010].fill(9); + after[5050] ^= 1; + after[17_000] ^= 1; + after.extend_from_slice(&[7; 3000]); + let writes = delta(&before, &after).unwrap(); + assert_eq!( + writes.len(), + 4, + "the header, two runs merged as one, one more, the tail" + ); + assert_eq!( + written(&before, after.len() as u64, &writes).unwrap(), + after + ); + let rewritten: Vec = before.iter().map(|byte| byte ^ 1).collect(); + assert!(delta(&before, &rewritten).is_none()); + assert!(delta(&before, &before[..10_000]).is_none()); + let outside = [Written { + offset: after.len() as u64, + bytes: vec![1], + }]; + assert!(written(&before, after.len() as u64, &outside).is_none()); + } +} diff --git a/crates/notebook/src/live/tests.rs b/crates/notebook/src/live/tests.rs index 77c6ded383b650541a905743e8e0f18f51be10d5..0d3cbbad3a7fa250b1ffec6c6bbb1db18c0c1bf7 100644 --- a/crates/notebook/src/live/tests.rs +++ b/crates/notebook/src/live/tests.rs @@ -520,3 +520,47 @@ fn a_malicious_relay_is_caught() { }); } } + +/// Ends of two versions each learn the other's version from the opening, before any key is +/// agreed, so each can say which should update. +#[test] +fn another_version_is_named_before_any_key() { + let listener = TcpListener::bind("127.0.0.1:0").unwrap(); + let address = listener.local_addr().unwrap(); + let older = thread::spawn(move || { + let mut stream = TcpStream::connect(address).unwrap(); + #[derive(minicbor::Encode)] + #[cbor(map)] + struct Open { + #[n(0)] + version: u16, + #[n(1)] + room: String, + #[cbor(n(2), with = "minicbor::bytes")] + pake: Vec, + } + let open = minicbor::to_vec(Open { + version: 1, + room: "room".into(), + pake: vec![0; 33], + }) + .unwrap(); + stream + .write_all(&[&(open.len() as u32).to_be_bytes()[..], &open].concat()) + .unwrap(); + let mut length = [0; 4]; + stream.read_exact(&mut length).unwrap(); + let mut answer = vec![0; u32::from_be_bytes(length) as usize]; + stream.read_exact(&mut answer).unwrap(); + answer + }); + let (mut stream, _) = listener.accept().unwrap(); + let Err(error) = wire::open(&mut stream, Side::Responder, "room", b"secret") else { + panic!("met a peer of another version"); + }; + assert_eq!(error.kind(), io::ErrorKind::Unsupported); + let version = error.get_ref().unwrap().downcast_ref::(); + assert_eq!(version, Some(&wire::Version(1))); + // The older end heard this one's opening, and with it its version. + assert!(!older.join().unwrap().is_empty()); +} diff --git a/crates/notebook/src/live/wire.rs b/crates/notebook/src/live/wire.rs index c4972b49b7bdd552d398fe43f40b405174411236..fa1697719761d8eced313906fd4b96376bd38a41 100644 --- a/crates/notebook/src/live/wire.rs +++ b/crates/notebook/src/live/wire.rs @@ -11,7 +11,23 @@ use spake2::{Ed25519Group, Identity, Password, Spake2}; use std::io::{self, Read, Write}; /// The opening's version. Frames after it never change shape; they grow by kinds and fields. -pub const VERSION: u16 = 1; +pub const VERSION: u16 = 2; + +/// The version of the opening a peer of another version sent, as `open` fails with it. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub struct Version(pub u16); + +impl std::fmt::Display for Version { + fn fmt(&self, f: &mut std::fmt::Formatter) -> std::fmt::Result { + write!( + f, + "The peer speaks Live Share version {}, this end {VERSION}", + self.0 + ) + } +} + +impl std::error::Error for Version {} /// The largest block either side reads. const MOST: usize = 16 << 20; @@ -24,6 +40,8 @@ pub mod kind { pub const PRESENCE: u16 = 16; /// A host's files changed: `Touched`. pub const TOUCHED: u16 = 18; + /// A section's bytes as a commit changed them: `Delta`. + pub const DELTA: u16 = 19; pub const WELCOME: u16 = 32; /// Storage requests to a host, each a `Request` answered by a `Reply`. pub const LIST: u16 = 257; @@ -55,6 +73,7 @@ pub const KNOWN: &[u16] = &[ kind::BYE, kind::PRESENCE, kind::TOUCHED, + kind::DELTA, kind::WELCOME, kind::LIST, kind::STAMP, @@ -208,6 +227,31 @@ pub struct Touched { pub paths: Vec, } +/// What a commit changed in a section a host serves: from the image with stamp `base`, the +/// image `length` long with `writes` in place. A guest holding the base needs read nothing. +#[derive(Clone, Debug, PartialEq, Encode, Decode)] +#[cbor(map)] +pub struct Delta { + #[n(0)] + pub path: String, + #[n(1)] + pub base: WireStamp, + #[n(2)] + pub length: u64, + #[n(3)] + pub writes: Vec, +} + +/// Bytes at an offset. +#[derive(Clone, Debug, PartialEq, Encode, Decode)] +#[cbor(map)] +pub struct Written { + #[n(0)] + pub offset: u64, + #[cbor(n(1), with = "minicbor::bytes")] + pub bytes: Vec, +} + /// A storage request; its kind names the verb, and the verb what it carries. #[derive(Clone, Debug, Default, PartialEq, Encode, Decode)] #[cbor(map)] @@ -490,12 +534,20 @@ pub fn open( } let theirs: Open = minicbor::decode(&read_block(stream)?).map_err(|_| invalid("A malformed opening"))?; - if theirs.version != VERSION || theirs.room != room { - return Err(invalid("The peer means another room or version")); - } + // A responder answers even a peer of another version, so that both ends can say which + // should update. if side == Side::Responder { write_block(stream, &minicbor::to_vec(&ours).map_err(io::Error::other)?)?; } + if theirs.version != VERSION { + return Err(io::Error::new( + io::ErrorKind::Unsupported, + Version(theirs.version), + )); + } + if theirs.room != room { + return Err(invalid("The peer means another room")); + } let key = pake .finish(&theirs.pake) .map_err(|_| invalid("A malformed key exchange"))?; diff --git a/crates/snowbound/src/live.rs b/crates/snowbound/src/live.rs index 6d4200fd36059ec51dab59fa6c22c1eb73d74d24..be4397ca7af4b37c3b42910c20bb06a8c8d8807e 100644 --- a/crates/snowbound/src/live.rs +++ b/crates/snowbound/src/live.rs @@ -451,18 +451,21 @@ impl State { }); } - /// `host` shares the notebook at `location`: its folder's changes reach its guests. + /// `host` shares the notebook at `location`: its folder's changes reach its guests, and + /// theirs its sections, at once. fn hosting(&mut self, location: String, host: Arc) { if let Some(background) = self .library_at(&location) .and_then(|library| library.background.clone()) { - let host = Arc::downgrade(&host); + background.set_settle(notebook::live::share::SETTLE); + let shared = Arc::downgrade(&host); background.on_touched(Some(Box::new(move |paths| { - if let Some(host) = host.upgrade() { + if let Some(host) = shared.upgrade() { host.touched(paths); } }))); + host.on_changed(Box::new(move |paths| background.touched(paths))); } self.peers.hosts.insert(location, host); } diff --git a/crates/snowbound/src/main.rs b/crates/snowbound/src/main.rs index ec7ead7967a2b4a3ce49dc9f907a3f5c184d4055..b0adff212fd0d868d9042ade3e7a07e875155eac 100644 --- a/crates/snowbound/src/main.rs +++ b/crates/snowbound/src/main.rs @@ -4226,6 +4226,16 @@ impl State { } } Event::Failed(error) => eprintln!("Synchronization stopped: {error}"), + // Its guests hear of it sooner than the folder's watch would tell them. + #[cfg(feature = "live")] + Event::Attempt { + status: notebook::EditStatus::Published { .. }, + .. + } => { + if let Some(host) = self.peers.hosts.get(&session.library.location) { + host.touched(&[session.tabs[session.tab].path.clone()]); + } + } Event::Attempt { .. } | Event::Unreachable(_) => {} } } diff --git a/crates/snowbound/src/share.rs b/crates/snowbound/src/share.rs index 240a4cf730deadeb1c96428b375643dea73a5b6b..990eef41dca10082a69c2e5d3e92807ab4910f3b 100644 --- a/crates/snowbound/src/share.rs +++ b/crates/snowbound/src/share.rs @@ -62,6 +62,14 @@ enum Status { fn refusal(refusal: &Refusal) -> String { match refusal { Refusal::Malformed => "Check the code. It looks like 7KQ-4MZ-9XR.".into(), + Refusal::Version { theirs_newer: true } => { + "The person sharing has a newer Snowbound. Update Snowbound, then try again.".into() + } + Refusal::Version { + theirs_newer: false, + } => "The person sharing has an older \ + Snowbound. Ask them to update it." + .into(), Refusal::Wrong => "That code or password doesn’t open a notebook. Check it with the \ person sharing." .into(), -- 2.54.0