From e01c7cf1781d73fb57bc66e80bb1b3da9cafa9d2 Mon Sep 17 00:00:00 2001 From: clover caruso Date: Fri, 2 Oct 2026 18:24:10 -0700 Subject: [PATCH] feat: Live Share serves a notebook's storage to guests who meet its host by code A host serves its notebook's storage verbs (notebook::live::share::Host) to the peers in the share's room; a guest runs the replica, queue and merge it runs on an SMB share against them (Notebook::open_hosted, Section::resume_hosted, Background::hosted), so offline queueing, rebases and conflict pages work as there and the host's files are written only by its own storage. A guest meets the host in a short code's room (SPAKE2 on its words and any password, numbered by the relay or by the host off it) and is welcomed with the share's room and its random secret; a new share has a new secret. Guests' commits are checked on the host's image first, paths outside the notebook and presence's secret are refused, and large bodies travel in answered 128 KiB chunks so a relay never queues much for a slow peer. A notebook opens from its last listing while the host is away. Presence meets in a room whose random secret is kept in .snowbound/live.json (Notebook::presence_room) in place of the TOC's identity, which every copy of a notebook holds for good. Live events now carry frames, meetings and departures, Live::leave says Bye, and the joining end of a code gives up after one wrong try, as the relay counts each one. onestore::Transaction gains a byte form to carry a commit to another machine. Assisted-by: claude-opus-5.5 --- crates/notebook/README.md | 46 +- crates/notebook/src/background.rs | 44 +- crates/notebook/src/live.rs | 470 ++++++--- crates/notebook/src/live/relay.rs | 62 +- crates/notebook/src/live/share.rs | 1342 ++++++++++++++++++++++++++ crates/notebook/src/live/tests.rs | 60 +- crates/notebook/src/live/wire.rs | 252 ++++- crates/notebook/src/live/words.txt | 1296 +++++++++++++++++++++++++ crates/notebook/src/session.rs | 76 +- crates/notebook/src/session/tests.rs | 9 + crates/notebook/src/sidecar.rs | 52 + crates/notebook/src/sidecar/tests.rs | 37 + crates/notebook/tests/live_share.rs | 431 +++++++++ crates/onestore/src/commit.rs | 84 ++ crates/snowbound/src/live.rs | 4 +- 15 files changed, 4055 insertions(+), 210 deletions(-) create mode 100644 crates/notebook/src/live/share.rs create mode 100644 crates/notebook/src/live/words.txt create mode 100644 crates/notebook/tests/live_share.rs 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() }) }) -- 2.54.0