authorgravatar for git@paperclover.netclover caruso <git@paperclover.net> 2026-10-02 23:06:34-07:00
committergravatar for git@paperclover.netclover caruso <git@paperclover.net> 2026-10-03 05:26:36-07:00
log5e6642680a2a6d3e0dcdc018cb7a021a7e010daa
tree193a6f17294a17cf69c19ef6fae5494855c5e5ae
parenta312c05b23b42586a70cb0be2c71d96e2b51eaf8
signature Signed by SSH key SHA256:52mNGHRsVFBDED9IAX5pe+LRWUefqTbxEReunq21QvU

feat: Live Share holds a room of thirty: the relay copies presence and a host's news, and guests catch up by changes

In a notebook's room a frame goes to the relay once and the relay copies it, to everyone or to the slots named (relay::GROUP, BROADCAST), sealed under a key of the sender's own and numbered, with its count of broadcasts, so a repeat, a reorder, a forgery or a lost broadcast still breaks the connection. Only a guest and its host keep a stream through the relay; hellos, presence and Touched go to everyone, deltas to the guests holding that section. Presence goes at most ten times a second. Touched carries each file's stamp hash, so a guest holding that image knows the stamp without asking, and a guest that missed a change reads it as writes from the image it holds. A room holds 64 peers and admits 120 a minute, a join told to wait a few seconds waits, a room's byte budget counts what the relay gives out, and /health counts bytes in and out. examples/live_crowd.rs runs 30 guests, 15 typing, through a local relay: relay traffic falls from about 2 MB/s each way to 0.15 in and 0.9 out, relay memory halves, and carets arrive in a few milliseconds. Keys into one section from fifteen typists still take seconds, as each commit waits its turn at the host. The relay must be redeployed before clients of this version meet through it. Assisted-by: claude-opus-5.5

14 files changed, 1511 insertions(+), 134 deletions(-)

crates/notebook/Cargo.toml+4
......@@ -63,6 +63,10 @@ tempfile = "3"
6363name = "live_latency"
6464required-features = ["live"]
6565
66[[example]]
67name = "live_crowd"
68required-features = ["live"]
69
6670[[example]]
6771name = "smb_offline_client"
6872required-features = ["smb"]
crates/notebook/examples/live_crowd.rs created+564
......@@ -0,0 +1,564 @@
1//! A crowd on one shared notebook: a host, and guests that join it through a relay, half of
2//! them typing on the same page. Reports how long keys and carets take to reach the others,
3//! what the relay passes on, and what the relay, the host and the guests cost.
4//!
5//! `cargo run --release -p notebook --features live --example live_crowd -- RELAY_BINARY
6//! [PEERS=30] [TYPISTS=PEERS/2] [SECONDS=60]`, the relay's binary built with
7//! `cargo build --release -p relay`. The relay, the host and the guests each run as a process
8//! of their own, so `ps` tells their costs apart.
9//!
10//! For watching a crowd in the app: `fixture FOLDER PARAGRAPHS` writes the notebook the
11//! host shares, and `guests URL CODE PEERS SECONDS` has that many join the app's share of it
12//! and move their carets about its page.
13
14use notebook::{
15 Replica,
16 live::{
17 Caret, Guid, Hello, Presence, Spot,
18 share::{self, Guest, Host, Sharing},
19 },
20 session::{Background, Notebook, Section},
21};
22use onestore::{
23 ExGuid,
24 op::{Edit, Op, PageOp},
25};
26use std::{
27 collections::HashMap,
28 io::{BufRead, BufReader, Read, Write},
29 net::TcpStream,
30 path::Path,
31 process::{Child, Command, Stdio},
32 sync::{
33 Arc, Mutex,
34 atomic::{AtomicBool, Ordering},
35 },
36 thread,
37 time::{Duration, Instant},
38};
39
40/// Keys each typist types a second, and how often observers look.
41const KEYS_PER_SECOND: f64 = 5.0;
42const LOOK: Duration = Duration::from_millis(2);
43/// Guests that time what they see.
44const OBSERVERS: usize = 3;
45
46fn until(what: &str, limit: Duration, done: impl Fn() -> bool) {
47 let deadline = Instant::now() + limit;
48 while !done() {
49 assert!(Instant::now() < deadline, "{what}");
50 thread::sleep(Duration::from_millis(5));
51 }
52}
53
54fn hello(name: &str) -> Hello {
55 Hello::new(name.into(), None).unwrap()
56}
57
58fn main() {
59 let args: Vec<String> = std::env::args().skip(1).collect();
60 match args.first().map(String::as_str) {
61 Some("host") => return host(&args[1], args[2].parse().unwrap()),
62 Some("fixture") => {
63 fixture(Path::new(&args[1]), args[2].parse().unwrap());
64 return;
65 }
66 Some("guests") => {
67 let seconds = Duration::from_secs(args[4].parse().unwrap());
68 return guests(&args[1], &args[2], args[3].parse().unwrap(), seconds);
69 }
70 _ => {}
71 }
72 let binary = args.first().expect("the relay's binary");
73 let number = |at: usize, default: usize| args.get(at).map_or(default, |n| n.parse().unwrap());
74 let peers = number(1, 30);
75 let typists = number(2, peers / 2);
76 let seconds = number(3, 60) as u64;
77 crowd(binary, peers, typists, Duration::from_secs(seconds));
78}
79
80/// Writes the notebook `Garden` in `directory`: a section whose page has `paragraphs`
81/// paragraphs, `Typed:` then `Note 1:` and on. Its folder.
82fn fixture(directory: &Path, paragraphs: usize) -> std::path::PathBuf {
83 let folder = directory.join("Garden");
84 std::fs::create_dir_all(&folder).unwrap();
85 let mut image = onestore::create_section("Garden.one", "Typed:", "Fixture").unwrap();
86 for n in (1..paragraphs).rev() {
87 let arena = onestore::Arena::default();
88 let mut section = onestore::Section::open(&arena, image).unwrap();
89 let (space, ..) = section.pages().unwrap()[0];
90 let page = section.page(space).unwrap();
91 let (text, _) = texts(&page).swap_remove(0);
92 let right = onestore::page::text::new_id().unwrap();
93 let ops = [
94 PageOp::Split {
95 text,
96 at: 6,
97 paragraph: onestore::page::text::new_id().unwrap(),
98 right,
99 lists: Vec::new(),
100 },
101 PageOp::Text {
102 text: right,
103 range: 0..0,
104 with: format!("Note {n}:"),
105 },
106 ];
107 let edit = Edit {
108 at: 134_000_000_000_000_000,
109 ops: ops.into_iter().map(|op| Op::Page { space, op }).collect(),
110 };
111 section.apply("Fixture", &edit).unwrap();
112 section.seal().unwrap();
113 image = section.image();
114 }
115 std::fs::write(folder.join("Garden.one"), image).unwrap();
116 folder
117}
118
119/// The host's process: shares a fixture notebook of a page with `paragraphs` paragraphs
120/// through `url`, says its code, and reports what the section grew to once its input closes.
121fn host(url: &str, paragraphs: usize) {
122 let directory = tempfile::tempdir().unwrap();
123 let folder = fixture(directory.path(), paragraphs);
124 let file = folder.join("Garden.one");
125 let storage = Notebook::open(&folder, directory.path().join("cache"))
126 .unwrap()
127 .into_storage();
128 let host = Host::start(
129 storage,
130 hello("Host"),
131 Sharing::new("").unwrap(),
132 "Garden",
133 None,
134 Some(url),
135 || {},
136 )
137 .unwrap();
138 until("no code", Duration::from_secs(30), || {
139 host.code().is_some_and(|code| share::code(&code).is_some())
140 });
141 println!("code {}", host.code().unwrap());
142 let _ = std::io::stdin().read_to_end(&mut Vec::new());
143 let image = std::fs::read(&file).unwrap();
144 let size = image.len();
145 let arena = onestore::Arena::default();
146 let mut section = onestore::Section::open(&arena, image).unwrap();
147 let pages = section.pages().unwrap().len();
148 let conflicts: usize = section
149 .conflicts()
150 .unwrap()
151 .iter()
152 .map(|(_, c)| c.len())
153 .sum();
154 println!("section {size} bytes, {pages} pages, {conflicts} conflict pages");
155}
156
157/// A guest as the app holds one: its share, the notebook's background, the section open
158/// and its file's identity.
159struct Open {
160 guest: Arc<Guest>,
161 _background: Background,
162 section: Section,
163 file: [u8; 16],
164}
165
166fn join(name: &str, code: &str, url: &str, cache: &Path) -> Open {
167 let welcome = share::join(hello(name), code, "", None, Some(url)).unwrap();
168 let guest = Guest::start(
169 hello(name),
170 welcome.share,
171 welcome.secret,
172 None,
173 Some(url),
174 || {},
175 )
176 .unwrap();
177 until("no host", Duration::from_secs(120), || {
178 guest.host().is_some()
179 });
180 let mut notebook = Notebook::open_hosted(Arc::clone(&guest), cache).unwrap();
181 let background = Background::hosted(Arc::clone(&guest), || {}).unwrap();
182 background.watch(notebook.replicas());
183 let replica = notebook.replica_path("Garden.one").unwrap();
184 std::fs::create_dir_all(replica.parent().unwrap()).unwrap();
185 let image = notebook.read_section("Garden.one").unwrap();
186 let file = onestore::Header::parse(&image[..1024]).unwrap().file_id;
187 let replica = Replica::open_or_create(&replica, None, || Ok(image)).unwrap();
188 let section =
189 Section::resume_hosted("Garden.one".into(), replica, Arc::clone(&guest), || {}).unwrap();
190 background.hold("Garden.one", &section);
191 Open {
192 guest,
193 _background: background,
194 section,
195 file,
196 }
197}
198
199/// Where in the open page's `paragraph`th text, `offset` units in, `open`'s guest is.
200fn presence(open: &Open, paragraph: usize, offset: u32) -> Option<Presence> {
201 let (space, texts) = texts_of(&open.section)?;
202 let (text, _) = texts.get(paragraph)?;
203 let at = Spot {
204 text: Guid {
205 guid: text.guid,
206 n: text.n,
207 },
208 offset,
209 };
210 Some(Presence {
211 section: Some(open.file),
212 page: Some(Guid {
213 guid: space.guid,
214 n: space.n,
215 }),
216 caret: Some(Caret {
217 anchor: at,
218 focus: at,
219 }),
220 })
221}
222
223/// `peers` guests join the share `code` names through `url` and, for `seconds`, put their
224/// carets in the page's paragraphs, each moving every few seconds.
225fn guests(url: &str, code: &str, peers: usize, seconds: Duration) {
226 let directory = tempfile::tempdir().unwrap();
227 let names = [
228 "Ada", "Grace", "Alan", "Barbara", "Edsger", "Frances", "Donald", "Margaret", "Ken",
229 "Radia", "Dennis", "Hedy", "Linus", "Karen", "John", "Sophie",
230 ];
231 let joining: Vec<_> = (0..peers)
232 .map(|n| {
233 let name = format!("{} {n}", names[n % names.len()]);
234 let (code, url, cache) = (
235 code.to_owned(),
236 url.to_owned(),
237 directory.path().join(n.to_string()),
238 );
239 thread::sleep(Duration::from_millis(100));
240 thread::spawn(move || join(&name, &code, &url, &cache))
241 })
242 .collect();
243 let guests: Vec<Open> = joining
244 .into_iter()
245 .map(|joined| joined.join().unwrap())
246 .collect();
247 let began = Instant::now();
248 let mut turn = 0;
249 while began.elapsed() < seconds {
250 for (n, open) in guests.iter().enumerate() {
251 let paragraphs = texts_of(&open.section).map_or(1, |(_, texts)| texts.len());
252 if (n + turn) % 4 == 0 || turn == 0 {
253 let at = (n * 7 + turn) % paragraphs;
254 if let Some(presence) = presence(open, at, (turn % 5) as u32) {
255 open.guest.set_presence(presence);
256 }
257 }
258 }
259 turn += 1;
260 thread::sleep(Duration::from_secs(1));
261 }
262}
263
264/// The page's object space and its paragraphs' texts: each one's object and what it reads.
265fn texts_of(section: &Section) -> Option<(ExGuid, Vec<(ExGuid, String)>)> {
266 let (space, ..) = *section.pages().ok()?.first()?;
267 let page = section.page(space).ok()?;
268 Some((space, texts(&page)))
269}
270
271fn texts(page: &onestore::page::Page) -> Vec<(ExGuid, String)> {
272 let outlines = page.objects.iter().filter_map(|object| match object {
273 onestore::page::PageObject::Outline(outline) => Some(outline),
274 _ => None,
275 });
276 outlines
277 .flat_map(|outline| &outline.paragraphs)
278 .filter_map(|p| p.text().map(|t| (t.id, t.text.text().to_owned())))
279 .collect()
280}
281
282/// The key typist `typist` types `n`th, unique to both.
283fn key(typist: usize, n: usize) -> char {
284 char::from_u32(0x4E00 + (typist * 600 + n) as u32).unwrap()
285}
286
287fn typist_of(key: char) -> Option<(usize, usize)> {
288 let at = (key as u32).checked_sub(0x4E00)? as usize;
289 Some((at / 600, at % 600))
290}
291
292/// Seconds of CPU and kilobytes resident `ps` reports for `pid`.
293fn cost(pid: u32) -> (f64, u64) {
294 let output = Command::new("ps")
295 .args(["-o", "time=,rss=", "-p", &pid.to_string()])
296 .output()
297 .unwrap();
298 let text = String::from_utf8_lossy(&output.stdout);
299 let mut words = text.split_whitespace();
300 let time = words.next().unwrap_or("0:0");
301 let seconds = time.split(':').fold(0.0, |total, part| {
302 total * 60.0 + part.parse::<f64>().unwrap_or(0.0)
303 });
304 (
305 seconds,
306 words.next().and_then(|rss| rss.parse().ok()).unwrap_or(0),
307 )
308}
309
310/// The relay's `/health`, as numbers by name.
311fn health(address: &str) -> HashMap<String, f64> {
312 let mut stream = TcpStream::connect(address).unwrap();
313 write!(
314 stream,
315 "GET /health HTTP/1.1\r\nHost: relay\r\nConnection: close\r\n\r\n"
316 )
317 .unwrap();
318 let mut text = String::new();
319 stream.read_to_string(&mut text).unwrap();
320 let body = text.split("\r\n\r\n").nth(1).unwrap_or_default();
321 body.trim()
322 .trim_matches(['{', '}'])
323 .split(',')
324 .filter_map(|pair| {
325 let (name, value) = pair.split_once(':')?;
326 Some((name.trim_matches('"').to_owned(), value.parse().ok()?))
327 })
328 .collect()
329}
330
331fn report(what: &str, mut times: Vec<Duration>) {
332 if times.is_empty() {
333 println!("{what}: none seen");
334 return;
335 }
336 times.sort();
337 let ms = |at: usize| times[at].as_secs_f64() * 1000.0;
338 println!(
339 "{what}: median {:.0} ms, p90 {:.0} ms, p99 {:.0} ms, worst {:.0} ms ({} seen)",
340 ms(times.len() / 2),
341 ms(times.len() * 9 / 10),
342 ms(times.len() * 99 / 100),
343 ms(times.len() - 1),
344 times.len()
345 );
346}
347
348struct Spawned(Child);
349
350impl Drop for Spawned {
351 fn drop(&mut self) {
352 let _ = self.0.kill();
353 let _ = self.0.wait();
354 }
355}
356
357fn crowd(binary: &str, peers: usize, typists: usize, window: Duration) {
358 let port = std::net::TcpListener::bind("127.0.0.1:0")
359 .unwrap()
360 .local_addr()
361 .unwrap()
362 .port();
363 let address = format!("127.0.0.1:{port}");
364 // Everyone comes from one address here, which a relay would take for one abuser.
365 let relay = Spawned(
366 Command::new(binary)
367 .args(["--listen", &address])
368 .args(["--max-connections-per-address", "1000"])
369 .args(["--joins-per-minute", "10000"])
370 .args(["--failures-per-minute", "10000"])
371 .stdout(Stdio::null())
372 .spawn()
373 .unwrap(),
374 );
375 until("no relay", Duration::from_secs(10), || {
376 TcpStream::connect(&address).is_ok()
377 });
378 let url = format!("ws://{address}");
379 let mut hosting = Spawned(
380 Command::new(std::env::current_exe().unwrap())
381 .args(["host", &url, &typists.max(1).to_string()])
382 .stdin(Stdio::piped())
383 .stdout(Stdio::piped())
384 .spawn()
385 .unwrap(),
386 );
387 let mut said = BufReader::new(hosting.0.stdout.take().unwrap()).lines();
388 let code = said.next().unwrap().unwrap();
389 let code = code.strip_prefix("code ").unwrap().to_owned();
390 println!("{peers} guests, {typists} typing {KEYS_PER_SECOND} keys a second, code {code}");
391
392 let directory = tempfile::tempdir().unwrap();
393 let began = Instant::now();
394 let joining: Vec<_> = (0..peers)
395 .map(|n| {
396 let (code, url) = (code.clone(), url.clone());
397 let cache = directory.path().join(format!("guest{n}"));
398 thread::sleep(Duration::from_millis(50));
399 thread::spawn(move || join(&format!("Guest {n}"), &code, &url, &cache))
400 })
401 .collect();
402 let guests: Vec<Arc<Open>> = joining
403 .into_iter()
404 .map(|joined| Arc::new(joined.join().unwrap()))
405 .collect();
406 println!("all joined in {:.1} s", began.elapsed().as_secs_f64());
407 until("not everyone met", Duration::from_secs(300), || {
408 guests.iter().all(|open| open.guest.peers().len() >= peers)
409 });
410 println!("all met in {:.1} s", began.elapsed().as_secs_f64());
411
412 let (_, texts) = texts_of(&guests[0].section).unwrap();
413 assert!(texts.len() >= typists, "a paragraph for each typist");
414 for open in &guests {
415 open.guest.set_presence(presence(open, 0, 0).unwrap());
416 }
417
418 // When each typist typed each key, and said its caret moved past it.
419 let typed: Arc<Mutex<HashMap<(usize, usize), Instant>>> = Arc::default();
420 let moved: Arc<Mutex<HashMap<(usize, usize), Instant>>> = Arc::default();
421 let seen_keys: Arc<Mutex<Vec<Duration>>> = Arc::default();
422 let seen_carets: Arc<Mutex<Vec<Duration>>> = Arc::default();
423 let done = Arc::new(AtomicBool::new(false));
424 let (relay_before, host_before, crowd_before) = (
425 cost(relay.0.id()),
426 cost(hosting.0.id()),
427 cost(std::process::id()),
428 );
429 let bytes_before = health(&address);
430 let started = Instant::now();
431
432 let mut threads = Vec::new();
433 for (n, open) in guests.iter().enumerate().take(typists) {
434 let (open, typed, done) = (Arc::clone(open), Arc::clone(&typed), Arc::clone(&done));
435 let moved = Arc::clone(&moved);
436 threads.push(thread::spawn(move || {
437 let pause = Duration::from_secs_f64(1.0 / KEYS_PER_SECOND);
438 thread::sleep(pause.mul_f64(n as f64 / typists as f64));
439 let mut count = 0;
440 while !done.load(Ordering::Acquire) && count < 600 {
441 open.section.events();
442 let Some((space, mut texts)) = texts_of(&open.section) else {
443 continue;
444 };
445 let (text, before) = texts.swap_remove(n);
446 let end = before.encode_utf16().count() as u32;
447 let edit = Edit {
448 at: 134_000_000_000_000_000,
449 ops: vec![Op::Page {
450 space,
451 op: PageOp::Text {
452 text,
453 range: end..end,
454 with: key(n, count).to_string(),
455 },
456 }],
457 };
458 typed.lock().unwrap().insert((n, count), Instant::now());
459 open.section.replica().apply("Typist", edit).unwrap();
460 count += 1;
461 if let Some(presence) = presence(&open, 0, count as u32) {
462 moved.lock().unwrap().insert((n, count), Instant::now());
463 open.guest.set_presence(presence);
464 }
465 thread::sleep(pause);
466 }
467 }));
468 }
469 for open in guests.iter().skip(typists).take(OBSERVERS) {
470 let (open, typed, done) = (Arc::clone(open), Arc::clone(&typed), Arc::clone(&done));
471 let moved = Arc::clone(&moved);
472 let (keys, carets) = (Arc::clone(&seen_keys), Arc::clone(&seen_carets));
473 threads.push(thread::spawn(move || {
474 let mut keys_seen = std::collections::HashSet::new();
475 let mut carets_seen = HashMap::new();
476 while !done.load(Ordering::Acquire) {
477 open.section.events();
478 let now = Instant::now();
479 if let Some((_, texts)) = texts_of(&open.section) {
480 let chars = texts.iter().flat_map(|(_, text)| text.chars());
481 for (typist, n) in chars.filter_map(typist_of) {
482 if keys_seen.insert((typist, n))
483 && let Some(at) = typed.lock().unwrap().get(&(typist, n))
484 {
485 keys.lock().unwrap().push(now - *at);
486 }
487 }
488 }
489 for peer in open.guest.peers() {
490 let Some(typist) = (peer.hello.name.strip_prefix("Guest "))
491 .and_then(|n| n.parse::<usize>().ok())
492 .filter(|n| *n < typists)
493 else {
494 continue;
495 };
496 let Some(count) = peer.presence.and_then(|p| p.caret).map(|c| c.focus.offset)
497 else {
498 continue;
499 };
500 let count = count as usize;
501 if count > 0
502 && carets_seen.insert(typist, count) != Some(count)
503 && let Some(at) = moved.lock().unwrap().get(&(typist, count))
504 {
505 carets.lock().unwrap().push(now - *at);
506 }
507 }
508 thread::sleep(LOOK);
509 }
510 }));
511 }
512 thread::sleep(window);
513 done.store(true, Ordering::Release);
514 let elapsed = started.elapsed().as_secs_f64();
515 let (relay_after, host_after, crowd_after) = (
516 cost(relay.0.id()),
517 cost(hosting.0.id()),
518 cost(std::process::id()),
519 );
520 let bytes_after = health(&address);
521 for thread in threads {
522 thread.join().unwrap();
523 }
524 let keys = typed.lock().unwrap().len();
525 report(
526 "keys, typist to observer",
527 std::mem::take(&mut seen_keys.lock().unwrap()),
528 );
529 report(
530 "carets, typist to observer",
531 std::mem::take(&mut seen_carets.lock().unwrap()),
532 );
533 let delta = |name: &str| bytes_after[name] - bytes_before.get(name).copied().unwrap_or(0.0);
534 let members = (peers + 1) as f64;
535 println!(
536 "relay passed on {:.0} KB/s in, {:.0} KB/s out: per peer {:.1} KB/s in, {:.1} KB/s out",
537 delta("bytes_in") / elapsed / 1000.0,
538 delta("bytes_out") / elapsed / 1000.0,
539 delta("bytes_in") / elapsed / 1000.0 / members,
540 delta("bytes_out") / elapsed / 1000.0 / members,
541 );
542 let load = |name: &str, (before, _): (f64, u64), (after, rss): (f64, u64)| {
543 println!(
544 "{name}: {:.0}% of a core, {:.1} MB resident",
545 (after - before) / elapsed * 100.0,
546 rss as f64 / 1024.0
547 );
548 };
549 load("relay", relay_before, relay_after);
550 load("host", host_before, host_after);
551 load(&format!("{peers} guests"), crowd_before, crowd_after);
552
553 // Every key typed should reach the host's file; conflict pages show as extra pages.
554 thread::sleep(Duration::from_secs(5));
555 let observed = texts_of(&guests[typists.min(peers - 1)].section).map_or(0, |(_, texts)| {
556 let chars = texts.iter().flat_map(|(_, text)| text.chars());
557 chars.filter_map(typist_of).count()
558 });
559 println!("{keys} keys typed, {observed} on an observer's page after 5 s");
560 drop(hosting.0.stdin.take());
561 for line in said.map_while(Result::ok) {
562 println!("host: {line}");
563 }
564}
crates/notebook/src/live.rs+151-24
......@@ -8,6 +8,7 @@
88//! from scratch.
99
1010pub use ::relay::code;
11mod group;
1112pub mod proxy;
1213mod relay;
1314pub mod share;
......@@ -38,6 +39,8 @@ const SERVICE: &str = "_snowbound._tcp.local.";
3839const PING: Duration = Duration::from_secs(15);
3940const GONE: Duration = Duration::from_secs(45);
4041const OPENING: Duration = Duration::from_secs(5);
42/// The most often presence goes to a peer.
43const PRESENCE_EVERY: Duration = Duration::from_millis(100);
4144/// The longest wait before meeting again.
4245const PATIENCE: Duration = Duration::from_secs(30);
4346/// Wrong tries of a code met off any relay before it admits no one new, as a relay burns one.
......@@ -137,16 +140,15 @@ pub struct Peer {
137140pub enum Event<'a> {
138141 /// The peers, their presence, the code or the relay's answer changed.
139142 Changed,
140 /// A peer was met, with the line to it.
143 /// A stream to a peer opened, with the line to it.
141144 Met(&'a Arc<Hello>, &'a Line),
142 /// A peer's connection ended.
145 /// A peer's stream ended.
143146 Left(&'a Arc<Hello>),
144 /// A frame of a kind presence doesn't read itself, with the line to answer on.
147 /// A frame of a kind presence doesn't read itself, on a stream or to the group.
145148 Frame {
146149 from: &'a Arc<Hello>,
147150 kind: u16,
148151 body: &'a [u8],
149 line: &'a Line,
150152 },
151153}
152154
......@@ -200,6 +202,48 @@ struct State {
200202 daemon: Option<ServiceDaemon>,
201203}
202204
205impl State {
206 /// The relay's group, in a notebook's room while the relay is reached.
207 fn group(&self) -> Option<&relay::Group> {
208 self.relay.as_ref()?.group.as_ref()
209 }
210}
211
212/// Presence sent at most every `PRESENCE_EVERY`, and only the newest.
213#[derive(Default)]
214struct Paced {
215 due: Option<Instant>,
216 last: Option<Instant>,
217 sent: Option<u64>,
218}
219
220impl Paced {
221 /// How long to wait for news: until presence is due, else `idle`.
222 fn wait(&self, idle: Duration) -> Duration {
223 self.due
224 .map_or(idle, |due| due.saturating_duration_since(Instant::now()))
225 }
226
227 /// Hears that presence changed.
228 fn changed(&mut self) {
229 let last = self.last;
230 self.due
231 .get_or_insert_with(|| last.map_or_else(Instant::now, |at| at + PRESENCE_EVERY));
232 }
233
234 /// The presence to send now, where it is due and new.
235 fn due(&mut self, state: &Mutex<State>) -> Option<Presence> {
236 self.due.filter(|due| *due <= Instant::now())?;
237 self.due = None;
238 let state = state.lock().unwrap();
239 if self.sent == Some(state.generation) {
240 return None;
241 }
242 (self.sent, self.last) = (Some(state.generation), Some(Instant::now()));
243 Some(state.presence.clone())
244 }
245}
246
203247struct Link {
204248 connection: u64,
205249 peer: Peer,
......@@ -401,7 +445,8 @@ impl Live {
401445 thread::spawn(move || shared.dial(address));
402446 }
403447
404 /// Says where this end is now; peers hear only the newest of quick changes.
448 /// Says where this end is now; peers hear only the newest of quick changes, at most every
449 /// tenth of a second.
405450 pub fn set_presence(&self, presence: Presence) {
406451 let mut state = self.shared.state.lock().unwrap();
407452 if state.presence == presence {
......@@ -410,14 +455,23 @@ impl Live {
410455 state.presence = presence;
411456 state.generation += 1;
412457 for link in state.peers.values() {
413 let _ = link.line.0.send(Out::Presence);
458 if self.shared.carries_presence(&*link.pipe) {
459 let _ = link.line.0.send(Out::Presence);
460 }
461 }
462 if let Some(group) = state.group() {
463 group.send(relay::Out::Presence);
414464 }
415465 }
416466
417 /// The peers connected now, by id.
467 /// The peers in the room now, by id.
418468 pub fn peers(&self) -> Vec<Peer> {
419 let state = self.shared.state.lock().unwrap();
420 state.peers.values().map(|link| link.peer.clone()).collect()
469 self.shared.peers()
470 }
471
472 /// A way to send to this room's peers that doesn't keep it open.
473 pub fn sender(&self) -> Sender {
474 Sender(Arc::downgrade(&self.shared))
421475 }
422476
423477 /// The line to `peer`, while it is connected.
......@@ -442,6 +496,45 @@ impl Live {
442496 }
443497}
444498
499/// Sends to a room's peers while the room is open.
500#[derive(Clone)]
501pub struct Sender(std::sync::Weak<Shared>);
502
503impl Sender {
504 /// Sends message `kind` holding `body` to the peers `to` names, or to everyone: once to
505 /// the room's group through the relay, and to each peer met directly.
506 pub fn send(&self, kind: u16, body: &impl Encode<()>, to: Option<&[[u8; 16]]>) {
507 let (Some(shared), Ok(body)) = (self.0.upgrade(), minicbor::to_vec(body)) else {
508 return;
509 };
510 let state = shared.state.lock().unwrap();
511 let group = state.group();
512 let named = |id: &[u8; 16]| to.is_none_or(|to| to.contains(id));
513 let mut slots = Vec::new();
514 for (id, link) in state.peers.iter().filter(|(id, _)| named(id)) {
515 match group {
516 Some(group) if !link.pipe.direct() => slots.extend(group.slot(id)),
517 _ => {
518 let _ = link.line.0.send(Out::Frame(kind, body.clone()));
519 }
520 }
521 }
522 let Some(group) = group else {
523 return;
524 };
525 match to {
526 None => group.send(relay::Out::Frame(kind, body, None)),
527 Some(to) => {
528 let unlinked = to.iter().filter(|id| !state.peers.contains_key(*id));
529 slots.extend(unlinked.filter_map(|id| group.slot(id)));
530 if !slots.is_empty() {
531 group.send(relay::Out::Frame(kind, body, Some(slots)));
532 }
533 }
534 }
535 }
536}
537
445538impl Drop for Live {
446539 fn drop(&mut self) {
447540 self.shared.stopped.store(true, Ordering::Release);
......@@ -461,6 +554,38 @@ impl Drop for Live {
461554}
462555
463556impl Shared {
557 /// The peers in the room: those met directly, and the rest as the relay's group or a
558 /// stream through the relay last heard of them.
559 fn peers(&self) -> Vec<Peer> {
560 let state = self.state.lock().unwrap();
561 let mut peers = BTreeMap::new();
562 if let Some(group) = state.group() {
563 for member in group.members.lock().unwrap().values() {
564 if let Some(peer) = &member.peer {
565 peers.insert(peer.hello.peer, peer.clone());
566 }
567 }
568 }
569 for (id, link) in &state.peers {
570 if link.pipe.direct() || !peers.contains_key(id) {
571 peers.insert(*id, link.peer.clone());
572 }
573 }
574 peers.into_values().collect()
575 }
576
577 /// Whether peer `id` is met directly, which is where it says everything.
578 fn direct(&self, id: &[u8; 16]) -> bool {
579 let state = self.state.lock().unwrap();
580 state.peers.get(id).is_some_and(|link| link.pipe.direct())
581 }
582
583 /// Whether presence goes on `pipe`: not on a stream through the relay in a notebook's
584 /// room, whose group carries it.
585 fn carries_presence(&self, pipe: &dyn Pipe) -> bool {
586 pipe.direct() || matches!(self.room, Room::Code { .. })
587 }
588
464589 /// The room's tag as it stands: a code's, once numbered.
465590 fn tag(&self) -> Option<String> {
466591 match &self.room {
......@@ -592,7 +717,9 @@ impl Shared {
592717 pipe.shutdown();
593718 return false;
594719 }
595 let _ = line.0.send(Out::Presence);
720 if self.carries_presence(&*pipe) {
721 let _ = line.0.send(Out::Presence);
722 }
596723 state.peers.insert(
597724 peer,
598725 Link {
......@@ -632,7 +759,6 @@ impl Shared {
632759 from: &hello,
633760 kind,
634761 body: &body,
635 line: &line,
636762 });
637763 // The peer hangs up after its bye.
638764 if kind == kind::BYE {
......@@ -660,22 +786,16 @@ impl Shared {
660786 true
661787 }
662788
663 /// Sends the newest presence whenever woken, frames as they come, and a ping when quiet.
789 /// Sends the newest presence when woken, at most every `PRESENCE_EVERY`, frames as they
790 /// come, and a ping when quiet.
664791 fn write(&self, pipe: &dyn Pipe, mut send: Sealer, outgoing: mpsc::Receiver<Out>) {
665792 let mut stream = pipe;
666 let mut sent = None;
793 let mut paced = Paced::default();
667794 loop {
668 let result = match outgoing.recv_timeout(PING) {
795 let result = match outgoing.recv_timeout(paced.wait(PING)) {
669796 Ok(Out::Presence) => {
670 let presence = {
671 let state = self.state.lock().unwrap();
672 if sent == Some(state.generation) {
673 continue;
674 }
675 sent = Some(state.generation);
676 state.presence.clone()
677 };
678 send.send(&mut stream, kind::PRESENCE, &presence)
797 paced.changed();
798 Ok(())
679799 }
680800 Ok(Out::Frame(kind, body)) => send.send_encoded(&mut stream, kind, &body),
681801 Ok(Out::Bye(body)) => {
......@@ -683,9 +803,16 @@ impl Shared {
683803 pipe.shutdown();
684804 return;
685805 }
686 Err(mpsc::RecvTimeoutError::Timeout) => send.send(&mut stream, kind::PING, &()),
806 Err(mpsc::RecvTimeoutError::Timeout) if paced.due.is_none() => {
807 send.send(&mut stream, kind::PING, &())
808 }
809 Err(mpsc::RecvTimeoutError::Timeout) => Ok(()),
687810 Err(mpsc::RecvTimeoutError::Disconnected) => return,
688811 };
812 let result = result.and_then(|()| match paced.due(&self.state) {
813 Some(presence) => send.send(&mut stream, kind::PRESENCE, &presence),
814 None => Ok(()),
815 });
689816 if result.is_err() {
690817 pipe.shutdown();
691818 return;
crates/notebook/src/live/group.rs created+204
......@@ -0,0 +1,204 @@
1//! Frames a peer in a notebook's room sends once through the relay, which copies them to
2//! everyone or to the peers named (`relay::GROUP`): presence, hellos, and a host's news of
3//! its files. Each is sealed under a key of its sender's own, derived from the room's secret
4//! and an id the sender picks for each connection, so the relay sees only ciphertext.
5//!
6//! A frame carries its number, which is its nonce, and the count of broadcasts its sender has
7//! sent: a frame numbered no later than the last, or a broadcast missing before it, means the
8//! relay lost, repeated, reordered or forged one, and this end meets the room again. A frame
9//! to some peers goes unseen by the others, so only what can tell itself stale is sent so.
10//! The room's members hold one key, so a frame proves its sender is in the room, not which
11//! member it is; members are trusted alike, as each may change the notebook itself.
12
13use super::{Hello, Peer};
14use aes_gcm::{
15 Aes256Gcm, KeyInit,
16 aead::{Aead, Payload},
17};
18use hmac::{Hmac, Mac};
19use sha2::Sha256;
20use std::io;
21
22/// A frame's sender id, number, broadcasts, and whether it is a broadcast.
23const HEADER: usize = 16 + 8 + 8 + 1;
24
25/// The room's group key.
26pub(super) struct Keys([u8; 32]);
27
28impl Keys {
29 pub(super) fn new(secret: &[u8]) -> Self {
30 Self(mac(secret, b"Snowbound live v2 group"))
31 }
32
33 fn cipher(&self, sender: &[u8; 16]) -> Aes256Gcm {
34 Aes256Gcm::new_from_slice(&mac(&self.0, sender)).expect("a 32-byte key")
35 }
36}
37
38fn mac(key: &[u8], message: &[u8]) -> [u8; 32] {
39 let mut mac = <Hmac<Sha256> as hmac::KeyInit>::new_from_slice(key).expect("any key length");
40 mac.update(message);
41 mac.finalize().into_bytes().into()
42}
43
44/// Seals this end's frames for one connection to the relay.
45pub(super) struct Sealer {
46 id: [u8; 16],
47 cipher: Aes256Gcm,
48 number: u64,
49 broadcasts: u64,
50}
51
52impl Sealer {
53 pub(super) fn new(keys: &Keys) -> io::Result<Self> {
54 let mut id = [0; 16];
55 getrandom::fill(&mut id).map_err(|_| io::Error::other("System random source failed"))?;
56 Ok(Self {
57 cipher: keys.cipher(&id),
58 id,
59 number: 0,
60 broadcasts: 0,
61 })
62 }
63
64 /// Message `kind` holding `body`, to everyone where `broadcast`.
65 pub(super) fn seal(&mut self, kind: u16, body: &[u8], broadcast: bool) -> io::Result<Vec<u8>> {
66 self.number += 1;
67 self.broadcasts += u64::from(broadcast);
68 let mut header = Vec::with_capacity(HEADER);
69 header.extend_from_slice(&self.id);
70 header.extend_from_slice(&self.number.to_be_bytes());
71 header.extend_from_slice(&self.broadcasts.to_be_bytes());
72 header.push(u8::from(broadcast));
73 let clear = [&kind.to_be_bytes()[..], body].concat();
74 let sealed = self
75 .cipher
76 .encrypt(
77 &nonce(self.number).into(),
78 Payload {
79 msg: &clear,
80 aad: &header,
81 },
82 )
83 .map_err(|_| io::Error::other("A frame could not be sealed"))?;
84 Ok([header, sealed].concat())
85 }
86}
87
88fn nonce(number: u64) -> [u8; 12] {
89 let mut nonce = [0; 12];
90 nonce[4..].copy_from_slice(&number.to_be_bytes());
91 nonce
92}
93
94/// A peer in the room as its frames say: the id its frames are sealed under, how far they
95/// have come, and the peer once its hello has.
96pub(super) struct Member {
97 id: [u8; 16],
98 cipher: Aes256Gcm,
99 number: u64,
100 broadcasts: u64,
101 pub(super) peer: Option<Peer>,
102}
103
104impl Member {
105 /// The member whose first frame is `frame`.
106 pub(super) fn new(keys: &Keys, frame: &[u8]) -> io::Result<Self> {
107 let id: [u8; 16] = frame
108 .get(..16)
109 .and_then(|id| id.try_into().ok())
110 .ok_or_else(|| broken("A group frame without its sender"))?;
111 Ok(Self {
112 cipher: keys.cipher(&id),
113 id,
114 number: 0,
115 broadcasts: 0,
116 peer: None,
117 })
118 }
119
120 /// The kind and body of `frame`, the next from this member; an error of kind
121 /// `InvalidData` means the relay tampered with its frames.
122 pub(super) fn open(&mut self, frame: &[u8]) -> io::Result<(u16, Vec<u8>)> {
123 let (header, sealed) = frame
124 .split_at_checked(HEADER)
125 .ok_or_else(|| broken("A group frame cut short"))?;
126 let field = |at: usize| u64::from_be_bytes(header[at..at + 8].try_into().unwrap());
127 let (number, broadcasts, broadcast) = (field(16), field(24), header[32] == 1);
128 if header[..16] != self.id {
129 return Err(broken("A group frame from another sender in its slot"));
130 }
131 let mut clear = self
132 .cipher
133 .decrypt(
134 &nonce(number).into(),
135 Payload {
136 msg: sealed,
137 aad: header,
138 },
139 )
140 .map_err(|_| broken("A group frame that does not open"))?;
141 // The first frame heard from a member is where its count starts.
142 let first = self.number == 0;
143 let due = self.broadcasts + u64::from(broadcast);
144 if !first && (number <= self.number || broadcasts != due) {
145 return Err(broken("A group frame lost, repeated or out of order"));
146 }
147 (self.number, self.broadcasts) = (number, broadcasts);
148 let body = clear.split_off(2.min(clear.len()));
149 let kind = u16::from_be_bytes(clear.try_into().map_err(|_| broken("An empty frame"))?);
150 Ok((kind, body))
151 }
152
153 pub(super) fn hello(&self) -> Option<&std::sync::Arc<Hello>> {
154 self.peer.as_ref().map(|peer| &peer.hello)
155 }
156}
157
158fn broken(message: &'static str) -> io::Error {
159 io::Error::new(io::ErrorKind::InvalidData, message)
160}
161
162#[cfg(test)]
163mod tests {
164 use super::*;
165
166 #[test]
167 fn group_frames_out_of_place_are_caught() {
168 let keys = Keys::new(b"room secret");
169 let mut ada = Sealer::new(&keys).unwrap();
170 let frames: Vec<_> = (0..4)
171 .map(|n| ada.seal(16, &[n], n != 1).unwrap())
172 .collect();
173 let member = || Member::new(&keys, &frames[0]).unwrap();
174
175 let mut grace = member();
176 for (n, frame) in frames.iter().enumerate() {
177 assert_eq!(grace.open(frame).unwrap(), (16, vec![n as u8]));
178 }
179 // A frame to others alone goes unseen without a break.
180 let mut unaddressed = member();
181 unaddressed.open(&frames[0]).unwrap();
182 assert_eq!(unaddressed.open(&frames[2]).unwrap(), (16, vec![2]));
183 // A repeat, a frame reordered, a lost broadcast, an altered frame, another key.
184 let mut repeated = member();
185 repeated.open(&frames[0]).unwrap();
186 assert!(repeated.open(&frames[0]).is_err());
187 let mut reordered = member();
188 reordered.open(&frames[1]).unwrap();
189 assert!(reordered.open(&frames[0]).is_err());
190 let mut lost = member();
191 lost.open(&frames[0]).unwrap();
192 assert!(lost.open(&frames[3]).is_err());
193 let mut altered = frames[3].clone();
194 *altered.last_mut().unwrap() ^= 1;
195 assert!(member().open(&altered).is_err());
196 let stranger = Keys::new(b"another secret");
197 assert!(
198 Member::new(&stranger, &frames[0])
199 .unwrap()
200 .open(&frames[0])
201 .is_err()
202 );
203 }
204}
crates/notebook/src/live/relay.rs+241-20
......@@ -1,15 +1,20 @@
1//! Meeting peers through a relay (`crates/relay`): one WebSocket to the room, carrying a
2//! stream to each peer in it, each opened with SPAKE2 and sealed end to end as on a LAN, so
3//! the relay sees only who talks to whom, when, and how much. A stream that breaks makes this
4//! end join the room again, which ends every stream it had there; peers then meet afresh.
1//! Meeting peers through a relay (`crates/relay`): one WebSocket to the room, carrying
2//! streams to peers in it, each opened with SPAKE2 and sealed end to end as on a LAN, so the
3//! relay sees only who talks to whom, when, and how much. In a code's room every two peers
4//! have a stream. In a notebook's room only a host and each guest do; everything else goes
5//! to the room's group (`group`), sent once and copied by the relay. A stream or group frame
6//! that breaks makes this end join the room again, which ends every stream it had there;
7//! peers then meet afresh.
58
69use super::{
7 Event, OPENING, PATIENCE, Pipe, Relayed as Answer, Shared, Side, code_parts,
10 Event, Hello, OPENING, PATIENCE, Paced, Peer, Pipe, Presence, Relayed as Answer, Room, Shared,
11 Side, code_parts, group,
812 transport::{self, Address, Failure, parse},
13 wire::kind,
914};
10use ::relay::{Notice, SLOT, Verdict, ws};
15use ::relay::{BROADCAST, GROUP, Notice, SLOT, Verdict, ws};
1116use std::{
12 collections::{HashMap, HashSet},
17 collections::{HashMap, HashSet, hash_map::Entry},
1318 io::{self, BufReader, Read, Write},
1419 sync::{Arc, Mutex, atomic::Ordering, mpsc},
1520 thread,
......@@ -67,7 +72,7 @@ fn keep(shared: &Arc<Shared>, address: &Address, port: u16) {
6772 None => format!("{}/v1/claim", address.path),
6873 },
6974 };
70 match connect(address, &path, owner) {
75 match connect(shared, address, &path, owner) {
7176 Ok((socket, reader)) => {
7277 shared.state.lock().unwrap().relay = Some(Arc::clone(&socket));
7378 answered(shared, Answer::Joined);
......@@ -99,14 +104,25 @@ fn keep(shared: &Arc<Shared>, address: &Address, port: u16) {
99104}
100105
101106/// Opens the relay's `path` at `address`: the connection, and what reads its messages.
102fn connect(address: &Address, path: &str, owner: bool) -> Result<(Arc<Socket>, Reader), Failure> {
107fn connect(
108 shared: &Shared,
109 address: &Address,
110 path: &str,
111 owner: bool,
112) -> Result<(Arc<Socket>, Reader), Failure> {
103113 let connection = transport::connect(address, path)?;
104114 let reader = ws::Reader::new(BufReader::new(connection.reader), MOST, false);
115 let group = matches!(shared.room, Room::Notebook(_)).then(|| Group {
116 keys: group::Keys::new(&shared.secret),
117 members: Mutex::default(),
118 out: Mutex::default(),
119 });
105120 let socket = Arc::new(Socket {
106121 send: Mutex::new(connection.writer),
107122 close: connection.close,
108123 owner,
109124 links: Mutex::default(),
125 group,
110126 });
111127 Ok((socket, reader))
112128}
......@@ -118,6 +134,42 @@ pub(super) struct Socket {
118134 /// Whether this end claimed the room for its code, and so tells the relay who knew it.
119135 owner: bool,
120136 links: Mutex<Links>,
137 /// A notebook's room's group.
138 pub(super) group: Option<Group>,
139}
140
141/// A notebook's room's group: its key, the peers heard in it, and the way to this end's
142/// thread sending to it.
143pub(super) struct Group {
144 keys: group::Keys,
145 /// The peers heard in the group, by slot.
146 pub(super) members: Mutex<HashMap<u32, group::Member>>,
147 /// To the thread sending this end's group frames, while connected.
148 out: Mutex<Option<mpsc::Sender<Out>>>,
149}
150
151/// What this end sends the group next.
152pub(super) enum Out {
153 /// The newest presence, to everyone.
154 Presence,
155 /// A message kind and its encoded body, to everyone or to the peers in these slots.
156 Frame(u16, Vec<u8>, Option<Vec<u32>>),
157}
158
159impl Group {
160 pub(super) fn send(&self, out: Out) {
161 if let Some(sender) = &*self.out.lock().unwrap() {
162 let _ = sender.send(out);
163 }
164 }
165
166 /// The slot of the member that is peer `id`.
167 pub(super) fn slot(&self, id: &[u8; 16]) -> Option<u32> {
168 let members = self.members.lock().unwrap();
169 members
170 .iter()
171 .find_map(|(slot, member)| (member.hello()?.peer == *id).then_some(*slot))
172 }
121173}
122174
123175/// The streams a socket carries, by the slot of the peer at the other end.
......@@ -169,6 +221,13 @@ fn session(
169221 }
170222 }
171223 });
224 if let Some(group) = &socket.group {
225 let (out, outgoing) = mpsc::channel();
226 *group.out.lock().unwrap() = Some(out);
227 let (shared, socket) = (Arc::clone(shared), Arc::clone(socket));
228 thread::spawn(move || speak(&shared, &socket, outgoing));
229 }
230 let notebook = socket.group.is_some();
172231 while let Ok(message) = reader.read() {
173232 match message {
174233 ws::Message::Text(text) => match text.parse() {
......@@ -191,12 +250,22 @@ fn session(
191250 }
192251 Ok(Notice::Welcome { you, members }) => {
193252 socket.links.lock().unwrap().me = you;
194 for slot in members {
253 for slot in members.into_iter().filter(|_| !notebook) {
195254 meet(shared, socket, tag.as_deref(), slot, None);
196255 }
197256 }
198 Ok(Notice::Joined(slot)) => meet(shared, socket, tag.as_deref(), slot, None),
199 Ok(Notice::Left(slot)) => socket.forget(slot),
257 Ok(Notice::Joined(slot)) if !notebook => {
258 meet(shared, socket, tag.as_deref(), slot, None);
259 }
260 Ok(Notice::Joined(_)) => {}
261 Ok(Notice::Left(slot)) => {
262 socket.forget(slot);
263 let left = (socket.group.as_ref())
264 .and_then(|group| group.members.lock().unwrap().remove(&slot));
265 if left.is_some() {
266 (shared.events)(Event::Changed);
267 }
268 }
200269 Ok(Notice::Burned) => {
201270 // Coming back, ask for another number.
202271 *nameplate = None;
......@@ -205,12 +274,17 @@ fn session(
205274 }
206275 Err(()) => {}
207276 },
208 ws::Message::Binary(data) => {
209 if let Some((slot, bytes)) = data.split_first_chunk::<SLOT>() {
277 ws::Message::Binary(data) => match data.split_first_chunk::<SLOT>() {
278 Some((slot, frame)) if u32::from_be_bytes(*slot) & GROUP != 0 => {
279 let slot = u32::from_be_bytes(*slot) & !GROUP;
280 heard(shared, socket, tag.as_deref(), slot, frame);
281 }
282 Some((slot, bytes)) => {
210283 let slot = u32::from_be_bytes(*slot);
211284 meet(shared, socket, tag.as_deref(), slot, Some(bytes.to_vec()));
212285 }
213 }
286 None => {}
287 },
214288 ws::Message::Ping(payload) => {
215289 let _ = socket.send(ws::PONG, &payload);
216290 }
......@@ -219,11 +293,155 @@ fn session(
219293 }
220294 }
221295 socket.links.lock().unwrap().inboxes.clear();
296 if let Some(group) = &socket.group {
297 *group.out.lock().unwrap() = None;
298 group.members.lock().unwrap().clear();
299 (shared.events)(Event::Changed);
300 }
222301}
223302
224/// Hands `bytes` from the peer in `slot` to its stream, or opens one: this end opens a stream
225/// to a peer with a lower slot when told of it, and answers one with a higher slot when its
226/// first bytes come.
303/// Sends this end's group frames: its hello, asking everyone for theirs, then its presence
304/// at most every `PRESENCE_EVERY`, and frames as they come.
305fn speak(shared: &Shared, socket: &Socket, outgoing: mpsc::Receiver<Out>) {
306 let Some(group) = &socket.group else {
307 return;
308 };
309 let Ok(mut sealer) = group::Sealer::new(&group.keys) else {
310 return;
311 };
312 let mut send = |kind: u16, body: &[u8], to: Option<&[u32]>| -> io::Result<()> {
313 let sealed = sealer.seal(kind, body, to.is_none())?;
314 let mut message = match to {
315 None => BROADCAST.to_be_bytes().to_vec(),
316 Some(slots) => {
317 let mut message = (GROUP | slots.len() as u32).to_be_bytes().to_vec();
318 slots
319 .iter()
320 .for_each(|slot| message.extend(slot.to_be_bytes()));
321 message
322 }
323 };
324 message.extend(sealed);
325 socket.send(ws::BINARY, &message)
326 };
327 let hello = minicbor::to_vec(&shared.me).unwrap_or_default();
328 if send(kind::HELLO, &hello, None).is_err() {
329 return;
330 }
331 let mut paced = Paced::default();
332 paced.changed();
333 loop {
334 let next = match paced.due {
335 Some(_) => outgoing.recv_timeout(paced.wait(Duration::ZERO)),
336 None => outgoing
337 .recv()
338 .map_err(|_| mpsc::RecvTimeoutError::Disconnected),
339 };
340 let result = match next {
341 Ok(Out::Presence) => {
342 paced.changed();
343 Ok(())
344 }
345 Ok(Out::Frame(kind, body, to)) => send(kind, &body, to.as_deref()),
346 Err(mpsc::RecvTimeoutError::Timeout) => Ok(()),
347 Err(mpsc::RecvTimeoutError::Disconnected) => return,
348 };
349 let result = result.and_then(|()| match paced.due(&shared.state) {
350 Some(presence) => send(
351 kind::PRESENCE,
352 &minicbor::to_vec(&presence).unwrap_or_default(),
353 None,
354 ),
355 None => Ok(()),
356 });
357 if result.is_err() {
358 return;
359 }
360 }
361}
362
363/// Takes in a group frame from the peer in `slot`: its hello, answered with this end's and
364/// meeting it where it serves; its presence; or what else it says.
365fn heard(shared: &Arc<Shared>, socket: &Arc<Socket>, tag: Option<&str>, slot: u32, frame: &[u8]) {
366 let Some(group) = &socket.group else {
367 return;
368 };
369 let opened = {
370 let mut members = group.members.lock().unwrap();
371 let member = match members.entry(slot) {
372 Entry::Occupied(member) => Ok(member.into_mut()),
373 Entry::Vacant(vacant) => {
374 group::Member::new(&group.keys, frame).map(|m| vacant.insert(m))
375 }
376 };
377 member.and_then(|member| {
378 let (kind, body) = member.open(frame)?;
379 Ok((kind, body, member.hello().cloned()))
380 })
381 };
382 let (kind, body, hello) = match opened {
383 Ok(opened) => opened,
384 Err(error) => {
385 eprintln!("Live: the relay broke the room's frames ({error}); meeting again");
386 socket.hang_up();
387 return;
388 }
389 };
390 match kind {
391 kind::HELLO | kind::HELLO_BACK => {
392 let Ok(hello) = minicbor::decode::<Hello>(&body) else {
393 return;
394 };
395 if hello.peer == shared.me.peer {
396 return;
397 }
398 let serves = hello.serves.is_some() && shared.me.serves.is_none();
399 if let Some(member) = group.members.lock().unwrap().get_mut(&slot) {
400 member.peer = Some(Peer {
401 hello: Arc::new(hello),
402 presence: None,
403 });
404 }
405 if kind == kind::HELLO {
406 let presence = minicbor::to_vec(&shared.state.lock().unwrap().presence);
407 let presence = presence.unwrap_or_default();
408 let me = minicbor::to_vec(&shared.me).unwrap_or_default();
409 group.send(Out::Frame(kind::HELLO_BACK, me, Some(vec![slot])));
410 group.send(Out::Frame(kind::PRESENCE, presence, Some(vec![slot])));
411 }
412 if serves {
413 meet(shared, socket, tag, slot, None);
414 }
415 (shared.events)(Event::Changed);
416 }
417 kind::PRESENCE => {
418 let Ok(presence) = minicbor::decode::<Presence>(&body) else {
419 return;
420 };
421 let mut members = group.members.lock().unwrap();
422 if let Some(peer) = members.get_mut(&slot).and_then(|m| m.peer.as_mut()) {
423 peer.presence = Some(presence);
424 drop(members);
425 (shared.events)(Event::Changed);
426 }
427 }
428 // A peer met directly says the rest there.
429 kind if kind < 256 => {
430 if let Some(hello) = hello.filter(|hello| !shared.direct(&hello.peer)) {
431 (shared.events)(Event::Frame {
432 from: &hello,
433 kind,
434 body: &body,
435 });
436 }
437 }
438 _ => {}
439 }
440}
441
442/// Hands `bytes` from the peer in `slot` to its stream, or opens one. In a code's room this
443/// end opens a stream to a peer with a lower slot when told of it, and answers one with a
444/// higher slot when its first bytes come; in a notebook's room a guest opens one to its host.
227445fn meet(
228446 shared: &Arc<Shared>,
229447 socket: &Arc<Socket>,
......@@ -238,9 +456,12 @@ fn meet(
238456 }
239457 return;
240458 }
459 let notebook = socket.group.is_some();
241460 let side = match bytes {
242 None if slot < links.me => Side::Initiator,
243 Some(_) if slot > links.me => Side::Responder,
461 None if notebook || slot < links.me => Side::Initiator,
462 Some(_) if (notebook && shared.me.serves.is_some()) || (!notebook && slot > links.me) => {
463 Side::Responder
464 }
244465 _ => return,
245466 };
246467 let Some(tag) = tag.filter(|_| !links.ended.contains(&slot)) else {
crates/notebook/src/live/share.rs+182-43
......@@ -8,19 +8,20 @@
88//! the next, so a relay never holds much for a slow peer.
99
1010use super::{
11 Event, Hello, Line, Live, Peer, Presence, Reach, Relayed, Room,
11 Event, Hello, Line, Live, Peer, Presence, Reach, Relayed, Room, Sender,
1212 wire::{
1313 self, Delta, Failure, Reply, Request, Touched, Welcome, WireEntry, WireStamp, Written, kind,
1414 },
1515};
1616use crate::{Error, Result, background::Reports, discover, session::Storage};
1717use onestore::{CommitError, CommitState, RevisionIndex, Stamp, Store, Transaction};
18use sha2::{Digest, Sha256};
1819use std::{
1920 collections::{BTreeMap, HashMap},
2021 io,
2122 path::PathBuf,
2223 sync::{
23 Arc, Condvar, Mutex,
24 Arc, Condvar, Mutex, OnceLock,
2425 atomic::{AtomicBool, AtomicU64, Ordering},
2526 mpsc,
2627 },
......@@ -42,9 +43,15 @@ const SNAPSHOT_AGE: Duration = Duration::from_secs(120);
4243/// Sections whose image a host keeps to check guests' commits on and tell them what changed,
4344/// and a guest keeps to take those changes on without reading.
4445const IMAGES: usize = 4;
46/// Earlier images of those a host keeps besides, to tell a guest that missed a change what
47/// it was, and the most bytes they hold.
48const VERSIONS: usize = 16;
49const VERSIONS_BYTES: usize = 64 << 20;
4550/// How long a report of a change settles in a guest's background: its host reports each
4651/// commit once, at once.
4752pub const SETTLE: Duration = Duration::from_millis(20);
53/// The longest a relay may ask a guest joining to wait before it is told so.
54const PATIENT: Duration = Duration::from_secs(10);
4855/// What a host says leaving as it stops sharing.
4956const STOPPED: &str = "stopped";
5057/// What a host says hanging up on a guest that asks too much too fast.
......@@ -172,7 +179,11 @@ pub fn join(
172179 match live.relayed() {
173180 Relayed::Refused(404, _) if settled => return Err(Refusal::NoOne),
174181 Relayed::Refused(410, _) => return Err(Refusal::Expired),
175 Relayed::Refused(429, wait) => return Err(Refusal::TooMany(wait)),
182 // A short wait, as a crowd joining at once meets, passes as the relay is asked
183 // again.
184 Relayed::Refused(429, wait) if wait.is_none_or(|wait| wait > PATIENT) => {
185 return Err(Refusal::TooMany(wait));
186 }
176187 Relayed::Refused(503, _) if settled => return Err(Refusal::Busy),
177188 Relayed::Unreachable(trouble) if waited > Duration::from_secs(10) => {
178189 return Err(Refusal::Unreachable(trouble));
......@@ -224,6 +235,7 @@ impl Host {
224235 puts: Mutex::default(),
225236 guests: Mutex::default(),
226237 host: Mutex::default(),
238 room: OnceLock::new(),
227239 });
228240 let serving = Hello {
229241 serves: Some(sharing.share),
......@@ -247,6 +259,7 @@ impl Host {
247259 _ => {}
248260 },
249261 )?;
262 let _ = served.room.set(room.sender());
250263 let host = Self {
251264 notebook: notebook.to_owned(),
252265 reach,
......@@ -382,15 +395,19 @@ struct Served {
382395 guests: Mutex<BTreeMap<[u8; 16], Admitted>>,
383396 /// Hears the paths guests changed, as the host's own notebook should.
384397 host: Mutex<Option<crate::session::Listener>>,
398 /// The share's room, to tell guests what changed.
399 room: OnceLock<Sender>,
385400}
386401
387/// A guest as its host serves it: the line to it, its requests waiting for its workers, and
388/// how many more it may start now.
402/// A guest as its host serves it: the line to it, its requests waiting for its workers, how
403/// many more it may start now, and the sections it read lately, newest first, as it keeps
404/// their images.
389405struct Admitted {
390406 line: Line,
391407 queue: mpsc::SyncSender<(u16, Vec<u8>)>,
392408 starts: f64,
393409 counted: Instant,
410 held: Vec<String>,
394411}
395412
396413/// A section's path, stamp and image.
......@@ -439,6 +456,7 @@ impl Served {
439456 queue,
440457 starts: BURST,
441458 counted: Instant::now(),
459 held: Vec::new(),
442460 };
443461 self.guests.lock().unwrap().insert(peer, guest);
444462 }
......@@ -486,16 +504,26 @@ impl Served {
486504 }
487505 }
488506
489 /// Tells every guest the files at `paths` changed.
507 /// Tells every guest the files at `paths` changed, and what they are now.
490508 fn tell(&self, paths: &[String]) {
491509 let touched = Touched {
492510 paths: paths.to_vec(),
511 stamps: (paths.iter())
512 .map(|path| self.storage.stamp(path).ok().map(|stamp| digest(&stamp)))
513 .collect(),
493514 };
494 for guest in self.guests.lock().unwrap().values() {
495 let _ = guest.line.send(kind::TOUCHED, &touched);
515 if let Some(room) = self.room.get() {
516 room.send(kind::TOUCHED, &touched, None);
496517 }
497518 }
498519
520 /// Notes that `guest` holds the section at `path`, as it now keeps its image.
521 fn holds(guest: &mut Admitted, path: &str) {
522 guest.held.retain(|held| held != path);
523 guest.held.insert(0, path.to_owned());
524 guest.held.truncate(IMAGES);
525 }
526
499527 fn handle(&self, peer: &[u8; 16], kind: u16, body: &[u8]) -> Reply {
500528 let request = match minicbor::decode::<Request>(body) {
501529 Ok(request) => request,
......@@ -574,7 +602,7 @@ impl Served {
574602 }
575603 kind::COMMIT => {
576604 let transaction = Transaction::from_bytes(&self.carried(peer, &request)?)?;
577 self.commit(path, &transaction)?;
605 self.commit(peer, path, &transaction)?;
578606 self.changed(&[path.to_owned()]);
579607 done
580608 }
......@@ -648,6 +676,15 @@ impl Served {
648676 let offset = request.offset.unwrap_or_default() as usize;
649677 let mut snapshots = self.snapshots.lock().unwrap();
650678 snapshots.retain(|_, (_, read)| read.elapsed() < SNAPSHOT_AGE);
679 if let (kind::READ, None, Some(base)) = (kind, request.handle, &request.stamp)
680 && let Ok(base) = Stamp::try_from(base)
681 && let Some(changes) = self.changes(&request.path, &base)?
682 {
683 if let Some(guest) = self.guests.lock().unwrap().get_mut(peer) {
684 Self::holds(guest, &request.path);
685 }
686 return Ok(changes);
687 }
651688 let (image, handle) = match request.handle {
652689 Some(handle) => {
653690 let (image, read) = snapshots
......@@ -658,6 +695,11 @@ impl Served {
658695 }
659696 None => {
660697 drop(snapshots);
698 if kind == kind::READ
699 && let Some(guest) = self.guests.lock().unwrap().get_mut(peer)
700 {
701 Self::holds(guest, &request.path);
702 }
661703 let limit = (request.limit.unwrap_or(LIMIT as u64) as usize).min(LIMIT);
662704 let image = match kind {
663705 kind::READ => self.image(&request.path)?,
......@@ -692,6 +734,32 @@ impl Served {
692734 })
693735 }
694736
737 /// What changed in the section at `path` since the image with stamp `base`, where that
738 /// image is kept and the changes are much smaller than the section.
739 fn changes(&self, path: &str, base: &Stamp) -> Result<Option<Reply>> {
740 let before = self
741 .images
742 .lock()
743 .unwrap()
744 .iter()
745 .find_map(|(held, at, image)| (held == path && at == base).then(|| Arc::clone(image)));
746 let Some(before) = before else {
747 return Ok(None);
748 };
749 let after = self.image(path)?;
750 let writes = match delta(&before, &after) {
751 Some(writes) => writes,
752 None if before == after => Vec::new(),
753 None => return Ok(None),
754 };
755 Ok(Some(Reply {
756 writes: Some(writes),
757 length: Some(after.len() as u64),
758 stamp: Some((&Stamp::of(&after)?).into()),
759 ..Reply::default()
760 }))
761 }
762
695763 /// The section or TOC at `path` as it stands: the image kept for it while its stamp holds.
696764 fn image(&self, path: &str) -> Result<Arc<Vec<u8>>> {
697765 let stamp = self.storage.stamp(path)?;
......@@ -711,18 +779,31 @@ impl Served {
711779 Ok(image)
712780 }
713781
782 /// Keeps `image` as the section at `path` now, with the newest of each of `IMAGES`
783 /// sections and earlier images up to `VERSIONS` and `VERSIONS_BYTES`.
714784 fn keep(&self, path: &str, image: Arc<Vec<u8>>) {
715785 let Ok(stamp) = Stamp::of(&image) else {
716786 return;
717787 };
718788 let mut images = self.images.lock().unwrap();
719 images.retain(|(held, ..)| held != path);
789 images.retain(|(held, at, _)| held != path || *at != stamp);
720790 images.insert(0, (path.to_owned(), stamp, image));
721 images.truncate(IMAGES);
791 let (mut newest, mut versions, mut bytes) = (Vec::new(), 0, 0);
792 images.retain(|(held, _, image)| {
793 if !newest.contains(held) {
794 newest.push(held.clone());
795 return newest.len() <= IMAGES;
796 }
797 versions += 1;
798 bytes += image.len();
799 newest.iter().take(IMAGES).any(|kept| kept == held)
800 && versions <= VERSIONS
801 && bytes <= VERSIONS_BYTES
802 });
722803 }
723804
724 /// Commits a guest's transaction once the section it makes parses.
725 fn commit(&self, path: &str, transaction: &Transaction) -> Result<()> {
805 /// Commits `guest`'s transaction once the section it makes parses.
806 fn commit(&self, guest: &[u8; 16], path: &str, transaction: &Transaction) -> Result<()> {
726807 let not_committed = |error: io::Error| {
727808 Error::Remote(CommitError {
728809 state: CommitState::NotCommitted,
......@@ -751,14 +832,15 @@ impl Served {
751832 )));
752833 }
753834 self.storage.commit(path, transaction)?;
754 self.tell_delta(path, &image, &next);
835 self.tell_delta(path, &image, &next, Some(guest));
755836 self.keep(path, Arc::new(next));
756837 Ok(())
757838 }
758839
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]) {
840 /// Tells the guests holding the section at `path` what a commit changed in it, from
841 /// `before` to `after`, where that is small enough to send: all but the guest whose
842 /// commit it was, which has it already.
843 fn tell_delta(&self, path: &str, before: &[u8], after: &[u8], committed: Option<&[u8; 16]>) {
762844 let (Ok(base), Some(writes)) = (Stamp::of(before), delta(before, after)) else {
763845 return;
764846 };
......@@ -768,8 +850,18 @@ impl Served {
768850 length: after.len() as u64,
769851 writes,
770852 };
771 for guest in self.guests.lock().unwrap().values() {
772 let _ = guest.line.send(kind::DELTA, &delta);
853 let mut guests = self.guests.lock().unwrap();
854 let holding: Vec<[u8; 16]> = guests
855 .iter_mut()
856 .filter(|(id, guest)| Some(*id) != committed && guest.held.iter().any(|h| h == path))
857 .map(|(id, guest)| {
858 Self::holds(guest, path);
859 *id
860 })
861 .collect();
862 drop(guests);
863 if let (Some(room), false) = (self.room.get(), holding.is_empty()) {
864 room.send(kind::DELTA, &delta, Some(&holding));
773865 }
774866 }
775867
......@@ -789,12 +881,21 @@ impl Served {
789881 return;
790882 }
791883 if let Ok(after) = self.storage.read(path) {
792 self.tell_delta(path, &before, &after);
884 self.tell_delta(path, &before, &after, None);
793885 self.keep(path, Arc::new(after));
794886 }
795887 }
796888}
797889
890/// A stamp as `Touched` names it.
891fn digest(stamp: &Stamp) -> u64 {
892 let hash = Sha256::new()
893 .chain_update(stamp.header)
894 .chain_update(stamp.length.to_le_bytes())
895 .finalize();
896 u64::from_le_bytes(hash[..8].try_into().expect("eight bytes"))
897}
898
798899/// The writes that make `after` of `before`, a commit's appended bytes, patches and header,
799900/// where they are much less than `after` itself.
800901fn delta(before: &[u8], after: &[u8]) -> Option<Vec<Written>> {
......@@ -977,6 +1078,9 @@ struct Inner {
9771078 asked: (Mutex<usize>, Condvar),
9781079 /// The sections read lately, kept as the host's deltas change them, newest first.
9791080 images: Mutex<Vec<(String, Arc<Vec<u8>>)>>,
1081 /// The stamps of sections whose image held is the host's now, as its last report of
1082 /// them said.
1083 current: Mutex<HashMap<String, Stamp>>,
9801084}
9811085
9821086impl Guest {
......@@ -999,6 +1103,7 @@ impl Guest {
9991103 stopped: AtomicBool::new(false),
10001104 asked: Default::default(),
10011105 images: Mutex::default(),
1106 current: Mutex::default(),
10021107 });
10031108 let heard = Arc::clone(&inner);
10041109 let live = Live::start(me, &Room::Notebook(secret), reach, relay, move |event| {
......@@ -1123,17 +1228,32 @@ impl Guest {
11231228 Ok(())
11241229 }
11251230
1126 /// Reads the file at `path` a chunk at a time, as `kind` reads it.
1231 /// Reads the file at `path` a chunk at a time, as `kind` reads it; a section held, as
1232 /// the changes to it.
11271233 fn read(&self, kind: u16, path: &str, limit: usize) -> io::Result<Vec<u8>> {
1234 let held = (kind == kind::READ)
1235 .then(|| self.inner.image(path))
1236 .flatten();
11281237 let first = self.ask(
11291238 kind,
11301239 Request {
11311240 path: path.to_owned(),
11321241 offset: Some(0),
11331242 limit: Some(limit as u64),
1243 stamp: (held.as_deref())
1244 .and_then(|image| Stamp::of(image).ok())
1245 .map(|stamp| (&stamp).into()),
11341246 ..Request::default()
11351247 },
11361248 )?;
1249 if let (Some(writes), Some(held)) = (&first.writes, held) {
1250 let stamp: Option<Stamp> = first.stamp.as_ref().and_then(|s| s.try_into().ok());
1251 let image = written(&held, first.length.unwrap_or_default(), writes)
1252 .filter(|image| stamp.is_some() && Stamp::of(image).ok() == stamp)
1253 .ok_or_else(|| io::Error::new(io::ErrorKind::InvalidData, "Changes that miss"))?;
1254 self.inner.hold(path, Arc::new(image.clone()));
1255 return Ok(image);
1256 }
11371257 let length = first.length.unwrap_or_default() as usize;
11381258 if length > limit {
11391259 return Err(io::ErrorKind::FileTooLarge.into());
......@@ -1166,9 +1286,8 @@ impl Guest {
11661286
11671287 /// The image of the section at `path` with stamp `stamp`, where this guest holds it.
11681288 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())
1289 let image = self.inner.image(path)?;
1290 (Stamp::of(&image).ok().as_ref() == Some(stamp)).then(|| image.to_vec())
11721291 }
11731292
11741293 pub(crate) fn entries(&self, folder: &str) -> io::Result<Vec<discover::Entry>> {
......@@ -1187,7 +1306,12 @@ impl Guest {
11871306 .collect())
11881307 }
11891308
1309 /// The stamp of the file at `path`: the host's last report of it where this guest holds
1310 /// that image, else the host's answer.
11901311 fn stamp(&self, path: &str) -> io::Result<Stamp> {
1312 if let Some(stamp) = self.inner.current.lock().unwrap().get(path) {
1313 return Ok(stamp.clone());
1314 }
11911315 let reply = self.ask(
11921316 kind::STAMP,
11931317 Request {
......@@ -1243,6 +1367,13 @@ impl Guest {
12431367}
12441368
12451369impl Inner {
1370 fn image(&self, path: &str) -> Option<Arc<Vec<u8>>> {
1371 let images = self.images.lock().unwrap();
1372 images
1373 .iter()
1374 .find_map(|(held, image)| (held == path).then(|| Arc::clone(image)))
1375 }
1376
12461377 fn hold(&self, path: &str, image: Arc<Vec<u8>>) {
12471378 let mut images = self.images.lock().unwrap();
12481379 images.retain(|(held, _)| held != path);
......@@ -1252,12 +1383,7 @@ impl Inner {
12521383
12531384 /// Takes on a delta to an image held.
12541385 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)));
1386 let held = self.image(&delta.path);
12611387 let base: Option<Stamp> = (&delta.base).try_into().ok();
12621388 if let Some(image) = held
12631389 && Stamp::of(&image).ok() == base
......@@ -1267,15 +1393,10 @@ impl Inner {
12671393 }
12681394 }
12691395
1270 /// Takes on this guest's own commit to an image held.
1396 /// Takes on this guest's own commit to an image held; the host's report, not this, says
1397 /// whether the image is the host's, as another's commit may have followed.
12711398 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 {
1399 if let Some(image) = self.image(path) {
12791400 let mut next = (*image).clone();
12801401 if transaction.apply(&mut next).is_ok() {
12811402 self.hold(path, Arc::new(next));
......@@ -1283,6 +1404,18 @@ impl Inner {
12831404 }
12841405 }
12851406
1407 /// Takes the host's report of what the files at `paths` are now.
1408 fn touched(&self, touched: &Touched) {
1409 let mut current = self.current.lock().unwrap();
1410 for (at, path) in touched.paths.iter().enumerate() {
1411 let stamp = self.image(path).and_then(|image| Stamp::of(&image).ok());
1412 match stamp.filter(|stamp| touched.stamps.get(at) == Some(&Some(digest(stamp)))) {
1413 Some(stamp) => current.insert(path.clone(), stamp),
1414 None => current.remove(path),
1415 };
1416 }
1417 }
1418
12861419 fn heard(&self, event: Event) {
12871420 let serves = |hello: &Hello| hello.serves == Some(self.share);
12881421 match event {
......@@ -1291,6 +1424,7 @@ impl Inner {
12911424 }
12921425 Event::Left(hello) if serves(hello) => {
12931426 *self.host.lock().unwrap() = None;
1427 self.current.lock().unwrap().clear();
12941428 // Each request waiting hears its answer was lost.
12951429 self.pending.lock().unwrap().clear();
12961430 if let Some(reports) = self.watch.lock().unwrap().take() {
......@@ -1313,10 +1447,11 @@ impl Inner {
13131447 }
13141448 }
13151449 kind::TOUCHED if serves(from) => {
1316 if let Ok(touched) = minicbor::decode::<Touched>(body)
1317 && let Some(reports) = &*self.watch.lock().unwrap()
1318 {
1319 reports.touched(&touched.paths);
1450 if let Ok(touched) = minicbor::decode::<Touched>(body) {
1451 self.touched(&touched);
1452 if let Some(reports) = &*self.watch.lock().unwrap() {
1453 reports.touched(&touched.paths);
1454 }
13201455 }
13211456 }
13221457 kind::BYE
......@@ -1367,7 +1502,11 @@ impl crate::Remote for HostedRemote {
13671502 }
13681503
13691504 fn publish(&mut self, transaction: &Transaction) -> std::result::Result<(), CommitError> {
1370 self.guest.commit(&self.path, transaction)?;
1505 if let Err(error) = self.guest.commit(&self.path, transaction) {
1506 // The host's image moved on, and its report may not have come yet.
1507 self.guest.inner.current.lock().unwrap().remove(&self.path);
1508 return Err(error);
1509 }
13711510 self.guest.inner.published(&self.path, transaction);
13721511 Ok(())
13731512 }
crates/notebook/src/live/tests.rs+19-11
......@@ -359,22 +359,25 @@ enum Tamper {
359359}
360360
361361/// A relay in the middle of Grace's connection that passes on what the real one at
362/// `upstream` says, except the message to her numbered `at` among those from peers, which it
362/// `upstream` says, except the first message to her from a peer once `armed`, which it
363363/// tampers with. Her next connection waits for `release`.
364364struct Malicious {
365365 url: String,
366 armed: Arc<std::sync::atomic::AtomicBool>,
366367 tampered: mpsc::Receiver<()>,
367368 rejoined: mpsc::Receiver<()>,
368369 release: mpsc::Sender<()>,
369370}
370371
371fn malicious(upstream: SocketAddr, tamper: Tamper, at: usize) -> Malicious {
372fn malicious(upstream: SocketAddr, tamper: Tamper) -> Malicious {
372373 use ::relay::ws::{self, Message};
373374 let listener = TcpListener::bind("127.0.0.1:0").unwrap();
374375 let url = format!("ws://{}", listener.local_addr().unwrap());
375376 let (tampered, told) = mpsc::channel();
376377 let (rejoined, heard) = mpsc::channel();
377378 let (release, released) = mpsc::channel();
379 let armed = Arc::new(std::sync::atomic::AtomicBool::new(false));
380 let arming = Arc::clone(&armed);
378381 thread::spawn(move || {
379382 for (index, client) in listener.incoming().enumerate() {
380383 let mut client = client.unwrap();
......@@ -389,19 +392,19 @@ fn malicious(upstream: SocketAddr, tamper: Tamper, at: usize) -> Malicious {
389392 let _ = io::copy(&mut up, &mut to_server);
390393 let _ = to_server.shutdown(Shutdown::Both);
391394 });
392 let tampered = tampered.clone();
395 let (tampered, armed) = (tampered.clone(), Arc::clone(&arming));
393396 thread::spawn(move || {
394397 let mut reading = io::BufReader::new(server);
395398 let head = ws::head(&mut reading).unwrap();
396399 client.write_all(head.as_bytes()).unwrap();
397400 let mut messages = ws::Reader::new(reading, 1 << 20, false);
398 let (mut count, mut held) = (0, None);
401 let mut held = None;
399402 while let Ok(message) = messages.read() {
400403 let frames: Vec<Vec<u8>> = match message {
401404 Message::Binary(mut data) => {
402 count += 1;
403405 let mut out = vec![];
404 if index == 0 && count - 1 == at {
406 if index == 0 && armed.swap(false, std::sync::atomic::Ordering::AcqRel)
407 {
405408 match tamper {
406409 Tamper::Drop => {}
407410 Tamper::Repeat => out = vec![data.clone(), data],
......@@ -411,7 +414,8 @@ fn malicious(upstream: SocketAddr, tamper: Tamper, at: usize) -> Malicious {
411414 out = vec![data];
412415 }
413416 Tamper::Inject => {
414 // A slot, length and number, then made-up bytes.
417 // A slot and part of the sender's id, then
418 // made-up bytes.
415419 let mut forged = data.clone();
416420 forged[16..].fill(7);
417421 out = vec![forged, data];
......@@ -441,6 +445,7 @@ fn malicious(upstream: SocketAddr, tamper: Tamper, at: usize) -> Malicious {
441445 });
442446 Malicious {
443447 url,
448 armed,
444449 tampered: told,
445450 rejoined: heard,
446451 release,
......@@ -460,9 +465,9 @@ fn recording(name: &str, room: &Room, relay: &str) -> (Live, Arc<Mutex<Vec<Optio
460465 let Some(shared) = watched.get().and_then(std::sync::Weak::upgrade) else {
461466 return;
462467 };
463 let entry = match shared.state.lock().unwrap().peers.values().next() {
468 let entry = match shared.peers().first() {
464469 None => Some(None),
465 Some(link) => link.peer.presence.clone().map(Some),
470 Some(peer) => peer.presence.clone().map(Some),
466471 };
467472 recorded.lock().unwrap().extend(entry);
468473 })
......@@ -487,14 +492,17 @@ fn a_malicious_relay_is_caught() {
487492 let (url, address) = relay(Default::default());
488493 let ada = Live::start(hello("Ada"), &room, None, Some(&url), |_| {}).unwrap();
489494 ada.set_presence(caret(1));
490 // Ada's opening, hello and first presence reach Grace; the next is tampered with.
491 let relay = malicious(address, tamper, 3);
495 let relay = malicious(address, tamper);
492496 let (grace, heard) = recording("Grace", &room, &relay.url);
493497 until(&grace, |peers| {
494498 peers
495499 .first()
496500 .is_some_and(|peer| peer.presence == Some(caret(1)))
497501 });
502 // Ada's next presence goes to everyone, so its loss is caught too.
503 relay
504 .armed
505 .store(true, std::sync::atomic::Ordering::Release);
498506 ada.set_presence(caret(2));
499507 relay
500508 .tampered
crates/notebook/src/live/wire.rs+12-1
......@@ -37,6 +37,8 @@ pub mod kind {
3737 /// Keeps a quiet connection open; it says nothing else.
3838 pub const PING: u16 = 2;
3939 pub const BYE: u16 = 3;
40 /// A `Hello` in answer to one heard in a notebook's room's group.
41 pub const HELLO_BACK: u16 = 4;
4042 pub const PRESENCE: u16 = 16;
4143 /// A host's files changed: `Touched`.
4244 pub const TOUCHED: u16 = 18;
......@@ -71,6 +73,7 @@ pub const KNOWN: &[u16] = &[
7173 kind::HELLO,
7274 kind::PING,
7375 kind::BYE,
76 kind::HELLO_BACK,
7477 kind::PRESENCE,
7578 kind::TOUCHED,
7679 kind::DELTA,
......@@ -225,6 +228,10 @@ pub struct Welcome {
225228pub struct Touched {
226229 #[n(0)]
227230 pub paths: Vec<String>,
231 /// Each path's stamp now, as `share::digest` hashes it, where it has one: a guest holding
232 /// that image knows the stamp without asking.
233 #[n(1)]
234 pub stamps: Vec<Option<u64>>,
228235}
229236
230237/// What a commit changed in a section a host serves: from the image with stamp `base`, the
......@@ -271,7 +278,8 @@ pub struct Request {
271278 pub limit: Option<u64>,
272279 #[cbor(n(5), with = "minicbor::bytes")]
273280 pub bytes: Option<Vec<u8>>,
274 /// A confirmation's or supersession's base.
281 /// A confirmation's or supersession's base; for a section's read, the image the guest
282 /// holds, which the reply may give the changes to.
275283 #[n(6)]
276284 pub stamp: Option<WireStamp>,
277285 #[cbor(n(7), with = "minicbor::bytes")]
......@@ -306,6 +314,9 @@ pub struct Reply {
306314 /// The snapshot a read's later chunks come from.
307315 #[n(7)]
308316 pub handle: Option<u64>,
317 /// A read's changes to the image with the request's stamp, in place of its bytes.
318 #[n(8)]
319 pub writes: Option<Vec<Written>>,
309320}
310321
311322/// Why a request failed: an `io::ErrorKind` as `error_kind` numbers it, and for a commit,
crates/notebook/tests/live_share.rs+6-3
......@@ -344,13 +344,16 @@ fn a_flooding_guest_is_hung_up_on() {
344344 },
345345 )
346346 .unwrap();
347 // The line to the host, while Mallory's stream to it is open.
347348 let served = |live: &Live| {
348 live.peers()
349 let host = live
350 .peers()
349351 .into_iter()
350 .find(|peer| peer.hello.serves == Some(welcome.share))
352 .find(|peer| peer.hello.serves == Some(welcome.share))?;
353 live.line(&host.hello.peer)
351354 };
352355 until("Mallory never met the host", || served(&mallory).is_some());
353 let line = mallory.line(&served(&mallory).unwrap().hello.peer).unwrap();
356 let line = served(&mallory).unwrap();
354357 const SENT: u64 = 5000;
355358 for id in 0..SENT {
356359 let request = Request {
crates/relay/README.md+9-5
......@@ -2,9 +2,11 @@
22
33The relay Live Share meets through when two Snowbounds aren't on one network. Peers join a
44room named by a tag (a hash of the notebook's secret, or a code's number), and the relay
5passes their messages between them. Every message after the opening is sealed end to end
6and numbered inside the seal, so the relay can't read, alter, drop or reorder one unseen;
7it learns who talks to whom, when, and how much. `src/lib.rs` describes the protocol.
5passes their messages between them, or copies one to everyone in the room (presence, a
6host's news of its files) so a sender with thirty peers sends it once. Every message after
7the opening is sealed end to end and numbered inside the seal, so the relay can't read,
8alter, drop or reorder one unseen; it learns who talks to whom, when, and how much.
9`src/lib.rs` describes the protocol.
810
911One static Linux executable with no configuration file and nothing on disk. It speaks plain
1012HTTP and WebSocket; a proxy in front of it terminates TLS.
......@@ -28,7 +30,7 @@ sudo install -m 755 /tmp/snowbound-relay /usr/local/bin/snowbound-relay
2830sudo install -m 644 /tmp/snowbound-relay.service /etc/systemd/system/
2931sudo systemctl daemon-reload
3032sudo systemctl enable --now snowbound-relay
31curl -s http://127.0.0.1:23592/health # {"rooms":0,"peers":0,"connections":1,"seconds":3}
33curl -s http://127.0.0.1:23592/health # {"rooms":0,"peers":0,"connections":1,"seconds":3,"bytes_in":0,"bytes_out":0}
3234```
3335
3436Then point a name at the server and put a TLS proxy in front. Caddy fetches its own
......@@ -115,7 +117,9 @@ snowbound.paperclover.net {
115117
116118## Options
117119
118`snowbound-relay --help` lists every option and its default. Each can also be set in the
120`snowbound-relay --help` lists every option and its default. A room holds 64 peers and
121admits 120 a minute, enough for a class joining at once; its bytes per second count what
122the relay gives out, so a copy to thirty peers costs thirty times its size. Each can also be set in the
119123unit's environment as `SNOWBOUND_RELAY_<OPTION>` (`SNOWBOUND_RELAY_MAX_ROOMS=64`).
120124
121125Memory stays under about `max-connections × (max-message + queue)` plus two small thread
crates/relay/src/lib.rs+10-1
......@@ -7,7 +7,13 @@
77//! the relay never knows; peers seal them end to end. Text messages are [`Notice`]s from the
88//! relay and [`Verdict`]s to it.
99//!
10//! Of each pair of peers the one with the higher slot opens the stream between them.
10//! A group message is sent once and copied by the relay: to `BROADCAST`, every other peer
11//! in the room the sender may reach; to `GROUP | n` followed by `n` slots, those peers.
12//! Received, its slot is the sender's marked with `GROUP`. Peers seal group messages under a
13//! key the room's secret gives them all.
14//!
15//! Of each pair of peers in a code's room the one with the higher slot opens the stream
16//! between them.
1117//! A peer joining a code's room waits, hearing only the room's owner, until the owner says
1218//! its opening `met`; a peer that `failed`, left first, or stayed silent too long counts a
1319//! wrong code against its address and the code. See `resources/live-share.md`.
......@@ -21,6 +27,9 @@ use std::{fmt, str::FromStr};
2127
2228/// The bytes before a binary message's payload: the slot it goes to or came from.
2329pub const SLOT: usize = 4;
30/// Marks a group message's slot: see the crate's documentation.
31pub const GROUP: u32 = 1 << 31;
32pub const BROADCAST: u32 = u32::MAX;
2433
2534/// What the relay tells a peer.
2635#[derive(Clone, Debug, PartialEq, Eq)]
crates/relay/src/main.rs+2-2
......@@ -14,13 +14,13 @@ underscores: SNOWBOUND_RELAY_LISTEN=127.0.0.1:23592. Flags win.
1414 --max-connections N (256)
1515 --max-connections-per-address N per IPv4 address or IPv6 /64 (16)
1616 --max-rooms N (128)
17 --max-room-peers N (16)
17 --max-room-peers N (64)
1818 --max-message BYTES (262144)
1919 --queue BYTES waiting to go to one peer before it is dropped (1048576)
2020 --idle SECONDS silence before a connection is closed (600)
2121 --room-bytes-per-second BYTES (4194304)
2222 --joins-per-minute N per address (30)
23 --room-joins-per-minute N (30)
23 --room-joins-per-minute N (120)
2424 --failures-per-minute N wrong codes per address before a lockout (10)
2525 --burn-after N wrong codes before a code admits no one new (5)
2626 --pending SECONDS for a peer joining a code to meet its owner (20)
crates/relay/src/server.rs+71-24
......@@ -2,7 +2,7 @@
22//! limits on everything a stranger can make it hold. A thread reads each connection and
33//! another writes it, from a queue capped in bytes.
44
5use crate::{Notice, SLOT, Verdict, ws};
5use crate::{BROADCAST, GROUP, Notice, SLOT, Verdict, ws};
66use std::{
77 collections::{BTreeMap, HashMap, VecDeque},
88 io::{BufReader, Write},
......@@ -10,7 +10,7 @@ use std::{
1010 ops::RangeInclusive,
1111 sync::{
1212 Arc, Condvar, Mutex,
13 atomic::{AtomicUsize, Ordering},
13 atomic::{AtomicU64, AtomicUsize, Ordering},
1414 },
1515 thread,
1616 time::{Duration, Instant},
......@@ -50,13 +50,13 @@ impl Default for Config {
5050 max_connections: 256,
5151 max_connections_per_address: 16,
5252 max_rooms: 128,
53 max_room_peers: 16,
53 max_room_peers: 64,
5454 max_message: 256 << 10,
5555 queue: 1 << 20,
5656 idle: Duration::from_secs(600),
5757 room_bytes_per_second: 4 << 20,
5858 joins_per_minute: 30,
59 room_joins_per_minute: 30,
59 room_joins_per_minute: 120,
6060 failures_per_minute: 10,
6161 burn_after: 5,
6262 pending: Duration::from_secs(20),
......@@ -84,6 +84,7 @@ pub fn serve(listener: TcpListener, config: Config) -> std::io::Result<()> {
8484 config,
8585 state: Mutex::default(),
8686 connections: AtomicUsize::new(0),
87 relayed: Default::default(),
8788 started: Instant::now(),
8889 });
8990 let sweeping = Arc::clone(&relay);
......@@ -120,6 +121,8 @@ struct Relay {
120121 state: Mutex<State>,
121122 /// Connections open, joined or not.
122123 connections: AtomicUsize,
124 /// Bytes of peers' messages taken in, and given out.
125 relayed: [AtomicU64; 2],
123126 started: Instant,
124127}
125128
......@@ -398,11 +401,15 @@ impl Relay {
398401 fn health(&self) -> String {
399402 let state = self.state.lock().unwrap();
400403 let peers: usize = state.rooms.values().map(|room| room.members.len()).sum();
404 let [taken, given] = &self.relayed;
401405 format!(
402 "{{\"rooms\":{},\"peers\":{peers},\"connections\":{},\"seconds\":{}}}\n",
406 "{{\"rooms\":{},\"peers\":{peers},\"connections\":{},\"seconds\":{},\
407 \"bytes_in\":{},\"bytes_out\":{}}}\n",
403408 state.rooms.len(),
404409 self.connections.load(Ordering::Acquire),
405 self.started.elapsed().as_secs()
410 self.started.elapsed().as_secs(),
411 taken.load(Ordering::Relaxed),
412 given.load(Ordering::Relaxed),
406413 )
407414 }
408415
......@@ -543,33 +550,72 @@ impl Relay {
543550 /// Acts on a message from `slot`: false once it is gone or broke the protocol.
544551 fn heard(&self, tag: &str, slot: u32, outbox: &Outbox, message: ws::Message) -> bool {
545552 match message {
546 ws::Message::Binary(mut data) => {
547 if data.len() < SLOT {
553 ws::Message::Binary(data) => {
554 let Some((to, rest)) = data.split_first_chunk::<SLOT>() else {
548555 return false;
549 }
550 let rate = self.config.room_bytes_per_second as f64;
551 let wait = match self.state.lock().unwrap().rooms.get_mut(tag) {
552 Some(room) => room.bytes.spend(data.len() as f64, rate, Instant::now()),
553 None => return false,
554556 };
555 // Waiting here slows the sender alone, as its socket fills.
556 thread::sleep(wait);
557 let to = u32::from_be_bytes(data[..SLOT].try_into().expect("a slot"));
557 let to = u32::from_be_bytes(*to);
558 // The slots a group message names, then what it carries.
559 let (named, payload) = match to {
560 BROADCAST => (None, rest),
561 to if to & GROUP != 0 => {
562 let Some((named, payload)) =
563 rest.split_at_checked((to & !GROUP) as usize * SLOT)
564 else {
565 return false;
566 };
567 let named = named.chunks_exact(SLOT);
568 (
569 Some(named.map(|slot| u32::from_be_bytes(slot.try_into().unwrap()))),
570 payload,
571 )
572 }
573 _ => (None, rest),
574 };
558575 let state = self.state.lock().unwrap();
559576 let Some(room) = state.rooms.get(tag) else {
560577 return false;
561578 };
562 let (Some(from), Some(target)) = (room.members.get(&slot), room.members.get(&to))
563 else {
564 return room.members.contains_key(&slot);
579 let Some(from) = room.members.get(&slot) else {
580 return false;
565581 };
566582 // One waiting for the code's owner talks to the owner alone.
567 if (from.pending.is_none() || room.owner == Some(to))
568 && (target.pending.is_none() || room.owner == Some(slot))
569 {
570 data[..SLOT].copy_from_slice(&slot.to_be_bytes());
571 target.outbox.push(ws::frame(ws::BINARY, &data, None));
583 let allowed = |to: u32| {
584 let target = room.members.get(&to)?;
585 (to != slot
586 && (from.pending.is_none() || room.owner == Some(to))
587 && (target.pending.is_none() || room.owner == Some(slot)))
588 .then_some(&target.outbox)
589 };
590 let (source, targets): (u32, Vec<&Arc<Outbox>>) = match (to, named) {
591 (BROADCAST, _) => (
592 slot | GROUP,
593 room.members.keys().filter_map(|to| allowed(*to)).collect(),
594 ),
595 (_, Some(named)) => (slot | GROUP, named.filter_map(allowed).collect()),
596 (to, None) => (slot, allowed(to).into_iter().collect()),
597 };
598 let message = ws::frame(
599 ws::BINARY,
600 &[&source.to_be_bytes()[..], payload].concat(),
601 None,
602 );
603 let targets: Vec<Arc<Outbox>> = targets.into_iter().cloned().collect();
604 drop(state);
605 // Waiting here slows the sender alone, as its socket fills.
606 let given = (SLOT + payload.len()) * targets.len();
607 let rate = self.config.room_bytes_per_second as f64;
608 let wait = match self.state.lock().unwrap().rooms.get_mut(tag) {
609 Some(room) => room.bytes.spend(given as f64, rate, Instant::now()),
610 None => return false,
611 };
612 thread::sleep(wait);
613 for target in targets {
614 target.push(message.clone());
572615 }
616 let [taken, out] = &self.relayed;
617 taken.fetch_add(data.len() as u64, Ordering::Relaxed);
618 out.fetch_add(given as u64, Ordering::Relaxed);
573619 true
574620 }
575621 ws::Message::Text(text) => {
......@@ -1043,6 +1089,7 @@ mod tests {
10431089 },
10441090 state: Mutex::default(),
10451091 connections: AtomicUsize::new(0),
1092 relayed: Default::default(),
10461093 started: Instant::now(),
10471094 };
10481095 let listener = TcpListener::bind("127.0.0.1:0").unwrap();
crates/relay/tests/relay.rs+36
......@@ -185,6 +185,42 @@ fn peers_in_a_room_pass_messages() {
185185 assert_eq!(ada.text(), "left 2");
186186}
187187
188/// A group message goes once to the relay, which copies it to everyone else in the room or
189/// to the slots it names, marked as the group's; the relay counts what it passed on.
190#[test]
191fn group_messages_are_copied() {
192 let relay = Relay::start(&[]);
193 let mut peers: Vec<Peer> = (1..=3)
194 .map(|_| relay.join("/v1/room/0123abcd").unwrap())
195 .collect();
196 assert_eq!(peers[0].text(), "welcome 1");
197 assert_eq!(peers[0].text(), "joined 2");
198 assert_eq!(peers[0].text(), "joined 3");
199 assert_eq!(peers[1].text(), "welcome 2 1");
200 assert_eq!(peers[1].text(), "joined 3");
201 assert_eq!(peers[2].text(), "welcome 3 1 2");
202 let from = |slot: u32, bytes: &[u8]| {
203 Message::Binary([&(slot | relay::GROUP).to_be_bytes()[..], bytes].concat())
204 };
205 peers[0].send_to(relay::BROADCAST, b"everyone");
206 assert_eq!(peers[1].hear(), from(1, b"everyone"));
207 assert_eq!(peers[2].hear(), from(1, b"everyone"));
208 let named = [&3u32.to_be_bytes()[..], b"just you"].concat();
209 peers[1].send_to(relay::GROUP | 1, &named);
210 assert_eq!(peers[2].hear(), from(2, b"just you"));
211 peers[0].send_to(2, b"stream");
212 assert_eq!(
213 peers[1].hear(),
214 Message::Binary([&1u32.to_be_bytes()[..], b"stream"].concat())
215 );
216 let health = relay.get("/health");
217 let given = 2 * (4 + 8) + (4 + 8) + (4 + 6);
218 assert!(
219 health.contains(&format!("\"bytes_out\":{given}")),
220 "{health}"
221 );
222}
223
188224/// One joining a code's room hears only its owner until the owner says it met it; then the
189225/// others in the room, and they it.
190226#[test]