diff --git a/crates/notebook/examples/live_crowd.rs b/crates/notebook/examples/live_crowd.rs index 362ea4ba5852cadaf3c8fa66e05cd14bb53f062d..7bc4e6912386a8ad36979240e8992185093e1548 100644 --- a/crates/notebook/examples/live_crowd.rs +++ b/crates/notebook/examples/live_crowd.rs @@ -133,6 +133,7 @@ fn host(url: &str, paragraphs: usize) { None, Some(url), || {}, + |_| Ok(()), ) .unwrap(); until("no code", Duration::from_secs(30), || { diff --git a/crates/notebook/examples/live_latency.rs b/crates/notebook/examples/live_latency.rs index 74bb8304fc2715ea34fea9c83acc8a984196dfb1..942a1c94a3e0959994a342a00c4b904ee188d7a2 100644 --- a/crates/notebook/examples/live_latency.rs +++ b/crates/notebook/examples/live_latency.rs @@ -160,6 +160,7 @@ fn main() { reach, relay, || {}, + |_| Ok(()), ) .unwrap(); until("no code", || { diff --git a/crates/notebook/src/live/share.rs b/crates/notebook/src/live/share.rs index 1bc7b04bbd6a9281f062119a896f8325137329f9..a8f5e15f1b6b0ee78bc54d5778879a73fe8dff32 100644 --- a/crates/notebook/src/live/share.rs +++ b/crates/notebook/src/live/share.rs @@ -1,9 +1,9 @@ //! Live Share: a notebook one Snowbound holds, opened on others through a short code. The //! host serves its notebook's storage verbs (`session::Storage`) and batches guests' ops into //! guarded publications. A guest runs the same replica and durable queue as on an SMB share; -//! protected sections and older peers use transactions. 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 +//! protected sections use transactions. A guest first meets the host in the code's room, +//! where approval grants a device's own access credential. Presence has a separate room +//! whose key changes when a device is removed. Large bodies travel a chunk at a time, each answered before //! the next, so a relay never holds much for a slow peer. use super::{ @@ -20,7 +20,7 @@ use std::{ io, path::PathBuf, sync::{ - Arc, Condvar, Mutex, OnceLock, + Arc, Condvar, Mutex, atomic::{AtomicBool, AtomicU64, Ordering}, mpsc, }, @@ -29,6 +29,8 @@ use std::{ }; mod batch; +mod membership; +pub use membership::Host; /// The most bytes one message of a read or an upload carries. const CHUNK: usize = 128 << 10; @@ -77,11 +79,22 @@ pub fn code(typed: &str) -> Option { #[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. + /// The presence key, rotated when a device is removed. pub secret: [u8; 16], /// The code, or its secret alone until it has a number (`super::code`). pub code: String, pub password: String, + #[serde(default)] + pub approve: bool, + #[serde(default)] + pub members: Vec, +} + +#[derive(Clone, Debug, PartialEq, Eq, serde::Serialize, serde::Deserialize)] +pub struct Device { + pub secret: [u8; 16], + pub name: String, + pub device: Option, } impl Sharing { @@ -95,6 +108,8 @@ impl Sharing { secret: random[16..].try_into().expect("16 bytes"), code: super::code::secret()?, password: password.to_owned(), + approve: false, + members: Vec::new(), }) } } @@ -111,7 +126,9 @@ pub enum Refusal { Malformed, /// The person sharing runs a Snowbound of another Live Share version: the newer one's /// `true` where it is theirs, so this one should update. - Version { theirs_newer: bool }, + Version { + theirs_newer: bool, + }, /// The code's secret or password is wrong. Wrong, /// No one shares with the code's number now. @@ -126,6 +143,9 @@ pub enum Refusal { Unreachable(super::Trouble), /// The relay let this end in, but no one answered. TimedOut, + Declined, + Cancelled, + NotAdmitted, } /// Meets the host sharing `code` (and `password`) as `me`, on the networks `reach` names and @@ -136,9 +156,23 @@ pub fn join( password: &str, reach: Option, relay: Option<&str>, +) -> std::result::Result { + join_while(me, code, password, reach, relay, |_| true) +} + +/// Joins while `waiting` returns true, reporting whether the host is deciding approval. +pub fn join_while( + me: Hello, + code: &str, + password: &str, + reach: Option, + relay: Option<&str>, + continue_joining: impl Fn(bool) -> bool, ) -> std::result::Result { let code = self::code(code).ok_or(Refusal::Malformed)?; let (welcomed, welcome) = mpsc::channel(); + let approving = Arc::new(AtomicBool::new(false)); + let approval = Arc::clone(&approving); let (changed, waiting) = mpsc::channel(); let live = Live::start( me, @@ -152,7 +186,24 @@ pub fn join( .. } => { if let Ok(body) = minicbor::decode::(body) { - let _ = welcomed.send(body); + let _ = welcomed.send(Ok(body)); + } + } + Event::Frame { + kind: kind::APPROVAL, + body, + .. + } => { + if let Ok(body) = minicbor::decode::(body) { + match body { + wire::Approval::Pending => approval.store(true, Ordering::Release), + wire::Approval::Declined => { + let _ = welcomed.send(Err(Refusal::Declined)); + } + wire::Approval::Failed => { + let _ = welcomed.send(Err(Refusal::NotAdmitted)); + } + } } } _ => { @@ -163,8 +214,11 @@ pub fn join( .map_err(|_| Refusal::Unreachable(super::Trouble::Other))?; let start = Instant::now(); loop { + if !continue_joining(approving.load(Ordering::Acquire)) { + return Err(Refusal::Cancelled); + } if let Ok(welcome) = welcome.try_recv() { - return Ok(welcome); + return welcome; } if let Some(version) = live.other_version() { return Err(Refusal::Version { @@ -192,199 +246,21 @@ pub fn join( 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), + _ if waited + > Duration::from_secs(if approving.load(Ordering::Acquire) { + 300 + } else { + 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, - /// The share's room and the code's, until it stops. - room: Mutex>, - 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(), - guests: Mutex::default(), - writers: Mutex::default(), - host: Mutex::default(), - room: OnceLock::new(), - }); - 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.admit(hello.peer, line), - Event::Left(hello) => heard.forget(&hello.peer), - Event::Frame { - from, kind, body, .. - } if wire::KNOWN.contains(&kind) && kind > 256 && kind != kind::REPLY => { - heard.queue(&from.peer, kind, body); - } - Event::Changed => told(), - _ => {} - }, - )?; - let _ = served.room.set(room.sender()); - let host = Self { - notebook: notebook.to_owned(), - reach, - relay: relay.map(str::to_owned), - pairing: Mutex::new(Some(pair(&me, &sharing, notebook, reach, relay, &events)?)), - me, - sharing: Mutex::new(sharing), - room: Mutex::new(Some(room)), - served, - events, - }; - Ok(host) - } - - /// The code guests type, once it has its number, and none once stopped. A code with too - /// many wrong tries is replaced by one with a new secret. - pub fn code(&self) -> Option { - let mut pairing = self.pairing.lock().unwrap(); - let pairing = pairing.as_mut()?; - if pairing.burned() { - let mut sharing = self.sharing.lock().unwrap(); - sharing.code = super::code::secret().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 { - let pairing = self.pairing.lock().unwrap(); - pairing.as_ref().map_or(Relayed::Unknown, Live::relayed) - } - - /// The peers in the share's room. - pub fn guests(&self) -> Vec { - let room = self.room.lock().unwrap(); - room.as_ref().map(Live::peers).unwrap_or_default() - } - - pub fn set_presence(&self, presence: Presence) { - if let Some(room) = &*self.room.lock().unwrap() { - room.set_presence(presence); - } - } - - /// Has `listener` hear the catalog paths guests change from now on, sooner than a watch - /// on the notebook's folder would. - pub fn on_changed(&self, listener: crate::session::Listener) { - *self.served.host.lock().unwrap() = Some(listener); - } - - /// Tells every guest the files at these catalog paths changed, with what changed in the - /// sections a guest read lately. - pub fn touched(&self, paths: &[String]) { - for path in paths { - self.served.changed_here(path); - } - self.served.tell(paths); - } - - /// Stops sharing: no one new is welcomed, and every guest hears so and is let go. - pub fn stop(&self) { - drop(self.pairing.lock().unwrap().take()); - let room = self.room.lock().unwrap().take(); - if let Some(room) = room { - 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, @@ -399,7 +275,7 @@ struct Served { /// Hears the paths guests changed, as the host's own notebook should. host: Mutex>, /// The share's room, to tell guests what changed. - room: OnceLock, + room: Mutex>, } /// A guest as its host serves it: the line to it, its requests waiting for its workers, how @@ -515,7 +391,7 @@ impl Served { .map(|path| self.storage.stamp(path).ok().map(|stamp| digest(&stamp))) .collect(), }; - if let Some(room) = self.room.get() { + if let Some(room) = self.room.lock().unwrap().as_ref() { room.send(kind::TOUCHED, &touched, None); } } @@ -538,6 +414,15 @@ impl Served { } }; let id = request.id; + if !self.guests.lock().unwrap().contains_key(peer) { + return failed( + id, + refused( + io::ErrorKind::PermissionDenied, + "The device is no longer connected", + ), + ); + } match self.answer(peer, kind, request) { Ok(reply) => Reply { id, ..reply }, Err(error) => failed(id, error), @@ -864,7 +749,7 @@ impl Served { }) .collect(); drop(guests); - if let (Some(room), false) = (self.room.get(), holding.is_empty()) { + if let (Some(room), false) = (self.room.lock().unwrap().as_ref(), holding.is_empty()) { room.send(kind::DELTA, &delta, Some(&holding)); } } @@ -1070,6 +955,12 @@ pub struct Guest { struct Inner { share: [u8; 16], + me: Hello, + reach: Option, + relay: Option, + events: Arc, + presence: Mutex>, + here: Mutex, /// The host and the line to it, while connected. host: Mutex, Line)>>, pending: Mutex>>, @@ -1077,7 +968,7 @@ struct Inner { /// Where the host's reports of changed files go, while a background watches. watch: Mutex>, /// The host stopped sharing. - stopped: AtomicBool, + ended: Mutex>, /// Bytes of chunks asked for and not yet given. asked: (Mutex, Condvar), /// The sections read lately, kept as the host's deltas change them, newest first. @@ -1087,6 +978,11 @@ struct Inner { current: Mutex>, } +enum Ended { + Stopped, + Removed, +} + 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. @@ -1100,11 +996,17 @@ impl Guest { ) -> io::Result> { let inner = Arc::new(Inner { share, + me: me.clone(), + reach, + relay: relay.map(str::to_owned), + events: Arc::new(events), + presence: Mutex::default(), + here: Mutex::default(), host: Mutex::default(), pending: Mutex::default(), next: AtomicU64::new(1), watch: Mutex::default(), - stopped: AtomicBool::new(false), + ended: Mutex::default(), asked: Default::default(), images: Mutex::default(), current: Mutex::default(), @@ -1112,7 +1014,7 @@ impl Guest { let heard = Arc::clone(&inner); let live = Live::start(me, &Room::Notebook(secret), reach, relay, move |event| { heard.heard(event); - events(); + (heard.events)(); })?; Ok(Arc::new(Self { live, inner })) } @@ -1130,16 +1032,25 @@ impl Guest { /// Whether the host said it stopped sharing. pub fn stopped(&self) -> bool { - self.inner.stopped.load(Ordering::Acquire) + self.inner.ended.lock().unwrap().is_some() } /// Everyone in the share's room: the host and the other guests. pub fn peers(&self) -> Vec { - self.live.peers() + self.inner + .presence + .lock() + .unwrap() + .as_ref() + .map(|(_, live)| live.peers()) + .unwrap_or_default() } pub fn set_presence(&self, presence: Presence) { - self.live.set_presence(presence); + *self.inner.here.lock().unwrap() = presence.clone(); + if let Some((_, live)) = &*self.inner.presence.lock().unwrap() { + live.set_presence(presence); + } } /// Sends the host's reports of changed files to `reports` while it stays connected. @@ -1152,10 +1063,15 @@ impl Guest { } 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" + if let Some(version) = self.live.other_version() { + return io::Error::new(io::ErrorKind::NotConnected, wire::Version(version)); + } + let message = match &*self.inner.ended.lock().unwrap() { + Some(Ended::Stopped) => "The host stopped sharing this notebook", + Some(Ended::Removed) => { + "This device was removed. Ask the person sharing for a new link or code." + } + None => "The computer sharing this notebook can’t be reached", }; io::Error::new(io::ErrorKind::NotConnected, message) } @@ -1420,14 +1336,23 @@ impl Inner { } } - fn heard(&self, event: Event) { - let serves = |hello: &Hello| hello.serves == Some(self.share); + fn heard(self: &Arc, event: Event) { + let serves = |hello: &Hello| { + hello.serves == Some(self.share) + || self + .host + .lock() + .unwrap() + .as_ref() + .is_some_and(|(host, _)| host.peer == hello.peer) + }; 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; + drop(self.presence.lock().unwrap().take()); self.current.lock().unwrap().clear(); // Each request waiting hears its answer was lost. self.pending.lock().unwrap().clear(); @@ -1438,6 +1363,11 @@ impl Inner { Event::Frame { from, kind, body, .. } => match kind { + kind::WELCOME if serves(from) => { + if let Ok(welcome) = minicbor::decode::(body) { + self.join_presence(welcome.room); + } + } kind::REPLY if serves(from) => { if let Ok(reply) = minicbor::decode::(body) && let Some(waiting) = self.pending.lock().unwrap().remove(&reply.id) @@ -1458,18 +1388,59 @@ impl Inner { } } } - kind::BYE - if serves(from) - && minicbor::decode::(body) - .is_ok_and(|bye| bye.reason == STOPPED) => - { - self.stopped.store(true, Ordering::Release); + kind::BYE if serves(from) => { + if let Ok(bye) = minicbor::decode::(body) { + let ended = match bye.reason.as_str() { + STOPPED => Some(Ended::Stopped), + "removed" => Some(Ended::Removed), + _ => None, + }; + if ended.is_some() { + *self.ended.lock().unwrap() = ended; + } + } } _ => {} }, _ => {} } } + + fn join_presence(self: &Arc, secret: [u8; 16]) { + let mut presence = self.presence.lock().unwrap(); + if presence.as_ref().is_some_and(|(held, _)| *held == secret) { + return; + } + let inner = Arc::downgrade(self); + let live = Live::start( + self.me.clone(), + &Room::Notebook(secret), + self.reach, + self.relay.as_deref(), + move |event| { + if let Some(inner) = inner.upgrade() { + if matches!(event, Event::Frame { .. }) { + inner.heard(event); + } + (inner.events)(); + } + }, + ); + match live { + Ok(live) => { + live.set_presence(self.here.lock().unwrap().clone()); + *presence = Some((secret, live)); + let mut current = self.current.lock().unwrap(); + let paths: Vec = current.keys().cloned().collect(); + current.clear(); + drop(current); + if let Some(reports) = &*self.watch.lock().unwrap() { + reports.touched(&paths); + } + } + Err(error) => eprintln!("Live Share: could not join presence: {error}"), + } + } } /// A section a Live Share host serves, as a guest's replica publishes to it. diff --git a/crates/notebook/src/live/share/batch.rs b/crates/notebook/src/live/share/batch.rs index 92dd7e92c4eda78e7f85eae4bffc08263a99a0fa..cc60919c8661d6a4b51febcd903288c74485f031 100644 --- a/crates/notebook/src/live/share/batch.rs +++ b/crates/notebook/src/live/share/batch.rs @@ -148,6 +148,12 @@ impl Served { image: &[u8], batch: &Waiting, ) -> Result> { + if !self.guests.lock().unwrap().contains_key(&batch.peer) { + return Err(refused( + io::ErrorKind::PermissionDenied, + "The device is no longer connected", + )); + } let stamp: Stamp = batch .request .stamp diff --git a/crates/notebook/src/live/share/membership.rs b/crates/notebook/src/live/share/membership.rs new file mode 100644 index 0000000000000000000000000000000000000000..8e0c550303dd5e4a63c8adfa81751eaf163ce29e --- /dev/null +++ b/crates/notebook/src/live/share/membership.rs @@ -0,0 +1,391 @@ +use super::*; + +pub struct Host { + members: Arc, + pairing: Mutex>, +} + +#[allow(clippy::type_complexity)] +struct Members { + me: Hello, + notebook: String, + reach: Option, + relay: Option, + sharing: Mutex, + changes: Mutex<()>, + pending: Mutex, Line)>>, + access: Mutex>, + room: Mutex>, + served: Arc, + events: Arc, + keep: Box io::Result<()> + Send + Sync>, +} + +impl Host { + /// Shares a notebook, keeping every credential change before granting or revoking access. + #[allow(clippy::too_many_arguments)] + pub fn start( + storage: Box, + me: Hello, + sharing: Sharing, + notebook: &str, + reach: Option, + relay: Option<&str>, + events: impl Fn() + Send + Sync + 'static, + keep: impl Fn(&Sharing) -> io::Result<()> + Send + Sync + 'static, + ) -> io::Result { + keep(&sharing)?; + let members = Arc::new(Members { + me, + notebook: notebook.to_owned(), + reach, + relay: relay.map(str::to_owned), + sharing: Mutex::new(sharing), + changes: Mutex::default(), + pending: Mutex::default(), + access: Mutex::default(), + room: Mutex::default(), + served: Arc::new(Served { + storage, + images: Mutex::default(), + snapshots: Mutex::default(), + puts: Mutex::default(), + guests: Mutex::default(), + writers: Mutex::default(), + host: Mutex::default(), + room: Mutex::default(), + }), + events: Arc::new(events), + keep: Box::new(keep), + }); + let sharing = members.sharing.lock().unwrap().clone(); + members.presence(sharing.secret)?; + for device in &sharing.members { + members.open(device.secret)?; + } + let pairing = Mutex::new(Some(members.pair(&sharing)?)); + Ok(Self { members, pairing }) + } + + pub fn code(&self) -> Option { + let mut pairing = self.pairing.lock().unwrap(); + let pairing = pairing.as_mut()?; + let mut sharing = self.members.sharing.lock().unwrap(); + if pairing.burned() { + let mut next = sharing.clone(); + next.code = super::super::code::secret().ok()?; + (self.members.keep)(&next).ok()?; + *pairing = self.members.pair(&next).ok()?; + *sharing = next; + } + let code = pairing.code(); + if let Some(code) = &code + && *code != sharing.code + { + let mut next = sharing.clone(); + next.code = code.clone(); + (self.members.keep)(&next).ok()?; + *sharing = next; + } + code + } + + pub fn sharing(&self) -> Sharing { + self.code(); + self.members.sharing.lock().unwrap().clone() + } + + pub fn relayed(&self) -> Relayed { + self.pairing + .lock() + .unwrap() + .as_ref() + .map_or(Relayed::Unknown, Live::relayed) + } + + pub fn guests(&self) -> Vec { + self.members + .room + .lock() + .unwrap() + .as_ref() + .map(Live::peers) + .unwrap_or_default() + } + + pub fn devices(&self) -> Vec<(Device, bool)> { + let sharing = self.members.sharing.lock().unwrap(); + let access = self.members.access.lock().unwrap(); + sharing + .members + .iter() + .map(|device| { + let connected = access + .get(&device.secret) + .is_some_and(|live| !live.peers().is_empty()); + (device.clone(), connected) + }) + .collect() + } + + pub fn requests(&self) -> Vec> { + self.members + .pending + .lock() + .unwrap() + .values() + .map(|(hello, _)| Arc::clone(hello)) + .collect() + } + + pub fn approve(&self, approve: bool) -> io::Result<()> { + let mut sharing = self.members.sharing.lock().unwrap(); + let mut next = sharing.clone(); + next.approve = approve; + (self.members.keep)(&next)?; + *sharing = next; + drop(sharing); + if !approve { + for hello in self.requests() { + self.allow(&hello.peer)?; + } + } + (self.members.events)(); + Ok(()) + } + + pub fn allow(&self, peer: &[u8; 16]) -> io::Result<()> { + let pending = self.members.pending.lock().unwrap().remove(peer); + if let Some((hello, line)) = pending { + if let Err(error) = self.members.grant(&hello, &line) { + self.members + .pending + .lock() + .unwrap() + .insert(*peer, (hello, line)); + return Err(error); + } + (self.members.events)(); + } + Ok(()) + } + + pub fn decline(&self, peer: &[u8; 16]) { + if let Some((_, line)) = self.members.pending.lock().unwrap().remove(peer) { + let _ = line.send(kind::APPROVAL, &wire::Approval::Declined); + (self.members.events)(); + } + } + + pub fn remove(&self, secret: &[u8; 16]) -> io::Result<()> { + let _change = self.members.changes.lock().unwrap(); + let mut sharing = self.members.sharing.lock().unwrap(); + if !sharing + .members + .iter() + .any(|device| device.secret == *secret) + { + return Ok(()); + } + let mut next = sharing.clone(); + next.members.retain(|device| device.secret != *secret); + getrandom::fill(&mut next.secret) + .map_err(|_| io::Error::other("System random source failed"))?; + next.code = super::super::code::secret()?; + (self.members.keep)(&next)?; + *sharing = next.clone(); + drop(sharing); + let removed = self.members.access.lock().unwrap().remove(secret); + if let Some(removed) = removed { + for peer in removed.peers() { + self.members.served.forget(&peer.hello.peer); + } + removed.leave("removed"); + } + self.members.presence(next.secret)?; + *self.pairing.lock().unwrap() = Some(self.members.pair(&next)?); + let access = self.members.access.lock().unwrap(); + for device in &next.members { + if let Some(live) = access.get(&device.secret) { + live.sender().send( + kind::WELCOME, + &self.members.welcome(&next, device.secret), + None, + ); + } + } + (self.members.events)(); + Ok(()) + } + + pub fn set_presence(&self, presence: Presence) { + if let Some(room) = &*self.members.room.lock().unwrap() { + room.set_presence(presence); + } + } + + pub fn on_changed(&self, listener: crate::session::Listener) { + *self.members.served.host.lock().unwrap() = Some(listener); + } + + pub fn touched(&self, paths: &[String]) { + for path in paths { + self.members.served.changed_here(path); + } + self.members.served.tell(paths); + } + + pub fn stop(&self) { + let _change = self.members.changes.lock().unwrap(); + drop(self.pairing.lock().unwrap().take()); + self.members.pending.lock().unwrap().clear(); + let access = std::mem::take(&mut *self.members.access.lock().unwrap()); + for (_, live) in access { + live.leave(STOPPED); + } + drop(self.members.room.lock().unwrap().take()); + } +} + +impl Members { + fn welcome(&self, sharing: &Sharing, secret: [u8; 16]) -> Welcome { + Welcome { + share: sharing.share, + secret, + room: sharing.secret, + notebook: self.notebook.clone(), + host: self.me.name.clone(), + } + } + + fn presence(self: &Arc, secret: [u8; 16]) -> io::Result<()> { + let told = Arc::clone(&self.events); + let live = Live::start( + self.me.clone(), + &Room::Notebook(secret), + self.reach, + self.relay.as_deref(), + move |_| told(), + )?; + *self.served.room.lock().unwrap() = Some(live.sender()); + *self.room.lock().unwrap() = Some(live); + Ok(()) + } + + fn open(self: &Arc, secret: [u8; 16]) -> io::Result<()> { + let members = Arc::downgrade(self); + let live = Live::start( + Hello { + serves: Some(self.sharing.lock().unwrap().share), + ..self.me.clone() + }, + &Room::Notebook(secret), + self.reach, + self.relay.as_deref(), + move |event| { + let Some(members) = members.upgrade() else { + return; + }; + match event { + Event::Met(hello, line) => { + let sharing = members.sharing.lock().unwrap(); + if !sharing.members.iter().any(|device| device.secret == secret) { + line.hang_up("removed"); + return; + } + members.served.admit(hello.peer, line); + let _ = line.send(kind::WELCOME, &members.welcome(&sharing, secret)); + } + Event::Left(hello) => members.served.forget(&hello.peer), + Event::Frame { from, kind, body } + if wire::KNOWN.contains(&kind) && kind > 256 && kind != kind::REPLY => + { + members.served.queue(&from.peer, kind, body) + } + Event::Changed => (members.events)(), + _ => {} + } + }, + )?; + self.access.lock().unwrap().insert(secret, live); + Ok(()) + } + + fn grant(self: &Arc, hello: &Hello, line: &Line) -> io::Result<()> { + let _change = self.changes.lock().unwrap(); + if self.room.lock().unwrap().is_none() { + return Err(io::ErrorKind::NotConnected.into()); + } + let mut secret = [0; 16]; + getrandom::fill(&mut secret) + .map_err(|_| io::Error::other("System random source failed"))?; + if self.sharing.lock().unwrap().members.len() >= 64 { + return Err(io::ErrorKind::ResourceBusy.into()); + } + self.open(secret)?; + let kept = (|| -> io::Result { + let mut sharing = self.sharing.lock().unwrap(); + let mut next = sharing.clone(); + next.members.push(Device { + secret, + name: hello.name.clone(), + device: hello.device.clone(), + }); + (self.keep)(&next)?; + *sharing = next; + Ok(self.welcome(&sharing, secret)) + })(); + let welcome = match kept { + Ok(welcome) => welcome, + Err(error) => { + drop(self.access.lock().unwrap().remove(&secret)); + return Err(error); + } + }; + line.send(kind::WELCOME, &welcome)?; + (self.events)(); + Ok(()) + } + + fn pair(self: &Arc, sharing: &Sharing) -> io::Result { + let members = Arc::downgrade(self); + Live::start( + Hello { + serves: None, + ..self.me.clone() + }, + &Room::share(&sharing.code, &sharing.password), + self.reach, + self.relay.as_deref(), + move |event| { + let Some(members) = members.upgrade() else { + return; + }; + match event { + Event::Met(hello, line) => { + if members.sharing.lock().unwrap().approve { + let mut pending = members.pending.lock().unwrap(); + if pending.len() < 32 { + pending.insert(hello.peer, (Arc::clone(hello), line.clone())); + let _ = line.send(kind::APPROVAL, &wire::Approval::Pending); + } else { + let _ = line.send(kind::APPROVAL, &wire::Approval::Failed); + } + drop(pending); + (members.events)(); + } else if let Err(error) = members.grant(hello, line) { + eprintln!("Live Share: could not admit a device: {error}"); + let _ = line.send(kind::APPROVAL, &wire::Approval::Failed); + } + } + Event::Left(hello) => { + members.pending.lock().unwrap().remove(&hello.peer); + (members.events)(); + } + Event::Changed => (members.events)(), + _ => {} + } + }, + ) + } +} diff --git a/crates/notebook/src/live/wire.rs b/crates/notebook/src/live/wire.rs index 15b1e98286b6665b12058e2a9c076b000d3fd34b..892135de87a079f6173267e61349916a4d8c8cee 100644 --- a/crates/notebook/src/live/wire.rs +++ b/crates/notebook/src/live/wire.rs @@ -1,7 +1,6 @@ //! What peers say to each other: an opening in the clear that meets through the secret both //! hold (SPAKE2), then frames sealed under the keys it agreed, each a message kind and a CBOR -//! map. A later version adds kinds and fields; a reader skips the kinds and fields it doesn't -//! know, so every version speaks to every other. +//! map. Readers skip unknown kinds and fields within the same opening version. use aes_gcm::{Aes256Gcm, KeyInit, aead::Aead}; use hmac::{Hmac, Mac}; @@ -10,8 +9,8 @@ use sha2::Sha256; use spake2::{Ed25519Group, Identity, Password, Spake2}; use std::io::{self, Read, Write}; -/// The opening's version. Frames after it never change shape; they grow by kinds and fields. -pub const VERSION: u16 = 2; +/// The opening's version; peers must agree on its authentication and access rules. +pub const VERSION: u16 = 3; /// The version of the opening a peer of another version sent, as `open` fails with it. #[derive(Clone, Copy, Debug, PartialEq, Eq)] @@ -45,6 +44,7 @@ pub mod kind { /// A section's bytes as a commit changed them: `Delta`. pub const DELTA: u16 = 19; pub const WELCOME: u16 = 32; + pub const APPROVAL: u16 = 33; /// Storage requests to a host, each a `Request` answered by a `Reply`. pub const LIST: u16 = 257; pub const STAMP: u16 = 258; @@ -79,6 +79,7 @@ pub const KNOWN: &[u16] = &[ kind::TOUCHED, kind::DELTA, kind::WELCOME, + kind::APPROVAL, kind::LIST, kind::STAMP, kind::READ, @@ -122,6 +123,8 @@ pub struct Hello { pub serves: Option<[u8; 16]>, #[n(6)] pub ops: Option, + #[n(7)] + pub device: Option, } impl Hello { @@ -137,6 +140,7 @@ impl Hello { kinds: KNOWN.to_vec(), serves: None, ops: Some(1), + device: None, }) } } @@ -211,8 +215,7 @@ pub struct Bye { 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. +/// A device's access credential, the current presence room, and the notebook's names. #[derive(Clone, Debug, PartialEq, Encode, Decode)] #[cbor(map)] pub struct Welcome { @@ -225,6 +228,19 @@ pub struct Welcome { /// The host's name for itself, as `Hello::name`. #[n(3)] pub host: String, + #[cbor(n(4), with = "minicbor::bytes")] + pub room: [u8; 16], +} + +#[derive(Clone, Debug, Encode, Decode)] +#[cbor(index_only)] +pub enum Approval { + #[n(0)] + Pending, + #[n(1)] + Declined, + #[n(2)] + Failed, } /// Paths a host's files changed at, by catalog path; `""` is the notebook's folder. diff --git a/crates/notebook/tests/live_membership.rs b/crates/notebook/tests/live_membership.rs new file mode 100644 index 0000000000000000000000000000000000000000..835d23afbb6c261384057fc70f97f0cc340bac23 --- /dev/null +++ b/crates/notebook/tests/live_membership.rs @@ -0,0 +1,190 @@ +#![cfg(feature = "live")] + +#[path = "support/live.rs"] +mod live; +use live::*; +use notebook::live::share::{self, Guest, Host, Sharing}; +use notebook::session::Notebook; +use std::{ + sync::{ + Arc, Mutex, + atomic::{AtomicBool, Ordering}, + mpsc, + }, + time::Duration, +}; + +#[test] +fn approval_keeps_credentials_before_welcome_and_a_decline_grants_nothing() { + let directory = tempfile::tempdir().unwrap(); + let folder = notebook(directory.path()); + let url = relay(Default::default()); + let saved = Arc::new(Mutex::new(None)); + let fail = Arc::new(AtomicBool::new(false)); + let (kept, failing) = (Arc::clone(&saved), Arc::clone(&fail)); + let mut sharing = Sharing::new("").unwrap(); + sharing.approve = true; + let host = Host::start( + Notebook::open(&folder, directory.path().join("host")) + .unwrap() + .into_storage(), + hello("Ada"), + sharing, + "Garden", + None, + Some(&url), + || {}, + move |sharing| { + if failing.load(Ordering::Acquire) { + return Err(std::io::ErrorKind::PermissionDenied.into()); + } + *kept.lock().unwrap() = Some(serde_json::to_vec(sharing).unwrap()); + Ok(()) + }, + ) + .unwrap(); + let code = self::code(&host); + let (reply, result) = mpsc::channel(); + let joining = url.clone(); + std::thread::spawn(move || { + let _ = reply.send(share::join(hello("Grace"), &code, "", None, Some(&joining))); + }); + until("the request never arrived", || host.requests().len() == 1); + assert!(result.try_recv().is_err()); + assert!(host.devices().is_empty()); + let peer = host.requests()[0].peer; + fail.store(true, Ordering::Release); + assert!(host.allow(&peer).is_err()); + assert!(host.devices().is_empty()); + assert_eq!(host.requests().len(), 1); + fail.store(false, Ordering::Release); + host.allow(&peer).unwrap(); + let welcome = result + .recv_timeout(Duration::from_secs(5)) + .unwrap() + .unwrap(); + let persisted: Sharing = + serde_json::from_slice(saved.lock().unwrap().as_ref().unwrap()).unwrap(); + assert_eq!(persisted.members[0].secret, welcome.secret); + assert_ne!(welcome.secret, welcome.room); + fail.store(true, Ordering::Release); + assert!(host.remove(&welcome.secret).is_err()); + assert_eq!(host.sharing(), persisted); + fail.store(false, Ordering::Release); + + let code = self::code(&host); + let (reply, result) = mpsc::channel(); + std::thread::spawn(move || { + let _ = reply.send(share::join(hello("Alan"), &code, "", None, Some(&url))); + }); + until("the second request never arrived", || { + host.requests().len() == 1 + }); + host.decline(&host.requests()[0].peer); + assert_eq!( + result.recv_timeout(Duration::from_secs(5)).unwrap(), + Err(share::Refusal::Declined) + ); + assert_eq!(host.devices().len(), 1); +} + +#[test] +fn removal_retires_one_credential_and_other_devices_reconnect_after_a_restart() { + 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 original_code = code(&host); + let (alice, _) = guest( + "Alice", + &original_code, + &url, + &directory.path().join("alice"), + ); + let (bob, _) = guest("Bob", &original_code, &url, &directory.path().join("bob")); + let before = host.sharing(); + let alice_key = before + .members + .iter() + .find(|device| device.name == "Alice") + .unwrap() + .secret; + let bob_key = before + .members + .iter() + .find(|device| device.name == "Bob") + .unwrap() + .secret; + host.remove(&bob_key).unwrap(); + until("Bob was not removed", || { + bob.stopped() && bob.host().is_none() + }); + until("presence did not move to the new room", || { + host.guests().len() == 1 + }); + let current = host.sharing(); + assert_ne!(before.secret, current.secret); + assert_eq!(current.members.len(), 1); + assert_ne!(original_code, code(&host)); + assert!(alice.host().is_some()); + assert!(!alice.stopped()); + let live = Notebook::open_hosted(Arc::clone(&alice), directory.path().join("alice")).unwrap(); + assert_eq!(live.catalog().sections.len(), 2); + let restarted: Sharing = + serde_json::from_slice(&serde_json::to_vec(¤t).unwrap()).unwrap(); + drop(host); + until("Alice's connection did not close", || { + alice.host().is_none() + }); + let host = self::host(&folder, &host_cache, &restarted, &url); + until("Alice did not reconnect", || alice.host().is_some()); + assert_eq!(host.devices()[0].0.secret, alice_key); + let forged = Guest::start( + hello("Alice"), + current.share, + bob_key, + None, + Some(&url), + || {}, + ) + .unwrap(); + std::thread::sleep(Duration::from_millis(500)); + assert!(forged.host().is_none()); + assert!(Notebook::open_hosted(forged, directory.path().join("forged")).is_err()); +} + +#[test] +fn cancelling_a_join_retires_its_pending_request() { + let directory = tempfile::tempdir().unwrap(); + let folder = notebook(directory.path()); + let url = relay(Default::default()); + let mut sharing = Sharing::new("").unwrap(); + sharing.approve = true; + let host = host(&folder, &directory.path().join("host"), &sharing, &url); + let code = code(&host); + let alive = Arc::new(AtomicBool::new(true)); + let continuing = Arc::clone(&alive); + let (reply, result) = mpsc::channel(); + std::thread::spawn(move || { + let _ = reply.send(share::join_while( + hello("Grace"), + &code, + "", + None, + Some(&url), + |_| continuing.load(Ordering::Acquire), + )); + }); + until("the request never arrived", || !host.requests().is_empty()); + alive.store(false, Ordering::Release); + assert_eq!( + result.recv_timeout(Duration::from_secs(5)).unwrap(), + Err(share::Refusal::Cancelled) + ); + until("the cancelled request stayed", || { + host.requests().is_empty() + }); + assert!(host.devices().is_empty()); +} diff --git a/crates/notebook/tests/live_share.rs b/crates/notebook/tests/live_share.rs index c6aafd8b8c63884b7677c7bab214dbbb5ca283b2..6d2967e243eb7435a2df0b602846b224a0e50cec 100644 --- a/crates/notebook/tests/live_share.rs +++ b/crates/notebook/tests/live_share.rs @@ -70,6 +70,7 @@ fn a_guest_queues_while_the_host_is_away() { let image = std::fs::read(&file).unwrap(); // Ada's computer goes to sleep. + let sharing = host.sharing(); drop(host); until("the host never left", || guest.host().is_none()); let id = replace(§ion, &image, 0..8, "Offline"); @@ -110,6 +111,7 @@ fn two_guests_on_one_page_conflict_as_on_a_share() { let file = folder.join("Garden.one"); let image = std::fs::read(&file).unwrap(); + let sharing = host.sharing(); drop(host); until("the host never left", || { grace.host().is_none() && alan.host().is_none() diff --git a/crates/notebook/tests/support/live.rs b/crates/notebook/tests/support/live.rs index 911b42f72d5abeb99cb1b7f9ed7c75b40c093645..96af2f076a6314ce6406dd63a9e4218e6189cdb2 100644 --- a/crates/notebook/tests/support/live.rs +++ b/crates/notebook/tests/support/live.rs @@ -83,6 +83,7 @@ pub fn host(folder: &Path, cache: &Path, sharing: &Sharing, url: &str) -> Host { None, Some(url), || {}, + |_| Ok(()), ) .unwrap() } diff --git a/crates/relay/src/main.rs b/crates/relay/src/main.rs index ad801093796efb29cbebd3f8cff6f556367f862b..069512f43a53f3f374bb30f43b54718658c6ac8c 100644 --- a/crates/relay/src/main.rs +++ b/crates/relay/src/main.rs @@ -12,14 +12,14 @@ underscores: SNOWBOUND_RELAY_LISTEN=127.0.0.1:23592. Flags win. --listen ADDRESS where to listen (127.0.0.1:23592) --trust-forwarded true|false count peers by X-Forwarded-For, behind a proxy (false) --max-connections N (256) - --max-connections-per-address N per IPv4 address or IPv6 /64 (16) + --max-connections-per-address N per IPv4 address or IPv6 /64 (128) --max-rooms N (128) --max-room-peers N (64) --max-message BYTES (262144) --queue BYTES waiting to go to one peer before it is dropped (1048576) --idle SECONDS silence before a connection is closed (600) --room-bytes-per-second BYTES (4194304) - --joins-per-minute N per address (30) + --joins-per-minute N per address (240) --room-joins-per-minute N (120) --failures-per-minute N wrong codes per address before a lockout (10) --burn-after N wrong codes before a code admits no one new (5) diff --git a/crates/relay/src/server.rs b/crates/relay/src/server.rs index 2b56d2a7f95866d2ed09181e354376f4c1d2fc43..1a5fd458a394ae48a4a361933444cd57f07b84cf 100644 --- a/crates/relay/src/server.rs +++ b/crates/relay/src/server.rs @@ -48,14 +48,14 @@ impl Default for Config { Self { trust_forwarded: false, max_connections: 256, - max_connections_per_address: 16, + max_connections_per_address: 128, max_rooms: 128, max_room_peers: 64, max_message: 256 << 10, queue: 1 << 20, idle: Duration::from_secs(600), room_bytes_per_second: 4 << 20, - joins_per_minute: 30, + joins_per_minute: 240, room_joins_per_minute: 120, failures_per_minute: 10, burn_after: 5, diff --git a/crates/snowbound/src/live.rs b/crates/snowbound/src/live.rs index d1d47f54b47970461d9aeb0e32ee9d6a59fccfb6..95adeb6cffe5c5f7afa08f9483c0327c43a0a217 100644 --- a/crates/snowbound/src/live.rs +++ b/crates/snowbound/src/live.rs @@ -90,7 +90,17 @@ fn hello() -> io::Result { let picture = (picture && name == platform::user_name()) .then(|| account_picture().clone()) .flatten(); - Hello::new(name, picture) + let mut hello = Hello::new(name, picture)?; + hello.device = std::env::var("COMPUTERNAME").ok().or_else(|| { + std::process::Command::new("hostname") + .output() + .ok() + .filter(|output| output.status.success()) + .and_then(|output| String::from_utf8(output.stdout).ok()) + .map(|name| name.trim().to_owned()) + .filter(|name| !name.is_empty()) + }); + Ok(hello) } /// The relay peers off this network meet through: `SNOWBOUND_LIVE_RELAY` (`off` for none), @@ -141,11 +151,13 @@ fn keep(file: &Path, value: &impl serde::Serialize) -> io::Result<()> { options.write(true).create(true).truncate(true); #[cfg(unix)] std::os::unix::fs::OpenOptionsExt::mode(&mut options, 0o600); - io::Write::write_all( - &mut options.open(&partial)?, - &serde_json::to_vec_pretty(value)?, - )?; - notebook::fs::rename(partial, file) + let mut output = options.open(&partial)?; + io::Write::write_all(&mut output, &serde_json::to_vec_pretty(value)?)?; + output.sync_all()?; + notebook::fs::rename(partial, file)?; + #[cfg(unix)] + std::fs::File::open(folder)?.sync_all()?; + Ok(()) } /// A notebook another computer shares, as this one joined it. @@ -157,6 +169,8 @@ struct Share { host: String, /// It has been listed here, so it opens without waiting for the host. listed: bool, + #[serde(default)] + protocol: u16, } const JOINED: &str = "joined.json"; @@ -186,6 +200,11 @@ impl Joined { )); }; let refused = |error: &dyn std::fmt::Display| (share.notebook.clone(), error.to_string()); + if share.protocol != live::wire::VERSION { + return Err(refused( + &"Open this notebook again with the person sharing’s current link or code.", + )); + } let guest = hello() .and_then(|me| { Guest::start( @@ -241,9 +260,10 @@ impl Joined { pub(crate) fn join( code: &str, password: &str, + waiting: impl Fn(bool) -> bool, ) -> Result { let me = hello().map_err(|_| live::share::Refusal::Unreachable(live::Trouble::Other))?; - live::share::join(me, code, password, reach(), relay().as_deref()) + live::share::join_while(me, code, password, reach(), relay().as_deref(), waiting) } /// Keeps `welcome` as the notebook this computer joined: its location. @@ -259,6 +279,7 @@ pub(crate) fn joined(cache: &Path, welcome: live::wire::Welcome) -> io::Result>, /// Shares as kept for the next launch. - sharing: Option>, + sharing: Option>>>, /// Shares starting on threads of their own, by location. pub(crate) starting: HashMap>, started: Option)>>, @@ -410,7 +431,9 @@ impl State { let sharing = self .peers .sharing - .get_or_insert_with(|| read_kept(&kept(&cache, HOSTING))) + .get_or_insert_with(|| Arc::new(Mutex::new(read_kept(&kept(&cache, HOSTING))))) + .lock() + .unwrap() .clone(); let open: Vec> = self.notebooks.clone(); for library in &open { @@ -443,25 +466,29 @@ impl State { for location in closed { self.stop_sharing(&location); } - let now: BTreeMap = (self.peers.hosts.iter()) - .map(|(location, host)| (location.clone(), host.sharing())) - .collect(); - let before: BTreeMap = (sharing.into_iter()) - .filter(|(location, _)| { - !self.peers.starting.contains_key(location) - && open.iter().any(|library| library.location == *location) + for host in self.peers.hosts.values() { + host.code(); + } + if self.peers.share.is_none() + && let Some(library) = open.iter().find(|library| { + self.peers + .hosts + .get(&library.location) + .is_some_and(|host| !host.requests().is_empty()) }) - .collect(); - if now != before { - if let Err(error) = keep(&kept(&cache, HOSTING), &now) { - eprintln!("Keeping what this computer shares: {error}"); - } - self.peers.sharing = Some(now); + { + self.open_live_share(Arc::clone(library)); } } /// Shares `library` as `sharing` says, on a thread of its own. pub(crate) fn start_sharing(&mut self, library: &Arc, sharing: Sharing) { + let file = kept(&self.cache, HOSTING); + let kept = Arc::clone( + self.peers + .sharing + .get_or_insert_with(|| Arc::new(Mutex::new(read_kept(&file)))), + ); self.peers.starting.insert(library.location.clone(), None); let (started, _) = self.peers.started.get_or_insert_with(mpsc::channel); let (started, library, redraw) = @@ -470,6 +497,7 @@ impl State { let host = (|| -> Result> { let storage = library.reopen()?.into_storage(); let told = redraw.clone(); + let location = library.location.clone(); Ok(Host::start( storage, hello()?, @@ -478,6 +506,14 @@ impl State { reach(), relay().as_deref(), move || told.wake_by_ref(), + move |sharing| { + let mut shares = kept.lock().unwrap(); + let mut next = shares.clone(); + next.insert(location.clone(), sharing.clone()); + keep(&file, &next)?; + *shares = next; + Ok(()) + }, )?) })(); let _ = started.send(( @@ -521,11 +557,12 @@ impl State { crate::spawn(move || host.stop()); } // Kept, the share would start again on the next frame, as after a relaunch. - if let Some(sharing) = &mut self.peers.sharing - && sharing.remove(location).is_some() - && let Err(error) = keep(&kept(&self.cache, HOSTING), sharing) - { - eprintln!("Keeping what this computer shares: {error}"); + if let Some(sharing) = &self.peers.sharing { + let mut sharing = sharing.lock().unwrap(); + sharing.remove(location); + if let Err(error) = keep(&kept(&self.cache, HOSTING), &*sharing) { + eprintln!("Keeping what this computer shares: {error}"); + } } } diff --git a/crates/snowbound/src/share.rs b/crates/snowbound/src/share.rs index 2dd2910aa0ec334d4313ae91ddf8707fbb515211..97933b38a88f644891307e17450d7a9bba3613bf 100644 --- a/crates/snowbound/src/share.rs +++ b/crates/snowbound/src/share.rs @@ -8,7 +8,11 @@ use notebook::live::{ share::{self, Refusal, Sharing}, wire::Welcome, }; -use std::sync::{Arc, mpsc}; +use std::sync::{ + Arc, + atomic::{AtomicBool, Ordering}, + mpsc, +}; use ui::{Anchor, Axis, Id, Spec, Ui, children, fill, fit, px}; use winit::keyboard::NamedKey; @@ -41,6 +45,9 @@ pub(crate) struct ShareDialog { library: Arc, protect: bool, password: String, + approve: bool, + replies: crate::live::Channel>, + error: Option, /// What was copied last: the code, or the link. copied: Option<&'static str>, } @@ -53,6 +60,14 @@ pub(crate) struct JoinDialog { password: String, status: Status, replies: crate::live::Channel>, + alive: Arc, + approval: Arc, +} + +impl Drop for JoinDialog { + fn drop(&mut self) { + self.alive.store(false, Ordering::Release); + } } enum Status { @@ -92,6 +107,9 @@ fn refusal(refusal: &Refusal) -> String { Refusal::Busy => "The Live Share relay is busy. Try again in a minute.".into(), Refusal::Unreachable(trouble) => unreachable(*trouble), Refusal::TimedOut => "The computer sharing didn’t answer. Try again.".into(), + Refusal::Declined => "The person sharing declined this request.".into(), + Refusal::Cancelled => "Joining cancelled.".into(), + Refusal::NotAdmitted => "The computer sharing couldn’t add this device. Ask the person sharing to check Live Share.".into(), } } @@ -267,6 +285,9 @@ impl State { library, protect: false, password: String::new(), + approve: true, + replies: mpsc::channel(), + error: None, copied: None, }); self.ui.open_popup(share_id()); @@ -289,6 +310,8 @@ impl State { password: String::new(), status: Status::Idle, replies: mpsc::channel(), + alive: Arc::new(AtomicBool::new(true)), + approval: Arc::new(AtomicBool::new(false)), }); self.ui.open_popup(join_id()); self.ui.focus_all(code_field()); @@ -307,6 +330,9 @@ impl State { return; } let location = dialog.library.location.clone(); + for result in dialog.replies.1.try_iter() { + dialog.error = result.err(); + } let host = self.peers.hosts.get(&location).cloned(); let starting = self.peers.starting.get(&location).cloned(); let ui = &mut self.ui; @@ -319,6 +345,7 @@ impl State { frame(ui, share_id(), "Live Share"); text(ui, "notebook", &dialog.library.name, true); let (mut start, mut stop, mut copy) = (false, false, None); + let (mut approval, mut allow, mut decline, mut remove) = (None, None, None, None); match (&host, &starting) { (Some(host), _) => { let code = host.code(); @@ -390,30 +417,85 @@ impl State { ), _ => {} } + if ui::check_box(ui, "approve", "Ask before joining", host.sharing().approve) + .clicked + { + approval = Some(!host.sharing().approve); + } + for request in host.requests() { + ui.open( + ("request", request.peer), + Spec { + size: [fill(), children()], + gap: 8.0, + ..Spec::default() + }, + ); + ui.leaf( + "name", + Spec { + size: [fill(), fit()], + text: Some(&format!("{} wants to join", request.name)), + overflow: ui::Overflow::Wrap, + ..Spec::default() + }, + ); + if ui::button(ui, "allow", "Allow").clicked { + allow = Some(request.peer); + } + if ui::button(ui, "decline", "Decline").clicked { + decline = Some(request.peer); + } + ui.close(); + } ui.leaf( "people", Spec { size: [fill(), px(theme.font_size * 2.0)], - text: Some("Connected"), + text: Some("Devices"), bold: true, role: Some(Role::Heading), ..Spec::default() }, ); - let guests = host.guests(); - if guests.is_empty() { + let devices = host.devices(); + if devices.is_empty() { text(ui, "nobody", "No one has joined yet.", true); } - for guest in &guests { + for (device, connected) in &devices { + ui.open( + ("device", device.secret), + Spec { + size: [fill(), children()], + gap: 8.0, + ..Spec::default() + }, + ); ui.leaf( - ("guest", guest.hello.peer), + "name", Spec { size: [fill(), px(theme.font_size * 1.6)], - text: Some(&guest.hello.name), + text: Some(&match &device.device { + Some(name) => format!("{} · {name}", device.name), + None => device.name.clone(), + }), role: Some(Role::ListItem), ..Spec::default() }, ); + ui.leaf( + "connection", + Spec { + size: [fit(), px(theme.font_size * 1.6)], + text: Some(if *connected { "Connected" } else { "Offline" }), + color: Some(theme.text_dim), + ..Spec::default() + }, + ); + if ui::button(ui, "remove", "Remove").clicked { + remove = Some(device.secret); + } + ui.close(); } buttons(ui); stop = ui::button(ui, "stop", "Stop Sharing").clicked; @@ -436,6 +518,9 @@ impl State { ui.focus_all(password_field()); } } + if ui::check_box(ui, "approve", "Ask before joining", dialog.approve).clicked { + dialog.approve = !dialog.approve; + } if dialog.protect { labelled(ui, "Password:", |ui| { let spec = field_spec(ui); @@ -455,6 +540,9 @@ impl State { let offered = host.is_none() && !matches!(starting, Some(None)); let done = ui::button(ui, "done", "Done").clicked || entered && !offered; ui.close(); + if let Some(error) = &dialog.error { + status(ui, error); + } ui.close(); if let Some((what, text)) = copy { dialog.copied = self.clipboard.set_text(text).is_ok().then_some(what); @@ -462,6 +550,27 @@ impl State { if stop { self.stop_sharing(&location); } + if let Some(host) = host { + if let Some(peer) = decline { + host.decline(&peer); + } + if approval.is_some() || allow.is_some() || remove.is_some() { + let replies = dialog.replies.0.clone(); + let redraw = self.redraw.clone(); + crate::spawn(move || { + let result = if let Some(approve) = approval { + host.approve(approve) + } else if let Some(peer) = allow { + host.allow(&peer) + } else { + host.remove(&remove.unwrap()) + }; + let _ = replies + .send(result.map_err(|error| format!("Couldn’t change access: {error}"))); + redraw.wake(); + }); + } + } if start && !(dialog.protect && dialog.password.is_empty()) { let password = if dialog.protect { dialog.password.clone() @@ -469,7 +578,10 @@ impl State { String::new() }; match Sharing::new(&password) { - Ok(sharing) => self.start_sharing(&dialog.library, sharing), + Ok(mut sharing) => { + sharing.approve = dialog.approve; + self.start_sharing(&dialog.library, sharing); + } Err(error) => { self.peers .starting @@ -539,7 +651,14 @@ impl State { }); match &dialog.status { Status::Idle => {} - Status::Waiting => status(ui, "Connecting…"), + Status::Waiting => status( + ui, + if dialog.approval.load(Ordering::Acquire) { + "Needs approval" + } else { + "Connecting…" + }, + ), Status::Failed(message) => status(ui, message), } buttons(ui); @@ -560,8 +679,15 @@ impl State { dialog.status = Status::Waiting; let (replies, redraw) = (dialog.replies.0.clone(), self.redraw.clone()); let password = dialog.password.clone(); + let (alive, approval) = + (Arc::clone(&dialog.alive), Arc::clone(&dialog.approval)); crate::spawn(move || { - let reply = crate::live::join(&code, &password); + let reply = crate::live::join(&code, &password, |waiting| { + if approval.swap(waiting, Ordering::AcqRel) != waiting { + redraw.wake_by_ref(); + } + alive.load(Ordering::Acquire) + }); let _ = replies.send(reply); redraw.wake(); }); diff --git a/crates/snowbound/tests/replay.rs b/crates/snowbound/tests/replay.rs index e8db0b0269ce65ef3e8f954f0fc30badb0373a05..87264e8d90751a5f58b70d5a0e1663540327fca5 100644 --- a/crates/snowbound/tests/replay.rs +++ b/crates/snowbound/tests/replay.rs @@ -723,11 +723,12 @@ fn stop_sharing_ends_the_share() { let mut steps = vec!["modifiers command shift", "key p", "modifiers", "settle"]; steps.extend(["type Live Share", "settle", "key Enter", "settle"]); steps.extend(["key Enter", "wait 1500", "accessibility shared"]); - // Copy, Copy Link, then Stop Sharing. + // Copy, Copy Link, Ask before joining, then Stop Sharing. steps.extend([ "key Tab", "key Tab", "key Tab", + "key Tab", "key Enter", "wait 1500", "accessibility stopped",