diff --git a/crates/notebook/README.md b/crates/notebook/README.md index eb3bdf08bbd503943cf90c4f473d17a209f502a0..8561bf8315d5b770f808d1649ad25a642044ef20 100644 --- a/crates/notebook/README.md +++ b/crates/notebook/README.md @@ -533,24 +533,38 @@ Device and simulator builds link for iOS. Native acceptance uses disposable OneNote 2010 clients and Samba; it does not establish on-device execution or physical power-loss durability. -## Live presence (feature `live`) +## Live presence and Live Share (feature `live`) -`live::Live::start(hello, room, reach, relay, notify)` listens on a TCP port and, with a +`live::Live::start(hello, room, reach, relay, events)` listens on a TCP port and, with a `Reach`, advertises `_snowbound._tcp` by mDNS on every network or on loopback alone, -connecting to the peers in the same `Room` that it finds: a notebook's identity, or a code -typed on both (`Room::Code("4-violet-otter")`). With a `relay` (`wss://live.example.net`, -`crates/relay`) it also joins the room there and meets its peers through it; a code's words -alone (`Room::Code("violet-otter")`) ask the relay for a number, and `code()` then has the -whole code. `connect(address)` meets a peer discovery did not find. Peers meet through -SPAKE2 on the room's secret, then every frame is AES-256-GCM under the keys it agreed: its -number, which is also its nonce, then a message kind and a CBOR map (`live::wire`). A frame -lost, repeated, reordered or forged on the way fails where it lands; the connection is -dropped as broken, nothing from it after the fault is applied, and the ends meet again from -scratch. A reader skips kinds and map keys it doesn't know, so later versions add both -freely. `set_presence` says which section, page and caret this end has (text object and -UTF-16 offset, as ops address text); a connection sends only the newest. `peers()` lists -each connected peer's `Hello` (name, picture) and presence, and `notify` runs whenever that -changes. Dropping the `Live` leaves. +connecting to the peers in the same `Room` that it finds: a notebook's or a share's random +secret (`Room::Notebook`), or a code typed on both (`Room::join("4-violet-otter", password)`). +With a `relay` (`wss://live.example.net`, `crates/relay`) it also joins the room there and +meets its peers through it; the end sharing a code (`Room::share`) has the relay number its +words, `code()` then has the whole code, and it burns a code after too many wrong tries. +`connect(address)` meets a peer discovery did not find. Peers meet through SPAKE2 on the +room's secret, then every frame is AES-256-GCM under the keys it agreed: its number, which is +also its nonce, then a message kind and a CBOR map (`live::wire`). A frame lost, repeated, +reordered or forged on the way fails where it lands; the connection is dropped as broken, +nothing from it after the fault is applied, and the ends meet again from scratch. A reader +skips kinds and map keys it doesn't know, so later versions add both freely. +`set_presence` says which section, page and caret this end has (text object and UTF-16 +offset, as ops address text); a connection sends only the newest. `peers()` lists each +connected peer's `Hello` (name, picture) and presence; `events` hears when they change, who +was met and left, and every other frame with the `Line` to answer on. `leave(reason)` says +`Bye` first; dropping the `Live` leaves at once. `Notebook::presence_room()` is the secret of +a notebook's presence room, kept in `.snowbound/live.json` and made where it has none. + +`live::share` is Live Share. `Host::start(storage, hello, sharing, name, reach, relay, events)` +serves `Notebook::into_storage()` to the peers in the share's room and welcomes whoever knows +`Sharing::code` (and its password) from the code's room; `code()` replaces a burned code with +new words, `guests()` lists who is connected, `touched(paths)` passes the host's own changes +on, and `stop()` lets every guest go. `join(hello, code, password, reach, relay)` returns the +`Welcome` (the share and its secret), or a `Refusal` saying why not: a wrong code, no one +sharing it, an expired code, too many wrong tries, or no relay. `Guest::start` joins the share; +`Notebook::open_hosted(guest, cache)`, `Background::hosted(guest, notify)` and +`Section::resume_hosted(path, replica, guest, notify)` then work as on a share, with +`Guest::host()` and `stopped()` for whether the host is there and still sharing. ## Queue measurement diff --git a/crates/notebook/src/background.rs b/crates/notebook/src/background.rs index 931b12540c81375be10bfacfff8506c1d2305331..5e1d130905e8120176e24a4b037986e5bbc94210 100644 --- a/crates/notebook/src/background.rs +++ b/crates/notebook/src/background.rs @@ -62,10 +62,15 @@ pub(crate) struct Shared { watched: Mutex, /// Asked for by `discard`, once the thread stops. discard: AtomicBool, + /// Hears every report of changed paths (`on_touched`). + listener: Mutex>, } +/// Hears the catalog paths a report names. +pub type Listener = Box; + /// A connection's report of its notebook folder's changes, from a watch armed on connecting. -#[cfg_attr(not(feature = "smb"), allow(dead_code))] +#[cfg_attr(not(any(feature = "smb", feature = "live")), allow(dead_code))] pub(crate) struct Reports { shared: Weak, connection: u64, @@ -135,6 +140,7 @@ impl Background { signal, watched: Mutex::default(), discard: AtomicBool::new(false), + listener: Mutex::default(), }); let weak = Arc::downgrade(&shared); let owner = Arc::clone(&shared); @@ -342,6 +348,35 @@ impl Background { ) } + /// Keeps the sections of a notebook a Live Share host serves in sync while they are not + /// open, with an offline copy of each; the host reports what changed. Paths are catalog + /// paths. + #[cfg(feature = "live")] + pub fn hosted( + guest: Arc, + notify: impl Fn() + Send + 'static, + ) -> Result { + Self::start( + true, + move |reports| { + guest.watch(reports)?; + let (bound, listing) = (Arc::clone(&guest), Arc::clone(&guest)); + let bind = move |path: &str| crate::live::share::HostedRemote::new(&bound, path); + let list = move |folder: &str| listing.entries(folder); + Ok(((bind, list), true)) + }, + notify, + ) + } + + /// Has `listener` hear every report of changed paths from now on, as a Live Share host + /// passes them to its guests; `None` stops. + pub fn on_touched(&self, listener: Option) { + if let Ok(mut held) = self.0.listener.lock() { + *held = listener; + } + } + /// Watches these sections (`Notebook::replicas`); a section not watched before is first /// checked soon after, the next one a little later. pub fn watch(&self, sections: Vec) { @@ -484,6 +519,11 @@ impl Shared { watched.touched(paths, Instant::now()); } self.signal.wake(); + if let Ok(listener) = self.listener.lock() + && let Some(listener) = &*listener + { + listener(paths); + } } /// The session holding the section at `path` let it go: it is checked now. @@ -499,7 +539,7 @@ impl Shared { } } -#[cfg_attr(not(feature = "smb"), allow(dead_code))] +#[cfg_attr(not(any(feature = "smb", feature = "live")), allow(dead_code))] impl Reports { /// The sections at or below these paths changed. pub(crate) fn touched(&self, paths: &[String]) { diff --git a/crates/notebook/src/live.rs b/crates/notebook/src/live.rs index 17347158e5fb8b25185b5eb2350b47cdbdb9fcb1..5fdce7106330e72794b9ddb5f5d7427945b45de4 100644 --- a/crates/notebook/src/live.rs +++ b/crates/notebook/src/live.rs @@ -1,16 +1,19 @@ -//! Live presence: who else has the notebook open, the page they are on and their caret, -//! straight from one Snowbound to another, or through a relay (`crates/relay`) where they -//! aren't on one network. Peers find each other with mDNS (`_snowbound._tcp`) or in the -//! relay's room, and meet through a secret both hold, a notebook's identity or a code typed -//! on both, which SPAKE2 turns into the keys every frame after the opening is sealed with. -//! Each connection has a thread reading and one writing. A connection whose frames arrive -//! out of order is dropped and met again from scratch. +//! Live presence and Live Share: who else has the notebook open, the page they are on and +//! their caret, straight from one Snowbound to another, or through a relay (`crates/relay`) +//! where they aren't on one network; and a notebook one machine holds opened on another +//! (`share`). Peers find each other with mDNS (`_snowbound._tcp`) or in the relay's room, and +//! meet through a secret both hold, a room's or a code typed on both, which SPAKE2 turns into +//! the keys every frame after the opening is sealed with. Each connection has a thread reading +//! and one writing. A connection whose frames arrive out of order is dropped and met again +//! from scratch. mod relay; +pub mod share; pub mod wire; pub use wire::{Caret, Guid, Hello, Presence, Spot}; use mdns_sd::{IfKind, ServiceDaemon, ServiceEvent, ServiceInfo}; +use minicbor::Encode; use sha2::{Digest, Sha256}; use std::{ collections::BTreeMap, @@ -22,7 +25,7 @@ use std::{ mpsc, }, thread, - time::Duration, + time::{Duration, Instant}, }; use wire::{Sealer, Side, kind}; @@ -33,37 +36,73 @@ const GONE: Duration = Duration::from_secs(45); const OPENING: Duration = Duration::from_secs(5); /// 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. +const TRIES: u32 = 5; /// The secret peers meet through. #[derive(Clone, Debug, PartialEq, Eq)] pub enum Room { - /// Everyone with the notebook: its table of contents' file identity, which only its files - /// hold. + /// A notebook's room, or a share's: a random secret only its members hold. Notebook([u8; 16]), - /// A code typed on both: `7-violet-otter`, whose number names it on the network and whose - /// words only the two people know. Its words alone (`violet-otter`) ask the relay for a - /// free number, and `Live::code` then has the whole code. - Code(String), + /// A code typed on both ends, `412-violet-otter`, with a password where one is set: its + /// number names it on the network, and its words and password only the two people know. + /// Its `owner`, the end sharing it, takes the number from a relay (the one the code has, + /// coming back; any free one for words alone) or picks one where it has no relay, and + /// judges every try, burning the code after too many wrong ones. + Code { + code: String, + password: String, + owner: bool, + }, } impl Room { - /// What names the room in the clear: a hash of a notebook's identity, a code's number; - /// none for a code the relay hasn't numbered. + /// The room a code typed on this end leads to. + pub fn join(code: &str, password: &str) -> Self { + Self::Code { + code: code.trim().to_lowercase(), + password: password.to_owned(), + owner: false, + } + } + + /// The room of a code this end shares: its words, or a whole code to keep its number. + pub fn share(code: &str, password: &str) -> Self { + Self::Code { + code: code.trim().to_lowercase(), + password: password.to_owned(), + owner: true, + } + } + + /// What names the room in the clear: a hash of a room's secret, a code's number; none for + /// a code without one yet. fn tag(&self) -> Option { match self { Room::Notebook(id) => Some(hex(&Sha256::digest( [&b"Snowbound room "[..], id].concat(), )[..8])), - Room::Code(code) => code_parts(code).0.map(|number| format!("code-{number}")), + Room::Code { code, .. } => code_parts(code).0.map(|number| format!("code-{number}")), } } fn secret(&self) -> Vec { match self { Room::Notebook(id) => id.to_vec(), - Room::Code(code) => code_parts(code).1.to_lowercase().into_bytes(), + Room::Code { code, password, .. } => { + let mut secret = code_parts(code).1.to_lowercase().into_bytes(); + if !password.is_empty() { + secret.push(b'\n'); + secret.extend_from_slice(password.as_bytes()); + } + secret + } } } + + fn owner(&self) -> bool { + matches!(self, Room::Code { owner: true, .. }) + } } /// A code's number, if it has one, and its words. @@ -91,19 +130,49 @@ pub struct Peer { pub presence: Option, } +/// What happened, as `Live::start`'s `events` hears it on a network thread. +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. + Met(&'a Arc, &'a Line), + /// A peer's connection ended. + Left(&'a Arc), + /// A frame of a kind presence doesn't read itself, with the line to answer on. + Frame { + from: &'a Arc, + kind: u16, + body: &'a [u8], + line: &'a Line, + }, +} + +/// How the relay last answered. +#[derive(Clone, Debug, PartialEq, Eq)] +pub enum Relayed { + /// Not asked yet, or no relay. + Unknown, + /// In the room. + Joined, + /// Answered with this HTTP status, and how long it asked to wait. + Refused(u16, Option), + /// Not reached. + Unreachable, +} + /// Presence on the network while it lives; dropping it leaves. pub struct Live { shared: Arc, address: SocketAddr, - daemon: Option, } struct Shared { me: Hello, room: Room, secret: Vec, + reach: Option, state: Mutex, - notify: Box, + events: Box, stopped: AtomicBool, connections: AtomicU64, } @@ -118,15 +187,43 @@ struct State { code: Option, /// The relay connection open now, to hang up on leaving. relay: Option>, + relayed: Option, + /// Meetings that failed on the secret: wrong codes tried here, or this end's. + failed: u32, + /// The code admits no one new. + burned: bool, + daemon: Option, } struct Link { connection: u64, peer: Peer, - wake: mpsc::Sender<()>, + line: Line, pipe: Arc, } +/// What a connection's writer sends next. +enum Out { + /// The newest presence. + Presence, + Frame(u16, Vec), + /// A `Bye`, after which it hangs up. + Bye(Vec), +} + +/// The way to one connected peer: frames sent on it go after those sent before. +#[derive(Clone)] +pub struct Line(mpsc::Sender); + +impl Line { + pub fn send(&self, kind: u16, body: &impl Encode<()>) -> io::Result<()> { + let body = minicbor::to_vec(body).map_err(io::Error::other)?; + self.0 + .send(Out::Frame(kind, body)) + .map_err(|_| io::ErrorKind::NotConnected.into()) + } +} + /// A stream to one peer, read by one thread and written by another: a TCP connection, or /// one carried through a relay. trait Pipe: Send + Sync { @@ -185,22 +282,14 @@ impl Write for &dyn Pipe { impl Live { /// Starts listening as `me` in `room`, advertised and looked for where `reach` says, or /// not at all with `None`, leaving peers to `connect`, and in the room at `relay` - /// (`wss://live.example.net`) where given. `notify` runs on a network thread whenever - /// `peers` or `code` changes. + /// (`wss://live.example.net`) where given. `events` runs on a network thread. pub fn start( me: Hello, room: &Room, reach: Option, relay: Option<&str>, - notify: impl Fn() + Send + Sync + 'static, + events: impl Fn(Event) + Send + Sync + 'static, ) -> io::Result { - let tag = room.tag(); - if tag.is_none() && relay.is_none() { - return Err(io::Error::new( - io::ErrorKind::InvalidInput, - "Only a relay can number a code", - )); - } let host = match reach { Some(Reach::Network) => IpAddr::V4(Ipv4Addr::UNSPECIFIED), Some(Reach::Loopback) | None => IpAddr::V4(Ipv4Addr::LOCALHOST), @@ -211,26 +300,36 @@ impl Live { address.set_ip(IpAddr::V4(Ipv4Addr::LOCALHOST)); } let code = match room { - Room::Code(code) if tag.is_some() => Some(code.trim().to_owned()), + Room::Code { code, .. } if room.tag().is_some() => Some(code.clone()), + // Off any relay, the end sharing words alone numbers them itself. + Room::Code { code, .. } if relay.is_none() => { + let mut number = [0; 2]; + getrandom::fill(&mut number) + .map_err(|_| io::Error::other("System random source failed"))?; + Some(format!( + "{}-{code}", + 1000 + u16::from_le_bytes(number) % 9000 + )) + } _ => None, }; let shared = Arc::new(Shared { me, room: room.clone(), secret: room.secret(), + reach, state: Mutex::new(State { code, ..State::default() }), - notify: Box::new(notify), + events: Box::new(events), stopped: AtomicBool::new(false), connections: AtomicU64::new(0), }); if let Some(relay) = relay { - relay::join(&shared, relay)?; + relay::join(&shared, relay, address.port())?; } let accepting = Arc::clone(&shared); - let accepted = tag.clone(); thread::Builder::new() .name("live accept".into()) .spawn(move || { @@ -238,23 +337,14 @@ impl Live { if accepting.stopped.load(Ordering::Acquire) { return; } - if let (Ok(stream), Some(tag)) = (stream, accepted.clone()) { + if let (Ok(stream), Some(tag)) = (stream, accepting.tag()) { let shared = Arc::clone(&accepting); thread::spawn(move || shared.run(Arc::new(stream), Side::Responder, &tag)); } } })?; - let daemon = match (reach, tag) { - (Some(reach), Some(tag)) => { - Some(advertise(&shared, reach, tag, address.port()).map_err(io::Error::other)?) - } - _ => None, - }; - Ok(Live { - shared, - address, - daemon, - }) + shared.advertise(address.port()); + Ok(Live { shared, address }) } /// Where this end listens. @@ -263,11 +353,28 @@ impl Live { } /// The code others type to meet this end: the one it was given, or the one the relay - /// numbered. + /// or this end numbered. pub fn code(&self) -> Option { self.shared.state.lock().unwrap().code.clone() } + /// How the relay last answered. + pub fn relayed(&self) -> Relayed { + let state = self.shared.state.lock().unwrap(); + state.relayed.clone().unwrap_or(Relayed::Unknown) + } + + /// Meetings that failed on the secret: wrong tries of this end's code, or this end's own + /// wrong code. + pub fn failed(&self) -> u32 { + self.shared.state.lock().unwrap().failed + } + + /// Whether this end's code had too many wrong tries and admits no one new. + pub fn burned(&self) -> bool { + self.shared.state.lock().unwrap().burned + } + /// Connects to a peer at `address` that discovery did not find. pub fn connect(&self, address: SocketAddr) { let shared = Arc::clone(&self.shared); @@ -283,7 +390,7 @@ impl Live { state.presence = presence; state.generation += 1; for link in state.peers.values() { - let _ = link.wake.send(()); + let _ = link.line.0.send(Out::Presence); } } @@ -292,17 +399,38 @@ impl Live { let state = self.shared.state.lock().unwrap(); state.peers.values().map(|link| link.peer.clone()).collect() } + + /// The line to `peer`, while it is connected. + pub fn line(&self, peer: &[u8; 16]) -> Option { + let state = self.shared.state.lock().unwrap(); + state.peers.get(peer).map(|link| link.line.clone()) + } + + /// Says `reason` to every peer and leaves, waiting a moment for them to hear it. + pub fn leave(self, reason: &str) { + let bye = minicbor::to_vec(wire::Bye { + reason: reason.into(), + }) + .unwrap_or_default(); + for link in self.shared.state.lock().unwrap().peers.values() { + let _ = link.line.0.send(Out::Bye(bye.clone())); + } + let deadline = Instant::now() + Duration::from_secs(1); + while !self.shared.state.lock().unwrap().peers.is_empty() && Instant::now() < deadline { + thread::sleep(Duration::from_millis(20)); + } + } } impl Drop for Live { fn drop(&mut self) { self.shared.stopped.store(true, Ordering::Release); - if let Some(daemon) = &self.daemon { - let _ = daemon.shutdown(); - } // Wakes the accepting thread to see it has stopped. let _ = TcpStream::connect_timeout(&self.address, OPENING); - let state = self.shared.state.lock().unwrap(); + let mut state = self.shared.state.lock().unwrap(); + if let Some(daemon) = state.daemon.take() { + let _ = daemon.shutdown(); + } for link in state.peers.values() { link.pipe.shutdown(); } @@ -312,69 +440,34 @@ impl Drop for Live { } } -/// Advertises `shared` as in room `tag` on `port` and connects to the peers in its room that -/// discovery finds with a higher id than its own, which leave the connecting to it. -fn advertise( - shared: &Arc, - reach: Reach, - tag: String, - port: u16, -) -> mdns_sd::Result { - let daemon = ServiceDaemon::new()?; - let id = hex(&shared.me.peer); - let properties = [("v", "1"), ("room", tag.as_str()), ("peer", id.as_str())]; - let host = format!("snowbound-{id}.local."); - let info = match reach { - Reach::Network => { - ServiceInfo::new(SERVICE, &id, &host, "", port, &properties[..])?.enable_addr_auto() - } - Reach::Loopback => { - daemon.disable_interface(IfKind::All)?; - daemon.enable_interface(IfKind::LoopbackV4)?; - ServiceInfo::new( - SERVICE, - &id, - &host, - IpAddr::V4(Ipv4Addr::LOCALHOST), - port, - &properties[..], - )? - } - }; - daemon.register(info)?; - let found = daemon.browse(SERVICE)?; - let shared = Arc::clone(shared); - thread::Builder::new() - .name("live discovery".into()) - .spawn(move || { - while let Ok(event) = found.recv() { - let ServiceEvent::ServiceResolved(service) = event else { - continue; - }; - let (Some(room), Some(peer)) = ( - service.get_property_val_str("room"), - service.get_property_val_str("peer"), - ) else { - continue; - }; - if room != tag || peer <= id.as_str() || shared.knows(peer) { - continue; - } - let mut addresses: Vec = - service.addresses.iter().map(|ip| ip.to_ip_addr()).collect(); - addresses.sort_by_key(|ip| (!ip.is_ipv4(), !ip.is_loopback())); - if let Some(ip) = addresses.first() { - let address = SocketAddr::new(*ip, service.port); - let shared = Arc::clone(&shared); - thread::spawn(move || shared.dial(address)); - } - } - }) - .map_err(|error| mdns_sd::Error::Msg(error.to_string()))?; - Ok(daemon) -} - impl Shared { + /// The room's tag as it stands: a code's, once numbered. + fn tag(&self) -> Option { + match &self.room { + Room::Code { .. } => code_parts(self.state.lock().unwrap().code.as_deref()?) + .0 + .map(|number| format!("code-{number}")), + room => room.tag(), + } + } + + /// Advertises this end on `port` once its room has a tag, where its reach says, and + /// connects to the peers in its room that discovery finds with a higher id than its own, + /// which leave the connecting to it. + fn advertise(self: &Arc, port: u16) { + let (Some(reach), Some(tag)) = (self.reach, self.tag()) else { + return; + }; + let mut state = self.state.lock().unwrap(); + if state.daemon.is_some() || self.stopped.load(Ordering::Acquire) { + return; + } + match discover(self, reach, tag, port) { + Ok(daemon) => state.daemon = Some(daemon), + Err(error) => eprintln!("Live: no discovery: {error}"), + } + } + fn knows(&self, peer: &str) -> bool { let state = self.state.lock().unwrap(); state.peers.keys().any(|id| hex(id) == peer) @@ -383,11 +476,11 @@ impl Shared { /// Connects to `address`, and again while the peer is there and the connection was the /// one this end kept. fn dial(self: Arc, address: SocketAddr) { - let Some(tag) = self.room.tag() else { - return; - }; let mut wait = Duration::from_secs(1); while let Ok(stream) = TcpStream::connect_timeout(&address, OPENING) { + let Some(tag) = self.tag() else { + return; + }; if !Arc::clone(&self).run(Arc::new(stream), Side::Initiator, &tag) || self.stopped.load(Ordering::Acquire) { @@ -398,10 +491,24 @@ impl Shared { } } + /// Counts a meeting that failed on the secret: the end sharing a code burns it after too + /// many, and an end that typed one gives up at once. + fn failed(&self) { + let mut state = self.state.lock().unwrap(); + state.failed += 1; + match &self.room { + Room::Code { owner: true, .. } if state.failed >= TRIES => state.burned = true, + Room::Code { owner: false, .. } => self.stopped.store(true, Ordering::Release), + _ => {} + } + drop(state); + (self.events)(Event::Changed); + } + /// Meets the peer at the other end of `pipe` in room `tag`, then reads from it until it /// goes: whether it was the connection kept to that peer. fn run(self: Arc, pipe: Arc, side: Side, tag: &str) -> bool { - if self.stopped.load(Ordering::Acquire) { + if self.stopped.load(Ordering::Acquire) || self.state.lock().unwrap().burned { pipe.shutdown(); return false; } @@ -424,14 +531,19 @@ impl Shared { Ok(met) => met, Err(error) => { eprintln!("Live: no meeting in {tag}: {error}"); + if error.kind() == io::ErrorKind::InvalidData { + self.failed(); + } pipe.shutdown(); return false; } }; let peer = hello.peer; let name = hello.name.clone(); + let hello = Arc::new(hello); let connection = self.connections.fetch_add(1, Ordering::Relaxed); - let (wake, woken) = mpsc::channel(); + let (out, outgoing) = mpsc::channel(); + let line = Line(out); { let mut state = self.state.lock().unwrap(); // A peer met both directly and through a relay keeps the direct connection, as @@ -450,39 +562,48 @@ impl Shared { pipe.shutdown(); return false; } - let _ = wake.send(()); + let _ = line.0.send(Out::Presence); state.peers.insert( peer, Link { connection, peer: Peer { - hello: Arc::new(hello), + hello: Arc::clone(&hello), presence: None, }, - wake, + line: line.clone(), pipe: Arc::clone(&pipe), }, ); } - (self.notify)(); let writing = Arc::clone(&pipe); let shared = Arc::clone(&self); - thread::spawn(move || shared.write(&*writing, send, woken)); + thread::spawn(move || shared.write(&*writing, send, outgoing)); + (self.events)(Event::Met(&hello, &line)); + (self.events)(Event::Changed); let ended = loop { let (message, body) = match receive.receive(&mut stream) { Ok(frame) => frame, Err(error) => break Some(error), }; - if message != kind::PRESENCE { - continue; + match message { + kind::HELLO | kind::PING => {} + kind::PRESENCE => { + let Ok(presence) = minicbor::decode::(&body) else { + break None; + }; + if let Some(link) = self.state.lock().unwrap().peers.get_mut(&peer) { + link.peer.presence = Some(presence); + } + (self.events)(Event::Changed); + } + kind => (self.events)(Event::Frame { + from: &hello, + kind, + body: &body, + line: &line, + }), } - let Ok(presence) = minicbor::decode::(&body) else { - break None; - }; - if let Some(link) = self.state.lock().unwrap().peers.get_mut(&peer) { - link.peer.presence = Some(presence); - } - (self.notify)(); }; if let Some(error) = ended.filter(|error| error.kind() == io::ErrorKind::InvalidData) { eprintln!("Live: the connection to {name} broke ({error}); meeting again"); @@ -497,18 +618,19 @@ impl Shared { { state.peers.remove(&peer); drop(state); - (self.notify)(); + (self.events)(Event::Left(&hello)); + (self.events)(Event::Changed); } true } - /// Sends the newest presence whenever woken, and a ping when quiet. - fn write(&self, pipe: &dyn Pipe, mut send: Sealer, woken: mpsc::Receiver<()>) { + /// Sends the newest presence whenever woken, 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; loop { - let result = match woken.recv_timeout(PING) { - Ok(()) => { + let result = match outgoing.recv_timeout(PING) { + Ok(Out::Presence) => { let presence = { let state = self.state.lock().unwrap(); if sent == Some(state.generation) { @@ -519,6 +641,12 @@ impl Shared { }; send.send(&mut stream, kind::PRESENCE, &presence) } + Ok(Out::Frame(kind, body)) => send.send_encoded(&mut stream, kind, &body), + Ok(Out::Bye(body)) => { + let _ = send.send_encoded(&mut stream, kind::BYE, &body); + pipe.shutdown(); + return; + } Err(mpsc::RecvTimeoutError::Timeout) => send.send(&mut stream, kind::PING, &()), Err(mpsc::RecvTimeoutError::Disconnected) => return, }; @@ -530,6 +658,70 @@ impl Shared { } } +/// Advertises `shared` as in room `tag` on `port` and connects to the peers in its room that +/// discovery finds with a higher id than its own. +fn discover( + shared: &Arc, + reach: Reach, + tag: String, + port: u16, +) -> mdns_sd::Result { + let daemon = ServiceDaemon::new()?; + let id = hex(&shared.me.peer); + let properties = [("v", "1"), ("room", tag.as_str()), ("peer", id.as_str())]; + let host = format!("snowbound-{id}.local."); + let info = match reach { + Reach::Network => { + ServiceInfo::new(SERVICE, &id, &host, "", port, &properties[..])?.enable_addr_auto() + } + Reach::Loopback => { + daemon.disable_interface(IfKind::All)?; + daemon.enable_interface(IfKind::LoopbackV4)?; + ServiceInfo::new( + SERVICE, + &id, + &host, + IpAddr::V4(Ipv4Addr::LOCALHOST), + port, + &properties[..], + )? + } + }; + daemon.register(info)?; + let found = daemon.browse(SERVICE)?; + let shared = Arc::downgrade(shared); + thread::Builder::new() + .name("live discovery".into()) + .spawn(move || { + while let Ok(event) = found.recv() { + let ServiceEvent::ServiceResolved(service) = event else { + continue; + }; + let Some(shared) = shared.upgrade() else { + return; + }; + let (Some(room), Some(peer)) = ( + service.get_property_val_str("room"), + service.get_property_val_str("peer"), + ) else { + continue; + }; + if room != tag || peer <= id.as_str() || shared.knows(peer) { + continue; + } + let mut addresses: Vec = + service.addresses.iter().map(|ip| ip.to_ip_addr()).collect(); + addresses.sort_by_key(|ip| (!ip.is_ipv4(), !ip.is_loopback())); + if let Some(ip) = addresses.first() { + let address = SocketAddr::new(*ip, service.port); + thread::spawn(move || shared.dial(address)); + } + } + }) + .map_err(|error| mdns_sd::Error::Msg(error.to_string()))?; + Ok(daemon) +} + fn hex(bytes: &[u8]) -> String { bytes.iter().map(|byte| format!("{byte:02x}")).collect() } diff --git a/crates/notebook/src/live/relay.rs b/crates/notebook/src/live/relay.rs index cfa5256e878dc2fb997bc58179249c78016fe9b6..b892514292c45361d9dd60005350782bebce751e 100644 --- a/crates/notebook/src/live/relay.rs +++ b/crates/notebook/src/live/relay.rs @@ -3,7 +3,7 @@ //! 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. -use super::{OPENING, PATIENCE, Pipe, Shared, Side, code_parts}; +use super::{Event, OPENING, PATIENCE, Pipe, Relayed as Answer, Shared, Side, code_parts}; use ::relay::{Notice, SLOT, Verdict, ws}; use base64::Engine; use rustls::{ClientConfig, ClientConnection, RootCertStore, pki_types::ServerName}; @@ -29,13 +29,14 @@ const MOST: usize = 1 << 20; type Reader = ws::Reader>>; -/// Joins `shared`'s room at the relay at `url`, and keeps joining whenever it falls out. -pub(super) fn join(shared: &Arc, url: &str) -> io::Result<()> { +/// Joins `shared`'s room at the relay at `url`, and keeps joining whenever it falls out; +/// `port` is where it listens, to advertise once the relay numbers its code. +pub(super) fn join(shared: &Arc, url: &str, port: u16) -> io::Result<()> { let address = parse(url)?; let shared = Arc::clone(shared); thread::Builder::new() .name("live relay".into()) - .spawn(move || keep(&shared, &address))?; + .spawn(move || keep(&shared, &address, port))?; Ok(()) } @@ -87,41 +88,56 @@ impl From for Failure { } } +/// Records how the relay answered, telling `shared`'s events where it changed. +fn answered(shared: &Shared, answer: Answer) { + let mut state = shared.state.lock().unwrap(); + if state.relayed.as_ref() != Some(&answer) { + state.relayed = Some(answer); + drop(state); + (shared.events)(Event::Changed); + } +} + /// Stays in the room at `address`, joining again whenever the connection ends, until /// `shared` stops. -fn keep(shared: &Arc, address: &Address) { +fn keep(shared: &Arc, address: &Address, port: u16) { let mut wait = Duration::from_secs(1); - // The number the relay gave this end's code, asked for again on joining again. - let mut nameplate = None; + let owner = shared.room.owner(); + // The number of this end's code, asked for again on joining again. + let mut nameplate = owner + .then(|| code_parts(shared.state.lock().unwrap().code.as_deref()?).0) + .flatten(); while !shared.stopped.load(Ordering::Acquire) { let began = Instant::now(); - let tag = shared.room.tag(); - let path = match &tag { - Some(tag) => format!("{}/v1/room/{tag}", address.path), - None => match nameplate { + let path = match (owner, shared.tag()) { + (false, Some(tag)) => format!("{}/v1/room/{tag}", address.path), + (false, None) => return, + (true, _) => match nameplate { Some(number) => format!("{}/v1/claim?nameplate={number}", address.path), None => format!("{}/v1/claim", address.path), }, }; - match connect(address, &path, tag.is_none()) { + match connect(address, &path, owner) { Ok((socket, reader)) => { shared.state.lock().unwrap().relay = Some(Arc::clone(&socket)); + answered(shared, Answer::Joined); if !shared.stopped.load(Ordering::Acquire) { - session(shared, &socket, reader, &mut nameplate); + session(shared, &socket, reader, &mut nameplate, port); } shared.state.lock().unwrap().relay = None; socket.hang_up(); } - Err(Failure::Refused(410, _)) => { - eprintln!("Live: the code has expired; ask for a new one"); - return; - } Err(Failure::Refused(status, retry)) => { eprintln!("Live: the relay refused to let this end in ({status})"); + answered(shared, Answer::Refused(status, retry)); + if status == 410 { + return; + } wait = wait.max(retry.unwrap_or_default()); } Err(Failure::Network(error)) => { eprintln!("Live: no relay at {}: {error}", address.authority); + answered(shared, Answer::Unreachable); } } if began.elapsed() >= STEADY { @@ -337,8 +353,9 @@ fn session( socket: &Arc, mut reader: Reader, nameplate: &mut Option, + port: u16, ) { - let mut tag = shared.room.tag(); + let mut tag = shared.tag(); // Pings while the session lasts: dropping `_beat` at its end stops them. let (_beat, beats) = mpsc::channel::<()>(); let pinging = Arc::clone(socket); @@ -355,16 +372,16 @@ fn session( Ok(Notice::Nameplate(number)) => { *nameplate = Some(number); tag = Some(format!("code-{number}")); - let super::Room::Code(words) = &shared.room else { + let super::Room::Code { code: words, .. } = &shared.room else { continue; }; let code = format!("{number}-{}", code_parts(words).1); let mut state = shared.state.lock().unwrap(); if state.code.as_ref() != Some(&code) { - eprintln!("Live: the code is {code}"); state.code = Some(code); drop(state); - (shared.notify)(); + shared.advertise(port); + (shared.events)(Event::Changed); } } Ok(Notice::Welcome { you, members }) => { @@ -378,7 +395,8 @@ fn session( Ok(Notice::Burned) => { // Coming back, ask for another number. *nameplate = None; - eprintln!("Live: too many wrong tries; the code admits no one new"); + shared.state.lock().unwrap().burned = true; + (shared.events)(Event::Changed); } Err(()) => {} }, diff --git a/crates/notebook/src/live/share.rs b/crates/notebook/src/live/share.rs new file mode 100644 index 0000000000000000000000000000000000000000..87c001fd3fd56d1be675a1200300aba76ed662c6 --- /dev/null +++ b/crates/notebook/src/live/share.rs @@ -0,0 +1,1342 @@ +//! Live Share: a notebook one Snowbound holds, opened on others through a short code. The +//! host serves its notebook's storage verbs (`session::Storage`) to each peer in the share's +//! room; a guest runs the replica, queue and merge it runs on an SMB share against those +//! verbs, so offline queueing, rebases and conflict pages work as there, and the host's files +//! stay what its own storage writes. A guest first meets the host in the code's room, where +//! the host welcomes it with the share's room and secret; a new share has a new secret, so +//! stopping retires every guest. Large bodies travel a chunk at a time, each answered before +//! the next, so a relay never holds much for a slow peer. +//! +//! The code's words come from the EFF's short word list +//! (, CC BY 3.0 US), `yo-yo` replaced by `yarn`. + +use super::{ + Event, Hello, Line, Live, Peer, Presence, Reach, Relayed, Room, + wire::{self, Failure, Reply, Request, Touched, Welcome, WireEntry, WireStamp, kind}, +}; +use crate::{Error, Result, background::Reports, discover, session::Storage}; +use onestore::{CommitError, CommitState, RevisionIndex, Stamp, Store, Transaction}; +use std::{ + collections::{BTreeMap, HashMap}, + io, + path::PathBuf, + sync::{ + Arc, Condvar, Mutex, + atomic::{AtomicBool, AtomicU64, Ordering}, + mpsc, + }, + thread, + time::{Duration, Instant}, +}; + +/// The most bytes one message of a read or an upload carries. +const CHUNK: usize = 128 << 10; +/// The most bytes of chunks a guest has asked for and not yet been given. +const WINDOW: usize = 512 << 10; +/// How long a request waits for its reply. +const TIMEOUT: Duration = Duration::from_secs(60); +/// The largest file read or written whole. +const LIMIT: usize = 256 << 20; +/// Snapshots of files being read, per guest, and how long one is kept unread. +const SNAPSHOTS: usize = 8; +const SNAPSHOT_AGE: Duration = Duration::from_secs(120); +/// Sections whose image a host keeps to check guests' commits on. +const IMAGES: usize = 4; +/// What a host says leaving as it stops sharing. +const STOPPED: &str = "stopped"; + +const WORDS: &str = include_str!("words.txt"); + +/// Two random words for a code, `violet-otter`. +pub fn words() -> io::Result { + let list: Vec<&str> = WORDS.lines().collect(); + let mut bytes = [0; 8]; + getrandom::fill(&mut bytes).map_err(|_| io::Error::other("System random source failed"))?; + let [first, second] = [&bytes[..4], &bytes[4..]] + .map(|bytes| u32::from_le_bytes(bytes.try_into().expect("4 bytes")) as usize); + let first = first % list.len(); + let mut second = second % (list.len() - 1); + if second >= first { + second += 1; + } + Ok(format!("{}-{}", list[first], list[second])) +} + +/// A code as typed, `412 Violet otter`, in the form it is met by, `412-violet-otter`; none +/// where it is not a number and two words. +pub fn code(typed: &str) -> Option { + let parts: Vec = typed + .split(|c: char| c.is_whitespace() || c == '-') + .filter(|part| !part.is_empty()) + .map(str::to_lowercase) + .collect(); + match &parts[..] { + [number, first, second] + if number.parse::().is_ok() + && [first, second] + .iter() + .all(|word| word.chars().all(|c| c.is_ascii_lowercase())) => + { + Some(parts.join("-")) + } + _ => None, + } +} + +/// What a host keeps of a share to take it up again after a relaunch. +#[derive(Clone, Debug, PartialEq, Eq, serde::Serialize, serde::Deserialize)] +pub struct Sharing { + pub share: [u8; 16], + /// The share room's secret, which every guest welcomed holds. + pub secret: [u8; 16], + /// The code: its words, with its number in front once one was given. + pub code: String, + pub password: String, +} + +impl Sharing { + /// A new share: its own id and secret, and new words. + pub fn new(password: &str) -> io::Result { + let mut random = [0; 32]; + getrandom::fill(&mut random) + .map_err(|_| io::Error::other("System random source failed"))?; + Ok(Self { + share: random[..16].try_into().expect("16 bytes"), + secret: random[16..].try_into().expect("16 bytes"), + code: words()?, + password: password.to_owned(), + }) + } +} + +/// Where a guest keeps a share's notebook: `live://` and the share's id. +pub fn location(share: &[u8; 16]) -> String { + format!("live://{}", super::hex(share)) +} + +/// Why joining failed. +#[derive(Clone, Debug, PartialEq, Eq)] +pub enum Refusal { + /// Not a number and two words. + Malformed, + /// The code's words or password are wrong. + Wrong, + /// No one shares with the code's number now. + NoOne, + /// The code had too many wrong tries. + Expired, + /// Too many wrong codes from this network; try again after the wait, where known. + TooMany(Option), + /// The relay is full. + Busy, + /// The relay couldn't be reached, and no one answered on this network. + Unreachable, + /// The relay let this end in, but no one answered. + TimedOut, +} + +/// Meets the host sharing `code` (and `password`) as `me`, on the networks `reach` names and +/// through `relay`: the share it welcomes this end to. +pub fn join( + me: Hello, + code: &str, + password: &str, + reach: Option, + relay: Option<&str>, +) -> std::result::Result { + let code = self::code(code).ok_or(Refusal::Malformed)?; + let (welcomed, welcome) = mpsc::channel(); + let (changed, waiting) = mpsc::channel(); + let live = Live::start( + me, + &Room::join(&code, password), + reach, + relay, + move |event| match event { + Event::Frame { + kind: kind::WELCOME, + body, + .. + } => { + if let Ok(body) = minicbor::decode::(body) { + let _ = welcomed.send(body); + } + } + _ => { + let _ = changed.send(()); + } + }, + ) + .map_err(|_| Refusal::Unreachable)?; + let start = Instant::now(); + loop { + if let Ok(welcome) = welcome.try_recv() { + return Ok(welcome); + } + if live.failed() > 0 { + return Err(Refusal::Wrong); + } + // A peer on this network may yet answer what the relay refused. + let waited = start.elapsed(); + let settled = waited > Duration::from_secs(3) || reach.is_none(); + 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)), + Relayed::Refused(503, _) if settled => return Err(Refusal::Busy), + Relayed::Unreachable if waited > Duration::from_secs(10) => { + return Err(Refusal::Unreachable); + } + Relayed::Unknown if relay.is_none() && waited > Duration::from_secs(10) => { + return Err(Refusal::NoOne); + } + _ if waited > Duration::from_secs(20) => return Err(Refusal::TimedOut), + _ => {} + } + let _ = waiting.recv_timeout(Duration::from_millis(100)); + } +} + +/// A notebook shared while it lives: the share's room, serving the notebook's storage to the +/// guests in it, and the code's room, welcoming whoever knows the code. +pub struct Host { + me: Hello, + notebook: String, + reach: Option, + relay: Option, + sharing: Mutex, + room: Live, + pairing: Mutex, + served: Arc, + events: Arc, +} + +impl Host { + /// Shares `storage`, the notebook named `notebook`, as `sharing` says, as `me`, where + /// `reach` and `relay` say. `events` runs on a network thread whenever the guests or the + /// code change. What guests change reaches the host as its own watch on the notebook's + /// folder reports it, and reaches the other guests at once. + pub fn start( + storage: Box, + me: Hello, + sharing: Sharing, + notebook: &str, + reach: Option, + relay: Option<&str>, + events: impl Fn() + Send + Sync + 'static, + ) -> io::Result { + let events: Arc = Arc::new(events); + let served = Arc::new(Served { + storage, + images: Mutex::default(), + snapshots: Mutex::default(), + puts: Mutex::default(), + lines: Mutex::default(), + }); + let serving = Hello { + serves: Some(sharing.share), + ..me.clone() + }; + let (heard, told) = (Arc::clone(&served), Arc::clone(&events)); + let room = Live::start( + serving, + &Room::Notebook(sharing.secret), + reach, + relay, + move |event| match event { + Event::Met(hello, line) => { + heard.lines.lock().unwrap().insert(hello.peer, line.clone()); + } + Event::Left(hello) => heard.forget(&hello.peer), + Event::Frame { + from, + kind, + body, + line, + } if wire::KNOWN.contains(&kind) && kind > 256 && kind != kind::REPLY => { + let (served, line, peer) = (Arc::clone(&heard), line.clone(), from.peer); + let body = body.to_vec(); + thread::spawn(move || { + let reply = served.handle(&peer, kind, &body); + let _ = line.send(kind::REPLY, &reply); + }); + } + Event::Changed => told(), + _ => {} + }, + )?; + let host = Self { + notebook: notebook.to_owned(), + reach, + relay: relay.map(str::to_owned), + pairing: Mutex::new(pair(&me, &sharing, notebook, reach, relay, &events)?), + me, + sharing: Mutex::new(sharing), + room, + served, + events, + }; + Ok(host) + } + + /// The code guests type, once it has its number. A code with too many wrong tries is + /// replaced by one with new words. + pub fn code(&self) -> Option { + let mut pairing = self.pairing.lock().unwrap(); + if pairing.burned() { + let mut sharing = self.sharing.lock().unwrap(); + sharing.code = words().ok()?; + let relay = self.relay.as_deref(); + *pairing = pair( + &self.me, + &sharing, + &self.notebook, + self.reach, + relay, + &self.events, + ) + .ok()?; + } + let code = pairing.code(); + if let Some(code) = &code { + self.sharing.lock().unwrap().code = code.clone(); + } + code + } + + /// The share as it stands, to take up again after a relaunch. + pub fn sharing(&self) -> Sharing { + self.code(); + self.sharing.lock().unwrap().clone() + } + + /// How the relay last answered the code's room. + pub fn relayed(&self) -> Relayed { + self.pairing.lock().unwrap().relayed() + } + + /// The peers in the share's room. + pub fn guests(&self) -> Vec { + self.room.peers() + } + + pub fn set_presence(&self, presence: Presence) { + self.room.set_presence(presence); + } + + /// Tells every guest the files at these catalog paths changed. + pub fn touched(&self, paths: &[String]) { + self.served.tell(paths); + } + + /// Stops sharing: every guest hears so and is let go. + pub fn stop(self) { + drop(self.pairing); + self.room.leave(STOPPED); + } +} + +/// The code's room, welcoming whoever knows the code to the share. +fn pair( + me: &Hello, + sharing: &Sharing, + notebook: &str, + reach: Option, + relay: Option<&str>, + events: &Arc, +) -> io::Result { + let welcome = Welcome { + share: sharing.share, + secret: sharing.secret, + notebook: notebook.to_owned(), + host: me.name.clone(), + }; + let told = Arc::clone(events); + Live::start( + Hello { + serves: None, + ..me.clone() + }, + &Room::share(&sharing.code, &sharing.password), + reach, + relay, + move |event| match event { + Event::Met(_, line) => { + let _ = line.send(kind::WELCOME, &welcome); + } + Event::Changed => told(), + _ => {} + }, + ) +} + +/// A host's side of its guests' storage requests. +struct Served { + storage: Box, + /// Sections' images by path, with their stamps, to check commits on. + images: Mutex>, + /// Files being read a chunk at a time, by guest and the read's first request. + snapshots: Mutex>, Instant)>>, + /// Bytes a later request carries, by guest and upload. + puts: Mutex>>, + lines: Mutex>, +} + +/// A section's path, stamp and image. +type Image = (String, Stamp, Arc>); +/// What a guest's requests left, by guest and request. +type ByGuest = HashMap<([u8; 16], u64), T>; + +/// Whether a guest may name `path`: a catalog path inside the notebook, and not presence's +/// own secret, which guests meet in the share's room instead. +fn allowed(path: &str) -> bool { + path.split('/').all(|part| { + !part.is_empty() && part != "." && part != ".." && !part.contains(['\\', '\0', ':']) + }) && !path.to_ascii_lowercase().starts_with(".snowbound/live") +} + +/// The folder holding catalog path `path`. +fn folder(path: &str) -> String { + path.rsplit_once('/') + .map_or(String::new(), |(folder, _)| folder.to_owned()) +} + +fn refused(kind: io::ErrorKind, message: &str) -> Error { + io::Error::new(kind, message.to_owned()).into() +} + +impl Served { + fn forget(&self, peer: &[u8; 16]) { + self.lines.lock().unwrap().remove(peer); + self.snapshots + .lock() + .unwrap() + .retain(|(guest, _), _| guest != peer); + self.puts + .lock() + .unwrap() + .retain(|(guest, _), _| guest != peer); + } + + /// Tells every guest the files at `paths` changed. + fn tell(&self, paths: &[String]) { + let touched = Touched { + paths: paths.to_vec(), + }; + for line in self.lines.lock().unwrap().values() { + let _ = line.send(kind::TOUCHED, &touched); + } + } + + fn handle(&self, peer: &[u8; 16], kind: u16, body: &[u8]) -> Reply { + let request = match minicbor::decode::(body) { + Ok(request) => request, + Err(_) => { + return failed( + 0, + refused(io::ErrorKind::InvalidData, "A malformed request"), + ); + } + }; + let id = request.id; + match self.answer(peer, kind, request) { + Ok(reply) => Reply { id, ..reply }, + Err(error) => failed(id, error), + } + } + + fn answer(&self, peer: &[u8; 16], kind: u16, request: Request) -> Result { + let path = request.path.as_str(); + if !(path.is_empty() && matches!(kind, kind::LIST | kind::PUT) || allowed(path)) + || request.to.as_deref().is_some_and(|to| !allowed(to)) + { + return Err(refused( + io::ErrorKind::PermissionDenied, + "Outside the notebook", + )); + } + let to = || { + request + .to + .as_deref() + .ok_or_else(|| refused(io::ErrorKind::InvalidInput, "No target")) + }; + let stamp = || -> Result { + Ok(request + .stamp + .as_ref() + .ok_or_else(|| refused(io::ErrorKind::InvalidInput, "No stamp"))? + .try_into()?) + }; + let done = Reply::default(); + Ok(match kind { + kind::LIST => Reply { + entries: Some( + self.storage + .entries(path)? + .into_iter() + .filter(|entry| allowed(&within(path, &entry.name))) + .map(|entry| wire_entry(&entry)) + .collect(), + ), + ..done + }, + kind::STAMP => Reply { + stamp: Some((&self.storage.stamp(path)?).into()), + ..done + }, + kind::EXISTS => Reply { + exists: Some(self.storage.exists(path)), + ..done + }, + kind::READ | kind::READ_FILE => self.read(peer, kind, &request)?, + kind::PUT => { + let handle = request.handle.unwrap_or_default(); + let bytes = request.bytes.unwrap_or_default(); + let mut puts = self.puts.lock().unwrap(); + let held = puts.entry((*peer, handle)).or_default(); + if request.offset != Some(held.len() as u64) || held.len() + bytes.len() > LIMIT { + return Err(refused( + io::ErrorKind::InvalidInput, + "An upload out of order", + )); + } + held.extend_from_slice(&bytes); + done + } + kind::COMMIT => { + let transaction = Transaction::from_bytes(&self.carried(peer, &request)?)?; + self.commit(path, &transaction)?; + self.tell(&[path.to_owned()]); + done + } + kind::CONFIRM => { + self.storage.confirm(path, &stamp()?)?; + done + } + kind::CREATE => { + self.storage.create(path, &self.carried(peer, &request)?)?; + self.tell(&[folder(path)]); + done + } + kind::CREATE_DIRECTORY => { + self.storage.create_directory(path)?; + self.tell(&[folder(path)]); + done + } + kind::HIDE => { + self.storage.hide(path)?; + done + } + kind::RENAME | kind::REPLACE => { + let to = to()?; + if kind == kind::RENAME { + self.storage.rename(path, to)?; + } else { + self.storage.replace(path, to)?; + } + self.tell(&[folder(path), folder(to)]); + done + } + kind::DELETE => { + self.storage.delete(path)?; + self.tell(&[folder(path)]); + done + } + kind::PLACE => { + let ancestor = request + .ancestor + .ok_or_else(|| refused(io::ErrorKind::InvalidInput, "No ancestor"))?; + let name = request.name.as_deref().unwrap_or_default(); + self.storage.place(path, ancestor, name)?; + self.tell(&[path.to_owned()]); + done + } + kind::SUPERSEDE => { + self.storage.supersede(path, &stamp()?, to()?)?; + self.tell(&[folder(path)]); + done + } + _ => return Err(refused(io::ErrorKind::Unsupported, "An unknown request")), + }) + } + + /// The bytes a request carries, itself or in its uploads. + fn carried(&self, peer: &[u8; 16], request: &Request) -> Result> { + match (&request.bytes, request.handle) { + (Some(bytes), _) => Ok(bytes.clone()), + (None, Some(handle)) => self + .puts + .lock() + .unwrap() + .remove(&(*peer, handle)) + .ok_or_else(|| refused(io::ErrorKind::InvalidInput, "No such upload")), + (None, None) => Ok(Vec::new()), + } + } + + /// A chunk of a file as it stood when its read began. + fn read(&self, peer: &[u8; 16], kind: u16, request: &Request) -> Result { + let offset = request.offset.unwrap_or_default() as usize; + let mut snapshots = self.snapshots.lock().unwrap(); + snapshots.retain(|_, (_, read)| read.elapsed() < SNAPSHOT_AGE); + let (image, handle) = match request.handle { + Some(handle) => { + let (image, read) = snapshots + .get_mut(&(*peer, handle)) + .ok_or_else(|| refused(io::ErrorKind::TimedOut, "The read went stale"))?; + *read = Instant::now(); + (Arc::clone(image), handle) + } + None => { + drop(snapshots); + let limit = (request.limit.unwrap_or(LIMIT as u64) as usize).min(LIMIT); + let image = match kind { + kind::READ => self.image(&request.path)?, + _ => Arc::new(self.storage.read_file(&request.path, limit)?), + }; + if image.len() > limit { + return Err(io::Error::from(io::ErrorKind::FileTooLarge).into()); + } + snapshots = self.snapshots.lock().unwrap(); + if snapshots.keys().filter(|(guest, _)| guest == peer).count() >= SNAPSHOTS { + snapshots.retain(|(guest, _), _| guest != peer); + } + (image, request.id) + } + }; + let end = image.len().min(offset.saturating_add(CHUNK)); + let bytes = image.get(offset..end).unwrap_or_default().to_vec(); + if end < image.len() { + snapshots.insert((*peer, handle), (Arc::clone(&image), Instant::now())); + } else { + snapshots.remove(&(*peer, handle)); + } + Ok(Reply { + bytes: Some(bytes), + length: Some(image.len() as u64), + handle: Some(handle), + stamp: (kind == kind::READ) + .then(|| Stamp::of(&image).ok()) + .flatten() + .map(|stamp| (&stamp).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)?; + let kept = self + .images + .lock() + .unwrap() + .iter() + .find_map(|(held, at, image)| { + (held == path && *at == stamp).then(|| Arc::clone(image)) + }); + if let Some(image) = kept { + return Ok(image); + } + let image = Arc::new(self.storage.read(path)?); + self.keep(path, Arc::clone(&image)); + Ok(image) + } + + 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.insert(0, (path.to_owned(), stamp, image)); + images.truncate(IMAGES); + } + + /// Commits a guest's transaction once the section it makes parses. + fn commit(&self, path: &str, transaction: &Transaction) -> Result<()> { + let not_committed = |error: io::Error| { + Error::Remote(CommitError { + state: CommitState::NotCommitted, + error, + }) + }; + let image = self.image(path).map_err(|error| match error { + Error::Io(error) => not_committed(error), + error => error, + })?; + if Stamp::of(&image).ok().as_ref() != Some(transaction.base()) { + return Err(not_committed(io::Error::new( + io::ErrorKind::ResourceBusy, + "The section changed since", + ))); + } + let mut next = (*image).clone(); + let checked = transaction.apply(&mut next).and_then(|()| { + let store = Store::parse(&next)?; + RevisionIndex::parse(&store).map(drop) + }); + if let Err(error) = checked { + return Err(not_committed(io::Error::new( + io::ErrorKind::InvalidData, + error.to_string(), + ))); + } + self.storage.commit(path, transaction)?; + self.keep(path, Arc::new(next)); + Ok(()) + } +} + +fn within(folder: &str, name: &str) -> String { + if folder.is_empty() { + name.to_owned() + } else { + format!("{folder}/{name}") + } +} + +fn wire_entry(entry: &discover::Entry) -> WireEntry { + WireEntry { + name: entry.name.clone(), + kind: match entry.kind { + discover::EntryKind::File => 0, + discover::EntryKind::Directory => 1, + discover::EntryKind::Other => 2, + discover::EntryKind::Evicted => 3, + }, + size: entry.listed.size, + modified: entry.listed.modified, + } +} + +fn entry(entry: &WireEntry) -> discover::Entry { + discover::Entry { + name: entry.name.clone(), + kind: match entry.kind { + 0 => discover::EntryKind::File, + 1 => discover::EntryKind::Directory, + 3 => discover::EntryKind::Evicted, + _ => discover::EntryKind::Other, + }, + listed: discover::Listed { + size: entry.size, + modified: entry.modified, + }, + } +} + +fn state_number(state: CommitState) -> u8 { + match state { + CommitState::NotCommitted => 0, + CommitState::Unknown => 1, + CommitState::Committed => 2, + } +} + +fn failed(id: u64, error: Error) -> Reply { + let (kind, state) = match &error { + Error::Io(error) | Error::RemoteIo(error) => (error.kind(), None), + Error::Remote(error) => (error.error.kind(), Some(state_number(error.state))), + Error::Document(_) => (io::ErrorKind::InvalidData, None), + _ => (io::ErrorKind::Other, None), + }; + let message = match &error { + Error::Remote(error) => error.error.to_string(), + error => error.to_string(), + }; + Reply { + id, + failure: Some(Failure { + kind: wire::error_number(kind), + message, + state, + }), + ..Reply::default() + } +} + +/// How a request failed. +enum Failed { + /// It never left this end. + Unsent(io::Error), + /// Its answer was lost: whatever it asked may have happened. + Lost(io::Error), + /// The host refused it. + Refused(Failure), +} + +impl Failed { + fn io(self) -> io::Error { + match self { + Failed::Unsent(error) | Failed::Lost(error) => error, + Failed::Refused(failure) => { + io::Error::new(wire::error_kind(failure.kind), failure.message) + } + } + } + + /// As a commit's failure: one never sent was not committed, one whose answer was lost may + /// have been. + fn commit(self) -> CommitError { + let state = match &self { + Failed::Unsent(_) => CommitState::NotCommitted, + Failed::Lost(_) => CommitState::Unknown, + Failed::Refused(failure) => match failure.state { + Some(1) => CommitState::Unknown, + Some(2) => CommitState::Committed, + _ => CommitState::NotCommitted, + }, + }; + CommitError { + state, + error: self.io(), + } + } +} + +/// A guest's way to a share: the share's room, and the requests waiting on its host. +pub struct Guest { + live: Live, + inner: Arc, +} + +struct Inner { + share: [u8; 16], + /// The host and the line to it, while connected. + host: Mutex, Line)>>, + pending: Mutex>>, + next: AtomicU64, + /// Where the host's reports of changed files go, while a background watches. + watch: Mutex>, + /// The host stopped sharing. + stopped: AtomicBool, + /// Bytes of chunks asked for and not yet given. + asked: (Mutex, Condvar), +} + +impl Guest { + /// Joins share `share` through its `secret` as `me`, where `reach` and `relay` say. + /// `events` runs on a network thread whenever the host comes or goes, or the peers change. + pub fn start( + me: Hello, + share: [u8; 16], + secret: [u8; 16], + reach: Option, + relay: Option<&str>, + events: impl Fn() + Send + Sync + 'static, + ) -> io::Result> { + let inner = Arc::new(Inner { + share, + host: Mutex::default(), + pending: Mutex::default(), + next: AtomicU64::new(1), + watch: Mutex::default(), + stopped: AtomicBool::new(false), + asked: Default::default(), + }); + let heard = Arc::clone(&inner); + let live = Live::start(me, &Room::Notebook(secret), reach, relay, move |event| { + heard.heard(event); + events(); + })?; + Ok(Arc::new(Self { live, inner })) + } + + /// Where the notebook is kept: `location(share)`. + pub fn location(&self) -> String { + location(&self.inner.share) + } + + /// The host, while connected. + pub fn host(&self) -> Option> { + let host = self.inner.host.lock().unwrap(); + host.as_ref().map(|(hello, _)| Arc::clone(hello)) + } + + /// Whether the host said it stopped sharing. + pub fn stopped(&self) -> bool { + self.inner.stopped.load(Ordering::Acquire) + } + + /// Everyone in the share's room: the host and the other guests. + pub fn peers(&self) -> Vec { + self.live.peers() + } + + pub fn set_presence(&self, presence: Presence) { + self.live.set_presence(presence); + } + + /// Sends the host's reports of changed files to `reports` while it stays connected. + pub(crate) fn watch(&self, reports: Reports) -> io::Result<()> { + if self.host().is_none() { + return Err(self.offline()); + } + *self.inner.watch.lock().unwrap() = Some(reports); + Ok(()) + } + + fn offline(&self) -> io::Error { + let message = if self.stopped() { + "The host stopped sharing this notebook" + } else { + "The computer sharing this notebook can’t be reached" + }; + io::Error::new(io::ErrorKind::NotConnected, message) + } + + /// Asks the host `kind` of `request`, waiting for its reply. + fn request(&self, kind: u16, mut request: Request) -> std::result::Result { + let Some((_, line)) = self.inner.host.lock().unwrap().clone() else { + return Err(Failed::Unsent(self.offline())); + }; + let chunked = matches!(kind, kind::READ | kind::READ_FILE | kind::PUT); + if chunked { + let (asked, room) = &self.inner.asked; + let mut asked = asked.lock().unwrap(); + while *asked + CHUNK > WINDOW { + asked = room.wait(asked).unwrap(); + } + *asked += CHUNK; + } + let id = self.inner.next.fetch_add(1, Ordering::Relaxed); + request.id = id; + let (answer, answered) = mpsc::channel(); + self.inner.pending.lock().unwrap().insert(id, answer); + let sent = line.send(kind, &request); + let reply = match sent { + Err(error) => Err(Failed::Unsent(error)), + Ok(()) => match answered.recv_timeout(TIMEOUT) { + Ok(reply) => Ok(reply), + Err(mpsc::RecvTimeoutError::Timeout) => Err(Failed::Lost(io::Error::new( + io::ErrorKind::TimedOut, + "The computer sharing this notebook didn’t answer", + ))), + Err(mpsc::RecvTimeoutError::Disconnected) => Err(Failed::Lost(self.offline())), + }, + }; + self.inner.pending.lock().unwrap().remove(&id); + if chunked { + let (asked, room) = &self.inner.asked; + *asked.lock().unwrap() -= CHUNK; + room.notify_all(); + } + let reply = reply?; + match reply.failure { + Some(failure) => Err(Failed::Refused(failure)), + None => Ok(reply), + } + } + + fn ask(&self, kind: u16, request: Request) -> io::Result { + self.request(kind, request).map_err(Failed::io) + } + + /// Puts `bytes` in `request`, or uploads them first where they are large. + fn carry(&self, request: &mut Request, bytes: Vec) -> std::result::Result<(), Failed> { + if bytes.len() <= CHUNK { + request.bytes = Some(bytes); + return Ok(()); + } + let handle = self.inner.next.fetch_add(1, Ordering::Relaxed); + for (index, chunk) in bytes.chunks(CHUNK).enumerate() { + let put = Request { + handle: Some(handle), + offset: Some((index * CHUNK) as u64), + bytes: Some(chunk.to_vec()), + ..Request::default() + }; + // Nothing the upload carries happens before the request that uses it. + self.request(kind::PUT, put) + .map_err(|failed| match failed { + Failed::Lost(error) => Failed::Unsent(error), + failed => failed, + })?; + } + request.handle = Some(handle); + Ok(()) + } + + /// Reads the file at `path` a chunk at a time, as `kind` reads it. + fn read(&self, kind: u16, path: &str, limit: usize) -> io::Result> { + let first = self.ask( + kind, + Request { + path: path.to_owned(), + offset: Some(0), + limit: Some(limit as u64), + ..Request::default() + }, + )?; + let length = first.length.unwrap_or_default() as usize; + if length > limit { + return Err(io::ErrorKind::FileTooLarge.into()); + } + let mut image = first.bytes.unwrap_or_default(); + while image.len() < length { + let chunk = self + .ask( + kind, + Request { + path: path.to_owned(), + offset: Some(image.len() as u64), + handle: first.handle, + ..Request::default() + }, + )? + .bytes + .unwrap_or_default(); + if chunk.is_empty() { + return Err(io::ErrorKind::UnexpectedEof.into()); + } + image.extend_from_slice(&chunk); + } + image.truncate(length); + Ok(image) + } + + pub(crate) fn entries(&self, folder: &str) -> io::Result> { + let reply = self.ask( + kind::LIST, + Request { + path: folder.to_owned(), + ..Request::default() + }, + )?; + Ok(reply + .entries + .unwrap_or_default() + .iter() + .map(entry) + .collect()) + } + + fn stamp(&self, path: &str) -> io::Result { + let reply = self.ask( + kind::STAMP, + Request { + path: path.to_owned(), + ..Request::default() + }, + )?; + reply + .stamp + .as_ref() + .ok_or_else(|| io::Error::from(io::ErrorKind::InvalidData))? + .try_into() + } + + fn commit( + &self, + path: &str, + transaction: &Transaction, + ) -> std::result::Result<(), CommitError> { + let mut request = Request { + path: path.to_owned(), + ..Request::default() + }; + self.carry(&mut request, transaction.to_bytes()) + .map_err(Failed::commit)?; + self.request(kind::COMMIT, request) + .map(drop) + .map_err(Failed::commit) + } + + fn confirm(&self, path: &str, base: &Stamp) -> std::result::Result<(), CommitError> { + let request = Request { + path: path.to_owned(), + stamp: Some(base.into()), + ..Request::default() + }; + self.request(kind::CONFIRM, request) + .map(drop) + .map_err(Failed::commit) + } + + /// A request on `path` that answers nothing but whether it happened. + fn verb(&self, kind: u16, path: &str, request: Request) -> Result<()> { + self.ask( + kind, + Request { + path: path.to_owned(), + ..request + }, + )?; + Ok(()) + } +} + +impl Inner { + fn heard(&self, event: Event) { + let serves = |hello: &Hello| hello.serves == Some(self.share); + match event { + Event::Met(hello, line) if serves(hello) => { + *self.host.lock().unwrap() = Some((Arc::clone(hello), line.clone())); + } + Event::Left(hello) if serves(hello) => { + *self.host.lock().unwrap() = None; + // Each request waiting hears its answer was lost. + self.pending.lock().unwrap().clear(); + if let Some(reports) = self.watch.lock().unwrap().take() { + reports.lost(); + } + } + Event::Frame { + from, kind, body, .. + } => match kind { + kind::REPLY if serves(from) => { + if let Ok(reply) = minicbor::decode::(body) + && let Some(waiting) = self.pending.lock().unwrap().remove(&reply.id) + { + let _ = waiting.send(reply); + } + } + kind::TOUCHED if serves(from) => { + if let Ok(touched) = minicbor::decode::(body) + && let Some(reports) = &*self.watch.lock().unwrap() + { + reports.touched(&touched.paths); + } + } + kind::BYE + if serves(from) + && minicbor::decode::(body) + .is_ok_and(|bye| bye.reason == STOPPED) => + { + self.stopped.store(true, Ordering::Release); + } + _ => {} + }, + _ => {} + } + } +} + +/// A section a Live Share host serves, as a guest's replica publishes to it. +pub struct HostedRemote { + guest: Arc, + path: String, +} + +impl HostedRemote { + pub fn new(guest: &Arc, path: &str) -> Self { + Self { + guest: Arc::clone(guest), + path: path.to_owned(), + } + } +} + +impl crate::Remote for HostedRemote { + fn read(&mut self) -> io::Result> { + self.guest.read(kind::READ, &self.path, LIMIT) + } + + fn stamp(&mut self) -> io::Result { + self.guest.stamp(&self.path) + } + + fn publish(&mut self, transaction: &Transaction) -> std::result::Result<(), CommitError> { + self.guest.commit(&self.path, transaction) + } + + fn confirm(&mut self, base: &Stamp) -> std::result::Result<(), CommitError> { + self.guest.confirm(&self.path, base) + } +} + +/// A notebook a Live Share host serves, as a guest's `session::Notebook` reaches it. Its +/// folders as last listed are kept at `listed`, so the notebook opens while the host can't +/// be reached. +pub struct Hosted { + guest: Arc, + listed: PathBuf, +} + +/// Each folder's entries as last listed: name, `WireEntry::kind`, size and modified time. +type Listings = BTreeMap>; + +impl Hosted { + pub(crate) fn new(guest: Arc, listed: PathBuf) -> Self { + Self { guest, listed } + } +} + +/// The host's folders, or as they were last listed while it can't be reached. +struct Source<'a> { + guest: &'a Guest, + kept: Listings, + listed: Listings, +} + +impl discover::Source for Source<'_> { + fn entries(&mut self, path: &str, limit: usize) -> io::Result> { + let entries = match self.guest.entries(path) { + Ok(entries) => entries, + Err(error) if error.kind() == io::ErrorKind::NotConnected => self + .kept + .get(path) + .ok_or(error)? + .iter() + .map(|(name, kind, size, modified)| { + entry(&WireEntry { + name: name.clone(), + kind: *kind, + size: *size, + modified: *modified, + }) + }) + .collect(), + Err(error) => return Err(error), + }; + if entries.len() > limit { + return Err(io::ErrorKind::FileTooLarge.into()); + } + self.listed.insert( + path.to_owned(), + entries + .iter() + .map(|listed| { + let wire = wire_entry(listed); + (wire.name, wire.kind, wire.size, wire.modified) + }) + .collect(), + ); + Ok(entries) + } + + fn read(&mut self, path: &str, limit: usize) -> io::Result> { + self.guest.read(kind::READ, path, limit) + } + + fn read_asset(&mut self, path: &str, limit: usize) -> io::Result> { + self.guest.read(kind::READ_FILE, path, limit) + } +} + +impl Storage for Hosted { + fn discover( + &self, + cache: &mut discover::Cache, + limits: discover::Limits, + ) -> Result { + let kept = crate::fs::read(&self.listed) + .ok() + .and_then(|bytes| serde_json::from_slice(&bytes).ok()) + .unwrap_or_default(); + let mut source = Source { + guest: &self.guest, + kept, + listed: Listings::new(), + }; + let folder = cache.discover(&mut source, limits)?; + if let Ok(json) = serde_json::to_vec(&source.listed) { + let _ = crate::fs::write(&self.listed, json); + } + Ok(folder) + } + + fn location(&self) -> String { + self.guest.location() + } + + fn entries(&self, folder: &str) -> io::Result> { + self.guest.entries(folder) + } + + fn exists(&self, path: &str) -> bool { + let request = Request { + path: path.to_owned(), + ..Request::default() + }; + self.guest + .ask(kind::EXISTS, request) + .is_ok_and(|reply| reply.exists == Some(true)) + } + + fn stamp(&self, path: &str) -> io::Result { + self.guest.stamp(path) + } + + fn read(&self, path: &str) -> Result> { + Ok(self.guest.read(kind::READ, path, LIMIT)?) + } + + fn read_file(&self, path: &str, limit: usize) -> Result> { + Ok(self.guest.read(kind::READ_FILE, path, limit)?) + } + + fn create(&self, path: &str, bytes: &[u8]) -> Result<()> { + let mut request = Request { + path: path.to_owned(), + ..Request::default() + }; + self.guest + .carry(&mut request, bytes.to_vec()) + .map_err(Failed::io)?; + self.guest.verb(kind::CREATE, path, request) + } + + fn create_directory(&self, path: &str) -> Result<()> { + self.guest + .verb(kind::CREATE_DIRECTORY, path, Request::default()) + } + + fn hide(&self, path: &str) -> Result<()> { + self.guest.verb(kind::HIDE, path, Request::default()) + } + + fn rename(&self, from: &str, to: &str) -> Result<()> { + let request = Request { + to: Some(to.to_owned()), + ..Request::default() + }; + self.guest.verb(kind::RENAME, from, request) + } + + fn rename_root(&self, _: &str, _: &[String]) -> Result { + Err(refused( + io::ErrorKind::Unsupported, + "Only the computer sharing this notebook can rename its folder", + )) + } + + fn replace(&self, from: &str, to: &str) -> Result<()> { + let request = Request { + to: Some(to.to_owned()), + ..Request::default() + }; + self.guest.verb(kind::REPLACE, from, request) + } + + fn delete(&self, path: &str) -> Result<()> { + self.guest.verb(kind::DELETE, path, Request::default()) + } + + fn place(&self, path: &str, ancestor: [u8; 16], name: &str) -> Result<()> { + let request = Request { + ancestor: Some(ancestor), + name: Some(name.to_owned()), + ..Request::default() + }; + self.guest.verb(kind::PLACE, path, request) + } + + fn commit(&self, path: &str, transaction: &Transaction) -> Result<()> { + Ok(self.guest.commit(path, transaction)?) + } + + fn confirm(&self, path: &str, base: &Stamp) -> std::result::Result<(), CommitError> { + self.guest.confirm(path, base) + } + + fn supersede(&self, path: &str, base: &Stamp, with: &str) -> Result<()> { + let request = Request { + to: Some(with.to_owned()), + stamp: Some(WireStamp::from(base)), + ..Request::default() + }; + self.guest.verb(kind::SUPERSEDE, path, request) + } +} diff --git a/crates/notebook/src/live/tests.rs b/crates/notebook/src/live/tests.rs index d25172625af521d330ee38f7e0f38c3af1be4c36..942f4ef1b9b70d9372a50c8ae073887f0fa3121e 100644 --- a/crates/notebook/src/live/tests.rs +++ b/crates/notebook/src/live/tests.rs @@ -43,9 +43,9 @@ fn caret(offset: u32) -> Presence { /// caret, and see the other leave. #[test] fn peers_meet_and_follow_presence() { - let room = Room::Code("7-violet-otter".into()); - let ada = Live::start(hello("Ada"), &room, None, None, || {}).unwrap(); - let grace = Live::start(hello("Grace"), &room, None, None, || {}).unwrap(); + let room = Room::join("7-violet-otter", ""); + let ada = Live::start(hello("Ada"), &room, None, None, |_| {}).unwrap(); + let grace = Live::start(hello("Grace"), &room, None, None, |_| {}).unwrap(); ada.set_presence(caret(1)); ada.connect(grace.address()); let seen = until(&grace, |peers| { @@ -70,18 +70,18 @@ fn peers_meet_and_follow_presence() { fn another_code_never_meets() { let ada = Live::start( hello("Ada"), - &Room::Code("7-violet-otter".into()), + &Room::join("7-violet-otter", ""), None, None, - || {}, + |_| {}, ) .unwrap(); let mallory = Live::start( hello("Mallory"), - &Room::Code("7-violet-ocelot".into()), + &Room::join("7-violet-ocelot", ""), None, None, - || {}, + |_| {}, ) .unwrap(); mallory.connect(ada.address()); @@ -115,8 +115,8 @@ fn later_fields_and_kinds_are_skipped() { } ); - let room = Room::Code("4-quiet-heron".into()); - let grace = Live::start(hello("Grace"), &room, None, None, || {}).unwrap(); + let room = Room::join("4-quiet-heron", ""); + let grace = Live::start(hello("Grace"), &room, None, None, |_| {}).unwrap(); // A later version: it greets, says something new, then where it is. let later = thread::spawn(move || { let mut stream = TcpStream::connect(grace.address()).unwrap(); @@ -146,8 +146,8 @@ fn later_fields_and_kinds_are_skipped() { #[ignore = "multicasts mDNS on the loopback interface"] fn peers_find_each_other_on_loopback() { let room = Room::Notebook([9; 16]); - let ada = Live::start(hello("Ada"), &room, Some(Reach::Loopback), None, || {}).unwrap(); - let grace = Live::start(hello("Grace"), &room, Some(Reach::Loopback), None, || {}).unwrap(); + let ada = Live::start(hello("Ada"), &room, Some(Reach::Loopback), None, |_| {}).unwrap(); + let grace = Live::start(hello("Grace"), &room, Some(Reach::Loopback), None, |_| {}).unwrap(); until(&ada, |peers| peers.len() == 1); until(&grace, |peers| peers.len() == 1); } @@ -272,9 +272,9 @@ fn status(address: SocketAddr, path: &str) -> String { fn peers_meet_through_a_relay() { let (url, _) = relay(Default::default()); let room = Room::Notebook([3; 16]); - let ada = Live::start(hello("Ada"), &room, None, Some(&url), || {}).unwrap(); + let ada = Live::start(hello("Ada"), &room, None, Some(&url), |_| {}).unwrap(); ada.set_presence(caret(1)); - let grace = Live::start(hello("Grace"), &room, None, Some(&url), || {}).unwrap(); + let grace = Live::start(hello("Grace"), &room, None, Some(&url), |_| {}).unwrap(); until(&grace, |peers| { peers.len() == 1 && peers[0].presence == Some(caret(1)) }); @@ -298,10 +298,10 @@ fn a_relay_numbers_a_code_and_burns_it_after_wrong_tries() { }); let host = Live::start( hello("Ada"), - &Room::Code("violet-otter".into()), + &Room::share("violet-otter", ""), None, Some(&url), - || {}, + |_| {}, ) .unwrap(); let deadline = Instant::now() + Duration::from_secs(10); @@ -316,23 +316,30 @@ fn a_relay_numbers_a_code_and_burns_it_after_wrong_tries() { assert_eq!(words, "violet-otter"); let guest = Live::start( hello("Grace"), - &Room::Code(code.clone()), + &Room::join(&code, ""), None, Some(&url), - || {}, + |_| {}, ) .unwrap(); until(&host, |peers| peers.len() == 1); until(&guest, |peers| peers.len() == 1); let path = format!("/v1/room/code-{number}"); - let wrong = Room::Code(format!("{number}-violet-ocelot")); - let mallory = Live::start(hello("Mallory"), &wrong, None, Some(&url), || {}).unwrap(); - // Mallory tries again a second later, and the second wrong try burns the code. Asking - // sooner would count as a try itself. - thread::sleep(Duration::from_secs(5)); + let wrong = Room::join(&format!("{number}-violet-ocelot"), ""); + // An end that typed a wrong code gives up at once; Mallory tries twice, and the second + // wrong try burns the code. + for _ in 0..2 { + let mallory = Live::start(hello("Mallory"), &wrong, None, Some(&url), |_| {}).unwrap(); + let deadline = Instant::now() + Duration::from_secs(10); + while mallory.failed() == 0 { + assert!(Instant::now() < deadline, "Mallory never tried"); + thread::sleep(Duration::from_millis(20)); + } + assert!(mallory.peers().is_empty()); + } + until(&host, |_| host.burned()); assert_eq!(status(address, &path), "HTTP/1.1 410 Gone"); - assert!(mallory.peers().is_empty()); assert_eq!(host.peers().len(), 1, "Grace stays"); } @@ -441,7 +448,10 @@ fn recording(name: &str, room: &Room, relay: &str) -> (Live, Arc>> = Arc::default(); let (recorded, watched) = (Arc::clone(&heard), Arc::clone(&shared)); - let live = Live::start(hello(name), room, None, Some(relay), move || { + let live = Live::start(hello(name), room, None, Some(relay), move |event| { + if !matches!(event, Event::Changed) { + return; + } let Some(shared) = watched.get().and_then(std::sync::Weak::upgrade) else { return; }; @@ -470,7 +480,7 @@ fn a_malicious_relay_is_caught() { Tamper::Inject, ] { let (url, address) = relay(Default::default()); - let ada = Live::start(hello("Ada"), &room, None, Some(&url), || {}).unwrap(); + 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); diff --git a/crates/notebook/src/live/wire.rs b/crates/notebook/src/live/wire.rs index bfdb8d4737da092c29c6ad106d77e9dff7c8e09b..c4972b49b7bdd552d398fe43f40b405174411236 100644 --- a/crates/notebook/src/live/wire.rs +++ b/crates/notebook/src/live/wire.rs @@ -20,11 +20,60 @@ pub mod kind { pub const HELLO: u16 = 1; /// Keeps a quiet connection open; it says nothing else. pub const PING: u16 = 2; + pub const BYE: u16 = 3; pub const PRESENCE: u16 = 16; + /// A host's files changed: `Touched`. + pub const TOUCHED: u16 = 18; + pub const WELCOME: u16 = 32; + /// Storage requests to a host, each a `Request` answered by a `Reply`. + pub const LIST: u16 = 257; + pub const STAMP: u16 = 258; + /// A section's consistent image, a chunk at a time. + pub const READ: u16 = 259; + pub const COMMIT: u16 = 260; + pub const CONFIRM: u16 = 261; + pub const CREATE: u16 = 262; + pub const CREATE_DIRECTORY: u16 = 263; + pub const HIDE: u16 = 264; + pub const RENAME: u16 = 265; + pub const REPLACE: u16 = 266; + pub const DELETE: u16 = 267; + pub const PLACE: u16 = 268; + pub const SUPERSEDE: u16 = 269; + /// Any other file as it stands, a chunk at a time. + pub const READ_FILE: u16 = 270; + pub const EXISTS: u16 = 271; + /// A chunk of the bytes a later request carries. + pub const PUT: u16 = 272; + pub const REPLY: u16 = 511; } /// The kinds this version reads, as `Hello::kinds` lists them. -pub const KNOWN: &[u16] = &[kind::HELLO, kind::PING, kind::PRESENCE]; +pub const KNOWN: &[u16] = &[ + kind::HELLO, + kind::PING, + kind::BYE, + kind::PRESENCE, + kind::TOUCHED, + kind::WELCOME, + kind::LIST, + kind::STAMP, + kind::READ, + kind::COMMIT, + kind::CONFIRM, + kind::CREATE, + kind::CREATE_DIRECTORY, + kind::HIDE, + kind::RENAME, + kind::REPLACE, + kind::DELETE, + kind::PLACE, + kind::SUPERSEDE, + kind::READ_FILE, + kind::EXISTS, + kind::PUT, + kind::REPLY, +]; /// The first message each way, and the only one with a name and a picture. #[derive(Clone, Debug, PartialEq, Encode, Decode)] @@ -44,6 +93,9 @@ pub struct Hello { /// The message kinds the sender reads; a request waits for its kind to be listed. #[n(4)] pub kinds: Vec, + /// The share this peer hosts, whose storage requests it answers. + #[cbor(n(5), with = "minicbor::bytes")] + pub serves: Option<[u8; 16]>, } impl Hello { @@ -57,6 +109,7 @@ impl Hello { picture, app: format!("Snowbound {}", env!("CARGO_PKG_VERSION")), kinds: KNOWN.to_vec(), + serves: None, }) } } @@ -123,6 +176,194 @@ impl From for onestore::ExGuid { } } +/// Why a peer leaves: `left`, or `stopped` for a host that stopped sharing. +#[derive(Clone, Debug, PartialEq, Encode, Decode)] +#[cbor(map)] +pub struct Bye { + #[n(0)] + pub reason: String, +} + +/// What the host of a share gives a peer that knew its code: the share's room, and names to +/// show it by. +#[derive(Clone, Debug, PartialEq, Encode, Decode)] +#[cbor(map)] +pub struct Welcome { + #[cbor(n(0), with = "minicbor::bytes")] + pub share: [u8; 16], + #[cbor(n(1), with = "minicbor::bytes")] + pub secret: [u8; 16], + #[n(2)] + pub notebook: String, + /// The host's name for itself, as `Hello::name`. + #[n(3)] + pub host: String, +} + +/// Paths a host's files changed at, by catalog path; `""` is the notebook's folder. +#[derive(Clone, Debug, Default, PartialEq, Encode, Decode)] +#[cbor(map)] +pub struct Touched { + #[n(0)] + pub paths: Vec, +} + +/// A storage request; its kind names the verb, and the verb what it carries. +#[derive(Clone, Debug, Default, PartialEq, Encode, Decode)] +#[cbor(map)] +pub struct Request { + /// Names the reply; a `PUT`'s names the bytes a later request carries. + #[n(0)] + pub id: u64, + /// A catalog path. + #[n(1)] + pub path: String, + /// A rename's or replacement's target, or the file a supersession puts in place. + #[n(2)] + pub to: Option, + #[n(3)] + pub offset: Option, + #[n(4)] + pub limit: Option, + #[cbor(n(5), with = "minicbor::bytes")] + pub bytes: Option>, + /// A confirmation's or supersession's base. + #[n(6)] + pub stamp: Option, + #[cbor(n(7), with = "minicbor::bytes")] + pub ancestor: Option<[u8; 16]>, + #[n(8)] + pub name: Option, + /// The `PUT`s, by id, whose bytes this request carries in place of `bytes`; a read's + /// snapshot after its first chunk. + #[n(9)] + pub handle: Option, +} + +/// A request's answer: a failure, or what the verb gives. +#[derive(Clone, Debug, Default, PartialEq, Encode, Decode)] +#[cbor(map)] +pub struct Reply { + #[n(0)] + pub id: u64, + #[n(1)] + pub failure: Option, + #[cbor(n(2), with = "minicbor::bytes")] + pub bytes: Option>, + /// A read's whole length. + #[n(3)] + pub length: Option, + #[n(4)] + pub stamp: Option, + #[n(5)] + pub entries: Option>, + #[n(6)] + pub exists: Option, + /// The snapshot a read's later chunks come from. + #[n(7)] + pub handle: Option, +} + +/// Why a request failed: an `io::ErrorKind` as `error_kind` numbers it, and for a commit, +/// how far it got (`commit_state`). +#[derive(Clone, Debug, PartialEq, Encode, Decode)] +#[cbor(map)] +pub struct Failure { + #[n(0)] + pub kind: u16, + #[n(1)] + pub message: String, + #[n(2)] + pub state: Option, +} + +/// An `onestore::Stamp`. +#[derive(Clone, Debug, PartialEq, Encode, Decode)] +#[cbor(map)] +pub struct WireStamp { + #[cbor(n(0), with = "minicbor::bytes")] + pub header: Vec, + #[n(1)] + pub length: u64, +} + +impl From<&onestore::Stamp> for WireStamp { + fn from(stamp: &onestore::Stamp) -> Self { + Self { + header: stamp.header.to_vec(), + length: stamp.length, + } + } +} + +impl TryFrom<&WireStamp> for onestore::Stamp { + type Error = io::Error; + + fn try_from(stamp: &WireStamp) -> io::Result { + Ok(Self { + header: stamp + .header + .as_slice() + .try_into() + .map_err(|_| invalid("A stamp's header is 1024 bytes"))?, + length: stamp.length, + }) + } +} + +/// A folder's entry, as `discover::Entry`. +#[derive(Clone, Debug, PartialEq, Encode, Decode)] +#[cbor(map)] +pub struct WireEntry { + #[n(0)] + pub name: String, + /// 0 a file, 1 a folder, 2 anything else, 3 a file kept elsewhere. + #[n(1)] + pub kind: u8, + #[n(2)] + pub size: u64, + #[n(3)] + pub modified: u64, +} + +/// The `io::ErrorKind`s a failure names, by number; any other is `Other`. +const ERROR_KINDS: [io::ErrorKind; 20] = [ + io::ErrorKind::Other, + io::ErrorKind::NotFound, + io::ErrorKind::PermissionDenied, + io::ErrorKind::AlreadyExists, + io::ErrorKind::InvalidInput, + io::ErrorKind::InvalidData, + io::ErrorKind::TimedOut, + io::ErrorKind::WouldBlock, + io::ErrorKind::ResourceBusy, + io::ErrorKind::Unsupported, + io::ErrorKind::FileTooLarge, + io::ErrorKind::NotConnected, + io::ErrorKind::ReadOnlyFilesystem, + io::ErrorKind::DirectoryNotEmpty, + io::ErrorKind::NotADirectory, + io::ErrorKind::IsADirectory, + io::ErrorKind::StorageFull, + io::ErrorKind::UnexpectedEof, + io::ErrorKind::Interrupted, + io::ErrorKind::BrokenPipe, +]; + +pub fn error_number(kind: io::ErrorKind) -> u16 { + ERROR_KINDS + .iter() + .position(|known| *known == kind) + .unwrap_or(0) as u16 +} + +pub fn error_kind(number: u16) -> io::ErrorKind { + ERROR_KINDS + .get(usize::from(number)) + .copied() + .unwrap_or(io::ErrorKind::Other) +} + /// The opening, sent in the clear by the side that connected and answered by the other. #[derive(Encode, Decode)] #[cbor(map)] @@ -161,8 +402,13 @@ impl Sealer { kind: u16, body: &impl Encode<()>, ) -> io::Result<()> { - let mut clear = kind.to_be_bytes().to_vec(); - minicbor::encode(body, &mut clear).map_err(io::Error::other)?; + let body = minicbor::to_vec(body).map_err(io::Error::other)?; + self.send_encoded(to, kind, &body) + } + + /// `send` for a body already encoded. + pub fn send_encoded(&mut self, to: &mut impl Write, kind: u16, body: &[u8]) -> io::Result<()> { + let clear = [&kind.to_be_bytes()[..], body].concat(); let number = self.count.to_be_bytes(); let sealed = self .cipher diff --git a/crates/notebook/src/live/words.txt b/crates/notebook/src/live/words.txt new file mode 100644 index 0000000000000000000000000000000000000000..b5d21df1340f8c8fc2aed10fa6c3ba2dbec5f1ce --- /dev/null +++ b/crates/notebook/src/live/words.txt @@ -0,0 +1,1296 @@ +acid +acorn +acre +acts +afar +affix +aged +agent +agile +aging +agony +ahead +aide +aids +aim +ajar +alarm +alias +alibi +alien +alike +alive +aloe +aloft +aloha +alone +amend +amino +ample +amuse +angel +anger +angle +ankle +apple +april +apron +aqua +area +arena +argue +arise +armed +armor +army +aroma +array +arson +art +ashen +ashes +atlas +atom +attic +audio +avert +avoid +awake +award +awoke +axis +bacon +badge +bagel +baggy +baked +baker +balmy +banjo +barge +barn +bash +basil +bask +batch +bath +baton +bats +blade +blank +blast +blaze +bleak +blend +bless +blimp +blink +bloat +blob +blog +blot +blunt +blurt +blush +boast +boat +body +boil +bok +bolt +boned +boney +bonus +bony +book +booth +boots +boss +botch +both +boxer +breed +bribe +brick +bride +brim +bring +brink +brisk +broad +broil +broke +brook +broom +brush +buck +bud +buggy +bulge +bulk +bully +bunch +bunny +bunt +bush +bust +busy +buzz +cable +cache +cadet +cage +cake +calm +cameo +canal +candy +cane +canon +cape +card +cargo +carol +carry +carve +case +cash +cause +cedar +chain +chair +chant +chaos +charm +chase +cheek +cheer +chef +chess +chest +chew +chief +chili +chill +chip +chomp +chop +chow +chuck +chump +chunk +churn +chute +cider +cinch +city +civic +civil +clad +claim +clamp +clap +clash +clasp +class +claw +clay +clean +clear +cleat +cleft +clerk +click +cling +clink +clip +cloak +clock +clone +cloth +cloud +clump +coach +coast +coat +cod +coil +coke +cola +cold +colt +coma +come +comic +comma +cone +cope +copy +coral +cork +cost +cot +couch +cough +cover +cozy +craft +cramp +crane +crank +crate +crave +crawl +crazy +creme +crepe +crept +crib +cried +crisp +crook +crop +cross +crowd +crown +crumb +crush +crust +cub +cult +cupid +cure +curl +curry +curse +curve +curvy +cushy +cut +cycle +dab +dad +daily +dairy +daisy +dance +dandy +darn +dart +dash +data +date +dawn +deaf +deal +dean +debit +debt +debug +decaf +decal +decay +deck +decor +decoy +deed +delay +denim +dense +dent +depth +derby +desk +dial +diary +dice +dig +dill +dime +dimly +diner +dingy +disco +dish +disk +ditch +ditzy +dizzy +dock +dodge +doing +doll +dome +donor +donut +dose +dot +dove +down +dowry +doze +drab +drama +drank +draw +dress +dried +drift +drill +drive +drone +droop +drove +drown +drum +dry +duck +duct +dude +dug +duke +duo +dusk +dust +duty +dwarf +dwell +eagle +early +earth +easel +east +eaten +eats +ebay +ebony +ebook +echo +edge +eel +eject +elbow +elder +elf +elk +elm +elope +elude +elves +email +emit +empty +emu +enter +entry +envoy +equal +erase +error +erupt +essay +etch +evade +even +evict +evil +evoke +exact +exit +fable +faced +fact +fade +fall +false +fancy +fang +fax +feast +feed +femur +fence +fend +ferry +fetal +fetch +fever +fiber +fifth +fifty +film +filth +final +finch +fit +five +flag +flaky +flame +flap +flask +fled +flick +fling +flint +flip +flirt +float +flock +flop +floss +flyer +foam +foe +fog +foil +folic +folk +food +fool +found +fox +foyer +frail +frame +fray +fresh +fried +frill +frisk +from +front +frost +froth +frown +froze +fruit +gag +gains +gala +game +gap +gas +gave +gear +gecko +geek +gem +genre +gift +gig +gills +given +giver +glad +glass +glide +gloss +glove +glow +glue +goal +going +golf +gong +good +gooey +goofy +gore +gown +grab +grain +grant +grape +graph +grasp +grass +grave +gravy +gray +green +greet +grew +grid +grief +grill +grip +grit +groom +grope +growl +grub +grunt +guide +gulf +gulp +gummy +guru +gush +gut +guy +habit +half +halo +halt +happy +harm +hash +hasty +hatch +hate +haven +hazel +hazy +heap +heat +heave +hedge +hefty +help +herbs +hers +hub +hug +hula +hull +human +humid +hump +hung +hunk +hunt +hurry +hurt +hush +hut +ice +icing +icon +icy +igloo +image +ion +iron +islam +issue +item +ivory +ivy +jab +jam +jaws +jazz +jeep +jelly +jet +jiffy +job +jog +jolly +jolt +jot +joy +judge +juice +juicy +july +jumbo +jump +junky +juror +jury +keep +keg +kept +kick +kilt +king +kite +kitty +kiwi +knee +knelt +koala +kung +ladle +lady +lair +lake +lance +land +lapel +large +lash +lasso +last +latch +late +lazy +left +legal +lemon +lend +lens +lent +level +lever +lid +life +lift +lilac +lily +limb +limes +line +lint +lion +lip +list +lived +liver +lunar +lunch +lung +lurch +lure +lurk +lying +lyric +mace +maker +malt +mama +mango +manor +many +map +march +mardi +marry +mash +match +mate +math +moan +mocha +moist +mold +mom +moody +mop +morse +most +motor +motto +mount +mouse +mousy +mouth +move +movie +mower +mud +mug +mulch +mule +mull +mumbo +mummy +mural +muse +music +musky +mute +nacho +nag +nail +name +nanny +nap +navy +near +neat +neon +nerd +nest +net +next +niece +ninth +nutty +oak +oasis +oat +ocean +oil +old +olive +omen +onion +only +ooze +opal +open +opera +opt +otter +ouch +ounce +outer +oval +oven +owl +ozone +pace +pagan +pager +palm +panda +panic +pants +panty +paper +park +party +pasta +patch +path +patio +payer +pecan +penny +pep +perch +perky +perm +pest +petal +petri +petty +photo +plank +plant +plaza +plead +plot +plow +pluck +plug +plus +poach +pod +poem +poet +pogo +point +poise +poker +polar +polio +polka +polo +pond +pony +poppy +pork +poser +pouch +pound +pout +power +prank +press +print +prior +prism +prize +probe +prong +proof +props +prude +prune +pry +pug +pull +pulp +pulse +puma +punch +punk +pupil +puppy +purr +purse +push +putt +quack +quake +query +quiet +quill +quilt +quit +quota +quote +rabid +race +rack +radar +radio +raft +rage +raid +rail +rake +rally +ramp +ranch +range +rank +rant +rash +raven +reach +react +ream +rebel +recap +relax +relay +relic +remix +repay +repel +reply +rerun +reset +rhyme +rice +rich +ride +rigid +rigor +rinse +riot +ripen +rise +risk +ritzy +rival +river +roast +robe +robin +rock +rogue +roman +romp +rope +rover +royal +ruby +rug +ruin +rule +runny +rush +rust +rut +sadly +sage +said +saint +salad +salon +salsa +salt +same +sandy +santa +satin +sauna +saved +savor +sax +say +scale +scam +scan +scare +scarf +scary +scoff +scold +scoop +scoot +scope +score +scorn +scout +scowl +scrap +scrub +scuba +scuff +sect +sedan +self +send +sepia +serve +set +seven +shack +shade +shady +shaft +shaky +sham +shape +share +sharp +shed +sheep +sheet +shelf +shell +shine +shiny +ship +shirt +shock +shop +shore +shout +shove +shown +showy +shred +shrug +shun +shush +shut +shy +sift +silk +silly +silo +sip +siren +sixth +size +skate +skew +skid +skier +skies +skip +skirt +skit +sky +slab +slack +slain +slam +slang +slash +slate +slaw +sled +sleek +sleep +sleet +slept +slice +slick +slimy +sling +slip +slit +slob +slot +slug +slum +slurp +slush +small +smash +smell +smile +smirk +smog +snack +snap +snare +snarl +sneak +sneer +sniff +snore +snort +snout +snowy +snub +snuff +speak +speed +spend +spent +spew +spied +spill +spiny +spoil +spoke +spoof +spool +spoon +sport +spot +spout +spray +spree +spur +squad +squat +squid +stack +staff +stage +stain +stall +stamp +stand +stank +stark +start +stash +state +stays +steam +steep +stem +step +stew +stick +sting +stir +stock +stole +stomp +stony +stood +stool +stoop +stop +storm +stout +stove +straw +stray +strut +stuck +stud +stuff +stump +stung +stunt +suds +sugar +sulk +surf +sushi +swab +swan +swarm +sway +swear +sweat +sweep +swell +swept +swim +swing +swipe +swirl +swoop +swore +syrup +tacky +taco +tag +take +tall +talon +tamer +tank +taper +taps +tarot +tart +task +taste +tasty +taunt +thank +thaw +theft +theme +thigh +thing +think +thong +thorn +those +throb +thud +thumb +thump +thus +tiara +tidal +tidy +tiger +tile +tilt +tint +tiny +trace +track +trade +train +trait +trap +trash +tray +treat +tree +trek +trend +trial +tribe +trick +trio +trout +truce +truck +trump +trunk +try +tug +tulip +tummy +turf +tusk +tutor +tutu +tux +tweak +tweet +twice +twine +twins +twirl +twist +uncle +uncut +undo +unify +union +unit +untie +upon +upper +urban +used +user +usher +utter +value +vapor +vegan +venue +verse +vest +veto +vice +video +view +viral +virus +visa +visor +vixen +vocal +voice +void +volt +voter +vowel +wad +wafer +wager +wages +wagon +wake +walk +wand +wasp +watch +water +wavy +wheat +whiff +whole +whoop +wick +widen +widow +width +wife +wifi +wilt +wimp +wind +wing +wink +wipe +wired +wiry +wise +wish +wispy +wok +wolf +womb +wool +woozy +word +work +worry +wound +woven +wrath +wreck +wrist +xerox +yahoo +yam +yard +year +yeast +yelp +yield +yarn +yodel +yoga +yoyo +yummy +zebra +zero +zesty +zippy +zone +zoom diff --git a/crates/notebook/src/session.rs b/crates/notebook/src/session.rs index 0bab77c9ff88a48f446e184b3488b0f4ca3f313f..f2fd32b69e4ae671a659f2a840109af9c5c0dfa1 100644 --- a/crates/notebook/src/session.rs +++ b/crates/notebook/src/session.rs @@ -1,7 +1,7 @@ //! The application's view of a notebook: sections opened through a local replica that //! publishes their edits to the section file in the background. -pub use crate::background::{Background, Known}; +pub use crate::background::{Background, Known, Listener}; use crate::{ EditStatus, Error, PendingEdit, Remote, Replica, Resolution, Result, SyncWorker, discover, fs, }; @@ -81,7 +81,11 @@ pub trait Storage: Send + Sync { ) -> Result; /// Where the notebook lives, which names its catalog's cache. fn location(&self) -> String; + /// A folder's entries, as discovery lists them. + fn entries(&self, folder: &str) -> io::Result>; fn exists(&self, path: &str) -> bool; + /// A section's or TOC's stamp, without reading its body or coordinating with writers. + fn stamp(&self, path: &str) -> io::Result; fn read(&self, path: &str) -> Result>; /// Reads a file of at most `limit` bytes as it stands, whatever it holds. fn read_file(&self, path: &str, limit: usize) -> Result>; @@ -103,6 +107,8 @@ pub trait Storage: Send + Sync { fn place(&self, path: &str, ancestor: [u8; 16], name: &str) -> Result<()>; /// Publishes a transaction made on the file's current image. fn commit(&self, path: &str, transaction: &Transaction) -> Result<()>; + /// Confirms that the file still has `base`'s stamp and is durable (`onestore::confirm`). + fn confirm(&self, path: &str, base: &Stamp) -> std::result::Result<(), CommitError>; /// Puts the file `with` in the place of the section or TOC at `path` under the coordination /// its writers take, provided `path` still has `base`'s stamp. fn supersede(&self, path: &str, base: &Stamp, with: &str) -> Result<()>; @@ -134,10 +140,19 @@ impl Storage for Directory { self.0.to_string_lossy().into_owned() } + fn entries(&self, folder: &str) -> io::Result> { + use discover::Source; + discover::Local::open(&self.0)?.entries(folder, LIMITS.entries) + } + fn exists(&self, path: &str) -> bool { fs::metadata(self.path(path)).is_ok() } + fn stamp(&self, path: &str) -> io::Result { + FileRemote(self.path(path)).stamp() + } + fn read(&self, path: &str) -> Result> { Ok(fs::read_file(self.path(path))?) } @@ -248,6 +263,10 @@ impl Storage for Directory { Ok(fs::commit_file(transaction, self.path(path))?) } + fn confirm(&self, path: &str, base: &Stamp) -> std::result::Result<(), CommitError> { + fs::confirm_file(self.path(path), base) + } + fn supersede(&self, path: &str, base: &Stamp, with: &str) -> Result<()> { Ok(fs::supersede_file(self.path(path), base, self.path(with))?) } @@ -288,6 +307,15 @@ impl Storage for Share { self.client.location(&self.root) } + fn entries(&self, folder: &str) -> io::Result> { + use discover::Source; + discover::Smb::new(&self.client, &self.root)?.entries(folder, LIMITS.entries) + } + + fn stamp(&self, path: &str) -> io::Result { + self.client.stamp(&self.path(path)) + } + fn exists(&self, path: &str) -> bool { let (folder, name) = split(path); self.client @@ -362,6 +390,10 @@ impl Storage for Share { .commit_transaction(&self.path(path), transaction)?) } + fn confirm(&self, path: &str, base: &Stamp) -> std::result::Result<(), CommitError> { + self.client.confirm(&self.path(path), base) + } + fn supersede(&self, path: &str, base: &Stamp, with: &str) -> Result<()> { Ok(self .client @@ -440,6 +472,27 @@ impl Notebook { ) } + /// Opens the notebook a Live Share host serves to `guest`. Sections open through + /// `Section::resume_hosted` with the catalog's paths; while the host can't be reached, + /// the notebook opens as its folders were last listed. + #[cfg(feature = "live")] + pub fn open_hosted( + guest: Arc, + cache: impl AsRef, + ) -> Result { + let listed = listing(cache.as_ref(), &guest.location()).with_extension("entries.json"); + Self::with( + Box::new(crate::live::share::Hosted::new(guest, listed)), + None, + cache, + ) + } + + /// The storage the notebook's files are in, for a Live Share host to serve. + pub fn into_storage(self) -> Box { + self.storage + } + fn with( storage: Box, root: Option, @@ -509,6 +562,11 @@ impl Notebook { crate::sidecar::mappings(&*self.storage) } + /// The secret of the notebook's live presence room, made where it has none. + pub fn presence_room(&self) -> Result<[u8; 16]> { + crate::sidecar::room(&*self.storage) + } + /// The picture a mapping names, once its bytes match its name. pub fn tag_art_file(&self, art: &str) -> Result> { crate::sidecar::art(&*self.storage, art) @@ -2023,6 +2081,22 @@ impl Section { ) } + /// Resumes a replica against a section a Live Share host serves at catalog `path`. + #[cfg(feature = "live")] + pub fn resume_hosted( + path: String, + replica: Replica, + guest: Arc, + notify: impl Fn() + Send + 'static, + ) -> Result { + Self::start( + PathBuf::from(&path), + replica, + move || Ok(crate::live::share::HostedRemote::new(&guest, &path)), + notify, + ) + } + fn start( file: PathBuf, replica: Replica, diff --git a/crates/notebook/src/session/tests.rs b/crates/notebook/src/session/tests.rs index af7fb86ffd3ec07d7eeb8547c8ca13929bec8769..06490e7945375ff2795823c74af7df3559751e40 100644 --- a/crates/notebook/src/session/tests.rs +++ b/crates/notebook/src/session/tests.rs @@ -28,9 +28,15 @@ impl Storage for Racing { fn location(&self) -> String { self.inner.location() } + fn entries(&self, folder: &str) -> io::Result> { + self.inner.entries(folder) + } fn exists(&self, path: &str) -> bool { self.inner.exists(path) } + fn stamp(&self, path: &str) -> io::Result { + self.inner.stamp(path) + } fn read(&self, path: &str) -> Result> { self.inner.read(path) } @@ -64,6 +70,9 @@ impl Storage for Racing { fn commit(&self, path: &str, transaction: &Transaction) -> Result<()> { self.inner.commit(path, transaction) } + fn confirm(&self, path: &str, base: &Stamp) -> std::result::Result<(), CommitError> { + self.inner.confirm(path, base) + } fn supersede(&self, path: &str, base: &Stamp, with: &str) -> Result<()> { self.race(path); self.inner.supersede(path, base, with) diff --git a/crates/notebook/src/sidecar.rs b/crates/notebook/src/sidecar.rs index c9a43c719d12b55373eeaf33e31daaefa2c49890..4bc4b591e9326de8960f4122a1499f0421f6916f 100644 --- a/crates/notebook/src/sidecar.rs +++ b/crates/notebook/src/sidecar.rs @@ -17,6 +17,9 @@ use std::io; const FOLDER: &str = ".snowbound"; const MAPPING: &str = ".snowbound/tags.json"; +/// Live presence's room secret: whoever reads the notebook's files may see who else has it +/// open. Live Share never serves it to guests, who meet in the share's own room. +const ROOM: &str = ".snowbound/live.json"; const ART: &str = ".snowbound/tags"; /// The most bytes read of a mapping or a picture. pub const LIMIT: usize = 1 << 20; @@ -155,6 +158,55 @@ pub(crate) fn map( Err(io::Error::from(io::ErrorKind::ResourceBusy).into()) } +/// The secret of the notebook's presence room (`live::Room::Notebook`), made where it has none. +/// The first writer's stays: a writer that finds one made meanwhile takes it. +pub(crate) fn room(storage: &dyn Storage) -> Result<[u8; 16]> { + #[derive(Serialize, Deserialize)] + struct Room { + room: String, + } + let read = || -> Result> { + let bytes = match storage.read_file(ROOM, LIMIT) { + Err(crate::Error::Io(error)) if error.kind() == io::ErrorKind::NotFound => { + return Ok(None); + } + bytes => bytes?, + }; + let room: Room = serde_json::from_slice(&bytes).map_err(io::Error::from)?; + let secret: Option> = (0..room.room.len()) + .step_by(2) + .map(|at| u8::from_str_radix(room.room.get(at..at + 2)?, 16).ok()) + .collect(); + Ok(Some( + secret + .and_then(|secret| secret.try_into().ok()) + .ok_or(io::Error::from(io::ErrorKind::InvalidData))?, + )) + }; + if let Some(secret) = read()? { + return Ok(secret); + } + match storage.create_directory(FOLDER) { + Err(crate::Error::Io(error)) if error.kind() == io::ErrorKind::AlreadyExists => {} + created => created?, + } + storage.hide(FOLDER)?; + let mut secret = [0; 16]; + getrandom::fill(&mut secret).map_err(|_| io::Error::other("System random source failed"))?; + let room = Room { + room: secret.iter().map(|byte| format!("{byte:02x}")).collect(), + }; + let written = temporary(ROOM); + storage.create( + &written, + &serde_json::to_vec(&room).map_err(io::Error::from)?, + )?; + if storage.rename(&written, ROOM).is_err() { + storage.delete(&written)?; + } + read()?.ok_or_else(|| io::Error::from(io::ErrorKind::NotFound).into()) +} + /// A name beside `path` no other writer picks. fn temporary(path: &str) -> String { use std::hash::{BuildHasher, RandomState}; diff --git a/crates/notebook/src/sidecar/tests.rs b/crates/notebook/src/sidecar/tests.rs index 0baf9b10e578b1533a87aea6e8361521f62814f7..178d6311c4a95b7f7178e276f98e1ad45b2335ef 100644 --- a/crates/notebook/src/sidecar/tests.rs +++ b/crates/notebook/src/sidecar/tests.rs @@ -83,6 +83,22 @@ impl Storage for Folder { "memory".into() } + fn entries(&self, _: &str) -> io::Result> { + unimplemented!() + } + + fn stamp(&self, _: &str) -> io::Result { + unimplemented!() + } + + fn confirm( + &self, + _: &str, + _: &onestore::Stamp, + ) -> std::result::Result<(), onestore::CommitError> { + unimplemented!() + } + fn exists(&self, path: &str) -> bool { self.files.lock().unwrap().contains_key(path) || self.folders.lock().unwrap().contains(path) } @@ -341,3 +357,24 @@ fn themes_merge_into_the_hidden_folder() { }; assert!(themes::write(&folder, built_in).is_err()); } + +/// The presence room's secret is made once, in the hidden folder, and every reader after +/// takes it; a writer that finds one made meanwhile takes that one. +#[test] +fn the_presence_room_is_made_once() { + let folder = Folder::default(); + let secret = room(&folder).unwrap(); + assert_eq!(room(&folder).unwrap(), secret); + assert!(folder.hidden.lock().unwrap().contains(FOLDER)); + let files: Vec<_> = folder.files.lock().unwrap().keys().cloned().collect(); + assert_eq!(files, [ROOM]); + let theirs = br#"{"room":"000102030405060708090a0b0c0d0e0f"}"#; + let other = Folder::default(); + other.folders.lock().unwrap().insert(FOLDER.into()); + other + .files + .lock() + .unwrap() + .insert(ROOM.into(), theirs.to_vec()); + assert_eq!(room(&other).unwrap(), std::array::from_fn(|at| at as u8)); +} diff --git a/crates/notebook/tests/live_share.rs b/crates/notebook/tests/live_share.rs new file mode 100644 index 0000000000000000000000000000000000000000..e24cf06409e17d365addb1dd8a5c7729817dbf23 --- /dev/null +++ b/crates/notebook/tests/live_share.rs @@ -0,0 +1,431 @@ +//! Live Share: guests open a notebook a host serves through a relay, edit it through the +//! replica and queue they use on a share, queue while the host is away, conflict as on a +//! share, pass protected sections through as ciphertext, and are refused with a wrong code. +#![cfg(feature = "live")] + +use notebook::{ + EditStatus, Replica, + live::{ + Hello, + share::{self, Guest, Host, Refusal, Sharing}, + }, + session::{Notebook, Section, SyncState}, +}; +use onestore::{ + Arena, ExGuid, + op::{Edit, Op, PageOp}, + protected::{Key, rekey}, +}; +use std::{ + net::TcpListener, + path::Path, + sync::Arc, + thread, + time::{Duration, Instant}, +}; + +#[path = "support/server.rs"] +mod server; + +const PASSWORD: &str = "fixture password"; + +/// A relay on this computer with `config`'s limits: its URL. +fn relay(config: relay::server::Config) -> String { + let listener = TcpListener::bind("127.0.0.1:0").unwrap(); + let url = format!("ws://{}", listener.local_addr().unwrap()); + thread::spawn(move || relay::server::serve(listener, config)); + url +} + +fn hello(name: &str) -> Hello { + Hello::new(name.into(), None).unwrap() +} + +fn until(what: &str, done: impl Fn() -> bool) { + let deadline = Instant::now() + Duration::from_secs(30); + while !done() { + assert!(Instant::now() < deadline, "{what}"); + thread::sleep(Duration::from_millis(20)); + } +} + +/// A notebook folder holding `Garden.one`, a page reading "Original text", and a protected +/// `Sealed.one` reading "Sealed text". +fn notebook(root: &Path) -> std::path::PathBuf { + let folder = root.join("Garden"); + std::fs::create_dir_all(&folder).unwrap(); + std::fs::write( + folder.join("Garden.one"), + onestore::create_section("Garden.one", "Original text", "Fixture").unwrap(), + ) + .unwrap(); + let plain = onestore::create_section("Sealed.one", "Sealed text", "Fixture").unwrap(); + let key = Key::new(PASSWORD).unwrap(); + std::fs::write( + folder.join("Sealed.one"), + rekey(&plain, None, Some(&key)).unwrap(), + ) + .unwrap(); + folder +} + +fn host(folder: &Path, cache: &Path, sharing: &Sharing, url: &str) -> Host { + let storage = Notebook::open(folder, cache).unwrap().into_storage(); + Host::start( + storage, + hello("Ada"), + sharing.clone(), + "Garden", + None, + Some(url), + || {}, + ) + .unwrap() +} + +/// The host's code once the relay has numbered it. +fn code(host: &Host) -> String { + until("the code was never numbered", || { + host.code().is_some_and(|code| share::code(&code).is_some()) + }); + host.code().unwrap() +} + +/// `name` joins with `code` and opens the notebook in `cache`. +fn guest(name: &str, code: &str, url: &str, cache: &Path) -> (Arc, Notebook) { + let welcome = share::join(hello(name), code, "", None, Some(url)).unwrap(); + assert_eq!( + (welcome.notebook.as_str(), welcome.host.as_str()), + ("Garden", "Ada") + ); + let guest = Guest::start( + hello(name), + welcome.share, + welcome.secret, + None, + Some(url), + || {}, + ) + .unwrap(); + until("the host was never met", || guest.host().is_some()); + let notebook = Notebook::open_hosted(Arc::clone(&guest), cache).unwrap(); + (guest, notebook) +} + +fn open(notebook: &Notebook, guest: &Arc, path: &str, key: Option<&Key>) -> Section { + let replica = notebook.replica_path(path).unwrap(); + std::fs::create_dir_all(replica.parent().unwrap()).unwrap(); + let replica = Replica::open_or_create(&replica, key, || notebook.read_section(path)).unwrap(); + Section::resume_hosted(path.into(), replica, Arc::clone(guest), || {}).unwrap() +} + +fn replace(section: &Section, image: &[u8], range: std::ops::Range, with: &str) -> u64 { + let (space, text, _) = server::text(image); + replaced(section, space, text, range, with) +} + +fn replaced( + section: &Section, + space: ExGuid, + text: ExGuid, + range: std::ops::Range, + with: &str, +) -> u64 { + let op = PageOp::Text { + text, + range, + with: with.into(), + }; + let edit = Edit { + at: 134_000_000_000_000_000, + ops: vec![Op::Page { space, op }], + }; + section.replica().apply("Guest", edit).unwrap() +} + +fn published(section: &Section, id: u64) { + until("the edit was never published", || { + matches!( + section.status(id).unwrap(), + Some(EditStatus::Published { .. }) + ) + }); +} + +/// A guest opens a section through the host, its edit lands in the host's file, and the +/// host's own edit reaches the guest. +#[test] +fn a_guest_edits_the_host_s_notebook() { + let directory = tempfile::tempdir().unwrap(); + let folder = notebook(directory.path()); + let url = relay(Default::default()); + let sharing = Sharing::new("").unwrap(); + let host = host(&folder, &directory.path().join("host"), &sharing, &url); + let (guest, notebook) = guest("Grace", &code(&host), &url, &directory.path().join("grace")); + let paths: Vec = notebook + .catalog() + .sections + .iter() + .map(|section| section.path.clone()) + .collect(); + assert_eq!(paths, ["Garden.one", "Sealed.one"]); + let section = open(¬ebook, &guest, "Garden.one", None); + let file = folder.join("Garden.one"); + let id = replace(§ion, &std::fs::read(&file).unwrap(), 0..8, "Grace's"); + published(§ion, id); + assert_eq!( + server::text(&std::fs::read(&file).unwrap()).2, + "Grace's text" + ); + assert_eq!(host.guests().len(), 1); + + // The host's own edit, as its app commits one. + let image = std::fs::read(&file).unwrap(); + let (space, text, _) = server::text(&image); + std::fs::write(&file, server::typed(&image, space, text, 0..7, "Ada's")).unwrap(); + host.touched(&["Garden.one".into()]); + until("the guest never saw the host's edit", || { + section + .page(space) + .is_ok_and(|page| server::page_texts(&page).contains(&"Ada's text".to_owned())) + }); +} + +/// While the host is away a guest's edits wait in its replica, the notebook opens from its +/// last listing, and once the host is back the edits reach its file. +#[test] +fn a_guest_queues_while_the_host_is_away() { + let directory = tempfile::tempdir().unwrap(); + let folder = notebook(directory.path()); + let url = relay(Default::default()); + let sharing = Sharing::new("").unwrap(); + let host_cache = directory.path().join("host"); + let host = host(&folder, &host_cache, &sharing, &url); + let cache = directory.path().join("grace"); + let (guest, notebook) = guest("Grace", &code(&host), &url, &cache); + let section = open(¬ebook, &guest, "Garden.one", None); + let file = folder.join("Garden.one"); + let image = std::fs::read(&file).unwrap(); + + // Ada's computer goes to sleep. + drop(host); + until("the host never left", || guest.host().is_none()); + let id = replace(§ion, &image, 0..8, "Offline"); + until("the section never said the host was away", || { + section + .sync_status() + .is_ok_and(|status| status.state() == SyncState::NotConnected && status.queued > 0) + }); + assert_eq!(std::fs::read(&file).unwrap(), image); + // The notebook opens from its last listing while the host is away. + let reopened = Notebook::open_hosted(Arc::clone(&guest), &cache).unwrap(); + assert_eq!(reopened.catalog().sections.len(), 2); + + let host = self::host(&folder, &host_cache, &sharing, &url); + published(§ion, id); + assert_eq!( + server::text(&std::fs::read(&file).unwrap()).2, + "Offline text" + ); + drop(host); +} + +/// Two guests that change the same words while apart: the second to publish gets OneNote's +/// conflict page, in its replica and in the host's file. +#[test] +fn two_guests_on_one_page_conflict_as_on_a_share() { + let directory = tempfile::tempdir().unwrap(); + let folder = notebook(directory.path()); + let url = relay(Default::default()); + let sharing = Sharing::new("").unwrap(); + let host_cache = directory.path().join("host"); + let host = host(&folder, &host_cache, &sharing, &url); + let code = code(&host); + let (grace, grace_notebook) = guest("Grace", &code, &url, &directory.path().join("grace")); + let (alan, alan_notebook) = guest("Alan", &code, &url, &directory.path().join("alan")); + let graces = open(&grace_notebook, &grace, "Garden.one", None); + let alans = open(&alan_notebook, &alan, "Garden.one", None); + let file = folder.join("Garden.one"); + let image = std::fs::read(&file).unwrap(); + + drop(host); + until("the host never left", || { + grace.host().is_none() && alan.host().is_none() + }); + let first = replace(&graces, &image, 0..8, "Grace's"); + let second = replace(&alans, &image, 0..8, "Alan's"); + let host = self::host(&folder, &host_cache, &sharing, &url); + published(&graces, first); + published(&alans, second); + let stored = std::fs::read(&file).unwrap(); + let conflicts = server::conflicts(&stored); + assert_eq!(conflicts.len(), 1, "one page holds a conflict page"); + let (user, kept) = &conflicts[0].1[0]; + assert_eq!(user, "Guest"); + let texts: Vec = server::pages(&stored) + .iter() + .flat_map(|(_, page)| server::page_texts(page)) + .chain(kept.iter().cloned()) + .collect(); + assert!( + texts.contains(&"Grace's text".to_owned()) && texts.contains(&"Alan's text".to_owned()), + "{texts:?}" + ); + assert!(!alans.conflicts().unwrap().is_empty() || !graces.conflicts().unwrap().is_empty()); + drop(host); +} + +/// A protected section's guest unlocks it with the password; the host only ever stores +/// what the guest sealed. +#[test] +fn a_protected_section_passes_through_as_ciphertext() { + let directory = tempfile::tempdir().unwrap(); + let folder = notebook(directory.path()); + let url = relay(Default::default()); + let sharing = Sharing::new("").unwrap(); + let host = host(&folder, &directory.path().join("host"), &sharing, &url); + let (guest, notebook) = guest("Grace", &code(&host), &url, &directory.path().join("grace")); + assert!(notebook.unlock("Sealed.one", "wrong password").is_err()); + let key = notebook.unlock("Sealed.one", PASSWORD).unwrap(); + let section = open(¬ebook, &guest, "Sealed.one", Some(&key)); + let file = folder.join("Sealed.one"); + let arena = Arena::default(); + let mut unlocked = + onestore::Section::unlock(&arena, std::fs::read(&file).unwrap(), &key).unwrap(); + let (space, ..) = unlocked.pages().unwrap()[0]; + let page = unlocked.page(space).unwrap(); + let text = page + .objects + .iter() + .find_map(|object| match object { + onestore::page::PageObject::Outline(outline) => outline + .paragraphs + .iter() + .find_map(|p| p.text().map(|t| t.id)), + _ => None, + }) + .unwrap(); + let id = replaced(§ion, space, text, 0..6, "Guarded"); + published(§ion, id); + let stored = std::fs::read(&file).unwrap(); + let marker = "Guarded" + .encode_utf16() + .flat_map(u16::to_le_bytes) + .collect::>(); + assert!( + !stored.windows(marker.len()).any(|window| window == marker), + "the host's file holds the new text in the clear" + ); + assert!(onestore::Section::open(&Arena::default(), stored.clone()).is_err()); + let arena = Arena::default(); + let reread = onestore::Section::unlock(&arena, stored, &key).unwrap(); + let texts = server::page_texts(&reread.page(space).unwrap()); + assert!(texts.contains(&"Guarded text".to_owned()), "{texts:?}"); +} + +/// A wrong code is refused without the guest learning anything; the host's code burns after +/// too many wrong tries and is replaced by new words; too many wrong codes from one network +/// lock it out; and a guest can't reach outside the notebook or presence's secret. +#[test] +fn wrong_codes_are_refused_and_counted() { + let directory = tempfile::tempdir().unwrap(); + let folder = notebook(directory.path()); + let url = relay(relay::server::Config { + burn_after: 2, + failures_per_minute: 3, + ..Default::default() + }); + let sharing = Sharing::new("").unwrap(); + let host = host(&folder, &directory.path().join("host"), &sharing, &url); + let code = code(&host); + let number = code.split('-').next().unwrap(); + let wrong = format!("{number}-violet-ocelot"); + let join = |code: &str| share::join(hello("Mallory"), code, "", None, Some(&url)); + assert_eq!(join("not a code").unwrap_err(), Refusal::Malformed); + assert_eq!(join(&wrong).unwrap_err(), Refusal::Wrong); + assert_eq!(join(&wrong).unwrap_err(), Refusal::Wrong); + // The code burned, and the host shares new words, under a number of their own. + until("the code never changed", || { + host.code() + .is_some_and(|now| now != code && share::code(&now).is_some()) + }); + let fresh = host.code().unwrap(); + assert!(matches!( + join(&code).unwrap_err(), + Refusal::Expired | Refusal::NoOne + )); + let number = fresh.split('-').next().unwrap(); + assert_eq!( + join(&format!("{number}-violet-ocelot")).unwrap_err(), + Refusal::Wrong + ); + // A third wrong code in a minute locks this network out, even from the right code. + assert!(matches!( + join(&fresh).unwrap_err(), + Refusal::TooMany(Some(_)) + )); + + // Joined, a guest still can't reach outside the notebook. + std::fs::create_dir_all(folder.join(".snowbound")).unwrap(); + std::fs::write(folder.join(".snowbound/live.json"), b"{}").unwrap(); + std::fs::write(directory.path().join("secret.txt"), b"secret").unwrap(); + let url = relay(Default::default()); + let host = self::host(&folder, &directory.path().join("host"), &sharing, &url); + let (_guest, notebook) = guest( + "Grace", + &self::code(&host), + &url, + &directory.path().join("g"), + ); + let storage = notebook.into_storage(); + for path in [ + "../secret.txt", + ".snowbound/live.json", + "/etc/hosts", + "a//b", + ] { + let error = storage.read_file(path, 1024).unwrap_err(); + assert!( + error.to_string().contains("Outside the notebook"), + "{path}: {error}" + ); + } + assert!(storage.rename_root("Elsewhere", &[]).is_err()); +} + +/// Files larger than one message go up and come back a chunk at a time, through a relay that +/// hangs up on a peer with more than its queue waiting. +#[test] +fn large_files_travel_in_chunks() { + let directory = tempfile::tempdir().unwrap(); + let folder = notebook(directory.path()); + let url = relay(relay::server::Config { + queue: 400 << 10, + ..Default::default() + }); + let sharing = Sharing::new("").unwrap(); + let host = host(&folder, &directory.path().join("host"), &sharing, &url); + let (_guest, notebook) = guest("Grace", &code(&host), &url, &directory.path().join("grace")); + let storage = notebook.into_storage(); + let bytes: Vec = (0..3_000_000u32).map(|at| (at * 7 % 251) as u8).collect(); + storage + .create("Garden_onefiles/big.bin", &bytes) + .unwrap_err(); + std::fs::create_dir(folder.join("Garden_onefiles")).unwrap(); + storage.create("Garden_onefiles/big.bin", &bytes).unwrap(); + assert_eq!( + std::fs::read(folder.join("Garden_onefiles/big.bin")).unwrap(), + bytes + ); + assert_eq!( + storage + .read_file("Garden_onefiles/big.bin", 4 << 20) + .unwrap(), + bytes + ); + assert!( + storage + .read_file("Garden_onefiles/big.bin", 1 << 20) + .is_err() + ); +} diff --git a/crates/onestore/src/commit.rs b/crates/onestore/src/commit.rs index a1c8544768ef2008adc12a79c6ade2487beb5d34..213e6f9d42e073ac8a6262577ca2d4bcf4a7c0af 100644 --- a/crates/onestore/src/commit.rs +++ b/crates/onestore/src/commit.rs @@ -497,6 +497,64 @@ impl Transaction { &self.base } + /// The transaction as bytes `from_bytes` reads back, to carry it to another machine to + /// commit: the base header and length, the new header, the appended bytes' length and + /// the bytes, then each patch's offset, length and bytes. Integers are little-endian. + pub fn to_bytes(&self) -> Vec { + let mut bytes = Vec::with_capacity(2064 + self.append.len()); + bytes.extend_from_slice(&self.base.header); + bytes.extend_from_slice(&self.base.length.to_le_bytes()); + bytes.extend_from_slice(&self.header); + bytes.extend_from_slice(&(self.append.len() as u64).to_le_bytes()); + bytes.extend_from_slice(&self.append); + for (offset, patch) in &self.patches { + bytes.extend_from_slice(&offset.to_le_bytes()); + bytes.extend_from_slice(&(patch.len() as u32).to_le_bytes()); + bytes.extend_from_slice(patch); + } + bytes + } + + /// Reads `to_bytes`'s form, refusing a patch outside the base's data area, which no + /// transaction writes. + pub fn from_bytes(bytes: &[u8]) -> Result { + let malformed = |message| crate::Error { offset: 0, message }; + fn take<'a>(rest: &mut &'a [u8], length: usize) -> Result<&'a [u8], crate::Error> { + let (taken, after) = rest.split_at_checked(length).ok_or(crate::Error { + offset: 0, + message: "A truncated transaction", + })?; + *rest = after; + Ok(taken) + } + let rest = &mut &bytes[..]; + let base_header: [u8; 1024] = take(rest, 1024)?.try_into().expect("1024 bytes"); + let length = u64::from_le_bytes(take(rest, 8)?.try_into().expect("8 bytes")); + let header: [u8; 1024] = take(rest, 1024)?.try_into().expect("1024 bytes"); + let appended = u64::from_le_bytes(take(rest, 8)?.try_into().expect("8 bytes")); + let appended = usize::try_from(appended).map_err(|_| malformed("A huge append"))?; + let append = take(rest, appended)?.to_vec(); + let mut patches = Vec::new(); + while !rest.is_empty() { + let offset = u64::from_le_bytes(take(rest, 8)?.try_into().expect("8 bytes")); + let size = u32::from_le_bytes(take(rest, 4)?.try_into().expect("4 bytes")); + let patch = take(rest, size as usize)?; + if offset < 1024 || offset.saturating_add(u64::from(size)) > length { + return Err(malformed("A patch outside the base's data")); + } + patches.push((offset, patch.to_vec())); + } + Ok(Self { + base: Stamp { + header: base_header, + length, + }, + append, + patches, + header, + }) + } + /// The bytes a commit writes, by offset, in the order `apply` writes them; the header, /// last, covers bytes 0..1024. pub fn writes(&self) -> impl Iterator { @@ -600,6 +658,32 @@ mod tests { } } + #[test] + fn transactions_read_back_from_bytes() { + let transaction = Transaction { + base: Stamp { + header: [1; 1024], + length: 4096, + }, + append: vec![2; 300], + patches: vec![(1024, vec![3; 8]), (4000, vec![4; 96])], + header: [5; 1024], + }; + let bytes = transaction.to_bytes(); + assert_eq!(Transaction::from_bytes(&bytes).unwrap(), transaction); + assert!(Transaction::from_bytes(&bytes[..bytes.len() - 1]).is_err()); + let outside = Transaction { + patches: vec![(4000, vec![4; 97])], + ..transaction.clone() + }; + assert!(Transaction::from_bytes(&outside.to_bytes()).is_err()); + let header = Transaction { + patches: vec![(1000, vec![4; 8])], + ..transaction + }; + assert!(Transaction::from_bytes(&header.to_bytes()).is_err()); + } + #[test] fn a_read_torn_by_a_commit_is_read_again_then_refused() { let section = crate::create_section("Torn.one", "Text", "Fixture").unwrap(); diff --git a/crates/snowbound/src/live.rs b/crates/snowbound/src/live.rs index 96ff14a75d91bbd24196eaf8b61e2bd8ea0382a9..d186b43f1a9aeab9e6c2b1f8c46d10093ea48b51 100644 --- a/crates/snowbound/src/live.rs +++ b/crates/snowbound/src/live.rs @@ -51,7 +51,7 @@ impl State { None => Some(DEFAULT_RELAY.to_owned()), }; let room = match std::env::var("SNOWBOUND_LIVE_CODE") { - Ok(code) => Some(Room::Code(code)), + Ok(code) => Some(Room::join(&code, "")), Err(_) => self .session .as_ref() @@ -67,7 +67,7 @@ impl State { .flatten(); let live = Hello::new(self.author.clone(), picture) .and_then(|me| { - live::Live::start(me, &room, reach, relay.as_deref(), move || { + live::Live::start(me, &room, reach, relay.as_deref(), move |_| { redraw.wake_by_ref() }) })