authorgravatar for git@paperclover.netclover caruso <git@paperclover.net> 2026-10-02 21:55:05-07:00
committergravatar for git@paperclover.netclover caruso <git@paperclover.net> 2026-10-03 05:26:36-07:00
log02acde3ba6bf8683a349deb02c298332a7c4f4ed
tree77bcf07d81cc0c05b30fbb6cd8084256e9819697
parentb528713ecc49c4ad7d46c039ac7d0acc2f9e72fe
signature Signed by SSH key SHA256:52mNGHRsVFBDED9IAX5pe+LRWUefqTbxEReunq21QvU

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

10 files changed, 617 insertions(+), 28 deletions(-)

crates/notebook/Cargo.toml+4
...@@ -59,6 +59,10 @@ libc = "0.2"...@@ -59,6 +59,10 @@ libc = "0.2"
59serde_json = "1"59serde_json = "1"
60tempfile = "3"60tempfile = "3"
6161
62[[example]]
63name = "live_latency"
64required-features = ["live"]
65
62[[example]]66[[example]]
63name = "smb_offline_client"67name = "smb_offline_client"
64required-features = ["smb"]68required-features = ["smb"]
crates/notebook/examples/live_latency.rs created+206
...@@ -0,0 +1,206 @@
1//! Measures Live Share's typing latency: how long after a keystroke lands in one guest's
2//! replica the other guest's section shows it, and the same from the host to a guest.
3//!
4//! `cargo run -p notebook --features live --example live_latency -- lan|relay [URL]`: `lan`
5//! meets by mDNS on this computer's loopback; `relay` through `URL`, or a relay this run
6//! starts on loopback.
7
8#[path = "../tests/support/server.rs"]
9mod server;
10
11use notebook::{
12 Replica,
13 live::{
14 Hello, Reach,
15 share::{self, Guest, Host, Sharing},
16 },
17 session::{Background, Notebook, Section},
18};
19use onestore::op::{Edit, Op, PageOp};
20use std::{
21 path::Path,
22 sync::Arc,
23 thread,
24 time::{Duration, Instant},
25};
26
27const KEYS: usize = 20;
28
29fn until(what: &str, done: impl Fn() -> bool) {
30 let deadline = Instant::now() + Duration::from_secs(60);
31 while !done() {
32 assert!(Instant::now() < deadline, "{what}");
33 thread::sleep(Duration::from_millis(1));
34 }
35}
36
37fn hello(name: &str) -> Hello {
38 Hello::new(name.into(), None).unwrap()
39}
40
41/// A guest as the app holds one: the notebook, its background, and the section open.
42struct Open {
43 _background: Background,
44 section: Section,
45}
46
47fn guest(
48 name: &str,
49 code: &str,
50 reach: Option<Reach>,
51 relay: Option<&str>,
52 cache: &Path,
53) -> (Arc<Guest>, Open) {
54 let welcome = share::join(hello(name), code, "", reach, relay).unwrap();
55 let guest = Guest::start(
56 hello(name),
57 welcome.share,
58 welcome.secret,
59 reach,
60 relay,
61 || {},
62 )
63 .unwrap();
64 until("no host", || guest.host().is_some());
65 let mut notebook = Notebook::open_hosted(Arc::clone(&guest), cache).unwrap();
66 let background = Background::hosted(Arc::clone(&guest), || {}).unwrap();
67 background.watch(notebook.replicas());
68 let replica = notebook.replica_path("Garden.one").unwrap();
69 std::fs::create_dir_all(replica.parent().unwrap()).unwrap();
70 let replica =
71 Replica::open_or_create(&replica, None, || notebook.read_section("Garden.one")).unwrap();
72 let section =
73 Section::resume_hosted("Garden.one".into(), replica, Arc::clone(&guest), || {}).unwrap();
74 background.hold("Garden.one", &section);
75 (
76 guest,
77 Open {
78 _background: background,
79 section,
80 },
81 )
82}
83
84/// Types `key` at the end of the page's text in `section`; returns the text it then reads.
85fn type_key(section: &Section, key: char) -> String {
86 let (space, text, before) = text_of(section);
87 let at = before.encode_utf16().count() as u32;
88 let op = PageOp::Text {
89 text,
90 range: at..at,
91 with: key.to_string(),
92 };
93 let edit = Edit {
94 at: 134_000_000_000_000_000,
95 ops: vec![Op::Page { space, op }],
96 };
97 section.replica().apply("Typist", edit).unwrap();
98 format!("{before}{key}")
99}
100
101fn text_of(section: &Section) -> (onestore::ExGuid, onestore::ExGuid, String) {
102 let (space, ..) = section.pages().unwrap()[0];
103 let page = section.page(space).unwrap();
104 page.objects
105 .iter()
106 .find_map(|object| match object {
107 onestore::page::PageObject::Outline(outline) => outline
108 .paragraphs
109 .iter()
110 .find_map(|p| p.text().map(|t| (space, t.id, t.text.text().to_owned()))),
111 _ => None,
112 })
113 .unwrap()
114}
115
116fn report(what: &str, mut times: Vec<Duration>) {
117 times.sort();
118 let ms = |at: usize| times[at].as_secs_f64() * 1000.0;
119 println!(
120 "{what}: median {:.0} ms, p90 {:.0} ms, worst {:.0} ms over {} keys",
121 ms(times.len() / 2),
122 ms(times.len() * 9 / 10),
123 ms(times.len() - 1),
124 times.len()
125 );
126}
127
128fn main() {
129 let args: Vec<String> = std::env::args().skip(1).collect();
130 let (reach, url) = match args.first().map(String::as_str) {
131 Some("lan") => (Some(Reach::Loopback), None),
132 Some("relay") => (
133 None,
134 Some(args.get(1).cloned().unwrap_or_else(|| {
135 let listener = std::net::TcpListener::bind("127.0.0.1:0").unwrap();
136 let url = format!("ws://{}", listener.local_addr().unwrap());
137 thread::spawn(move || relay::server::serve(listener, Default::default()));
138 url
139 })),
140 ),
141 _ => panic!("lan or relay [URL]"),
142 };
143 let relay = url.as_deref();
144 let directory = tempfile::tempdir().unwrap();
145 let folder = directory.path().join("Garden");
146 std::fs::create_dir_all(&folder).unwrap();
147 std::fs::write(
148 folder.join("Garden.one"),
149 onestore::create_section("Garden.one", "Typed:", "Fixture").unwrap(),
150 )
151 .unwrap();
152 let storage = Notebook::open(&folder, directory.path().join("host"))
153 .unwrap()
154 .into_storage();
155 let host = Host::start(
156 storage,
157 hello("Ada"),
158 Sharing::new("").unwrap(),
159 "Garden",
160 reach,
161 relay,
162 || {},
163 )
164 .unwrap();
165 until("no code", || {
166 host.code().is_some_and(|code| share::code(&code).is_some())
167 });
168 let code = host.code().unwrap();
169 let (_grace, grace) = guest("Grace", &code, reach, relay, &directory.path().join("g"));
170 let (_alan, alan) = guest("Alan", &code, reach, relay, &directory.path().join("a"));
171
172 let mut times = Vec::new();
173 for key in "abcdefghijklmnopqrstuvwxyz".chars().take(KEYS) {
174 let typed = Instant::now();
175 let expected = type_key(&grace.section, key);
176 until("the other guest never saw a key", || {
177 alan.section.events();
178 text_of(&alan.section).2 == expected
179 });
180 times.push(typed.elapsed());
181 thread::sleep(Duration::from_millis(150));
182 }
183 report("guest to guest", times);
184
185 // The host's own keystrokes, published to its file as its app does, then reported to
186 // the share as its app reports them.
187 let file = folder.join("Garden.one");
188 let mut times = Vec::new();
189 for key in "ABCDEFGHIJKLMNOPQRSTUVWXYZ".chars().take(KEYS) {
190 let image = std::fs::read(&file).unwrap();
191 let (space, text, before) = server::text(&image);
192 let at = before.encode_utf16().count() as u32;
193 let typed = Instant::now();
194 let next = server::typed(&image, space, text, at..at, &key.to_string());
195 std::fs::write(&file, &next).unwrap();
196 host.touched(&["Garden.one".into()]);
197 let expected = format!("{before}{key}");
198 until("the guest never saw the host's key", || {
199 alan.section.events();
200 text_of(&alan.section).2 == expected
201 });
202 times.push(typed.elapsed());
203 thread::sleep(Duration::from_millis(150));
204 }
205 report("host to guest", times);
206}
crates/notebook/src/background.rs+17-7
...@@ -91,6 +91,8 @@ struct Watched {...@@ -91,6 +91,8 @@ struct Watched {
91 lost: bool,91 lost: bool,
92 /// The current connection's watch reports the folder's changes.92 /// The current connection's watch reports the folder's changes.
93 reported: bool,93 reported: bool,
94 /// How long a report settles, where not `SETTLE` (`Background::set_settle`).
95 settle: Option<Duration>,
94}96}
9597
96struct Watch {98struct Watch {
...@@ -356,7 +358,7 @@ impl Background {...@@ -356,7 +358,7 @@ impl Background {
356 guest: Arc<crate::live::share::Guest>,358 guest: Arc<crate::live::share::Guest>,
357 notify: impl Fn() + Send + 'static,359 notify: impl Fn() + Send + 'static,
358 ) -> Result<Self> {360 ) -> Result<Self> {
359 Self::start(361 let background = Self::start(
360 true,362 true,
361 move |reports| {363 move |reports| {
362 guest.watch(reports)?;364 guest.watch(reports)?;
...@@ -366,7 +368,17 @@ impl Background {...@@ -366,7 +368,17 @@ impl Background {
366 Ok(((bind, list), true))368 Ok(((bind, list), true))
367 },369 },
368 notify,370 notify,
369 )371 )?;
372 background.set_settle(crate::live::share::SETTLE);
373 Ok(background)
374 }
375
376 /// Checks what a report names after `settle` rather than a second, where reports come
377 /// one to a commit, as a Live Share host sends them.
378 pub fn set_settle(&self, settle: Duration) {
379 if let Ok(mut watched) = self.0.watched.lock() {
380 watched.settle = Some(settle);
381 }
370 }382 }
371383
372 /// Has `listener` hear every report of changed paths from now on, as a Live Share host384 /// Has `listener` hear every report of changed paths from now on, as a Live Share host
...@@ -638,6 +650,7 @@ impl Watched {...@@ -638,6 +650,7 @@ impl Watched {
638 }650 }
639651
640 fn touched(&mut self, paths: &[String], now: Instant) {652 fn touched(&mut self, paths: &[String], now: Instant) {
653 let settled = now + self.settle.unwrap_or(SETTLE);
641 for path in paths {654 for path in paths {
642 let path = path.trim_matches('/');655 let path = path.trim_matches('/');
643 let section = self656 let section = self
...@@ -647,16 +660,13 @@ impl Watched {...@@ -647,16 +660,13 @@ impl Watched {
647 match section {660 match section {
648 Some(index) => {661 Some(index) => {
649 let watch = &mut self.sections[index];662 let watch = &mut self.sections[index];
650 watch.due = watch.due.min(now + SETTLE);663 watch.due = watch.due.min(settled);
651 watch.current = false;664 watch.current = false;
652 watch.reported = true;665 watch.reported = true;
653 }666 }
654 None if self.sections.iter().any(|watch| within(&watch.path, path)) => {667 None if self.sections.iter().any(|watch| within(&watch.path, path)) => {
655 self.folders.insert(path.to_owned());668 self.folders.insert(path.to_owned());
656 self.relist = Some(669 self.relist = Some(self.relist.map_or(settled, |due| due.min(settled)));
657 self.relist
658 .map_or(now + SETTLE, |due| due.min(now + SETTLE)),
659 );
660 }670 }
661 None => {}671 None => {}
662 }672 }
crates/notebook/src/live.rs+20-2
...@@ -192,6 +192,8 @@ struct State {...@@ -192,6 +192,8 @@ struct State {
192 failed: u32,192 failed: u32,
193 /// The code admits no one new.193 /// The code admits no one new.
194 burned: bool,194 burned: bool,
195 /// The Live Share version of the last peer met that speaks another.
196 outdated: Option<u16>,
195 daemon: Option<ServiceDaemon>,197 daemon: Option<ServiceDaemon>,
196}198}
197199
...@@ -379,6 +381,12 @@ impl Live {...@@ -379,6 +381,12 @@ impl Live {
379 self.shared.state.lock().unwrap().failed381 self.shared.state.lock().unwrap().failed
380 }382 }
381383
384 /// The Live Share version of the last peer met that speaks another, which one of the two
385 /// must update to meet.
386 pub fn other_version(&self) -> Option<u16> {
387 self.shared.state.lock().unwrap().outdated
388 }
389
382 /// Whether this end's code had too many wrong tries and admits no one new.390 /// Whether this end's code had too many wrong tries and admits no one new.
383 pub fn burned(&self) -> bool {391 pub fn burned(&self) -> bool {
384 self.shared.state.lock().unwrap().burned392 self.shared.state.lock().unwrap().burned
...@@ -535,12 +543,22 @@ impl Shared {...@@ -535,12 +543,22 @@ impl Shared {
535 stream.set_read_timeout(GONE)?;543 stream.set_read_timeout(GONE)?;
536 Ok((send, receive, hello))544 Ok((send, receive, hello))
537 })();545 })();
538 pipe.met(met.is_ok());546 let other = (met.as_ref().err())
547 .and_then(|error| error.get_ref()?.downcast_ref::<wire::Version>())
548 .copied();
549 // A peer of another version guessed nothing, the keys never being agreed.
550 pipe.met(met.is_ok() || other.is_some());
539 let (send, mut receive, hello) = match met {551 let (send, mut receive, hello) = match met {
540 Ok(met) => met,552 Ok(met) => met,
541 Err(error) => {553 Err(error) => {
542 eprintln!("Live: no meeting in {tag}: {error}");554 eprintln!("Live: no meeting in {tag}: {error}");
543 if error.kind() == io::ErrorKind::InvalidData {555 if let Some(wire::Version(version)) = other {
556 self.state.lock().unwrap().outdated = Some(version);
557 if matches!(self.room, Room::Code { owner: false, .. }) {
558 self.stopped.store(true, Ordering::Release);
559 }
560 (self.events)(Event::Changed);
561 } else if error.kind() == io::ErrorKind::InvalidData {
544 self.failed();562 self.failed();
545 }563 }
546 pipe.shutdown();564 pipe.shutdown();
crates/notebook/src/live/share.rs+246-12
...@@ -9,7 +9,9 @@...@@ -9,7 +9,9 @@
99
10use super::{10use super::{
11 Event, Hello, Line, Live, Peer, Presence, Reach, Relayed, Room,11 Event, Hello, Line, Live, Peer, Presence, Reach, Relayed, Room,
12 wire::{self, Failure, Reply, Request, Touched, Welcome, WireEntry, WireStamp, kind},12 wire::{
13 self, Delta, Failure, Reply, Request, Touched, Welcome, WireEntry, WireStamp, Written, kind,
14 },
13};15};
14use crate::{Error, Result, background::Reports, discover, session::Storage};16use crate::{Error, Result, background::Reports, discover, session::Storage};
15use onestore::{CommitError, CommitState, RevisionIndex, Stamp, Store, Transaction};17use onestore::{CommitError, CommitState, RevisionIndex, Stamp, Store, Transaction};
...@@ -37,8 +39,12 @@ const LIMIT: usize = 256 << 20;...@@ -37,8 +39,12 @@ const LIMIT: usize = 256 << 20;
37/// Snapshots of files being read, per guest, and how long one is kept unread.39/// Snapshots of files being read, per guest, and how long one is kept unread.
38const SNAPSHOTS: usize = 8;40const SNAPSHOTS: usize = 8;
39const SNAPSHOT_AGE: Duration = Duration::from_secs(120);41const SNAPSHOT_AGE: Duration = Duration::from_secs(120);
40/// Sections whose image a host keeps to check guests' commits on.42/// Sections whose image a host keeps to check guests' commits on and tell them what changed,
43/// and a guest keeps to take those changes on without reading.
41const IMAGES: usize = 4;44const IMAGES: usize = 4;
45/// How long a report of a change settles in a guest's background: its host reports each
46/// commit once, at once.
47pub const SETTLE: Duration = Duration::from_millis(20);
42/// What a host says leaving as it stops sharing.48/// What a host says leaving as it stops sharing.
43const STOPPED: &str = "stopped";49const STOPPED: &str = "stopped";
44/// What a host says hanging up on a guest that asks too much too fast.50/// 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 {...@@ -95,6 +101,9 @@ pub fn location(share: &[u8; 16]) -> String {
95pub enum Refusal {101pub enum Refusal {
96 /// Not a code, or one mistyped, which spends none of the relay's tries.102 /// Not a code, or one mistyped, which spends none of the relay's tries.
97 Malformed,103 Malformed,
104 /// The person sharing runs a Snowbound of another Live Share version: the newer one's
105 /// `true` where it is theirs, so this one should update.
106 Version { theirs_newer: bool },
98 /// The code's secret or password is wrong.107 /// The code's secret or password is wrong.
99 Wrong,108 Wrong,
100 /// No one shares with the code's number now.109 /// No one shares with the code's number now.
...@@ -149,6 +158,11 @@ pub fn join(...@@ -149,6 +158,11 @@ pub fn join(
149 if let Ok(welcome) = welcome.try_recv() {158 if let Ok(welcome) = welcome.try_recv() {
150 return Ok(welcome);159 return Ok(welcome);
151 }160 }
161 if let Some(version) = live.other_version() {
162 return Err(Refusal::Version {
163 theirs_newer: version > wire::VERSION,
164 });
165 }
152 if live.failed() > 0 {166 if live.failed() > 0 {
153 return Err(Refusal::Wrong);167 return Err(Refusal::Wrong);
154 }168 }
...@@ -209,6 +223,7 @@ impl Host {...@@ -209,6 +223,7 @@ impl Host {
209 snapshots: Mutex::default(),223 snapshots: Mutex::default(),
210 puts: Mutex::default(),224 puts: Mutex::default(),
211 guests: Mutex::default(),225 guests: Mutex::default(),
226 host: Mutex::default(),
212 });227 });
213 let serving = Hello {228 let serving = Hello {
214 serves: Some(sharing.share),229 serves: Some(sharing.share),
...@@ -296,8 +311,18 @@ impl Host {...@@ -296,8 +311,18 @@ impl Host {
296 }311 }
297 }312 }
298313
299 /// Tells every guest the files at these catalog paths changed.314 /// Has `listener` hear the catalog paths guests change from now on, sooner than a watch
315 /// on the notebook's folder would.
316 pub fn on_changed(&self, listener: crate::session::Listener) {
317 *self.served.host.lock().unwrap() = Some(listener);
318 }
319
320 /// Tells every guest the files at these catalog paths changed, with what changed in the
321 /// sections a guest read lately.
300 pub fn touched(&self, paths: &[String]) {322 pub fn touched(&self, paths: &[String]) {
323 for path in paths {
324 self.served.changed_here(path);
325 }
301 self.served.tell(paths);326 self.served.tell(paths);
302 }327 }
303328
...@@ -355,6 +380,8 @@ struct Served {...@@ -355,6 +380,8 @@ struct Served {
355 /// Bytes a later request carries, by guest and upload.380 /// Bytes a later request carries, by guest and upload.
356 puts: Mutex<ByGuest<Vec<u8>>>,381 puts: Mutex<ByGuest<Vec<u8>>>,
357 guests: Mutex<BTreeMap<[u8; 16], Admitted>>,382 guests: Mutex<BTreeMap<[u8; 16], Admitted>>,
383 /// Hears the paths guests changed, as the host's own notebook should.
384 host: Mutex<Option<crate::session::Listener>>,
358}385}
359386
360/// A guest as its host serves it: the line to it, its requests waiting for its workers, and387/// A guest as its host serves it: the line to it, its requests waiting for its workers, and
...@@ -451,6 +478,14 @@ impl Served {...@@ -451,6 +478,14 @@ impl Served {
451 .retain(|(guest, _), _| guest != peer);478 .retain(|(guest, _), _| guest != peer);
452 }479 }
453480
481 /// Tells every guest, and the host, that a guest changed the files at `paths`.
482 fn changed(&self, paths: &[String]) {
483 self.tell(paths);
484 if let Some(listener) = &*self.host.lock().unwrap() {
485 listener(paths);
486 }
487 }
488
454 /// Tells every guest the files at `paths` changed.489 /// Tells every guest the files at `paths` changed.
455 fn tell(&self, paths: &[String]) {490 fn tell(&self, paths: &[String]) {
456 let touched = Touched {491 let touched = Touched {
...@@ -540,7 +575,7 @@ impl Served {...@@ -540,7 +575,7 @@ impl Served {
540 kind::COMMIT => {575 kind::COMMIT => {
541 let transaction = Transaction::from_bytes(&self.carried(peer, &request)?)?;576 let transaction = Transaction::from_bytes(&self.carried(peer, &request)?)?;
542 self.commit(path, &transaction)?;577 self.commit(path, &transaction)?;
543 self.tell(&[path.to_owned()]);578 self.changed(&[path.to_owned()]);
544 done579 done
545 }580 }
546 kind::CONFIRM => {581 kind::CONFIRM => {
...@@ -549,12 +584,12 @@ impl Served {...@@ -549,12 +584,12 @@ impl Served {
549 }584 }
550 kind::CREATE => {585 kind::CREATE => {
551 self.storage.create(path, &self.carried(peer, &request)?)?;586 self.storage.create(path, &self.carried(peer, &request)?)?;
552 self.tell(&[folder(path)]);587 self.changed(&[folder(path)]);
553 done588 done
554 }589 }
555 kind::CREATE_DIRECTORY => {590 kind::CREATE_DIRECTORY => {
556 self.storage.create_directory(path)?;591 self.storage.create_directory(path)?;
557 self.tell(&[folder(path)]);592 self.changed(&[folder(path)]);
558 done593 done
559 }594 }
560 kind::HIDE => {595 kind::HIDE => {
...@@ -568,12 +603,12 @@ impl Served {...@@ -568,12 +603,12 @@ impl Served {
568 } else {603 } else {
569 self.storage.replace(path, to)?;604 self.storage.replace(path, to)?;
570 }605 }
571 self.tell(&[folder(path), folder(to)]);606 self.changed(&[folder(path), folder(to)]);
572 done607 done
573 }608 }
574 kind::DELETE => {609 kind::DELETE => {
575 self.storage.delete(path)?;610 self.storage.delete(path)?;
576 self.tell(&[folder(path)]);611 self.changed(&[folder(path)]);
577 done612 done
578 }613 }
579 kind::PLACE => {614 kind::PLACE => {
...@@ -582,12 +617,12 @@ impl Served {...@@ -582,12 +617,12 @@ impl Served {
582 .ok_or_else(|| refused(io::ErrorKind::InvalidInput, "No ancestor"))?;617 .ok_or_else(|| refused(io::ErrorKind::InvalidInput, "No ancestor"))?;
583 let name = request.name.as_deref().unwrap_or_default();618 let name = request.name.as_deref().unwrap_or_default();
584 self.storage.place(path, ancestor, name)?;619 self.storage.place(path, ancestor, name)?;
585 self.tell(&[path.to_owned()]);620 self.changed(&[path.to_owned()]);
586 done621 done
587 }622 }
588 kind::SUPERSEDE => {623 kind::SUPERSEDE => {
589 self.storage.supersede(path, &stamp()?, to()?)?;624 self.storage.supersede(path, &stamp()?, to()?)?;
590 self.tell(&[folder(path)]);625 self.changed(&[folder(path)]);
591 done626 done
592 }627 }
593 _ => return Err(refused(io::ErrorKind::Unsupported, "An unknown request")),628 _ => return Err(refused(io::ErrorKind::Unsupported, "An unknown request")),
...@@ -716,9 +751,103 @@ impl Served {...@@ -716,9 +751,103 @@ impl Served {
716 )));751 )));
717 }752 }
718 self.storage.commit(path, transaction)?;753 self.storage.commit(path, transaction)?;
754 self.tell_delta(path, &image, &next);
719 self.keep(path, Arc::new(next));755 self.keep(path, Arc::new(next));
720 Ok(())756 Ok(())
721 }757 }
758
759 /// Tells every guest what a commit changed in the section at `path`, from `before` to
760 /// `after`, where that is small enough to send.
761 fn tell_delta(&self, path: &str, before: &[u8], after: &[u8]) {
762 let (Ok(base), Some(writes)) = (Stamp::of(before), delta(before, after)) else {
763 return;
764 };
765 let delta = Delta {
766 path: path.to_owned(),
767 base: (&base).into(),
768 length: after.len() as u64,
769 writes,
770 };
771 for guest in self.guests.lock().unwrap().values() {
772 let _ = guest.line.send(kind::DELTA, &delta);
773 }
774 }
775
776 /// The host's own change to the section at `path`, which guests hear as a delta where
777 /// the host kept the image before it.
778 fn changed_here(&self, path: &str) {
779 let kept = self
780 .images
781 .lock()
782 .unwrap()
783 .iter()
784 .find_map(|(held, at, image)| (held == path).then(|| (at.clone(), Arc::clone(image))));
785 let Some((at, before)) = kept else {
786 return;
787 };
788 if self.storage.stamp(path).is_ok_and(|now| now == at) {
789 return;
790 }
791 if let Ok(after) = self.storage.read(path) {
792 self.tell_delta(path, &before, &after);
793 self.keep(path, Arc::new(after));
794 }
795 }
796}
797
798/// The writes that make `after` of `before`, a commit's appended bytes, patches and header,
799/// where they are much less than `after` itself.
800fn delta(before: &[u8], after: &[u8]) -> Option<Vec<Written>> {
801 const BLOCK: usize = 4096;
802 if before.len() < 1024 || after.len() < before.len() {
803 return None;
804 }
805 let mut writes = vec![Written {
806 offset: 0,
807 bytes: after[..1024].to_vec(),
808 }];
809 let mut at = 1024;
810 while at < before.len() {
811 let end = (at + BLOCK).min(before.len());
812 if before[at..end] != after[at..end] {
813 let first = (at..end).find(|&i| before[i] != after[i]).unwrap_or(at);
814 let last = (at..end).rfind(|&i| before[i] != after[i]).unwrap_or(first) + 1;
815 match writes.last_mut() {
816 Some(write) if write.offset as usize + write.bytes.len() + 64 >= first => {
817 let from = write.offset as usize;
818 write.bytes = after[from..last].to_vec();
819 }
820 _ => writes.push(Written {
821 offset: first as u64,
822 bytes: after[first..last].to_vec(),
823 }),
824 }
825 }
826 at = end;
827 }
828 if after.len() > before.len() {
829 writes.push(Written {
830 offset: before.len() as u64,
831 bytes: after[before.len()..].to_vec(),
832 });
833 }
834 let sent: usize = writes.iter().map(|write| write.bytes.len()).sum();
835 (sent <= after.len() / 2).then_some(writes)
836}
837
838/// `image` with `writes`, `length` long: none where a write falls outside it.
839fn written(image: &[u8], length: u64, writes: &[Written]) -> Option<Vec<u8>> {
840 let length = usize::try_from(length)
841 .ok()
842 .filter(|length| *length >= image.len())?;
843 let mut next = image.to_vec();
844 next.resize(length, 0);
845 for write in writes {
846 let offset = usize::try_from(write.offset).ok()?;
847 next.get_mut(offset..offset.checked_add(write.bytes.len())?)?
848 .copy_from_slice(&write.bytes);
849 }
850 Some(next)
722}851}
723852
724fn within(folder: &str, name: &str) -> String {853fn within(folder: &str, name: &str) -> String {
...@@ -846,6 +975,8 @@ struct Inner {...@@ -846,6 +975,8 @@ struct Inner {
846 stopped: AtomicBool,975 stopped: AtomicBool,
847 /// Bytes of chunks asked for and not yet given.976 /// Bytes of chunks asked for and not yet given.
848 asked: (Mutex<usize>, Condvar),977 asked: (Mutex<usize>, Condvar),
978 /// The sections read lately, kept as the host's deltas change them, newest first.
979 images: Mutex<Vec<(String, Arc<Vec<u8>>)>>,
849}980}
850981
851impl Guest {982impl Guest {
...@@ -867,6 +998,7 @@ impl Guest {...@@ -867,6 +998,7 @@ impl Guest {
867 watch: Mutex::default(),998 watch: Mutex::default(),
868 stopped: AtomicBool::new(false),999 stopped: AtomicBool::new(false),
869 asked: Default::default(),1000 asked: Default::default(),
1001 images: Mutex::default(),
870 });1002 });
871 let heard = Arc::clone(&inner);1003 let heard = Arc::clone(&inner);
872 let live = Live::start(me, &Room::Notebook(secret), reach, relay, move |event| {1004 let live = Live::start(me, &Room::Notebook(secret), reach, relay, move |event| {
...@@ -1026,9 +1158,19 @@ impl Guest {...@@ -1026,9 +1158,19 @@ impl Guest {
1026 image.extend_from_slice(&chunk);1158 image.extend_from_slice(&chunk);
1027 }1159 }
1028 image.truncate(length);1160 image.truncate(length);
1161 if kind == kind::READ {
1162 self.inner.hold(path, Arc::new(image.clone()));
1163 }
1029 Ok(image)1164 Ok(image)
1030 }1165 }
10311166
1167 /// The image of the section at `path` with stamp `stamp`, where this guest holds it.
1168 fn held(&self, path: &str, stamp: &Stamp) -> Option<Vec<u8>> {
1169 let images = self.inner.images.lock().unwrap();
1170 let (_, image) = images.iter().find(|(held, _)| held == path)?;
1171 (Stamp::of(image).ok().as_ref() == Some(stamp)).then(|| image.to_vec())
1172 }
1173
1032 pub(crate) fn entries(&self, folder: &str) -> io::Result<Vec<discover::Entry>> {1174 pub(crate) fn entries(&self, folder: &str) -> io::Result<Vec<discover::Entry>> {
1033 let reply = self.ask(1175 let reply = self.ask(
1034 kind::LIST,1176 kind::LIST,
...@@ -1101,6 +1243,46 @@ impl Guest {...@@ -1101,6 +1243,46 @@ impl Guest {
1101}1243}
11021244
1103impl Inner {1245impl Inner {
1246 fn hold(&self, path: &str, image: Arc<Vec<u8>>) {
1247 let mut images = self.images.lock().unwrap();
1248 images.retain(|(held, _)| held != path);
1249 images.insert(0, (path.to_owned(), image));
1250 images.truncate(IMAGES);
1251 }
1252
1253 /// Takes on a delta to an image held.
1254 fn apply(&self, delta: &Delta) {
1255 let held = self
1256 .images
1257 .lock()
1258 .unwrap()
1259 .iter()
1260 .find_map(|(path, image)| (*path == delta.path).then(|| Arc::clone(image)));
1261 let base: Option<Stamp> = (&delta.base).try_into().ok();
1262 if let Some(image) = held
1263 && Stamp::of(&image).ok() == base
1264 && let Some(next) = written(&image, delta.length, &delta.writes)
1265 {
1266 self.hold(&delta.path, Arc::new(next));
1267 }
1268 }
1269
1270 /// Takes on this guest's own commit to an image held.
1271 fn published(&self, path: &str, transaction: &Transaction) {
1272 let held = self
1273 .images
1274 .lock()
1275 .unwrap()
1276 .iter()
1277 .find_map(|(held, image)| (held == path).then(|| Arc::clone(image)));
1278 if let Some(image) = held {
1279 let mut next = (*image).clone();
1280 if transaction.apply(&mut next).is_ok() {
1281 self.hold(path, Arc::new(next));
1282 }
1283 }
1284 }
1285
1104 fn heard(&self, event: Event) {1286 fn heard(&self, event: Event) {
1105 let serves = |hello: &Hello| hello.serves == Some(self.share);1287 let serves = |hello: &Hello| hello.serves == Some(self.share);
1106 match event {1288 match event {
...@@ -1125,6 +1307,11 @@ impl Inner {...@@ -1125,6 +1307,11 @@ impl Inner {
1125 let _ = waiting.send(reply);1307 let _ = waiting.send(reply);
1126 }1308 }
1127 }1309 }
1310 kind::DELTA if serves(from) => {
1311 if let Ok(delta) = minicbor::decode::<Delta>(body) {
1312 self.apply(&delta);
1313 }
1314 }
1128 kind::TOUCHED if serves(from) => {1315 kind::TOUCHED if serves(from) => {
1129 if let Ok(touched) = minicbor::decode::<Touched>(body)1316 if let Ok(touched) = minicbor::decode::<Touched>(body)
1130 && let Some(reports) = &*self.watch.lock().unwrap()1317 && let Some(reports) = &*self.watch.lock().unwrap()
...@@ -1150,6 +1337,8 @@ impl Inner {...@@ -1150,6 +1337,8 @@ impl Inner {
1150pub struct HostedRemote {1337pub struct HostedRemote {
1151 guest: Arc<Guest>,1338 guest: Arc<Guest>,
1152 path: String,1339 path: String,
1340 /// The stamp last asked for, which an image the guest holds may already have.
1341 seen: Option<Stamp>,
1153}1342}
11541343
1155impl HostedRemote {1344impl HostedRemote {
...@@ -1157,21 +1346,30 @@ impl HostedRemote {...@@ -1157,21 +1346,30 @@ impl HostedRemote {
1157 Self {1346 Self {
1158 guest: Arc::clone(guest),1347 guest: Arc::clone(guest),
1159 path: path.to_owned(),1348 path: path.to_owned(),
1349 seen: None,
1160 }1350 }
1161 }1351 }
1162}1352}
11631353
1164impl crate::Remote for HostedRemote {1354impl crate::Remote for HostedRemote {
1165 fn read(&mut self) -> io::Result<Vec<u8>> {1355 fn read(&mut self) -> io::Result<Vec<u8>> {
1356 if let Some(image) = (self.seen.as_ref()).and_then(|seen| self.guest.held(&self.path, seen))
1357 {
1358 return Ok(image);
1359 }
1166 self.guest.read(kind::READ, &self.path, LIMIT)1360 self.guest.read(kind::READ, &self.path, LIMIT)
1167 }1361 }
11681362
1169 fn stamp(&mut self) -> io::Result<Stamp> {1363 fn stamp(&mut self) -> io::Result<Stamp> {
1170 self.guest.stamp(&self.path)1364 let stamp = self.guest.stamp(&self.path)?;
1365 self.seen = Some(stamp.clone());
1366 Ok(stamp)
1171 }1367 }
11721368
1173 fn publish(&mut self, transaction: &Transaction) -> std::result::Result<(), CommitError> {1369 fn publish(&mut self, transaction: &Transaction) -> std::result::Result<(), CommitError> {
1174 self.guest.commit(&self.path, transaction)1370 self.guest.commit(&self.path, transaction)?;
1371 self.guest.inner.published(&self.path, transaction);
1372 Ok(())
1175 }1373 }
11761374
1177 fn confirm(&mut self, base: &Stamp) -> std::result::Result<(), CommitError> {1375 fn confirm(&mut self, base: &Stamp) -> std::result::Result<(), CommitError> {
...@@ -1373,3 +1571,39 @@ impl Storage for Hosted {...@@ -1373,3 +1571,39 @@ impl Storage for Hosted {
1373 self.guest.verb(kind::SUPERSEDE, path, request)1571 self.guest.verb(kind::SUPERSEDE, path, request)
1374 }1572 }
1375}1573}
1574
1575#[cfg(test)]
1576mod tests {
1577 use super::*;
1578
1579 /// A commit's writes, found by comparing images, rebuild the image after it from the one
1580 /// before, and a change touching most of the file is left to be read whole.
1581 #[test]
1582 fn deltas_rebuild_the_image_after_a_commit() {
1583 let before: Vec<u8> = (0..20_000u32).map(|at| (at % 251) as u8).collect();
1584 let mut after = before.clone();
1585 after[3] ^= 1;
1586 after[5000..5010].fill(9);
1587 after[5050] ^= 1;
1588 after[17_000] ^= 1;
1589 after.extend_from_slice(&[7; 3000]);
1590 let writes = delta(&before, &after).unwrap();
1591 assert_eq!(
1592 writes.len(),
1593 4,
1594 "the header, two runs merged as one, one more, the tail"
1595 );
1596 assert_eq!(
1597 written(&before, after.len() as u64, &writes).unwrap(),
1598 after
1599 );
1600 let rewritten: Vec<u8> = before.iter().map(|byte| byte ^ 1).collect();
1601 assert!(delta(&before, &rewritten).is_none());
1602 assert!(delta(&before, &before[..10_000]).is_none());
1603 let outside = [Written {
1604 offset: after.len() as u64,
1605 bytes: vec![1],
1606 }];
1607 assert!(written(&before, after.len() as u64, &outside).is_none());
1608 }
1609}
crates/notebook/src/live/tests.rs+44
...@@ -520,3 +520,47 @@ fn a_malicious_relay_is_caught() {...@@ -520,3 +520,47 @@ fn a_malicious_relay_is_caught() {
520 });520 });
521 }521 }
522}522}
523
524/// Ends of two versions each learn the other's version from the opening, before any key is
525/// agreed, so each can say which should update.
526#[test]
527fn another_version_is_named_before_any_key() {
528 let listener = TcpListener::bind("127.0.0.1:0").unwrap();
529 let address = listener.local_addr().unwrap();
530 let older = thread::spawn(move || {
531 let mut stream = TcpStream::connect(address).unwrap();
532 #[derive(minicbor::Encode)]
533 #[cbor(map)]
534 struct Open {
535 #[n(0)]
536 version: u16,
537 #[n(1)]
538 room: String,
539 #[cbor(n(2), with = "minicbor::bytes")]
540 pake: Vec<u8>,
541 }
542 let open = minicbor::to_vec(Open {
543 version: 1,
544 room: "room".into(),
545 pake: vec![0; 33],
546 })
547 .unwrap();
548 stream
549 .write_all(&[&(open.len() as u32).to_be_bytes()[..], &open].concat())
550 .unwrap();
551 let mut length = [0; 4];
552 stream.read_exact(&mut length).unwrap();
553 let mut answer = vec![0; u32::from_be_bytes(length) as usize];
554 stream.read_exact(&mut answer).unwrap();
555 answer
556 });
557 let (mut stream, _) = listener.accept().unwrap();
558 let Err(error) = wire::open(&mut stream, Side::Responder, "room", b"secret") else {
559 panic!("met a peer of another version");
560 };
561 assert_eq!(error.kind(), io::ErrorKind::Unsupported);
562 let version = error.get_ref().unwrap().downcast_ref::<wire::Version>();
563 assert_eq!(version, Some(&wire::Version(1)));
564 // The older end heard this one's opening, and with it its version.
565 assert!(!older.join().unwrap().is_empty());
566}
crates/notebook/src/live/wire.rs+56-4
...@@ -11,7 +11,23 @@ use spake2::{Ed25519Group, Identity, Password, Spake2};...@@ -11,7 +11,23 @@ use spake2::{Ed25519Group, Identity, Password, Spake2};
11use std::io::{self, Read, Write};11use std::io::{self, Read, Write};
1212
13/// The opening's version. Frames after it never change shape; they grow by kinds and fields.13/// The opening's version. Frames after it never change shape; they grow by kinds and fields.
14pub const VERSION: u16 = 1;14pub const VERSION: u16 = 2;
15
16/// The version of the opening a peer of another version sent, as `open` fails with it.
17#[derive(Clone, Copy, Debug, PartialEq, Eq)]
18pub struct Version(pub u16);
19
20impl std::fmt::Display for Version {
21 fn fmt(&self, f: &mut std::fmt::Formatter) -> std::fmt::Result {
22 write!(
23 f,
24 "The peer speaks Live Share version {}, this end {VERSION}",
25 self.0
26 )
27 }
28}
29
30impl std::error::Error for Version {}
15/// The largest block either side reads.31/// The largest block either side reads.
16const MOST: usize = 16 << 20;32const MOST: usize = 16 << 20;
1733
...@@ -24,6 +40,8 @@ pub mod kind {...@@ -24,6 +40,8 @@ pub mod kind {
24 pub const PRESENCE: u16 = 16;40 pub const PRESENCE: u16 = 16;
25 /// A host's files changed: `Touched`.41 /// A host's files changed: `Touched`.
26 pub const TOUCHED: u16 = 18;42 pub const TOUCHED: u16 = 18;
43 /// A section's bytes as a commit changed them: `Delta`.
44 pub const DELTA: u16 = 19;
27 pub const WELCOME: u16 = 32;45 pub const WELCOME: u16 = 32;
28 /// Storage requests to a host, each a `Request` answered by a `Reply`.46 /// Storage requests to a host, each a `Request` answered by a `Reply`.
29 pub const LIST: u16 = 257;47 pub const LIST: u16 = 257;
...@@ -55,6 +73,7 @@ pub const KNOWN: &[u16] = &[...@@ -55,6 +73,7 @@ pub const KNOWN: &[u16] = &[
55 kind::BYE,73 kind::BYE,
56 kind::PRESENCE,74 kind::PRESENCE,
57 kind::TOUCHED,75 kind::TOUCHED,
76 kind::DELTA,
58 kind::WELCOME,77 kind::WELCOME,
59 kind::LIST,78 kind::LIST,
60 kind::STAMP,79 kind::STAMP,
...@@ -208,6 +227,31 @@ pub struct Touched {...@@ -208,6 +227,31 @@ pub struct Touched {
208 pub paths: Vec<String>,227 pub paths: Vec<String>,
209}228}
210229
230/// What a commit changed in a section a host serves: from the image with stamp `base`, the
231/// image `length` long with `writes` in place. A guest holding the base needs read nothing.
232#[derive(Clone, Debug, PartialEq, Encode, Decode)]
233#[cbor(map)]
234pub struct Delta {
235 #[n(0)]
236 pub path: String,
237 #[n(1)]
238 pub base: WireStamp,
239 #[n(2)]
240 pub length: u64,
241 #[n(3)]
242 pub writes: Vec<Written>,
243}
244
245/// Bytes at an offset.
246#[derive(Clone, Debug, PartialEq, Encode, Decode)]
247#[cbor(map)]
248pub struct Written {
249 #[n(0)]
250 pub offset: u64,
251 #[cbor(n(1), with = "minicbor::bytes")]
252 pub bytes: Vec<u8>,
253}
254
211/// A storage request; its kind names the verb, and the verb what it carries.255/// A storage request; its kind names the verb, and the verb what it carries.
212#[derive(Clone, Debug, Default, PartialEq, Encode, Decode)]256#[derive(Clone, Debug, Default, PartialEq, Encode, Decode)]
213#[cbor(map)]257#[cbor(map)]
...@@ -490,12 +534,20 @@ pub fn open(...@@ -490,12 +534,20 @@ pub fn open(
490 }534 }
491 let theirs: Open =535 let theirs: Open =
492 minicbor::decode(&read_block(stream)?).map_err(|_| invalid("A malformed opening"))?;536 minicbor::decode(&read_block(stream)?).map_err(|_| invalid("A malformed opening"))?;
493 if theirs.version != VERSION || theirs.room != room {537 // A responder answers even a peer of another version, so that both ends can say which
494 return Err(invalid("The peer means another room or version"));538 // should update.
495 }
496 if side == Side::Responder {539 if side == Side::Responder {
497 write_block(stream, &minicbor::to_vec(&ours).map_err(io::Error::other)?)?;540 write_block(stream, &minicbor::to_vec(&ours).map_err(io::Error::other)?)?;
498 }541 }
542 if theirs.version != VERSION {
543 return Err(io::Error::new(
544 io::ErrorKind::Unsupported,
545 Version(theirs.version),
546 ));
547 }
548 if theirs.room != room {
549 return Err(invalid("The peer means another room"));
550 }
499 let key = pake551 let key = pake
500 .finish(&theirs.pake)552 .finish(&theirs.pake)
501 .map_err(|_| invalid("A malformed key exchange"))?;553 .map_err(|_| invalid("A malformed key exchange"))?;
crates/snowbound/src/live.rs+6-3
...@@ -451,18 +451,21 @@ impl State {...@@ -451,18 +451,21 @@ impl State {
451 });451 });
452 }452 }
453453
454 /// `host` shares the notebook at `location`: its folder's changes reach its guests.454 /// `host` shares the notebook at `location`: its folder's changes reach its guests, and
455 /// theirs its sections, at once.
455 fn hosting(&mut self, location: String, host: Arc<Host>) {456 fn hosting(&mut self, location: String, host: Arc<Host>) {
456 if let Some(background) = self457 if let Some(background) = self
457 .library_at(&location)458 .library_at(&location)
458 .and_then(|library| library.background.clone())459 .and_then(|library| library.background.clone())
459 {460 {
460 let host = Arc::downgrade(&host);461 background.set_settle(notebook::live::share::SETTLE);
462 let shared = Arc::downgrade(&host);
461 background.on_touched(Some(Box::new(move |paths| {463 background.on_touched(Some(Box::new(move |paths| {
462 if let Some(host) = host.upgrade() {464 if let Some(host) = shared.upgrade() {
463 host.touched(paths);465 host.touched(paths);
464 }466 }
465 })));467 })));
468 host.on_changed(Box::new(move |paths| background.touched(paths)));
466 }469 }
467 self.peers.hosts.insert(location, host);470 self.peers.hosts.insert(location, host);
468 }471 }
crates/snowbound/src/main.rs+10
...@@ -4226,6 +4226,16 @@ impl State {...@@ -4226,6 +4226,16 @@ impl State {
4226 }4226 }
4227 }4227 }
4228 Event::Failed(error) => eprintln!("Synchronization stopped: {error}"),4228 Event::Failed(error) => eprintln!("Synchronization stopped: {error}"),
4229 // Its guests hear of it sooner than the folder's watch would tell them.
4230 #[cfg(feature = "live")]
4231 Event::Attempt {
4232 status: notebook::EditStatus::Published { .. },
4233 ..
4234 } => {
4235 if let Some(host) = self.peers.hosts.get(&session.library.location) {
4236 host.touched(&[session.tabs[session.tab].path.clone()]);
4237 }
4238 }
4229 Event::Attempt { .. } | Event::Unreachable(_) => {}4239 Event::Attempt { .. } | Event::Unreachable(_) => {}
4230 }4240 }
4231 }4241 }
crates/snowbound/src/share.rs+8
...@@ -62,6 +62,14 @@ enum Status {...@@ -62,6 +62,14 @@ enum Status {
62fn refusal(refusal: &Refusal) -> String {62fn refusal(refusal: &Refusal) -> String {
63 match refusal {63 match refusal {
64 Refusal::Malformed => "Check the code. It looks like 7KQ-4MZ-9XR.".into(),64 Refusal::Malformed => "Check the code. It looks like 7KQ-4MZ-9XR.".into(),
65 Refusal::Version { theirs_newer: true } => {
66 "The person sharing has a newer Snowbound. Update Snowbound, then try again.".into()
67 }
68 Refusal::Version {
69 theirs_newer: false,
70 } => "The person sharing has an older \
71 Snowbound. Ask them to update it."
72 .into(),
65 Refusal::Wrong => "That code or password doesn’t open a notebook. Check it with the \73 Refusal::Wrong => "That code or password doesn’t open a notebook. Check it with the \
66 person sharing."74 person sharing."
67 .into(),75 .into(),