From 5e6642680a2a6d3e0dcdc018cb7a021a7e010daa Mon Sep 17 00:00:00 2001 From: clover caruso Date: Fri, 2 Oct 2026 23:06:34 -0700 Subject: [PATCH] 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 --- crates/notebook/Cargo.toml | 4 + crates/notebook/examples/live_crowd.rs | 564 +++++++++++++++++++++++++ crates/notebook/src/live.rs | 175 ++++++-- crates/notebook/src/live/group.rs | 204 +++++++++ crates/notebook/src/live/relay.rs | 261 +++++++++++- crates/notebook/src/live/share.rs | 225 ++++++++-- crates/notebook/src/live/tests.rs | 30 +- crates/notebook/src/live/wire.rs | 13 +- crates/notebook/tests/live_share.rs | 9 +- crates/relay/README.md | 14 +- crates/relay/src/lib.rs | 11 +- crates/relay/src/main.rs | 4 +- crates/relay/src/server.rs | 99 +++-- crates/relay/tests/relay.rs | 36 ++ 14 files changed, 1513 insertions(+), 136 deletions(-) create mode 100644 crates/notebook/examples/live_crowd.rs create mode 100644 crates/notebook/src/live/group.rs diff --git a/crates/notebook/Cargo.toml b/crates/notebook/Cargo.toml index b28ab9b2c05445d9a2d0109fe2e36663a6b1efe8..89071756ff0d4478dcb0700aca8cd864cd152d80 100644 --- a/crates/notebook/Cargo.toml +++ b/crates/notebook/Cargo.toml @@ -63,6 +63,10 @@ tempfile = "3" name = "live_latency" required-features = ["live"] +[[example]] +name = "live_crowd" +required-features = ["live"] + [[example]] name = "smb_offline_client" required-features = ["smb"] diff --git a/crates/notebook/examples/live_crowd.rs b/crates/notebook/examples/live_crowd.rs new file mode 100644 index 0000000000000000000000000000000000000000..57ed5d7282faa91063a3ed7e68e08cdf46c3c10e --- /dev/null +++ b/crates/notebook/examples/live_crowd.rs @@ -0,0 +1,564 @@ +//! A crowd on one shared notebook: a host, and guests that join it through a relay, half of +//! them typing on the same page. Reports how long keys and carets take to reach the others, +//! what the relay passes on, and what the relay, the host and the guests cost. +//! +//! `cargo run --release -p notebook --features live --example live_crowd -- RELAY_BINARY +//! [PEERS=30] [TYPISTS=PEERS/2] [SECONDS=60]`, the relay's binary built with +//! `cargo build --release -p relay`. The relay, the host and the guests each run as a process +//! of their own, so `ps` tells their costs apart. +//! +//! For watching a crowd in the app: `fixture FOLDER PARAGRAPHS` writes the notebook the +//! host shares, and `guests URL CODE PEERS SECONDS` has that many join the app's share of it +//! and move their carets about its page. + +use notebook::{ + Replica, + live::{ + Caret, Guid, Hello, Presence, Spot, + share::{self, Guest, Host, Sharing}, + }, + session::{Background, Notebook, Section}, +}; +use onestore::{ + ExGuid, + op::{Edit, Op, PageOp}, +}; +use std::{ + collections::HashMap, + io::{BufRead, BufReader, Read, Write}, + net::TcpStream, + path::Path, + process::{Child, Command, Stdio}, + sync::{ + Arc, Mutex, + atomic::{AtomicBool, Ordering}, + }, + thread, + time::{Duration, Instant}, +}; + +/// Keys each typist types a second, and how often observers look. +const KEYS_PER_SECOND: f64 = 5.0; +const LOOK: Duration = Duration::from_millis(2); +/// Guests that time what they see. +const OBSERVERS: usize = 3; + +fn until(what: &str, limit: Duration, done: impl Fn() -> bool) { + let deadline = Instant::now() + limit; + while !done() { + assert!(Instant::now() < deadline, "{what}"); + thread::sleep(Duration::from_millis(5)); + } +} + +fn hello(name: &str) -> Hello { + Hello::new(name.into(), None).unwrap() +} + +fn main() { + let args: Vec = std::env::args().skip(1).collect(); + match args.first().map(String::as_str) { + Some("host") => return host(&args[1], args[2].parse().unwrap()), + Some("fixture") => { + fixture(Path::new(&args[1]), args[2].parse().unwrap()); + return; + } + Some("guests") => { + let seconds = Duration::from_secs(args[4].parse().unwrap()); + return guests(&args[1], &args[2], args[3].parse().unwrap(), seconds); + } + _ => {} + } + let binary = args.first().expect("the relay's binary"); + let number = |at: usize, default: usize| args.get(at).map_or(default, |n| n.parse().unwrap()); + let peers = number(1, 30); + let typists = number(2, peers / 2); + let seconds = number(3, 60) as u64; + crowd(binary, peers, typists, Duration::from_secs(seconds)); +} + +/// Writes the notebook `Garden` in `directory`: a section whose page has `paragraphs` +/// paragraphs, `Typed:` then `Note 1:` and on. Its folder. +fn fixture(directory: &Path, paragraphs: usize) -> std::path::PathBuf { + let folder = directory.join("Garden"); + std::fs::create_dir_all(&folder).unwrap(); + let mut image = onestore::create_section("Garden.one", "Typed:", "Fixture").unwrap(); + for n in (1..paragraphs).rev() { + let arena = onestore::Arena::default(); + let mut section = onestore::Section::open(&arena, image).unwrap(); + let (space, ..) = section.pages().unwrap()[0]; + let page = section.page(space).unwrap(); + let (text, _) = texts(&page).swap_remove(0); + let right = onestore::page::text::new_id().unwrap(); + let ops = [ + PageOp::Split { + text, + at: 6, + paragraph: onestore::page::text::new_id().unwrap(), + right, + lists: Vec::new(), + }, + PageOp::Text { + text: right, + range: 0..0, + with: format!("Note {n}:"), + }, + ]; + let edit = Edit { + at: 134_000_000_000_000_000, + ops: ops.into_iter().map(|op| Op::Page { space, op }).collect(), + }; + section.apply("Fixture", &edit).unwrap(); + section.seal().unwrap(); + image = section.image(); + } + std::fs::write(folder.join("Garden.one"), image).unwrap(); + folder +} + +/// The host's process: shares a fixture notebook of a page with `paragraphs` paragraphs +/// through `url`, says its code, and reports what the section grew to once its input closes. +fn host(url: &str, paragraphs: usize) { + let directory = tempfile::tempdir().unwrap(); + let folder = fixture(directory.path(), paragraphs); + let file = folder.join("Garden.one"); + let storage = Notebook::open(&folder, directory.path().join("cache")) + .unwrap() + .into_storage(); + let host = Host::start( + storage, + hello("Host"), + Sharing::new("").unwrap(), + "Garden", + None, + Some(url), + || {}, + ) + .unwrap(); + until("no code", Duration::from_secs(30), || { + host.code().is_some_and(|code| share::code(&code).is_some()) + }); + println!("code {}", host.code().unwrap()); + let _ = std::io::stdin().read_to_end(&mut Vec::new()); + let image = std::fs::read(&file).unwrap(); + let size = image.len(); + let arena = onestore::Arena::default(); + let mut section = onestore::Section::open(&arena, image).unwrap(); + let pages = section.pages().unwrap().len(); + let conflicts: usize = section + .conflicts() + .unwrap() + .iter() + .map(|(_, c)| c.len()) + .sum(); + println!("section {size} bytes, {pages} pages, {conflicts} conflict pages"); +} + +/// A guest as the app holds one: its share, the notebook's background, the section open +/// and its file's identity. +struct Open { + guest: Arc, + _background: Background, + section: Section, + file: [u8; 16], +} + +fn join(name: &str, code: &str, url: &str, cache: &Path) -> Open { + let welcome = share::join(hello(name), code, "", None, Some(url)).unwrap(); + let guest = Guest::start( + hello(name), + welcome.share, + welcome.secret, + None, + Some(url), + || {}, + ) + .unwrap(); + until("no host", Duration::from_secs(120), || { + guest.host().is_some() + }); + let mut notebook = Notebook::open_hosted(Arc::clone(&guest), cache).unwrap(); + let background = Background::hosted(Arc::clone(&guest), || {}).unwrap(); + background.watch(notebook.replicas()); + let replica = notebook.replica_path("Garden.one").unwrap(); + std::fs::create_dir_all(replica.parent().unwrap()).unwrap(); + let image = notebook.read_section("Garden.one").unwrap(); + let file = onestore::Header::parse(&image[..1024]).unwrap().file_id; + let replica = Replica::open_or_create(&replica, None, || Ok(image)).unwrap(); + let section = + Section::resume_hosted("Garden.one".into(), replica, Arc::clone(&guest), || {}).unwrap(); + background.hold("Garden.one", §ion); + Open { + guest, + _background: background, + section, + file, + } +} + +/// Where in the open page's `paragraph`th text, `offset` units in, `open`'s guest is. +fn presence(open: &Open, paragraph: usize, offset: u32) -> Option { + let (space, texts) = texts_of(&open.section)?; + let (text, _) = texts.get(paragraph)?; + let at = Spot { + text: Guid { + guid: text.guid, + n: text.n, + }, + offset, + }; + Some(Presence { + section: Some(open.file), + page: Some(Guid { + guid: space.guid, + n: space.n, + }), + caret: Some(Caret { + anchor: at, + focus: at, + }), + }) +} + +/// `peers` guests join the share `code` names through `url` and, for `seconds`, put their +/// carets in the page's paragraphs, each moving every few seconds. +fn guests(url: &str, code: &str, peers: usize, seconds: Duration) { + let directory = tempfile::tempdir().unwrap(); + let names = [ + "Ada", "Grace", "Alan", "Barbara", "Edsger", "Frances", "Donald", "Margaret", "Ken", + "Radia", "Dennis", "Hedy", "Linus", "Karen", "John", "Sophie", + ]; + let joining: Vec<_> = (0..peers) + .map(|n| { + let name = format!("{} {n}", names[n % names.len()]); + let (code, url, cache) = ( + code.to_owned(), + url.to_owned(), + directory.path().join(n.to_string()), + ); + thread::sleep(Duration::from_millis(100)); + thread::spawn(move || join(&name, &code, &url, &cache)) + }) + .collect(); + let guests: Vec = joining + .into_iter() + .map(|joined| joined.join().unwrap()) + .collect(); + let began = Instant::now(); + let mut turn = 0; + while began.elapsed() < seconds { + for (n, open) in guests.iter().enumerate() { + let paragraphs = texts_of(&open.section).map_or(1, |(_, texts)| texts.len()); + if (n + turn) % 4 == 0 || turn == 0 { + let at = (n * 7 + turn) % paragraphs; + if let Some(presence) = presence(open, at, (turn % 5) as u32) { + open.guest.set_presence(presence); + } + } + } + turn += 1; + thread::sleep(Duration::from_secs(1)); + } +} + +/// The page's object space and its paragraphs' texts: each one's object and what it reads. +fn texts_of(section: &Section) -> Option<(ExGuid, Vec<(ExGuid, String)>)> { + let (space, ..) = *section.pages().ok()?.first()?; + let page = section.page(space).ok()?; + Some((space, texts(&page))) +} + +fn texts(page: &onestore::page::Page) -> Vec<(ExGuid, String)> { + let outlines = page.objects.iter().filter_map(|object| match object { + onestore::page::PageObject::Outline(outline) => Some(outline), + _ => None, + }); + outlines + .flat_map(|outline| &outline.paragraphs) + .filter_map(|p| p.text().map(|t| (t.id, t.text.text().to_owned()))) + .collect() +} + +/// The key typist `typist` types `n`th, unique to both. +fn key(typist: usize, n: usize) -> char { + char::from_u32(0x4E00 + (typist * 600 + n) as u32).unwrap() +} + +fn typist_of(key: char) -> Option<(usize, usize)> { + let at = (key as u32).checked_sub(0x4E00)? as usize; + Some((at / 600, at % 600)) +} + +/// Seconds of CPU and kilobytes resident `ps` reports for `pid`. +fn cost(pid: u32) -> (f64, u64) { + let output = Command::new("ps") + .args(["-o", "time=,rss=", "-p", &pid.to_string()]) + .output() + .unwrap(); + let text = String::from_utf8_lossy(&output.stdout); + let mut words = text.split_whitespace(); + let time = words.next().unwrap_or("0:0"); + let seconds = time.split(':').fold(0.0, |total, part| { + total * 60.0 + part.parse::().unwrap_or(0.0) + }); + ( + seconds, + words.next().and_then(|rss| rss.parse().ok()).unwrap_or(0), + ) +} + +/// The relay's `/health`, as numbers by name. +fn health(address: &str) -> HashMap { + let mut stream = TcpStream::connect(address).unwrap(); + write!( + stream, + "GET /health HTTP/1.1\r\nHost: relay\r\nConnection: close\r\n\r\n" + ) + .unwrap(); + let mut text = String::new(); + stream.read_to_string(&mut text).unwrap(); + let body = text.split("\r\n\r\n").nth(1).unwrap_or_default(); + body.trim() + .trim_matches(['{', '}']) + .split(',') + .filter_map(|pair| { + let (name, value) = pair.split_once(':')?; + Some((name.trim_matches('"').to_owned(), value.parse().ok()?)) + }) + .collect() +} + +fn report(what: &str, mut times: Vec) { + if times.is_empty() { + println!("{what}: none seen"); + return; + } + times.sort(); + let ms = |at: usize| times[at].as_secs_f64() * 1000.0; + println!( + "{what}: median {:.0} ms, p90 {:.0} ms, p99 {:.0} ms, worst {:.0} ms ({} seen)", + ms(times.len() / 2), + ms(times.len() * 9 / 10), + ms(times.len() * 99 / 100), + ms(times.len() - 1), + times.len() + ); +} + +struct Spawned(Child); + +impl Drop for Spawned { + fn drop(&mut self) { + let _ = self.0.kill(); + let _ = self.0.wait(); + } +} + +fn crowd(binary: &str, peers: usize, typists: usize, window: Duration) { + let port = std::net::TcpListener::bind("127.0.0.1:0") + .unwrap() + .local_addr() + .unwrap() + .port(); + let address = format!("127.0.0.1:{port}"); + // Everyone comes from one address here, which a relay would take for one abuser. + let relay = Spawned( + Command::new(binary) + .args(["--listen", &address]) + .args(["--max-connections-per-address", "1000"]) + .args(["--joins-per-minute", "10000"]) + .args(["--failures-per-minute", "10000"]) + .stdout(Stdio::null()) + .spawn() + .unwrap(), + ); + until("no relay", Duration::from_secs(10), || { + TcpStream::connect(&address).is_ok() + }); + let url = format!("ws://{address}"); + let mut hosting = Spawned( + Command::new(std::env::current_exe().unwrap()) + .args(["host", &url, &typists.max(1).to_string()]) + .stdin(Stdio::piped()) + .stdout(Stdio::piped()) + .spawn() + .unwrap(), + ); + let mut said = BufReader::new(hosting.0.stdout.take().unwrap()).lines(); + let code = said.next().unwrap().unwrap(); + let code = code.strip_prefix("code ").unwrap().to_owned(); + println!("{peers} guests, {typists} typing {KEYS_PER_SECOND} keys a second, code {code}"); + + let directory = tempfile::tempdir().unwrap(); + let began = Instant::now(); + let joining: Vec<_> = (0..peers) + .map(|n| { + let (code, url) = (code.clone(), url.clone()); + let cache = directory.path().join(format!("guest{n}")); + thread::sleep(Duration::from_millis(50)); + thread::spawn(move || join(&format!("Guest {n}"), &code, &url, &cache)) + }) + .collect(); + let guests: Vec> = joining + .into_iter() + .map(|joined| Arc::new(joined.join().unwrap())) + .collect(); + println!("all joined in {:.1} s", began.elapsed().as_secs_f64()); + until("not everyone met", Duration::from_secs(300), || { + guests.iter().all(|open| open.guest.peers().len() >= peers) + }); + println!("all met in {:.1} s", began.elapsed().as_secs_f64()); + + let (_, texts) = texts_of(&guests[0].section).unwrap(); + assert!(texts.len() >= typists, "a paragraph for each typist"); + for open in &guests { + open.guest.set_presence(presence(open, 0, 0).unwrap()); + } + + // When each typist typed each key, and said its caret moved past it. + let typed: Arc>> = Arc::default(); + let moved: Arc>> = Arc::default(); + let seen_keys: Arc>> = Arc::default(); + let seen_carets: Arc>> = Arc::default(); + let done = Arc::new(AtomicBool::new(false)); + let (relay_before, host_before, crowd_before) = ( + cost(relay.0.id()), + cost(hosting.0.id()), + cost(std::process::id()), + ); + let bytes_before = health(&address); + let started = Instant::now(); + + let mut threads = Vec::new(); + for (n, open) in guests.iter().enumerate().take(typists) { + let (open, typed, done) = (Arc::clone(open), Arc::clone(&typed), Arc::clone(&done)); + let moved = Arc::clone(&moved); + threads.push(thread::spawn(move || { + let pause = Duration::from_secs_f64(1.0 / KEYS_PER_SECOND); + thread::sleep(pause.mul_f64(n as f64 / typists as f64)); + let mut count = 0; + while !done.load(Ordering::Acquire) && count < 600 { + open.section.events(); + let Some((space, mut texts)) = texts_of(&open.section) else { + continue; + }; + let (text, before) = texts.swap_remove(n); + let end = before.encode_utf16().count() as u32; + let edit = Edit { + at: 134_000_000_000_000_000, + ops: vec![Op::Page { + space, + op: PageOp::Text { + text, + range: end..end, + with: key(n, count).to_string(), + }, + }], + }; + typed.lock().unwrap().insert((n, count), Instant::now()); + open.section.replica().apply("Typist", edit).unwrap(); + count += 1; + if let Some(presence) = presence(&open, 0, count as u32) { + moved.lock().unwrap().insert((n, count), Instant::now()); + open.guest.set_presence(presence); + } + thread::sleep(pause); + } + })); + } + for open in guests.iter().skip(typists).take(OBSERVERS) { + let (open, typed, done) = (Arc::clone(open), Arc::clone(&typed), Arc::clone(&done)); + let moved = Arc::clone(&moved); + let (keys, carets) = (Arc::clone(&seen_keys), Arc::clone(&seen_carets)); + threads.push(thread::spawn(move || { + let mut keys_seen = std::collections::HashSet::new(); + let mut carets_seen = HashMap::new(); + while !done.load(Ordering::Acquire) { + open.section.events(); + let now = Instant::now(); + if let Some((_, texts)) = texts_of(&open.section) { + let chars = texts.iter().flat_map(|(_, text)| text.chars()); + for (typist, n) in chars.filter_map(typist_of) { + if keys_seen.insert((typist, n)) + && let Some(at) = typed.lock().unwrap().get(&(typist, n)) + { + keys.lock().unwrap().push(now - *at); + } + } + } + for peer in open.guest.peers() { + let Some(typist) = (peer.hello.name.strip_prefix("Guest ")) + .and_then(|n| n.parse::().ok()) + .filter(|n| *n < typists) + else { + continue; + }; + let Some(count) = peer.presence.and_then(|p| p.caret).map(|c| c.focus.offset) + else { + continue; + }; + let count = count as usize; + if count > 0 + && carets_seen.insert(typist, count) != Some(count) + && let Some(at) = moved.lock().unwrap().get(&(typist, count)) + { + carets.lock().unwrap().push(now - *at); + } + } + thread::sleep(LOOK); + } + })); + } + thread::sleep(window); + done.store(true, Ordering::Release); + let elapsed = started.elapsed().as_secs_f64(); + let (relay_after, host_after, crowd_after) = ( + cost(relay.0.id()), + cost(hosting.0.id()), + cost(std::process::id()), + ); + let bytes_after = health(&address); + for thread in threads { + thread.join().unwrap(); + } + let keys = typed.lock().unwrap().len(); + report( + "keys, typist to observer", + std::mem::take(&mut seen_keys.lock().unwrap()), + ); + report( + "carets, typist to observer", + std::mem::take(&mut seen_carets.lock().unwrap()), + ); + let delta = |name: &str| bytes_after[name] - bytes_before.get(name).copied().unwrap_or(0.0); + let members = (peers + 1) as f64; + println!( + "relay passed on {:.0} KB/s in, {:.0} KB/s out: per peer {:.1} KB/s in, {:.1} KB/s out", + delta("bytes_in") / elapsed / 1000.0, + delta("bytes_out") / elapsed / 1000.0, + delta("bytes_in") / elapsed / 1000.0 / members, + delta("bytes_out") / elapsed / 1000.0 / members, + ); + let load = |name: &str, (before, _): (f64, u64), (after, rss): (f64, u64)| { + println!( + "{name}: {:.0}% of a core, {:.1} MB resident", + (after - before) / elapsed * 100.0, + rss as f64 / 1024.0 + ); + }; + load("relay", relay_before, relay_after); + load("host", host_before, host_after); + load(&format!("{peers} guests"), crowd_before, crowd_after); + + // Every key typed should reach the host's file; conflict pages show as extra pages. + thread::sleep(Duration::from_secs(5)); + let observed = texts_of(&guests[typists.min(peers - 1)].section).map_or(0, |(_, texts)| { + let chars = texts.iter().flat_map(|(_, text)| text.chars()); + chars.filter_map(typist_of).count() + }); + println!("{keys} keys typed, {observed} on an observer's page after 5 s"); + drop(hosting.0.stdin.take()); + for line in said.map_while(Result::ok) { + println!("host: {line}"); + } +} diff --git a/crates/notebook/src/live.rs b/crates/notebook/src/live.rs index b04de9bc9999cdd860bc0edbcdd8b29ba5583db2..c47e7ed5a054c3f9efcefcefc23c19a68c41583f 100644 --- a/crates/notebook/src/live.rs +++ b/crates/notebook/src/live.rs @@ -8,6 +8,7 @@ //! from scratch. pub use ::relay::code; +mod group; pub mod proxy; mod relay; pub mod share; @@ -38,6 +39,8 @@ const SERVICE: &str = "_snowbound._tcp.local."; const PING: Duration = Duration::from_secs(15); const GONE: Duration = Duration::from_secs(45); const OPENING: Duration = Duration::from_secs(5); +/// The most often presence goes to a peer. +const PRESENCE_EVERY: Duration = Duration::from_millis(100); /// The longest wait before meeting again. const PATIENCE: Duration = Duration::from_secs(30); /// 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 { pub enum Event<'a> { /// The peers, their presence, the code or the relay's answer changed. Changed, - /// A peer was met, with the line to it. + /// A stream to a peer opened, with the line to it. Met(&'a Arc, &'a Line), - /// A peer's connection ended. + /// A peer's stream ended. Left(&'a Arc), - /// A frame of a kind presence doesn't read itself, with the line to answer on. + /// A frame of a kind presence doesn't read itself, on a stream or to the group. Frame { from: &'a Arc, kind: u16, body: &'a [u8], - line: &'a Line, }, } @@ -200,6 +202,48 @@ struct State { daemon: Option, } +impl State { + /// The relay's group, in a notebook's room while the relay is reached. + fn group(&self) -> Option<&relay::Group> { + self.relay.as_ref()?.group.as_ref() + } +} + +/// Presence sent at most every `PRESENCE_EVERY`, and only the newest. +#[derive(Default)] +struct Paced { + due: Option, + last: Option, + sent: Option, +} + +impl Paced { + /// How long to wait for news: until presence is due, else `idle`. + fn wait(&self, idle: Duration) -> Duration { + self.due + .map_or(idle, |due| due.saturating_duration_since(Instant::now())) + } + + /// Hears that presence changed. + fn changed(&mut self) { + let last = self.last; + self.due + .get_or_insert_with(|| last.map_or_else(Instant::now, |at| at + PRESENCE_EVERY)); + } + + /// The presence to send now, where it is due and new. + fn due(&mut self, state: &Mutex) -> Option { + self.due.filter(|due| *due <= Instant::now())?; + self.due = None; + let state = state.lock().unwrap(); + if self.sent == Some(state.generation) { + return None; + } + (self.sent, self.last) = (Some(state.generation), Some(Instant::now())); + Some(state.presence.clone()) + } +} + struct Link { connection: u64, peer: Peer, @@ -401,7 +445,8 @@ impl Live { thread::spawn(move || shared.dial(address)); } - /// Says where this end is now; peers hear only the newest of quick changes. + /// Says where this end is now; peers hear only the newest of quick changes, at most every + /// tenth of a second. pub fn set_presence(&self, presence: Presence) { let mut state = self.shared.state.lock().unwrap(); if state.presence == presence { @@ -410,14 +455,23 @@ impl Live { state.presence = presence; state.generation += 1; for link in state.peers.values() { - let _ = link.line.0.send(Out::Presence); + if self.shared.carries_presence(&*link.pipe) { + let _ = link.line.0.send(Out::Presence); + } + } + if let Some(group) = state.group() { + group.send(relay::Out::Presence); } } - /// The peers connected now, by id. + /// The peers in the room now, by id. pub fn peers(&self) -> Vec { - let state = self.shared.state.lock().unwrap(); - state.peers.values().map(|link| link.peer.clone()).collect() + self.shared.peers() + } + + /// A way to send to this room's peers that doesn't keep it open. + pub fn sender(&self) -> Sender { + Sender(Arc::downgrade(&self.shared)) } /// The line to `peer`, while it is connected. @@ -442,6 +496,45 @@ impl Live { } } +/// Sends to a room's peers while the room is open. +#[derive(Clone)] +pub struct Sender(std::sync::Weak); + +impl Sender { + /// Sends message `kind` holding `body` to the peers `to` names, or to everyone: once to + /// the room's group through the relay, and to each peer met directly. + pub fn send(&self, kind: u16, body: &impl Encode<()>, to: Option<&[[u8; 16]]>) { + let (Some(shared), Ok(body)) = (self.0.upgrade(), minicbor::to_vec(body)) else { + return; + }; + let state = shared.state.lock().unwrap(); + let group = state.group(); + let named = |id: &[u8; 16]| to.is_none_or(|to| to.contains(id)); + let mut slots = Vec::new(); + for (id, link) in state.peers.iter().filter(|(id, _)| named(id)) { + match group { + Some(group) if !link.pipe.direct() => slots.extend(group.slot(id)), + _ => { + let _ = link.line.0.send(Out::Frame(kind, body.clone())); + } + } + } + let Some(group) = group else { + return; + }; + match to { + None => group.send(relay::Out::Frame(kind, body, None)), + Some(to) => { + let unlinked = to.iter().filter(|id| !state.peers.contains_key(*id)); + slots.extend(unlinked.filter_map(|id| group.slot(id))); + if !slots.is_empty() { + group.send(relay::Out::Frame(kind, body, Some(slots))); + } + } + } + } +} + impl Drop for Live { fn drop(&mut self) { self.shared.stopped.store(true, Ordering::Release); @@ -461,6 +554,38 @@ impl Drop for Live { } impl Shared { + /// The peers in the room: those met directly, and the rest as the relay's group or a + /// stream through the relay last heard of them. + fn peers(&self) -> Vec { + let state = self.state.lock().unwrap(); + let mut peers = BTreeMap::new(); + if let Some(group) = state.group() { + for member in group.members.lock().unwrap().values() { + if let Some(peer) = &member.peer { + peers.insert(peer.hello.peer, peer.clone()); + } + } + } + for (id, link) in &state.peers { + if link.pipe.direct() || !peers.contains_key(id) { + peers.insert(*id, link.peer.clone()); + } + } + peers.into_values().collect() + } + + /// Whether peer `id` is met directly, which is where it says everything. + fn direct(&self, id: &[u8; 16]) -> bool { + let state = self.state.lock().unwrap(); + state.peers.get(id).is_some_and(|link| link.pipe.direct()) + } + + /// Whether presence goes on `pipe`: not on a stream through the relay in a notebook's + /// room, whose group carries it. + fn carries_presence(&self, pipe: &dyn Pipe) -> bool { + pipe.direct() || matches!(self.room, Room::Code { .. }) + } + /// The room's tag as it stands: a code's, once numbered. fn tag(&self) -> Option { match &self.room { @@ -592,7 +717,9 @@ impl Shared { pipe.shutdown(); return false; } - let _ = line.0.send(Out::Presence); + if self.carries_presence(&*pipe) { + let _ = line.0.send(Out::Presence); + } state.peers.insert( peer, Link { @@ -632,7 +759,6 @@ impl Shared { from: &hello, kind, body: &body, - line: &line, }); // The peer hangs up after its bye. if kind == kind::BYE { @@ -660,22 +786,16 @@ impl Shared { true } - /// Sends the newest presence whenever woken, frames as they come, and a ping when quiet. + /// Sends the newest presence when woken, at most every `PRESENCE_EVERY`, frames as they + /// come, and a ping when quiet. fn write(&self, pipe: &dyn Pipe, mut send: Sealer, outgoing: mpsc::Receiver) { let mut stream = pipe; - let mut sent = None; + let mut paced = Paced::default(); loop { - let result = match outgoing.recv_timeout(PING) { + let result = match outgoing.recv_timeout(paced.wait(PING)) { Ok(Out::Presence) => { - let presence = { - let state = self.state.lock().unwrap(); - if sent == Some(state.generation) { - continue; - } - sent = Some(state.generation); - state.presence.clone() - }; - send.send(&mut stream, kind::PRESENCE, &presence) + paced.changed(); + Ok(()) } Ok(Out::Frame(kind, body)) => send.send_encoded(&mut stream, kind, &body), Ok(Out::Bye(body)) => { @@ -683,9 +803,16 @@ impl Shared { pipe.shutdown(); return; } - Err(mpsc::RecvTimeoutError::Timeout) => send.send(&mut stream, kind::PING, &()), + Err(mpsc::RecvTimeoutError::Timeout) if paced.due.is_none() => { + send.send(&mut stream, kind::PING, &()) + } + Err(mpsc::RecvTimeoutError::Timeout) => Ok(()), Err(mpsc::RecvTimeoutError::Disconnected) => return, }; + let result = result.and_then(|()| match paced.due(&self.state) { + Some(presence) => send.send(&mut stream, kind::PRESENCE, &presence), + None => Ok(()), + }); if result.is_err() { pipe.shutdown(); return; diff --git a/crates/notebook/src/live/group.rs b/crates/notebook/src/live/group.rs new file mode 100644 index 0000000000000000000000000000000000000000..e0e2ea6e3849a201d2640e999dcbf827c00dcd50 --- /dev/null +++ b/crates/notebook/src/live/group.rs @@ -0,0 +1,204 @@ +//! Frames a peer in a notebook's room sends once through the relay, which copies them to +//! everyone or to the peers named (`relay::GROUP`): presence, hellos, and a host's news of +//! its files. Each is sealed under a key of its sender's own, derived from the room's secret +//! and an id the sender picks for each connection, so the relay sees only ciphertext. +//! +//! A frame carries its number, which is its nonce, and the count of broadcasts its sender has +//! sent: a frame numbered no later than the last, or a broadcast missing before it, means the +//! relay lost, repeated, reordered or forged one, and this end meets the room again. A frame +//! to some peers goes unseen by the others, so only what can tell itself stale is sent so. +//! The room's members hold one key, so a frame proves its sender is in the room, not which +//! member it is; members are trusted alike, as each may change the notebook itself. + +use super::{Hello, Peer}; +use aes_gcm::{ + Aes256Gcm, KeyInit, + aead::{Aead, Payload}, +}; +use hmac::{Hmac, Mac}; +use sha2::Sha256; +use std::io; + +/// A frame's sender id, number, broadcasts, and whether it is a broadcast. +const HEADER: usize = 16 + 8 + 8 + 1; + +/// The room's group key. +pub(super) struct Keys([u8; 32]); + +impl Keys { + pub(super) fn new(secret: &[u8]) -> Self { + Self(mac(secret, b"Snowbound live v2 group")) + } + + fn cipher(&self, sender: &[u8; 16]) -> Aes256Gcm { + Aes256Gcm::new_from_slice(&mac(&self.0, sender)).expect("a 32-byte key") + } +} + +fn mac(key: &[u8], message: &[u8]) -> [u8; 32] { + let mut mac = as hmac::KeyInit>::new_from_slice(key).expect("any key length"); + mac.update(message); + mac.finalize().into_bytes().into() +} + +/// Seals this end's frames for one connection to the relay. +pub(super) struct Sealer { + id: [u8; 16], + cipher: Aes256Gcm, + number: u64, + broadcasts: u64, +} + +impl Sealer { + pub(super) fn new(keys: &Keys) -> io::Result { + let mut id = [0; 16]; + getrandom::fill(&mut id).map_err(|_| io::Error::other("System random source failed"))?; + Ok(Self { + cipher: keys.cipher(&id), + id, + number: 0, + broadcasts: 0, + }) + } + + /// Message `kind` holding `body`, to everyone where `broadcast`. + pub(super) fn seal(&mut self, kind: u16, body: &[u8], broadcast: bool) -> io::Result> { + self.number += 1; + self.broadcasts += u64::from(broadcast); + let mut header = Vec::with_capacity(HEADER); + header.extend_from_slice(&self.id); + header.extend_from_slice(&self.number.to_be_bytes()); + header.extend_from_slice(&self.broadcasts.to_be_bytes()); + header.push(u8::from(broadcast)); + let clear = [&kind.to_be_bytes()[..], body].concat(); + let sealed = self + .cipher + .encrypt( + &nonce(self.number).into(), + Payload { + msg: &clear, + aad: &header, + }, + ) + .map_err(|_| io::Error::other("A frame could not be sealed"))?; + Ok([header, sealed].concat()) + } +} + +fn nonce(number: u64) -> [u8; 12] { + let mut nonce = [0; 12]; + nonce[4..].copy_from_slice(&number.to_be_bytes()); + nonce +} + +/// A peer in the room as its frames say: the id its frames are sealed under, how far they +/// have come, and the peer once its hello has. +pub(super) struct Member { + id: [u8; 16], + cipher: Aes256Gcm, + number: u64, + broadcasts: u64, + pub(super) peer: Option, +} + +impl Member { + /// The member whose first frame is `frame`. + pub(super) fn new(keys: &Keys, frame: &[u8]) -> io::Result { + let id: [u8; 16] = frame + .get(..16) + .and_then(|id| id.try_into().ok()) + .ok_or_else(|| broken("A group frame without its sender"))?; + Ok(Self { + cipher: keys.cipher(&id), + id, + number: 0, + broadcasts: 0, + peer: None, + }) + } + + /// The kind and body of `frame`, the next from this member; an error of kind + /// `InvalidData` means the relay tampered with its frames. + pub(super) fn open(&mut self, frame: &[u8]) -> io::Result<(u16, Vec)> { + let (header, sealed) = frame + .split_at_checked(HEADER) + .ok_or_else(|| broken("A group frame cut short"))?; + let field = |at: usize| u64::from_be_bytes(header[at..at + 8].try_into().unwrap()); + let (number, broadcasts, broadcast) = (field(16), field(24), header[32] == 1); + if header[..16] != self.id { + return Err(broken("A group frame from another sender in its slot")); + } + let mut clear = self + .cipher + .decrypt( + &nonce(number).into(), + Payload { + msg: sealed, + aad: header, + }, + ) + .map_err(|_| broken("A group frame that does not open"))?; + // The first frame heard from a member is where its count starts. + let first = self.number == 0; + let due = self.broadcasts + u64::from(broadcast); + if !first && (number <= self.number || broadcasts != due) { + return Err(broken("A group frame lost, repeated or out of order")); + } + (self.number, self.broadcasts) = (number, broadcasts); + let body = clear.split_off(2.min(clear.len())); + let kind = u16::from_be_bytes(clear.try_into().map_err(|_| broken("An empty frame"))?); + Ok((kind, body)) + } + + pub(super) fn hello(&self) -> Option<&std::sync::Arc> { + self.peer.as_ref().map(|peer| &peer.hello) + } +} + +fn broken(message: &'static str) -> io::Error { + io::Error::new(io::ErrorKind::InvalidData, message) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn group_frames_out_of_place_are_caught() { + let keys = Keys::new(b"room secret"); + let mut ada = Sealer::new(&keys).unwrap(); + let frames: Vec<_> = (0..4) + .map(|n| ada.seal(16, &[n], n != 1).unwrap()) + .collect(); + let member = || Member::new(&keys, &frames[0]).unwrap(); + + let mut grace = member(); + for (n, frame) in frames.iter().enumerate() { + assert_eq!(grace.open(frame).unwrap(), (16, vec![n as u8])); + } + // A frame to others alone goes unseen without a break. + let mut unaddressed = member(); + unaddressed.open(&frames[0]).unwrap(); + assert_eq!(unaddressed.open(&frames[2]).unwrap(), (16, vec![2])); + // A repeat, a frame reordered, a lost broadcast, an altered frame, another key. + let mut repeated = member(); + repeated.open(&frames[0]).unwrap(); + assert!(repeated.open(&frames[0]).is_err()); + let mut reordered = member(); + reordered.open(&frames[1]).unwrap(); + assert!(reordered.open(&frames[0]).is_err()); + let mut lost = member(); + lost.open(&frames[0]).unwrap(); + assert!(lost.open(&frames[3]).is_err()); + let mut altered = frames[3].clone(); + *altered.last_mut().unwrap() ^= 1; + assert!(member().open(&altered).is_err()); + let stranger = Keys::new(b"another secret"); + assert!( + Member::new(&stranger, &frames[0]) + .unwrap() + .open(&frames[0]) + .is_err() + ); + } +} diff --git a/crates/notebook/src/live/relay.rs b/crates/notebook/src/live/relay.rs index 553b41ae388f9ad0c51a8129cf68bf54b37575f4..74c5aa8d5f616ea21a8365fb3532fbf564d2daa2 100644 --- a/crates/notebook/src/live/relay.rs +++ b/crates/notebook/src/live/relay.rs @@ -1,15 +1,20 @@ -//! Meeting peers through a relay (`crates/relay`): one WebSocket to the room, carrying a -//! stream to each peer in it, each opened with SPAKE2 and sealed end to end as on a LAN, so -//! the relay sees only who talks to whom, when, and how much. A stream that breaks makes this -//! end join the room again, which ends every stream it had there; peers then meet afresh. +//! Meeting peers through a relay (`crates/relay`): one WebSocket to the room, carrying +//! streams to peers in it, each opened with SPAKE2 and sealed end to end as on a LAN, so the +//! relay sees only who talks to whom, when, and how much. In a code's room every two peers +//! have a stream. In a notebook's room only a host and each guest do; everything else goes +//! to the room's group (`group`), sent once and copied by the relay. A stream or group frame +//! that breaks makes this end join the room again, which ends every stream it had there; +//! peers then meet afresh. use super::{ - Event, OPENING, PATIENCE, Pipe, Relayed as Answer, Shared, Side, code_parts, + Event, Hello, OPENING, PATIENCE, Paced, Peer, Pipe, Presence, Relayed as Answer, Room, Shared, + Side, code_parts, group, transport::{self, Address, Failure, parse}, + wire::kind, }; -use ::relay::{Notice, SLOT, Verdict, ws}; +use ::relay::{BROADCAST, GROUP, Notice, SLOT, Verdict, ws}; use std::{ - collections::{HashMap, HashSet}, + collections::{HashMap, HashSet, hash_map::Entry}, io::{self, BufReader, Read, Write}, sync::{Arc, Mutex, atomic::Ordering, mpsc}, thread, @@ -67,7 +72,7 @@ fn keep(shared: &Arc, address: &Address, port: u16) { None => format!("{}/v1/claim", address.path), }, }; - match connect(address, &path, owner) { + match connect(shared, address, &path, owner) { Ok((socket, reader)) => { shared.state.lock().unwrap().relay = Some(Arc::clone(&socket)); answered(shared, Answer::Joined); @@ -99,14 +104,25 @@ fn keep(shared: &Arc, address: &Address, port: u16) { } /// Opens the relay's `path` at `address`: the connection, and what reads its messages. -fn connect(address: &Address, path: &str, owner: bool) -> Result<(Arc, Reader), Failure> { +fn connect( + shared: &Shared, + address: &Address, + path: &str, + owner: bool, +) -> Result<(Arc, Reader), Failure> { let connection = transport::connect(address, path)?; let reader = ws::Reader::new(BufReader::new(connection.reader), MOST, false); + let group = matches!(shared.room, Room::Notebook(_)).then(|| Group { + keys: group::Keys::new(&shared.secret), + members: Mutex::default(), + out: Mutex::default(), + }); let socket = Arc::new(Socket { send: Mutex::new(connection.writer), close: connection.close, owner, links: Mutex::default(), + group, }); Ok((socket, reader)) } @@ -118,6 +134,42 @@ pub(super) struct Socket { /// Whether this end claimed the room for its code, and so tells the relay who knew it. owner: bool, links: Mutex, + /// A notebook's room's group. + pub(super) group: Option, +} + +/// A notebook's room's group: its key, the peers heard in it, and the way to this end's +/// thread sending to it. +pub(super) struct Group { + keys: group::Keys, + /// The peers heard in the group, by slot. + pub(super) members: Mutex>, + /// To the thread sending this end's group frames, while connected. + out: Mutex>>, +} + +/// What this end sends the group next. +pub(super) enum Out { + /// The newest presence, to everyone. + Presence, + /// A message kind and its encoded body, to everyone or to the peers in these slots. + Frame(u16, Vec, Option>), +} + +impl Group { + pub(super) fn send(&self, out: Out) { + if let Some(sender) = &*self.out.lock().unwrap() { + let _ = sender.send(out); + } + } + + /// The slot of the member that is peer `id`. + pub(super) fn slot(&self, id: &[u8; 16]) -> Option { + let members = self.members.lock().unwrap(); + members + .iter() + .find_map(|(slot, member)| (member.hello()?.peer == *id).then_some(*slot)) + } } /// The streams a socket carries, by the slot of the peer at the other end. @@ -169,6 +221,13 @@ fn session( } } }); + if let Some(group) = &socket.group { + let (out, outgoing) = mpsc::channel(); + *group.out.lock().unwrap() = Some(out); + let (shared, socket) = (Arc::clone(shared), Arc::clone(socket)); + thread::spawn(move || speak(&shared, &socket, outgoing)); + } + let notebook = socket.group.is_some(); while let Ok(message) = reader.read() { match message { ws::Message::Text(text) => match text.parse() { @@ -191,12 +250,22 @@ fn session( } Ok(Notice::Welcome { you, members }) => { socket.links.lock().unwrap().me = you; - for slot in members { + for slot in members.into_iter().filter(|_| !notebook) { meet(shared, socket, tag.as_deref(), slot, None); } } - Ok(Notice::Joined(slot)) => meet(shared, socket, tag.as_deref(), slot, None), - Ok(Notice::Left(slot)) => socket.forget(slot), + Ok(Notice::Joined(slot)) if !notebook => { + meet(shared, socket, tag.as_deref(), slot, None); + } + Ok(Notice::Joined(_)) => {} + Ok(Notice::Left(slot)) => { + socket.forget(slot); + let left = (socket.group.as_ref()) + .and_then(|group| group.members.lock().unwrap().remove(&slot)); + if left.is_some() { + (shared.events)(Event::Changed); + } + } Ok(Notice::Burned) => { // Coming back, ask for another number. *nameplate = None; @@ -205,12 +274,17 @@ fn session( } Err(()) => {} }, - ws::Message::Binary(data) => { - if let Some((slot, bytes)) = data.split_first_chunk::() { + ws::Message::Binary(data) => match data.split_first_chunk::() { + Some((slot, frame)) if u32::from_be_bytes(*slot) & GROUP != 0 => { + let slot = u32::from_be_bytes(*slot) & !GROUP; + heard(shared, socket, tag.as_deref(), slot, frame); + } + Some((slot, bytes)) => { let slot = u32::from_be_bytes(*slot); meet(shared, socket, tag.as_deref(), slot, Some(bytes.to_vec())); } - } + None => {} + }, ws::Message::Ping(payload) => { let _ = socket.send(ws::PONG, &payload); } @@ -219,11 +293,155 @@ fn session( } } socket.links.lock().unwrap().inboxes.clear(); + if let Some(group) = &socket.group { + *group.out.lock().unwrap() = None; + group.members.lock().unwrap().clear(); + (shared.events)(Event::Changed); + } } -/// Hands `bytes` from the peer in `slot` to its stream, or opens one: this end opens a stream -/// to a peer with a lower slot when told of it, and answers one with a higher slot when its -/// first bytes come. +/// Sends this end's group frames: its hello, asking everyone for theirs, then its presence +/// at most every `PRESENCE_EVERY`, and frames as they come. +fn speak(shared: &Shared, socket: &Socket, outgoing: mpsc::Receiver) { + let Some(group) = &socket.group else { + return; + }; + let Ok(mut sealer) = group::Sealer::new(&group.keys) else { + return; + }; + let mut send = |kind: u16, body: &[u8], to: Option<&[u32]>| -> io::Result<()> { + let sealed = sealer.seal(kind, body, to.is_none())?; + let mut message = match to { + None => BROADCAST.to_be_bytes().to_vec(), + Some(slots) => { + let mut message = (GROUP | slots.len() as u32).to_be_bytes().to_vec(); + slots + .iter() + .for_each(|slot| message.extend(slot.to_be_bytes())); + message + } + }; + message.extend(sealed); + socket.send(ws::BINARY, &message) + }; + let hello = minicbor::to_vec(&shared.me).unwrap_or_default(); + if send(kind::HELLO, &hello, None).is_err() { + return; + } + let mut paced = Paced::default(); + paced.changed(); + loop { + let next = match paced.due { + Some(_) => outgoing.recv_timeout(paced.wait(Duration::ZERO)), + None => outgoing + .recv() + .map_err(|_| mpsc::RecvTimeoutError::Disconnected), + }; + let result = match next { + Ok(Out::Presence) => { + paced.changed(); + Ok(()) + } + Ok(Out::Frame(kind, body, to)) => send(kind, &body, to.as_deref()), + Err(mpsc::RecvTimeoutError::Timeout) => Ok(()), + Err(mpsc::RecvTimeoutError::Disconnected) => return, + }; + let result = result.and_then(|()| match paced.due(&shared.state) { + Some(presence) => send( + kind::PRESENCE, + &minicbor::to_vec(&presence).unwrap_or_default(), + None, + ), + None => Ok(()), + }); + if result.is_err() { + return; + } + } +} + +/// Takes in a group frame from the peer in `slot`: its hello, answered with this end's and +/// meeting it where it serves; its presence; or what else it says. +fn heard(shared: &Arc, socket: &Arc, tag: Option<&str>, slot: u32, frame: &[u8]) { + let Some(group) = &socket.group else { + return; + }; + let opened = { + let mut members = group.members.lock().unwrap(); + let member = match members.entry(slot) { + Entry::Occupied(member) => Ok(member.into_mut()), + Entry::Vacant(vacant) => { + group::Member::new(&group.keys, frame).map(|m| vacant.insert(m)) + } + }; + member.and_then(|member| { + let (kind, body) = member.open(frame)?; + Ok((kind, body, member.hello().cloned())) + }) + }; + let (kind, body, hello) = match opened { + Ok(opened) => opened, + Err(error) => { + eprintln!("Live: the relay broke the room's frames ({error}); meeting again"); + socket.hang_up(); + return; + } + }; + match kind { + kind::HELLO | kind::HELLO_BACK => { + let Ok(hello) = minicbor::decode::(&body) else { + return; + }; + if hello.peer == shared.me.peer { + return; + } + let serves = hello.serves.is_some() && shared.me.serves.is_none(); + if let Some(member) = group.members.lock().unwrap().get_mut(&slot) { + member.peer = Some(Peer { + hello: Arc::new(hello), + presence: None, + }); + } + if kind == kind::HELLO { + let presence = minicbor::to_vec(&shared.state.lock().unwrap().presence); + let presence = presence.unwrap_or_default(); + let me = minicbor::to_vec(&shared.me).unwrap_or_default(); + group.send(Out::Frame(kind::HELLO_BACK, me, Some(vec![slot]))); + group.send(Out::Frame(kind::PRESENCE, presence, Some(vec![slot]))); + } + if serves { + meet(shared, socket, tag, slot, None); + } + (shared.events)(Event::Changed); + } + kind::PRESENCE => { + let Ok(presence) = minicbor::decode::(&body) else { + return; + }; + let mut members = group.members.lock().unwrap(); + if let Some(peer) = members.get_mut(&slot).and_then(|m| m.peer.as_mut()) { + peer.presence = Some(presence); + drop(members); + (shared.events)(Event::Changed); + } + } + // A peer met directly says the rest there. + kind if kind < 256 => { + if let Some(hello) = hello.filter(|hello| !shared.direct(&hello.peer)) { + (shared.events)(Event::Frame { + from: &hello, + kind, + body: &body, + }); + } + } + _ => {} + } +} + +/// Hands `bytes` from the peer in `slot` to its stream, or opens one. In a code's room this +/// end opens a stream to a peer with a lower slot when told of it, and answers one with a +/// higher slot when its first bytes come; in a notebook's room a guest opens one to its host. fn meet( shared: &Arc, socket: &Arc, @@ -238,9 +456,12 @@ fn meet( } return; } + let notebook = socket.group.is_some(); let side = match bytes { - None if slot < links.me => Side::Initiator, - Some(_) if slot > links.me => Side::Responder, + None if notebook || slot < links.me => Side::Initiator, + Some(_) if (notebook && shared.me.serves.is_some()) || (!notebook && slot > links.me) => { + Side::Responder + } _ => return, }; let Some(tag) = tag.filter(|_| !links.ended.contains(&slot)) else { diff --git a/crates/notebook/src/live/share.rs b/crates/notebook/src/live/share.rs index ff7f9a8a8581860d9491332cf6c74c3fc42dcfcf..bea88c2a0668175924ab2988e0a1479fe7b8fefc 100644 --- a/crates/notebook/src/live/share.rs +++ b/crates/notebook/src/live/share.rs @@ -8,19 +8,20 @@ //! the next, so a relay never holds much for a slow peer. use super::{ - Event, Hello, Line, Live, Peer, Presence, Reach, Relayed, Room, + Event, Hello, Line, Live, Peer, Presence, Reach, Relayed, Room, Sender, wire::{ self, Delta, Failure, Reply, Request, Touched, Welcome, WireEntry, WireStamp, Written, kind, }, }; use crate::{Error, Result, background::Reports, discover, session::Storage}; use onestore::{CommitError, CommitState, RevisionIndex, Stamp, Store, Transaction}; +use sha2::{Digest, Sha256}; use std::{ collections::{BTreeMap, HashMap}, io, path::PathBuf, sync::{ - Arc, Condvar, Mutex, + Arc, Condvar, Mutex, OnceLock, atomic::{AtomicBool, AtomicU64, Ordering}, mpsc, }, @@ -42,9 +43,15 @@ const SNAPSHOT_AGE: Duration = Duration::from_secs(120); /// Sections whose image a host keeps to check guests' commits on and tell them what changed, /// and a guest keeps to take those changes on without reading. const IMAGES: usize = 4; +/// Earlier images of those a host keeps besides, to tell a guest that missed a change what +/// it was, and the most bytes they hold. +const VERSIONS: usize = 16; +const VERSIONS_BYTES: usize = 64 << 20; /// How long a report of a change settles in a guest's background: its host reports each /// commit once, at once. pub const SETTLE: Duration = Duration::from_millis(20); +/// The longest a relay may ask a guest joining to wait before it is told so. +const PATIENT: Duration = Duration::from_secs(10); /// What a host says leaving as it stops sharing. const STOPPED: &str = "stopped"; /// What a host says hanging up on a guest that asks too much too fast. @@ -172,7 +179,11 @@ pub fn join( match live.relayed() { Relayed::Refused(404, _) if settled => return Err(Refusal::NoOne), Relayed::Refused(410, _) => return Err(Refusal::Expired), - Relayed::Refused(429, wait) => return Err(Refusal::TooMany(wait)), + // A short wait, as a crowd joining at once meets, passes as the relay is asked + // again. + Relayed::Refused(429, wait) if wait.is_none_or(|wait| wait > PATIENT) => { + return Err(Refusal::TooMany(wait)); + } Relayed::Refused(503, _) if settled => return Err(Refusal::Busy), Relayed::Unreachable(trouble) if waited > Duration::from_secs(10) => { return Err(Refusal::Unreachable(trouble)); @@ -224,6 +235,7 @@ impl Host { puts: Mutex::default(), guests: Mutex::default(), host: Mutex::default(), + room: OnceLock::new(), }); let serving = Hello { serves: Some(sharing.share), @@ -247,6 +259,7 @@ impl Host { _ => {} }, )?; + let _ = served.room.set(room.sender()); let host = Self { notebook: notebook.to_owned(), reach, @@ -382,15 +395,19 @@ struct Served { guests: Mutex>, /// Hears the paths guests changed, as the host's own notebook should. host: Mutex>, + /// The share's room, to tell guests what changed. + room: OnceLock, } -/// A guest as its host serves it: the line to it, its requests waiting for its workers, and -/// how many more it may start now. +/// A guest as its host serves it: the line to it, its requests waiting for its workers, how +/// many more it may start now, and the sections it read lately, newest first, as it keeps +/// their images. struct Admitted { line: Line, queue: mpsc::SyncSender<(u16, Vec)>, starts: f64, counted: Instant, + held: Vec, } /// A section's path, stamp and image. @@ -439,6 +456,7 @@ impl Served { queue, starts: BURST, counted: Instant::now(), + held: Vec::new(), }; self.guests.lock().unwrap().insert(peer, guest); } @@ -486,16 +504,26 @@ impl Served { } } - /// Tells every guest the files at `paths` changed. + /// Tells every guest the files at `paths` changed, and what they are now. fn tell(&self, paths: &[String]) { let touched = Touched { paths: paths.to_vec(), + stamps: (paths.iter()) + .map(|path| self.storage.stamp(path).ok().map(|stamp| digest(&stamp))) + .collect(), }; - for guest in self.guests.lock().unwrap().values() { - let _ = guest.line.send(kind::TOUCHED, &touched); + if let Some(room) = self.room.get() { + room.send(kind::TOUCHED, &touched, None); } } + /// Notes that `guest` holds the section at `path`, as it now keeps its image. + fn holds(guest: &mut Admitted, path: &str) { + guest.held.retain(|held| held != path); + guest.held.insert(0, path.to_owned()); + guest.held.truncate(IMAGES); + } + fn handle(&self, peer: &[u8; 16], kind: u16, body: &[u8]) -> Reply { let request = match minicbor::decode::(body) { Ok(request) => request, @@ -574,7 +602,7 @@ impl Served { } kind::COMMIT => { let transaction = Transaction::from_bytes(&self.carried(peer, &request)?)?; - self.commit(path, &transaction)?; + self.commit(peer, path, &transaction)?; self.changed(&[path.to_owned()]); done } @@ -648,6 +676,15 @@ impl Served { let offset = request.offset.unwrap_or_default() as usize; let mut snapshots = self.snapshots.lock().unwrap(); snapshots.retain(|_, (_, read)| read.elapsed() < SNAPSHOT_AGE); + if let (kind::READ, None, Some(base)) = (kind, request.handle, &request.stamp) + && let Ok(base) = Stamp::try_from(base) + && let Some(changes) = self.changes(&request.path, &base)? + { + if let Some(guest) = self.guests.lock().unwrap().get_mut(peer) { + Self::holds(guest, &request.path); + } + return Ok(changes); + } let (image, handle) = match request.handle { Some(handle) => { let (image, read) = snapshots @@ -658,6 +695,11 @@ impl Served { } None => { drop(snapshots); + if kind == kind::READ + && let Some(guest) = self.guests.lock().unwrap().get_mut(peer) + { + Self::holds(guest, &request.path); + } let limit = (request.limit.unwrap_or(LIMIT as u64) as usize).min(LIMIT); let image = match kind { kind::READ => self.image(&request.path)?, @@ -692,6 +734,32 @@ impl Served { }) } + /// What changed in the section at `path` since the image with stamp `base`, where that + /// image is kept and the changes are much smaller than the section. + fn changes(&self, path: &str, base: &Stamp) -> Result> { + let before = self + .images + .lock() + .unwrap() + .iter() + .find_map(|(held, at, image)| (held == path && at == base).then(|| Arc::clone(image))); + let Some(before) = before else { + return Ok(None); + }; + let after = self.image(path)?; + let writes = match delta(&before, &after) { + Some(writes) => writes, + None if before == after => Vec::new(), + None => return Ok(None), + }; + Ok(Some(Reply { + writes: Some(writes), + length: Some(after.len() as u64), + stamp: Some((&Stamp::of(&after)?).into()), + ..Reply::default() + })) + } + /// The section or TOC at `path` as it stands: the image kept for it while its stamp holds. fn image(&self, path: &str) -> Result>> { let stamp = self.storage.stamp(path)?; @@ -711,18 +779,31 @@ impl Served { Ok(image) } + /// Keeps `image` as the section at `path` now, with the newest of each of `IMAGES` + /// sections and earlier images up to `VERSIONS` and `VERSIONS_BYTES`. fn keep(&self, path: &str, image: Arc>) { let Ok(stamp) = Stamp::of(&image) else { return; }; let mut images = self.images.lock().unwrap(); - images.retain(|(held, ..)| held != path); + images.retain(|(held, at, _)| held != path || *at != stamp); images.insert(0, (path.to_owned(), stamp, image)); - images.truncate(IMAGES); + let (mut newest, mut versions, mut bytes) = (Vec::new(), 0, 0); + images.retain(|(held, _, image)| { + if !newest.contains(held) { + newest.push(held.clone()); + return newest.len() <= IMAGES; + } + versions += 1; + bytes += image.len(); + newest.iter().take(IMAGES).any(|kept| kept == held) + && versions <= VERSIONS + && bytes <= VERSIONS_BYTES + }); } - /// Commits a guest's transaction once the section it makes parses. - fn commit(&self, path: &str, transaction: &Transaction) -> Result<()> { + /// Commits `guest`'s transaction once the section it makes parses. + fn commit(&self, guest: &[u8; 16], path: &str, transaction: &Transaction) -> Result<()> { let not_committed = |error: io::Error| { Error::Remote(CommitError { state: CommitState::NotCommitted, @@ -751,14 +832,15 @@ impl Served { ))); } self.storage.commit(path, transaction)?; - self.tell_delta(path, &image, &next); + self.tell_delta(path, &image, &next, Some(guest)); self.keep(path, Arc::new(next)); Ok(()) } - /// Tells every guest what a commit changed in the section at `path`, from `before` to - /// `after`, where that is small enough to send. - fn tell_delta(&self, path: &str, before: &[u8], after: &[u8]) { + /// Tells the guests holding the section at `path` what a commit changed in it, from + /// `before` to `after`, where that is small enough to send: all but the guest whose + /// commit it was, which has it already. + fn tell_delta(&self, path: &str, before: &[u8], after: &[u8], committed: Option<&[u8; 16]>) { let (Ok(base), Some(writes)) = (Stamp::of(before), delta(before, after)) else { return; }; @@ -768,8 +850,18 @@ impl Served { length: after.len() as u64, writes, }; - for guest in self.guests.lock().unwrap().values() { - let _ = guest.line.send(kind::DELTA, &delta); + let mut guests = self.guests.lock().unwrap(); + let holding: Vec<[u8; 16]> = guests + .iter_mut() + .filter(|(id, guest)| Some(*id) != committed && guest.held.iter().any(|h| h == path)) + .map(|(id, guest)| { + Self::holds(guest, path); + *id + }) + .collect(); + drop(guests); + if let (Some(room), false) = (self.room.get(), holding.is_empty()) { + room.send(kind::DELTA, &delta, Some(&holding)); } } @@ -789,12 +881,21 @@ impl Served { return; } if let Ok(after) = self.storage.read(path) { - self.tell_delta(path, &before, &after); + self.tell_delta(path, &before, &after, None); self.keep(path, Arc::new(after)); } } } +/// A stamp as `Touched` names it. +fn digest(stamp: &Stamp) -> u64 { + let hash = Sha256::new() + .chain_update(stamp.header) + .chain_update(stamp.length.to_le_bytes()) + .finalize(); + u64::from_le_bytes(hash[..8].try_into().expect("eight bytes")) +} + /// The writes that make `after` of `before`, a commit's appended bytes, patches and header, /// where they are much less than `after` itself. fn delta(before: &[u8], after: &[u8]) -> Option> { @@ -977,6 +1078,9 @@ struct Inner { asked: (Mutex, Condvar), /// The sections read lately, kept as the host's deltas change them, newest first. images: Mutex>)>>, + /// The stamps of sections whose image held is the host's now, as its last report of + /// them said. + current: Mutex>, } impl Guest { @@ -999,6 +1103,7 @@ impl Guest { stopped: AtomicBool::new(false), asked: Default::default(), images: Mutex::default(), + current: Mutex::default(), }); let heard = Arc::clone(&inner); let live = Live::start(me, &Room::Notebook(secret), reach, relay, move |event| { @@ -1123,17 +1228,32 @@ impl Guest { Ok(()) } - /// Reads the file at `path` a chunk at a time, as `kind` reads it. + /// Reads the file at `path` a chunk at a time, as `kind` reads it; a section held, as + /// the changes to it. fn read(&self, kind: u16, path: &str, limit: usize) -> io::Result> { + let held = (kind == kind::READ) + .then(|| self.inner.image(path)) + .flatten(); let first = self.ask( kind, Request { path: path.to_owned(), offset: Some(0), limit: Some(limit as u64), + stamp: (held.as_deref()) + .and_then(|image| Stamp::of(image).ok()) + .map(|stamp| (&stamp).into()), ..Request::default() }, )?; + if let (Some(writes), Some(held)) = (&first.writes, held) { + let stamp: Option = first.stamp.as_ref().and_then(|s| s.try_into().ok()); + let image = written(&held, first.length.unwrap_or_default(), writes) + .filter(|image| stamp.is_some() && Stamp::of(image).ok() == stamp) + .ok_or_else(|| io::Error::new(io::ErrorKind::InvalidData, "Changes that miss"))?; + self.inner.hold(path, Arc::new(image.clone())); + return Ok(image); + } let length = first.length.unwrap_or_default() as usize; if length > limit { return Err(io::ErrorKind::FileTooLarge.into()); @@ -1166,9 +1286,8 @@ impl Guest { /// The image of the section at `path` with stamp `stamp`, where this guest holds it. fn held(&self, path: &str, stamp: &Stamp) -> Option> { - let images = self.inner.images.lock().unwrap(); - let (_, image) = images.iter().find(|(held, _)| held == path)?; - (Stamp::of(image).ok().as_ref() == Some(stamp)).then(|| image.to_vec()) + let image = self.inner.image(path)?; + (Stamp::of(&image).ok().as_ref() == Some(stamp)).then(|| image.to_vec()) } pub(crate) fn entries(&self, folder: &str) -> io::Result> { @@ -1187,7 +1306,12 @@ impl Guest { .collect()) } + /// The stamp of the file at `path`: the host's last report of it where this guest holds + /// that image, else the host's answer. fn stamp(&self, path: &str) -> io::Result { + if let Some(stamp) = self.inner.current.lock().unwrap().get(path) { + return Ok(stamp.clone()); + } let reply = self.ask( kind::STAMP, Request { @@ -1243,6 +1367,13 @@ impl Guest { } impl Inner { + fn image(&self, path: &str) -> Option>> { + let images = self.images.lock().unwrap(); + images + .iter() + .find_map(|(held, image)| (held == path).then(|| Arc::clone(image))) + } + fn hold(&self, path: &str, image: Arc>) { let mut images = self.images.lock().unwrap(); images.retain(|(held, _)| held != path); @@ -1252,12 +1383,7 @@ impl Inner { /// Takes on a delta to an image held. fn apply(&self, delta: &Delta) { - let held = self - .images - .lock() - .unwrap() - .iter() - .find_map(|(path, image)| (*path == delta.path).then(|| Arc::clone(image))); + let held = self.image(&delta.path); let base: Option = (&delta.base).try_into().ok(); if let Some(image) = held && Stamp::of(&image).ok() == base @@ -1267,15 +1393,10 @@ impl Inner { } } - /// Takes on this guest's own commit to an image held. + /// Takes on this guest's own commit to an image held; the host's report, not this, says + /// whether the image is the host's, as another's commit may have followed. fn published(&self, path: &str, transaction: &Transaction) { - let held = self - .images - .lock() - .unwrap() - .iter() - .find_map(|(held, image)| (held == path).then(|| Arc::clone(image))); - if let Some(image) = held { + if let Some(image) = self.image(path) { let mut next = (*image).clone(); if transaction.apply(&mut next).is_ok() { self.hold(path, Arc::new(next)); @@ -1283,6 +1404,18 @@ impl Inner { } } + /// Takes the host's report of what the files at `paths` are now. + fn touched(&self, touched: &Touched) { + let mut current = self.current.lock().unwrap(); + for (at, path) in touched.paths.iter().enumerate() { + let stamp = self.image(path).and_then(|image| Stamp::of(&image).ok()); + match stamp.filter(|stamp| touched.stamps.get(at) == Some(&Some(digest(stamp)))) { + Some(stamp) => current.insert(path.clone(), stamp), + None => current.remove(path), + }; + } + } + fn heard(&self, event: Event) { let serves = |hello: &Hello| hello.serves == Some(self.share); match event { @@ -1291,6 +1424,7 @@ impl Inner { } Event::Left(hello) if serves(hello) => { *self.host.lock().unwrap() = None; + self.current.lock().unwrap().clear(); // Each request waiting hears its answer was lost. self.pending.lock().unwrap().clear(); if let Some(reports) = self.watch.lock().unwrap().take() { @@ -1313,10 +1447,11 @@ impl Inner { } } kind::TOUCHED if serves(from) => { - if let Ok(touched) = minicbor::decode::(body) - && let Some(reports) = &*self.watch.lock().unwrap() - { - reports.touched(&touched.paths); + if let Ok(touched) = minicbor::decode::(body) { + self.touched(&touched); + if let Some(reports) = &*self.watch.lock().unwrap() { + reports.touched(&touched.paths); + } } } kind::BYE @@ -1367,7 +1502,11 @@ impl crate::Remote for HostedRemote { } fn publish(&mut self, transaction: &Transaction) -> std::result::Result<(), CommitError> { - self.guest.commit(&self.path, transaction)?; + if let Err(error) = self.guest.commit(&self.path, transaction) { + // The host's image moved on, and its report may not have come yet. + self.guest.inner.current.lock().unwrap().remove(&self.path); + return Err(error); + } self.guest.inner.published(&self.path, transaction); Ok(()) } diff --git a/crates/notebook/src/live/tests.rs b/crates/notebook/src/live/tests.rs index 0d3cbbad3a7fa250b1ffec6c6bbb1db18c0c1bf7..39696aa0d7ab896ba77421948f835d1039252296 100644 --- a/crates/notebook/src/live/tests.rs +++ b/crates/notebook/src/live/tests.rs @@ -359,22 +359,25 @@ enum Tamper { } /// A relay in the middle of Grace's connection that passes on what the real one at -/// `upstream` says, except the message to her numbered `at` among those from peers, which it +/// `upstream` says, except the first message to her from a peer once `armed`, which it /// tampers with. Her next connection waits for `release`. struct Malicious { url: String, + armed: Arc, tampered: mpsc::Receiver<()>, rejoined: mpsc::Receiver<()>, release: mpsc::Sender<()>, } -fn malicious(upstream: SocketAddr, tamper: Tamper, at: usize) -> Malicious { +fn malicious(upstream: SocketAddr, tamper: Tamper) -> Malicious { use ::relay::ws::{self, Message}; let listener = TcpListener::bind("127.0.0.1:0").unwrap(); let url = format!("ws://{}", listener.local_addr().unwrap()); let (tampered, told) = mpsc::channel(); let (rejoined, heard) = mpsc::channel(); let (release, released) = mpsc::channel(); + let armed = Arc::new(std::sync::atomic::AtomicBool::new(false)); + let arming = Arc::clone(&armed); thread::spawn(move || { for (index, client) in listener.incoming().enumerate() { let mut client = client.unwrap(); @@ -389,19 +392,19 @@ fn malicious(upstream: SocketAddr, tamper: Tamper, at: usize) -> Malicious { let _ = io::copy(&mut up, &mut to_server); let _ = to_server.shutdown(Shutdown::Both); }); - let tampered = tampered.clone(); + let (tampered, armed) = (tampered.clone(), Arc::clone(&arming)); thread::spawn(move || { let mut reading = io::BufReader::new(server); let head = ws::head(&mut reading).unwrap(); client.write_all(head.as_bytes()).unwrap(); let mut messages = ws::Reader::new(reading, 1 << 20, false); - let (mut count, mut held) = (0, None); + let mut held = None; while let Ok(message) = messages.read() { let frames: Vec> = match message { Message::Binary(mut data) => { - count += 1; let mut out = vec![]; - if index == 0 && count - 1 == at { + if index == 0 && armed.swap(false, std::sync::atomic::Ordering::AcqRel) + { match tamper { Tamper::Drop => {} Tamper::Repeat => out = vec![data.clone(), data], @@ -411,7 +414,8 @@ fn malicious(upstream: SocketAddr, tamper: Tamper, at: usize) -> Malicious { out = vec![data]; } Tamper::Inject => { - // A slot, length and number, then made-up bytes. + // A slot and part of the sender's id, then + // made-up bytes. let mut forged = data.clone(); forged[16..].fill(7); out = vec![forged, data]; @@ -441,6 +445,7 @@ fn malicious(upstream: SocketAddr, tamper: Tamper, at: usize) -> Malicious { }); Malicious { url, + armed, tampered: told, rejoined: heard, release, @@ -460,9 +465,9 @@ fn recording(name: &str, room: &Room, relay: &str) -> (Live, Arc Some(None), - Some(link) => link.peer.presence.clone().map(Some), + Some(peer) => peer.presence.clone().map(Some), }; recorded.lock().unwrap().extend(entry); }) @@ -487,14 +492,17 @@ fn a_malicious_relay_is_caught() { let (url, address) = relay(Default::default()); let ada = Live::start(hello("Ada"), &room, None, Some(&url), |_| {}).unwrap(); ada.set_presence(caret(1)); - // Ada's opening, hello and first presence reach Grace; the next is tampered with. - let relay = malicious(address, tamper, 3); + let relay = malicious(address, tamper); let (grace, heard) = recording("Grace", &room, &relay.url); until(&grace, |peers| { peers .first() .is_some_and(|peer| peer.presence == Some(caret(1))) }); + // Ada's next presence goes to everyone, so its loss is caught too. + relay + .armed + .store(true, std::sync::atomic::Ordering::Release); ada.set_presence(caret(2)); relay .tampered diff --git a/crates/notebook/src/live/wire.rs b/crates/notebook/src/live/wire.rs index fa1697719761d8eced313906fd4b96376bd38a41..9cf8d2c396b7a7be1354bf954020f5e24f40f4c2 100644 --- a/crates/notebook/src/live/wire.rs +++ b/crates/notebook/src/live/wire.rs @@ -37,6 +37,8 @@ pub mod kind { /// Keeps a quiet connection open; it says nothing else. pub const PING: u16 = 2; pub const BYE: u16 = 3; + /// A `Hello` in answer to one heard in a notebook's room's group. + pub const HELLO_BACK: u16 = 4; pub const PRESENCE: u16 = 16; /// A host's files changed: `Touched`. pub const TOUCHED: u16 = 18; @@ -71,6 +73,7 @@ pub const KNOWN: &[u16] = &[ kind::HELLO, kind::PING, kind::BYE, + kind::HELLO_BACK, kind::PRESENCE, kind::TOUCHED, kind::DELTA, @@ -225,6 +228,10 @@ pub struct Welcome { pub struct Touched { #[n(0)] pub paths: Vec, + /// Each path's stamp now, as `share::digest` hashes it, where it has one: a guest holding + /// that image knows the stamp without asking. + #[n(1)] + pub stamps: Vec>, } /// What a commit changed in a section a host serves: from the image with stamp `base`, the @@ -271,7 +278,8 @@ pub struct Request { pub limit: Option, #[cbor(n(5), with = "minicbor::bytes")] pub bytes: Option>, - /// A confirmation's or supersession's base. + /// A confirmation's or supersession's base; for a section's read, the image the guest + /// holds, which the reply may give the changes to. #[n(6)] pub stamp: Option, #[cbor(n(7), with = "minicbor::bytes")] @@ -306,6 +314,9 @@ pub struct Reply { /// The snapshot a read's later chunks come from. #[n(7)] pub handle: Option, + /// A read's changes to the image with the request's stamp, in place of its bytes. + #[n(8)] + pub writes: Option>, } /// Why a request failed: an `io::ErrorKind` as `error_kind` numbers it, and for a commit, diff --git a/crates/notebook/tests/live_share.rs b/crates/notebook/tests/live_share.rs index 929ad88fd02d4d1f72118fb54635356d4d9f8ab4..95eadeb8882740a24e87d901e0aa8e0c7178738f 100644 --- a/crates/notebook/tests/live_share.rs +++ b/crates/notebook/tests/live_share.rs @@ -344,13 +344,16 @@ fn a_flooding_guest_is_hung_up_on() { }, ) .unwrap(); + // The line to the host, while Mallory's stream to it is open. let served = |live: &Live| { - live.peers() + let host = live + .peers() .into_iter() - .find(|peer| peer.hello.serves == Some(welcome.share)) + .find(|peer| peer.hello.serves == Some(welcome.share))?; + live.line(&host.hello.peer) }; until("Mallory never met the host", || served(&mallory).is_some()); - let line = mallory.line(&served(&mallory).unwrap().hello.peer).unwrap(); + let line = served(&mallory).unwrap(); const SENT: u64 = 5000; for id in 0..SENT { let request = Request { diff --git a/crates/relay/README.md b/crates/relay/README.md index 5237d02a6ade9aa33c0e8a58629b0b71dd9ac37e..a95a496196a2f6157f5a9aba333cf4c578a48617 100644 --- a/crates/relay/README.md +++ b/crates/relay/README.md @@ -2,9 +2,11 @@ The relay Live Share meets through when two Snowbounds aren't on one network. Peers join a room named by a tag (a hash of the notebook's secret, or a code's number), and the relay -passes their messages between them. Every message after the opening is sealed end to end -and numbered inside the seal, so the relay can't read, alter, drop or reorder one unseen; -it learns who talks to whom, when, and how much. `src/lib.rs` describes the protocol. +passes their messages between them, or copies one to everyone in the room (presence, a +host's news of its files) so a sender with thirty peers sends it once. Every message after +the opening is sealed end to end and numbered inside the seal, so the relay can't read, +alter, drop or reorder one unseen; it learns who talks to whom, when, and how much. +`src/lib.rs` describes the protocol. One static Linux executable with no configuration file and nothing on disk. It speaks plain HTTP 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 sudo install -m 644 /tmp/snowbound-relay.service /etc/systemd/system/ sudo systemctl daemon-reload sudo systemctl enable --now snowbound-relay -curl -s http://127.0.0.1:23592/health # {"rooms":0,"peers":0,"connections":1,"seconds":3} +curl -s http://127.0.0.1:23592/health # {"rooms":0,"peers":0,"connections":1,"seconds":3,"bytes_in":0,"bytes_out":0} ``` Then point a name at the server and put a TLS proxy in front. Caddy fetches its own @@ -115,7 +117,9 @@ snowbound.paperclover.net { ## Options -`snowbound-relay --help` lists every option and its default. Each can also be set in the +`snowbound-relay --help` lists every option and its default. A room holds 64 peers and +admits 120 a minute, enough for a class joining at once; its bytes per second count what +the relay gives out, so a copy to thirty peers costs thirty times its size. Each can also be set in the unit's environment as `SNOWBOUND_RELAY_