| author | |
| committer | |
| log | 725072f4398ea34fd120598c30e2c7c90714ffd3 |
| tree | 48448889f0b2f9ee899bfe71e5c347e88412b824 |
| parent | 55d5656dafafcc0610b573a1d68d6051fa92c260 |
| signature | Signed by SSH key SHA256:52mNGHRsVFBDED9IAX5pe+LRWUefqTbxEReunq21QvU |
Give each admitted device its own persisted storage credential, and rotate the
presence key and invite code on removal. Persist grants before sending credentials
and reject queued requests after withdrawal. Show pending requests and connected
or offline devices in Live Share; joining can wait for approval or be cancelled.
Raise relay address limits to accommodate the private device rooms behind a
shared network. Full CI, admission and revocation tests, native UI replay, and a
30-guest typing run passed.
Assisted-by: gpt-6.1-sol14 files changed, 1007 insertions(+), 264 deletions(-)
crates/notebook/examples/live_crowd.rs+1| ... | ... | @@ -133,6 +133,7 @@ fn host(url: &str, paragraphs: usize) { |
| 133 | 133 | None, |
| 134 | 134 | Some(url), |
| 135 | 135 | || {}, |
| 136 | |_| Ok(()), | |
| 136 | 137 | ) |
| 137 | 138 | .unwrap(); |
| 138 | 139 | until("no code", Duration::from_secs(30), || { |
crates/notebook/examples/live_latency.rs+1| ... | ... | @@ -160,6 +160,7 @@ fn main() { |
| 160 | 160 | reach, |
| 161 | 161 | relay, |
| 162 | 162 | || {}, |
| 163 | |_| Ok(()), | |
| 163 | 164 | ) |
| 164 | 165 | .unwrap(); |
| 165 | 166 | until("no code", || { |
crates/notebook/src/live/share.rs+187-216| ... | ... | @@ -1,9 +1,9 @@ |
| 1 | 1 | //! Live Share: a notebook one Snowbound holds, opened on others through a short code. The |
| 2 | 2 | //! host serves its notebook's storage verbs (`session::Storage`) and batches guests' ops into |
| 3 | 3 | //! guarded publications. A guest runs the same replica and durable queue as on an SMB share; |
| 4 | //! protected sections and older peers use transactions. A guest first meets the host in the code's room, where | |
| 5 | //! the host welcomes it with the share's room and secret; a new share has a new secret, so | |
| 6 | //! stopping retires every guest. Large bodies travel a chunk at a time, each answered before | |
| 4 | //! protected sections use transactions. A guest first meets the host in the code's room, | |
| 5 | //! where approval grants a device's own access credential. Presence has a separate room | |
| 6 | //! whose key changes when a device is removed. Large bodies travel a chunk at a time, each answered before | |
| 7 | 7 | //! the next, so a relay never holds much for a slow peer. |
| 8 | 8 | |
| 9 | 9 | use super::{ |
| ... | ... | @@ -20,7 +20,7 @@ use std::{ |
| 20 | 20 | io, |
| 21 | 21 | path::PathBuf, |
| 22 | 22 | sync::{ |
| 23 | Arc, Condvar, Mutex, OnceLock, | |
| 23 | Arc, Condvar, Mutex, | |
| 24 | 24 | atomic::{AtomicBool, AtomicU64, Ordering}, |
| 25 | 25 | mpsc, |
| 26 | 26 | }, |
| ... | ... | @@ -29,6 +29,8 @@ use std::{ |
| 29 | 29 | }; |
| 30 | 30 | |
| 31 | 31 | mod batch; |
| 32 | mod membership; | |
| 33 | pub use membership::Host; | |
| 32 | 34 | |
| 33 | 35 | /// The most bytes one message of a read or an upload carries. |
| 34 | 36 | const CHUNK: usize = 128 << 10; |
| ... | ... | @@ -77,11 +79,22 @@ pub fn code(typed: &str) -> Option<String> { |
| 77 | 79 | #[derive(Clone, Debug, PartialEq, Eq, serde::Serialize, serde::Deserialize)] |
| 78 | 80 | pub struct Sharing { |
| 79 | 81 | pub share: [u8; 16], |
| 80 | /// The share room's secret, which every guest welcomed holds. | |
| 82 | /// The presence key, rotated when a device is removed. | |
| 81 | 83 | pub secret: [u8; 16], |
| 82 | 84 | /// The code, or its secret alone until it has a number (`super::code`). |
| 83 | 85 | pub code: String, |
| 84 | 86 | pub password: String, |
| 87 | #[serde(default)] | |
| 88 | pub approve: bool, | |
| 89 | #[serde(default)] | |
| 90 | pub members: Vec<Device>, | |
| 91 | } | |
| 92 | ||
| 93 | #[derive(Clone, Debug, PartialEq, Eq, serde::Serialize, serde::Deserialize)] | |
| 94 | pub struct Device { | |
| 95 | pub secret: [u8; 16], | |
| 96 | pub name: String, | |
| 97 | pub device: Option<String>, | |
| 85 | 98 | } |
| 86 | 99 | |
| 87 | 100 | impl Sharing { |
| ... | ... | @@ -95,6 +108,8 @@ impl Sharing { |
| 95 | 108 | secret: random[16..].try_into().expect("16 bytes"), |
| 96 | 109 | code: super::code::secret()?, |
| 97 | 110 | password: password.to_owned(), |
| 111 | approve: false, | |
| 112 | members: Vec::new(), | |
| 98 | 113 | }) |
| 99 | 114 | } |
| 100 | 115 | } |
| ... | ... | @@ -111,7 +126,9 @@ pub enum Refusal { |
| 111 | 126 | Malformed, |
| 112 | 127 | /// The person sharing runs a Snowbound of another Live Share version: the newer one's |
| 113 | 128 | /// `true` where it is theirs, so this one should update. |
| 114 | Version { theirs_newer: bool }, | |
| 129 | Version { | |
| 130 | theirs_newer: bool, | |
| 131 | }, | |
| 115 | 132 | /// The code's secret or password is wrong. |
| 116 | 133 | Wrong, |
| 117 | 134 | /// No one shares with the code's number now. |
| ... | ... | @@ -126,6 +143,9 @@ pub enum Refusal { |
| 126 | 143 | Unreachable(super::Trouble), |
| 127 | 144 | /// The relay let this end in, but no one answered. |
| 128 | 145 | TimedOut, |
| 146 | Declined, | |
| 147 | Cancelled, | |
| 148 | NotAdmitted, | |
| 129 | 149 | } |
| 130 | 150 | |
| 131 | 151 | /// Meets the host sharing `code` (and `password`) as `me`, on the networks `reach` names and |
| ... | ... | @@ -136,9 +156,23 @@ pub fn join( |
| 136 | 156 | password: &str, |
| 137 | 157 | reach: Option<Reach>, |
| 138 | 158 | relay: Option<&str>, |
| 159 | ) -> std::result::Result<Welcome, Refusal> { | |
| 160 | join_while(me, code, password, reach, relay, |_| true) | |
| 161 | } | |
| 162 | ||
| 163 | /// Joins while `waiting` returns true, reporting whether the host is deciding approval. | |
| 164 | pub fn join_while( | |
| 165 | me: Hello, | |
| 166 | code: &str, | |
| 167 | password: &str, | |
| 168 | reach: Option<Reach>, | |
| 169 | relay: Option<&str>, | |
| 170 | continue_joining: impl Fn(bool) -> bool, | |
| 139 | 171 | ) -> std::result::Result<Welcome, Refusal> { |
| 140 | 172 | let code = self::code(code).ok_or(Refusal::Malformed)?; |
| 141 | 173 | let (welcomed, welcome) = mpsc::channel(); |
| 174 | let approving = Arc::new(AtomicBool::new(false)); | |
| 175 | let approval = Arc::clone(&approving); | |
| 142 | 176 | let (changed, waiting) = mpsc::channel(); |
| 143 | 177 | let live = Live::start( |
| 144 | 178 | me, |
| ... | ... | @@ -152,7 +186,24 @@ pub fn join( |
| 152 | 186 | .. |
| 153 | 187 | } => { |
| 154 | 188 | if let Ok(body) = minicbor::decode::<Welcome>(body) { |
| 155 | let _ = welcomed.send(body); | |
| 189 | let _ = welcomed.send(Ok(body)); | |
| 190 | } | |
| 191 | } | |
| 192 | Event::Frame { | |
| 193 | kind: kind::APPROVAL, | |
| 194 | body, | |
| 195 | .. | |
| 196 | } => { | |
| 197 | if let Ok(body) = minicbor::decode::<wire::Approval>(body) { | |
| 198 | match body { | |
| 199 | wire::Approval::Pending => approval.store(true, Ordering::Release), | |
| 200 | wire::Approval::Declined => { | |
| 201 | let _ = welcomed.send(Err(Refusal::Declined)); | |
| 202 | } | |
| 203 | wire::Approval::Failed => { | |
| 204 | let _ = welcomed.send(Err(Refusal::NotAdmitted)); | |
| 205 | } | |
| 206 | } | |
| 156 | 207 | } |
| 157 | 208 | } |
| 158 | 209 | _ => { |
| ... | ... | @@ -163,8 +214,11 @@ pub fn join( |
| 163 | 214 | .map_err(|_| Refusal::Unreachable(super::Trouble::Other))?; |
| 164 | 215 | let start = Instant::now(); |
| 165 | 216 | loop { |
| 217 | if !continue_joining(approving.load(Ordering::Acquire)) { | |
| 218 | return Err(Refusal::Cancelled); | |
| 219 | } | |
| 166 | 220 | if let Ok(welcome) = welcome.try_recv() { |
| 167 | return Ok(welcome); | |
| 221 | return welcome; | |
| 168 | 222 | } |
| 169 | 223 | if let Some(version) = live.other_version() { |
| 170 | 224 | return Err(Refusal::Version { |
| ... | ... | @@ -192,199 +246,21 @@ pub fn join( |
| 192 | 246 | Relayed::Unknown if relay.is_none() && waited > Duration::from_secs(10) => { |
| 193 | 247 | return Err(Refusal::NoOne); |
| 194 | 248 | } |
| 195 | _ if waited > Duration::from_secs(20) => return Err(Refusal::TimedOut), | |
| 249 | _ if waited | |
| 250 | > Duration::from_secs(if approving.load(Ordering::Acquire) { | |
| 251 | 300 | |
| 252 | } else { | |
| 253 | 20 | |
| 254 | }) => | |
| 255 | { | |
| 256 | return Err(Refusal::TimedOut); | |
| 257 | } | |
| 196 | 258 | _ => {} |
| 197 | 259 | } |
| 198 | 260 | let _ = waiting.recv_timeout(Duration::from_millis(100)); |
| 199 | 261 | } |
| 200 | 262 | } |
| 201 | 263 | |
| 202 | /// A notebook shared while it lives: the share's room, serving the notebook's storage to the | |
| 203 | /// guests in it, and the code's room, welcoming whoever knows the code. | |
| 204 | pub struct Host { | |
| 205 | me: Hello, | |
| 206 | notebook: String, | |
| 207 | reach: Option<Reach>, | |
| 208 | relay: Option<String>, | |
| 209 | sharing: Mutex<Sharing>, | |
| 210 | /// The share's room and the code's, until it stops. | |
| 211 | room: Mutex<Option<Live>>, | |
| 212 | pairing: Mutex<Option<Live>>, | |
| 213 | served: Arc<Served>, | |
| 214 | events: Arc<dyn Fn() + Send + Sync>, | |
| 215 | } | |
| 216 | ||
| 217 | impl Host { | |
| 218 | /// Shares `storage`, the notebook named `notebook`, as `sharing` says, as `me`, where | |
| 219 | /// `reach` and `relay` say. `events` runs on a network thread whenever the guests or the | |
| 220 | /// code change. What guests change reaches the host as its own watch on the notebook's | |
| 221 | /// folder reports it, and reaches the other guests at once. | |
| 222 | pub fn start( | |
| 223 | storage: Box<dyn Storage>, | |
| 224 | me: Hello, | |
| 225 | sharing: Sharing, | |
| 226 | notebook: &str, | |
| 227 | reach: Option<Reach>, | |
| 228 | relay: Option<&str>, | |
| 229 | events: impl Fn() + Send + Sync + 'static, | |
| 230 | ) -> io::Result<Self> { | |
| 231 | let events: Arc<dyn Fn() + Send + Sync> = Arc::new(events); | |
| 232 | let served = Arc::new(Served { | |
| 233 | storage, | |
| 234 | images: Mutex::default(), | |
| 235 | snapshots: Mutex::default(), | |
| 236 | puts: Mutex::default(), | |
| 237 | guests: Mutex::default(), | |
| 238 | writers: Mutex::default(), | |
| 239 | host: Mutex::default(), | |
| 240 | room: OnceLock::new(), | |
| 241 | }); | |
| 242 | let serving = Hello { | |
| 243 | serves: Some(sharing.share), | |
| 244 | ..me.clone() | |
| 245 | }; | |
| 246 | let (heard, told) = (Arc::clone(&served), Arc::clone(&events)); | |
| 247 | let room = Live::start( | |
| 248 | serving, | |
| 249 | &Room::Notebook(sharing.secret), | |
| 250 | reach, | |
| 251 | relay, | |
| 252 | move |event| match event { | |
| 253 | Event::Met(hello, line) => heard.admit(hello.peer, line), | |
| 254 | Event::Left(hello) => heard.forget(&hello.peer), | |
| 255 | Event::Frame { | |
| 256 | from, kind, body, .. | |
| 257 | } if wire::KNOWN.contains(&kind) && kind > 256 && kind != kind::REPLY => { | |
| 258 | heard.queue(&from.peer, kind, body); | |
| 259 | } | |
| 260 | Event::Changed => told(), | |
| 261 | _ => {} | |
| 262 | }, | |
| 263 | )?; | |
| 264 | let _ = served.room.set(room.sender()); | |
| 265 | let host = Self { | |
| 266 | notebook: notebook.to_owned(), | |
| 267 | reach, | |
| 268 | relay: relay.map(str::to_owned), | |
| 269 | pairing: Mutex::new(Some(pair(&me, &sharing, notebook, reach, relay, &events)?)), | |
| 270 | me, | |
| 271 | sharing: Mutex::new(sharing), | |
| 272 | room: Mutex::new(Some(room)), | |
| 273 | served, | |
| 274 | events, | |
| 275 | }; | |
| 276 | Ok(host) | |
| 277 | } | |
| 278 | ||
| 279 | /// The code guests type, once it has its number, and none once stopped. A code with too | |
| 280 | /// many wrong tries is replaced by one with a new secret. | |
| 281 | pub fn code(&self) -> Option<String> { | |
| 282 | let mut pairing = self.pairing.lock().unwrap(); | |
| 283 | let pairing = pairing.as_mut()?; | |
| 284 | if pairing.burned() { | |
| 285 | let mut sharing = self.sharing.lock().unwrap(); | |
| 286 | sharing.code = super::code::secret().ok()?; | |
| 287 | let relay = self.relay.as_deref(); | |
| 288 | *pairing = pair( | |
| 289 | &self.me, | |
| 290 | &sharing, | |
| 291 | &self.notebook, | |
| 292 | self.reach, | |
| 293 | relay, | |
| 294 | &self.events, | |
| 295 | ) | |
| 296 | .ok()?; | |
| 297 | } | |
| 298 | let code = pairing.code(); | |
| 299 | if let Some(code) = &code { | |
| 300 | self.sharing.lock().unwrap().code = code.clone(); | |
| 301 | } | |
| 302 | code | |
| 303 | } | |
| 304 | ||
| 305 | /// The share as it stands, to take up again after a relaunch. | |
| 306 | pub fn sharing(&self) -> Sharing { | |
| 307 | self.code(); | |
| 308 | self.sharing.lock().unwrap().clone() | |
| 309 | } | |
| 310 | ||
| 311 | /// How the relay last answered the code's room. | |
| 312 | pub fn relayed(&self) -> Relayed { | |
| 313 | let pairing = self.pairing.lock().unwrap(); | |
| 314 | pairing.as_ref().map_or(Relayed::Unknown, Live::relayed) | |
| 315 | } | |
| 316 | ||
| 317 | /// The peers in the share's room. | |
| 318 | pub fn guests(&self) -> Vec<Peer> { | |
| 319 | let room = self.room.lock().unwrap(); | |
| 320 | room.as_ref().map(Live::peers).unwrap_or_default() | |
| 321 | } | |
| 322 | ||
| 323 | pub fn set_presence(&self, presence: Presence) { | |
| 324 | if let Some(room) = &*self.room.lock().unwrap() { | |
| 325 | room.set_presence(presence); | |
| 326 | } | |
| 327 | } | |
| 328 | ||
| 329 | /// Has `listener` hear the catalog paths guests change from now on, sooner than a watch | |
| 330 | /// on the notebook's folder would. | |
| 331 | pub fn on_changed(&self, listener: crate::session::Listener) { | |
| 332 | *self.served.host.lock().unwrap() = Some(listener); | |
| 333 | } | |
| 334 | ||
| 335 | /// Tells every guest the files at these catalog paths changed, with what changed in the | |
| 336 | /// sections a guest read lately. | |
| 337 | pub fn touched(&self, paths: &[String]) { | |
| 338 | for path in paths { | |
| 339 | self.served.changed_here(path); | |
| 340 | } | |
| 341 | self.served.tell(paths); | |
| 342 | } | |
| 343 | ||
| 344 | /// Stops sharing: no one new is welcomed, and every guest hears so and is let go. | |
| 345 | pub fn stop(&self) { | |
| 346 | drop(self.pairing.lock().unwrap().take()); | |
| 347 | let room = self.room.lock().unwrap().take(); | |
| 348 | if let Some(room) = room { | |
| 349 | room.leave(STOPPED); | |
| 350 | } | |
| 351 | } | |
| 352 | } | |
| 353 | ||
| 354 | /// The code's room, welcoming whoever knows the code to the share. | |
| 355 | fn pair( | |
| 356 | me: &Hello, | |
| 357 | sharing: &Sharing, | |
| 358 | notebook: &str, | |
| 359 | reach: Option<Reach>, | |
| 360 | relay: Option<&str>, | |
| 361 | events: &Arc<dyn Fn() + Send + Sync>, | |
| 362 | ) -> io::Result<Live> { | |
| 363 | let welcome = Welcome { | |
| 364 | share: sharing.share, | |
| 365 | secret: sharing.secret, | |
| 366 | notebook: notebook.to_owned(), | |
| 367 | host: me.name.clone(), | |
| 368 | }; | |
| 369 | let told = Arc::clone(events); | |
| 370 | Live::start( | |
| 371 | Hello { | |
| 372 | serves: None, | |
| 373 | ..me.clone() | |
| 374 | }, | |
| 375 | &Room::share(&sharing.code, &sharing.password), | |
| 376 | reach, | |
| 377 | relay, | |
| 378 | move |event| match event { | |
| 379 | Event::Met(_, line) => { | |
| 380 | let _ = line.send(kind::WELCOME, &welcome); | |
| 381 | } | |
| 382 | Event::Changed => told(), | |
| 383 | _ => {} | |
| 384 | }, | |
| 385 | ) | |
| 386 | } | |
| 387 | ||
| 388 | 264 | /// A host's side of its guests' storage requests. |
| 389 | 265 | struct Served { |
| 390 | 266 | storage: Box<dyn Storage>, |
| ... | ... | @@ -399,7 +275,7 @@ struct Served { |
| 399 | 275 | /// Hears the paths guests changed, as the host's own notebook should. |
| 400 | 276 | host: Mutex<Option<crate::session::Listener>>, |
| 401 | 277 | /// The share's room, to tell guests what changed. |
| 402 | room: OnceLock<Sender>, | |
| 278 | room: Mutex<Option<Sender>>, | |
| 403 | 279 | } |
| 404 | 280 | |
| 405 | 281 | /// A guest as its host serves it: the line to it, its requests waiting for its workers, how |
| ... | ... | @@ -515,7 +391,7 @@ impl Served { |
| 515 | 391 | .map(|path| self.storage.stamp(path).ok().map(|stamp| digest(&stamp))) |
| 516 | 392 | .collect(), |
| 517 | 393 | }; |
| 518 | if let Some(room) = self.room.get() { | |
| 394 | if let Some(room) = self.room.lock().unwrap().as_ref() { | |
| 519 | 395 | room.send(kind::TOUCHED, &touched, None); |
| 520 | 396 | } |
| 521 | 397 | } |
| ... | ... | @@ -538,6 +414,15 @@ impl Served { |
| 538 | 414 | } |
| 539 | 415 | }; |
| 540 | 416 | let id = request.id; |
| 417 | if !self.guests.lock().unwrap().contains_key(peer) { | |
| 418 | return failed( | |
| 419 | id, | |
| 420 | refused( | |
| 421 | io::ErrorKind::PermissionDenied, | |
| 422 | "The device is no longer connected", | |
| 423 | ), | |
| 424 | ); | |
| 425 | } | |
| 541 | 426 | match self.answer(peer, kind, request) { |
| 542 | 427 | Ok(reply) => Reply { id, ..reply }, |
| 543 | 428 | Err(error) => failed(id, error), |
| ... | ... | @@ -864,7 +749,7 @@ impl Served { |
| 864 | 749 | }) |
| 865 | 750 | .collect(); |
| 866 | 751 | drop(guests); |
| 867 | if let (Some(room), false) = (self.room.get(), holding.is_empty()) { | |
| 752 | if let (Some(room), false) = (self.room.lock().unwrap().as_ref(), holding.is_empty()) { | |
| 868 | 753 | room.send(kind::DELTA, &delta, Some(&holding)); |
| 869 | 754 | } |
| 870 | 755 | } |
| ... | ... | @@ -1070,6 +955,12 @@ pub struct Guest { |
| 1070 | 955 | |
| 1071 | 956 | struct Inner { |
| 1072 | 957 | share: [u8; 16], |
| 958 | me: Hello, | |
| 959 | reach: Option<Reach>, | |
| 960 | relay: Option<String>, | |
| 961 | events: Arc<dyn Fn() + Send + Sync>, | |
| 962 | presence: Mutex<Option<([u8; 16], Live)>>, | |
| 963 | here: Mutex<Presence>, | |
| 1073 | 964 | /// The host and the line to it, while connected. |
| 1074 | 965 | host: Mutex<Option<(Arc<Hello>, Line)>>, |
| 1075 | 966 | pending: Mutex<HashMap<u64, mpsc::Sender<Reply>>>, |
| ... | ... | @@ -1077,7 +968,7 @@ struct Inner { |
| 1077 | 968 | /// Where the host's reports of changed files go, while a background watches. |
| 1078 | 969 | watch: Mutex<Option<Reports>>, |
| 1079 | 970 | /// The host stopped sharing. |
| 1080 | stopped: AtomicBool, | |
| 971 | ended: Mutex<Option<Ended>>, | |
| 1081 | 972 | /// Bytes of chunks asked for and not yet given. |
| 1082 | 973 | asked: (Mutex<usize>, Condvar), |
| 1083 | 974 | /// The sections read lately, kept as the host's deltas change them, newest first. |
| ... | ... | @@ -1087,6 +978,11 @@ struct Inner { |
| 1087 | 978 | current: Mutex<HashMap<String, Stamp>>, |
| 1088 | 979 | } |
| 1089 | 980 | |
| 981 | enum Ended { | |
| 982 | Stopped, | |
| 983 | Removed, | |
| 984 | } | |
| 985 | ||
| 1090 | 986 | impl Guest { |
| 1091 | 987 | /// Joins share `share` through its `secret` as `me`, where `reach` and `relay` say. |
| 1092 | 988 | /// `events` runs on a network thread whenever the host comes or goes, or the peers change. |
| ... | ... | @@ -1100,11 +996,17 @@ impl Guest { |
| 1100 | 996 | ) -> io::Result<Arc<Self>> { |
| 1101 | 997 | let inner = Arc::new(Inner { |
| 1102 | 998 | share, |
| 999 | me: me.clone(), | |
| 1000 | reach, | |
| 1001 | relay: relay.map(str::to_owned), | |
| 1002 | events: Arc::new(events), | |
| 1003 | presence: Mutex::default(), | |
| 1004 | here: Mutex::default(), | |
| 1103 | 1005 | host: Mutex::default(), |
| 1104 | 1006 | pending: Mutex::default(), |
| 1105 | 1007 | next: AtomicU64::new(1), |
| 1106 | 1008 | watch: Mutex::default(), |
| 1107 | stopped: AtomicBool::new(false), | |
| 1009 | ended: Mutex::default(), | |
| 1108 | 1010 | asked: Default::default(), |
| 1109 | 1011 | images: Mutex::default(), |
| 1110 | 1012 | current: Mutex::default(), |
| ... | ... | @@ -1112,7 +1014,7 @@ impl Guest { |
| 1112 | 1014 | let heard = Arc::clone(&inner); |
| 1113 | 1015 | let live = Live::start(me, &Room::Notebook(secret), reach, relay, move |event| { |
| 1114 | 1016 | heard.heard(event); |
| 1115 | events(); | |
| 1017 | (heard.events)(); | |
| 1116 | 1018 | })?; |
| 1117 | 1019 | Ok(Arc::new(Self { live, inner })) |
| 1118 | 1020 | } |
| ... | ... | @@ -1130,16 +1032,25 @@ impl Guest { |
| 1130 | 1032 | |
| 1131 | 1033 | /// Whether the host said it stopped sharing. |
| 1132 | 1034 | pub fn stopped(&self) -> bool { |
| 1133 | self.inner.stopped.load(Ordering::Acquire) | |
| 1035 | self.inner.ended.lock().unwrap().is_some() | |
| 1134 | 1036 | } |
| 1135 | 1037 | |
| 1136 | 1038 | /// Everyone in the share's room: the host and the other guests. |
| 1137 | 1039 | pub fn peers(&self) -> Vec<Peer> { |
| 1138 | self.live.peers() | |
| 1040 | self.inner | |
| 1041 | .presence | |
| 1042 | .lock() | |
| 1043 | .unwrap() | |
| 1044 | .as_ref() | |
| 1045 | .map(|(_, live)| live.peers()) | |
| 1046 | .unwrap_or_default() | |
| 1139 | 1047 | } |
| 1140 | 1048 | |
| 1141 | 1049 | pub fn set_presence(&self, presence: Presence) { |
| 1142 | self.live.set_presence(presence); | |
| 1050 | *self.inner.here.lock().unwrap() = presence.clone(); | |
| 1051 | if let Some((_, live)) = &*self.inner.presence.lock().unwrap() { | |
| 1052 | live.set_presence(presence); | |
| 1053 | } | |
| 1143 | 1054 | } |
| 1144 | 1055 | |
| 1145 | 1056 | /// Sends the host's reports of changed files to `reports` while it stays connected. |
| ... | ... | @@ -1152,10 +1063,15 @@ impl Guest { |
| 1152 | 1063 | } |
| 1153 | 1064 | |
| 1154 | 1065 | fn offline(&self) -> io::Error { |
| 1155 | let message = if self.stopped() { | |
| 1156 | "The host stopped sharing this notebook" | |
| 1157 | } else { | |
| 1158 | "The computer sharing this notebook can’t be reached" | |
| 1066 | if let Some(version) = self.live.other_version() { | |
| 1067 | return io::Error::new(io::ErrorKind::NotConnected, wire::Version(version)); | |
| 1068 | } | |
| 1069 | let message = match &*self.inner.ended.lock().unwrap() { | |
| 1070 | Some(Ended::Stopped) => "The host stopped sharing this notebook", | |
| 1071 | Some(Ended::Removed) => { | |
| 1072 | "This device was removed. Ask the person sharing for a new link or code." | |
| 1073 | } | |
| 1074 | None => "The computer sharing this notebook can’t be reached", | |
| 1159 | 1075 | }; |
| 1160 | 1076 | io::Error::new(io::ErrorKind::NotConnected, message) |
| 1161 | 1077 | } |
| ... | ... | @@ -1420,14 +1336,23 @@ impl Inner { |
| 1420 | 1336 | } |
| 1421 | 1337 | } |
| 1422 | 1338 | |
| 1423 | fn heard(&self, event: Event) { | |
| 1424 | let serves = |hello: &Hello| hello.serves == Some(self.share); | |
| 1339 | fn heard(self: &Arc<Self>, event: Event) { | |
| 1340 | let serves = |hello: &Hello| { | |
| 1341 | hello.serves == Some(self.share) | |
| 1342 | || self | |
| 1343 | .host | |
| 1344 | .lock() | |
| 1345 | .unwrap() | |
| 1346 | .as_ref() | |
| 1347 | .is_some_and(|(host, _)| host.peer == hello.peer) | |
| 1348 | }; | |
| 1425 | 1349 | match event { |
| 1426 | 1350 | Event::Met(hello, line) if serves(hello) => { |
| 1427 | 1351 | *self.host.lock().unwrap() = Some((Arc::clone(hello), line.clone())); |
| 1428 | 1352 | } |
| 1429 | 1353 | Event::Left(hello) if serves(hello) => { |
| 1430 | 1354 | *self.host.lock().unwrap() = None; |
| 1355 | drop(self.presence.lock().unwrap().take()); | |
| 1431 | 1356 | self.current.lock().unwrap().clear(); |
| 1432 | 1357 | // Each request waiting hears its answer was lost. |
| 1433 | 1358 | self.pending.lock().unwrap().clear(); |
| ... | ... | @@ -1438,6 +1363,11 @@ impl Inner { |
| 1438 | 1363 | Event::Frame { |
| 1439 | 1364 | from, kind, body, .. |
| 1440 | 1365 | } => match kind { |
| 1366 | kind::WELCOME if serves(from) => { | |
| 1367 | if let Ok(welcome) = minicbor::decode::<Welcome>(body) { | |
| 1368 | self.join_presence(welcome.room); | |
| 1369 | } | |
| 1370 | } | |
| 1441 | 1371 | kind::REPLY if serves(from) => { |
| 1442 | 1372 | if let Ok(reply) = minicbor::decode::<Reply>(body) |
| 1443 | 1373 | && let Some(waiting) = self.pending.lock().unwrap().remove(&reply.id) |
| ... | ... | @@ -1458,18 +1388,59 @@ impl Inner { |
| 1458 | 1388 | } |
| 1459 | 1389 | } |
| 1460 | 1390 | } |
| 1461 | kind::BYE | |
| 1462 | if serves(from) | |
| 1463 | && minicbor::decode::<wire::Bye>(body) | |
| 1464 | .is_ok_and(|bye| bye.reason == STOPPED) => | |
| 1465 | { | |
| 1466 | self.stopped.store(true, Ordering::Release); | |
| 1391 | kind::BYE if serves(from) => { | |
| 1392 | if let Ok(bye) = minicbor::decode::<wire::Bye>(body) { | |
| 1393 | let ended = match bye.reason.as_str() { | |
| 1394 | STOPPED => Some(Ended::Stopped), | |
| 1395 | "removed" => Some(Ended::Removed), | |
| 1396 | _ => None, | |
| 1397 | }; | |
| 1398 | if ended.is_some() { | |
| 1399 | *self.ended.lock().unwrap() = ended; | |
| 1400 | } | |
| 1401 | } | |
| 1467 | 1402 | } |
| 1468 | 1403 | _ => {} |
| 1469 | 1404 | }, |
| 1470 | 1405 | _ => {} |
| 1471 | 1406 | } |
| 1472 | 1407 | } |
| 1408 | ||
| 1409 | fn join_presence(self: &Arc<Self>, secret: [u8; 16]) { | |
| 1410 | let mut presence = self.presence.lock().unwrap(); | |
| 1411 | if presence.as_ref().is_some_and(|(held, _)| *held == secret) { | |
| 1412 | return; | |
| 1413 | } | |
| 1414 | let inner = Arc::downgrade(self); | |
| 1415 | let live = Live::start( | |
| 1416 | self.me.clone(), | |
| 1417 | &Room::Notebook(secret), | |
| 1418 | self.reach, | |
| 1419 | self.relay.as_deref(), | |
| 1420 | move |event| { | |
| 1421 | if let Some(inner) = inner.upgrade() { | |
| 1422 | if matches!(event, Event::Frame { .. }) { | |
| 1423 | inner.heard(event); | |
| 1424 | } | |
| 1425 | (inner.events)(); | |
| 1426 | } | |
| 1427 | }, | |
| 1428 | ); | |
| 1429 | match live { | |
| 1430 | Ok(live) => { | |
| 1431 | live.set_presence(self.here.lock().unwrap().clone()); | |
| 1432 | *presence = Some((secret, live)); | |
| 1433 | let mut current = self.current.lock().unwrap(); | |
| 1434 | let paths: Vec<String> = current.keys().cloned().collect(); | |
| 1435 | current.clear(); | |
| 1436 | drop(current); | |
| 1437 | if let Some(reports) = &*self.watch.lock().unwrap() { | |
| 1438 | reports.touched(&paths); | |
| 1439 | } | |
| 1440 | } | |
| 1441 | Err(error) => eprintln!("Live Share: could not join presence: {error}"), | |
| 1442 | } | |
| 1443 | } | |
| 1473 | 1444 | } |
| 1474 | 1445 | |
| 1475 | 1446 | /// A section a Live Share host serves, as a guest's replica publishes to it. |
crates/notebook/src/live/share/batch.rs+6| ... | ... | @@ -148,6 +148,12 @@ impl Served { |
| 148 | 148 | image: &[u8], |
| 149 | 149 | batch: &Waiting, |
| 150 | 150 | ) -> Result<Option<Transaction>> { |
| 151 | if !self.guests.lock().unwrap().contains_key(&batch.peer) { | |
| 152 | return Err(refused( | |
| 153 | io::ErrorKind::PermissionDenied, | |
| 154 | "The device is no longer connected", | |
| 155 | )); | |
| 156 | } | |
| 151 | 157 | let stamp: Stamp = batch |
| 152 | 158 | .request |
| 153 | 159 | .stamp |
crates/notebook/src/live/share/membership.rs created+391| ... | ... | @@ -0,0 +1,391 @@ |
| 1 | use super::*; | |
| 2 | ||
| 3 | pub struct Host { | |
| 4 | members: Arc<Members>, | |
| 5 | pairing: Mutex<Option<Live>>, | |
| 6 | } | |
| 7 | ||
| 8 | #[allow(clippy::type_complexity)] | |
| 9 | struct Members { | |
| 10 | me: Hello, | |
| 11 | notebook: String, | |
| 12 | reach: Option<Reach>, | |
| 13 | relay: Option<String>, | |
| 14 | sharing: Mutex<Sharing>, | |
| 15 | changes: Mutex<()>, | |
| 16 | pending: Mutex<BTreeMap<[u8; 16], (Arc<Hello>, Line)>>, | |
| 17 | access: Mutex<HashMap<[u8; 16], Live>>, | |
| 18 | room: Mutex<Option<Live>>, | |
| 19 | served: Arc<Served>, | |
| 20 | events: Arc<dyn Fn() + Send + Sync>, | |
| 21 | keep: Box<dyn Fn(&Sharing) -> io::Result<()> + Send + Sync>, | |
| 22 | } | |
| 23 | ||
| 24 | impl Host { | |
| 25 | /// Shares a notebook, keeping every credential change before granting or revoking access. | |
| 26 | #[allow(clippy::too_many_arguments)] | |
| 27 | pub fn start( | |
| 28 | storage: Box<dyn Storage>, | |
| 29 | me: Hello, | |
| 30 | sharing: Sharing, | |
| 31 | notebook: &str, | |
| 32 | reach: Option<Reach>, | |
| 33 | relay: Option<&str>, | |
| 34 | events: impl Fn() + Send + Sync + 'static, | |
| 35 | keep: impl Fn(&Sharing) -> io::Result<()> + Send + Sync + 'static, | |
| 36 | ) -> io::Result<Self> { | |
| 37 | keep(&sharing)?; | |
| 38 | let members = Arc::new(Members { | |
| 39 | me, | |
| 40 | notebook: notebook.to_owned(), | |
| 41 | reach, | |
| 42 | relay: relay.map(str::to_owned), | |
| 43 | sharing: Mutex::new(sharing), | |
| 44 | changes: Mutex::default(), | |
| 45 | pending: Mutex::default(), | |
| 46 | access: Mutex::default(), | |
| 47 | room: Mutex::default(), | |
| 48 | served: Arc::new(Served { | |
| 49 | storage, | |
| 50 | images: Mutex::default(), | |
| 51 | snapshots: Mutex::default(), | |
| 52 | puts: Mutex::default(), | |
| 53 | guests: Mutex::default(), | |
| 54 | writers: Mutex::default(), | |
| 55 | host: Mutex::default(), | |
| 56 | room: Mutex::default(), | |
| 57 | }), | |
| 58 | events: Arc::new(events), | |
| 59 | keep: Box::new(keep), | |
| 60 | }); | |
| 61 | let sharing = members.sharing.lock().unwrap().clone(); | |
| 62 | members.presence(sharing.secret)?; | |
| 63 | for device in &sharing.members { | |
| 64 | members.open(device.secret)?; | |
| 65 | } | |
| 66 | let pairing = Mutex::new(Some(members.pair(&sharing)?)); | |
| 67 | Ok(Self { members, pairing }) | |
| 68 | } | |
| 69 | ||
| 70 | pub fn code(&self) -> Option<String> { | |
| 71 | let mut pairing = self.pairing.lock().unwrap(); | |
| 72 | let pairing = pairing.as_mut()?; | |
| 73 | let mut sharing = self.members.sharing.lock().unwrap(); | |
| 74 | if pairing.burned() { | |
| 75 | let mut next = sharing.clone(); | |
| 76 | next.code = super::super::code::secret().ok()?; | |
| 77 | (self.members.keep)(&next).ok()?; | |
| 78 | *pairing = self.members.pair(&next).ok()?; | |
| 79 | *sharing = next; | |
| 80 | } | |
| 81 | let code = pairing.code(); | |
| 82 | if let Some(code) = &code | |
| 83 | && *code != sharing.code | |
| 84 | { | |
| 85 | let mut next = sharing.clone(); | |
| 86 | next.code = code.clone(); | |
| 87 | (self.members.keep)(&next).ok()?; | |
| 88 | *sharing = next; | |
| 89 | } | |
| 90 | code | |
| 91 | } | |
| 92 | ||
| 93 | pub fn sharing(&self) -> Sharing { | |
| 94 | self.code(); | |
| 95 | self.members.sharing.lock().unwrap().clone() | |
| 96 | } | |
| 97 | ||
| 98 | pub fn relayed(&self) -> Relayed { | |
| 99 | self.pairing | |
| 100 | .lock() | |
| 101 | .unwrap() | |
| 102 | .as_ref() | |
| 103 | .map_or(Relayed::Unknown, Live::relayed) | |
| 104 | } | |
| 105 | ||
| 106 | pub fn guests(&self) -> Vec<Peer> { | |
| 107 | self.members | |
| 108 | .room | |
| 109 | .lock() | |
| 110 | .unwrap() | |
| 111 | .as_ref() | |
| 112 | .map(Live::peers) | |
| 113 | .unwrap_or_default() | |
| 114 | } | |
| 115 | ||
| 116 | pub fn devices(&self) -> Vec<(Device, bool)> { | |
| 117 | let sharing = self.members.sharing.lock().unwrap(); | |
| 118 | let access = self.members.access.lock().unwrap(); | |
| 119 | sharing | |
| 120 | .members | |
| 121 | .iter() | |
| 122 | .map(|device| { | |
| 123 | let connected = access | |
| 124 | .get(&device.secret) | |
| 125 | .is_some_and(|live| !live.peers().is_empty()); | |
| 126 | (device.clone(), connected) | |
| 127 | }) | |
| 128 | .collect() | |
| 129 | } | |
| 130 | ||
| 131 | pub fn requests(&self) -> Vec<Arc<Hello>> { | |
| 132 | self.members | |
| 133 | .pending | |
| 134 | .lock() | |
| 135 | .unwrap() | |
| 136 | .values() | |
| 137 | .map(|(hello, _)| Arc::clone(hello)) | |
| 138 | .collect() | |
| 139 | } | |
| 140 | ||
| 141 | pub fn approve(&self, approve: bool) -> io::Result<()> { | |
| 142 | let mut sharing = self.members.sharing.lock().unwrap(); | |
| 143 | let mut next = sharing.clone(); | |
| 144 | next.approve = approve; | |
| 145 | (self.members.keep)(&next)?; | |
| 146 | *sharing = next; | |
| 147 | drop(sharing); | |
| 148 | if !approve { | |
| 149 | for hello in self.requests() { | |
| 150 | self.allow(&hello.peer)?; | |
| 151 | } | |
| 152 | } | |
| 153 | (self.members.events)(); | |
| 154 | Ok(()) | |
| 155 | } | |
| 156 | ||
| 157 | pub fn allow(&self, peer: &[u8; 16]) -> io::Result<()> { | |
| 158 | let pending = self.members.pending.lock().unwrap().remove(peer); | |
| 159 | if let Some((hello, line)) = pending { | |
| 160 | if let Err(error) = self.members.grant(&hello, &line) { | |
| 161 | self.members | |
| 162 | .pending | |
| 163 | .lock() | |
| 164 | .unwrap() | |
| 165 | .insert(*peer, (hello, line)); | |
| 166 | return Err(error); | |
| 167 | } | |
| 168 | (self.members.events)(); | |
| 169 | } | |
| 170 | Ok(()) | |
| 171 | } | |
| 172 | ||
| 173 | pub fn decline(&self, peer: &[u8; 16]) { | |
| 174 | if let Some((_, line)) = self.members.pending.lock().unwrap().remove(peer) { | |
| 175 | let _ = line.send(kind::APPROVAL, &wire::Approval::Declined); | |
| 176 | (self.members.events)(); | |
| 177 | } | |
| 178 | } | |
| 179 | ||
| 180 | pub fn remove(&self, secret: &[u8; 16]) -> io::Result<()> { | |
| 181 | let _change = self.members.changes.lock().unwrap(); | |
| 182 | let mut sharing = self.members.sharing.lock().unwrap(); | |
| 183 | if !sharing | |
| 184 | .members | |
| 185 | .iter() | |
| 186 | .any(|device| device.secret == *secret) | |
| 187 | { | |
| 188 | return Ok(()); | |
| 189 | } | |
| 190 | let mut next = sharing.clone(); | |
| 191 | next.members.retain(|device| device.secret != *secret); | |
| 192 | getrandom::fill(&mut next.secret) | |
| 193 | .map_err(|_| io::Error::other("System random source failed"))?; | |
| 194 | next.code = super::super::code::secret()?; | |
| 195 | (self.members.keep)(&next)?; | |
| 196 | *sharing = next.clone(); | |
| 197 | drop(sharing); | |
| 198 | let removed = self.members.access.lock().unwrap().remove(secret); | |
| 199 | if let Some(removed) = removed { | |
| 200 | for peer in removed.peers() { | |
| 201 | self.members.served.forget(&peer.hello.peer); | |
| 202 | } | |
| 203 | removed.leave("removed"); | |
| 204 | } | |
| 205 | self.members.presence(next.secret)?; | |
| 206 | *self.pairing.lock().unwrap() = Some(self.members.pair(&next)?); | |
| 207 | let access = self.members.access.lock().unwrap(); | |
| 208 | for device in &next.members { | |
| 209 | if let Some(live) = access.get(&device.secret) { | |
| 210 | live.sender().send( | |
| 211 | kind::WELCOME, | |
| 212 | &self.members.welcome(&next, device.secret), | |
| 213 | None, | |
| 214 | ); | |
| 215 | } | |
| 216 | } | |
| 217 | (self.members.events)(); | |
| 218 | Ok(()) | |
| 219 | } | |
| 220 | ||
| 221 | pub fn set_presence(&self, presence: Presence) { | |
| 222 | if let Some(room) = &*self.members.room.lock().unwrap() { | |
| 223 | room.set_presence(presence); | |
| 224 | } | |
| 225 | } | |
| 226 | ||
| 227 | pub fn on_changed(&self, listener: crate::session::Listener) { | |
| 228 | *self.members.served.host.lock().unwrap() = Some(listener); | |
| 229 | } | |
| 230 | ||
| 231 | pub fn touched(&self, paths: &[String]) { | |
| 232 | for path in paths { | |
| 233 | self.members.served.changed_here(path); | |
| 234 | } | |
| 235 | self.members.served.tell(paths); | |
| 236 | } | |
| 237 | ||
| 238 | pub fn stop(&self) { | |
| 239 | let _change = self.members.changes.lock().unwrap(); | |
| 240 | drop(self.pairing.lock().unwrap().take()); | |
| 241 | self.members.pending.lock().unwrap().clear(); | |
| 242 | let access = std::mem::take(&mut *self.members.access.lock().unwrap()); | |
| 243 | for (_, live) in access { | |
| 244 | live.leave(STOPPED); | |
| 245 | } | |
| 246 | drop(self.members.room.lock().unwrap().take()); | |
| 247 | } | |
| 248 | } | |
| 249 | ||
| 250 | impl Members { | |
| 251 | fn welcome(&self, sharing: &Sharing, secret: [u8; 16]) -> Welcome { | |
| 252 | Welcome { | |
| 253 | share: sharing.share, | |
| 254 | secret, | |
| 255 | room: sharing.secret, | |
| 256 | notebook: self.notebook.clone(), | |
| 257 | host: self.me.name.clone(), | |
| 258 | } | |
| 259 | } | |
| 260 | ||
| 261 | fn presence(self: &Arc<Self>, secret: [u8; 16]) -> io::Result<()> { | |
| 262 | let told = Arc::clone(&self.events); | |
| 263 | let live = Live::start( | |
| 264 | self.me.clone(), | |
| 265 | &Room::Notebook(secret), | |
| 266 | self.reach, | |
| 267 | self.relay.as_deref(), | |
| 268 | move |_| told(), | |
| 269 | )?; | |
| 270 | *self.served.room.lock().unwrap() = Some(live.sender()); | |
| 271 | *self.room.lock().unwrap() = Some(live); | |
| 272 | Ok(()) | |
| 273 | } | |
| 274 | ||
| 275 | fn open(self: &Arc<Self>, secret: [u8; 16]) -> io::Result<()> { | |
| 276 | let members = Arc::downgrade(self); | |
| 277 | let live = Live::start( | |
| 278 | Hello { | |
| 279 | serves: Some(self.sharing.lock().unwrap().share), | |
| 280 | ..self.me.clone() | |
| 281 | }, | |
| 282 | &Room::Notebook(secret), | |
| 283 | self.reach, | |
| 284 | self.relay.as_deref(), | |
| 285 | move |event| { | |
| 286 | let Some(members) = members.upgrade() else { | |
| 287 | return; | |
| 288 | }; | |
| 289 | match event { | |
| 290 | Event::Met(hello, line) => { | |
| 291 | let sharing = members.sharing.lock().unwrap(); | |
| 292 | if !sharing.members.iter().any(|device| device.secret == secret) { | |
| 293 | line.hang_up("removed"); | |
| 294 | return; | |
| 295 | } | |
| 296 | members.served.admit(hello.peer, line); | |
| 297 | let _ = line.send(kind::WELCOME, &members.welcome(&sharing, secret)); | |
| 298 | } | |
| 299 | Event::Left(hello) => members.served.forget(&hello.peer), | |
| 300 | Event::Frame { from, kind, body } | |
| 301 | if wire::KNOWN.contains(&kind) && kind > 256 && kind != kind::REPLY => | |
| 302 | { | |
| 303 | members.served.queue(&from.peer, kind, body) | |
| 304 | } | |
| 305 | Event::Changed => (members.events)(), | |
| 306 | _ => {} | |
| 307 | } | |
| 308 | }, | |
| 309 | )?; | |
| 310 | self.access.lock().unwrap().insert(secret, live); | |
| 311 | Ok(()) | |
| 312 | } | |
| 313 | ||
| 314 | fn grant(self: &Arc<Self>, hello: &Hello, line: &Line) -> io::Result<()> { | |
| 315 | let _change = self.changes.lock().unwrap(); | |
| 316 | if self.room.lock().unwrap().is_none() { | |
| 317 | return Err(io::ErrorKind::NotConnected.into()); | |
| 318 | } | |
| 319 | let mut secret = [0; 16]; | |
| 320 | getrandom::fill(&mut secret) | |
| 321 | .map_err(|_| io::Error::other("System random source failed"))?; | |
| 322 | if self.sharing.lock().unwrap().members.len() >= 64 { | |
| 323 | return Err(io::ErrorKind::ResourceBusy.into()); | |
| 324 | } | |
| 325 | self.open(secret)?; | |
| 326 | let kept = (|| -> io::Result<Welcome> { | |
| 327 | let mut sharing = self.sharing.lock().unwrap(); | |
| 328 | let mut next = sharing.clone(); | |
| 329 | next.members.push(Device { | |
| 330 | secret, | |
| 331 | name: hello.name.clone(), | |
| 332 | device: hello.device.clone(), | |
| 333 | }); | |
| 334 | (self.keep)(&next)?; | |
| 335 | *sharing = next; | |
| 336 | Ok(self.welcome(&sharing, secret)) | |
| 337 | })(); | |
| 338 | let welcome = match kept { | |
| 339 | Ok(welcome) => welcome, | |
| 340 | Err(error) => { | |
| 341 | drop(self.access.lock().unwrap().remove(&secret)); | |
| 342 | return Err(error); | |
| 343 | } | |
| 344 | }; | |
| 345 | line.send(kind::WELCOME, &welcome)?; | |
| 346 | (self.events)(); | |
| 347 | Ok(()) | |
| 348 | } | |
| 349 | ||
| 350 | fn pair(self: &Arc<Self>, sharing: &Sharing) -> io::Result<Live> { | |
| 351 | let members = Arc::downgrade(self); | |
| 352 | Live::start( | |
| 353 | Hello { | |
| 354 | serves: None, | |
| 355 | ..self.me.clone() | |
| 356 | }, | |
| 357 | &Room::share(&sharing.code, &sharing.password), | |
| 358 | self.reach, | |
| 359 | self.relay.as_deref(), | |
| 360 | move |event| { | |
| 361 | let Some(members) = members.upgrade() else { | |
| 362 | return; | |
| 363 | }; | |
| 364 | match event { | |
| 365 | Event::Met(hello, line) => { | |
| 366 | if members.sharing.lock().unwrap().approve { | |
| 367 | let mut pending = members.pending.lock().unwrap(); | |
| 368 | if pending.len() < 32 { | |
| 369 | pending.insert(hello.peer, (Arc::clone(hello), line.clone())); | |
| 370 | let _ = line.send(kind::APPROVAL, &wire::Approval::Pending); | |
| 371 | } else { | |
| 372 | let _ = line.send(kind::APPROVAL, &wire::Approval::Failed); | |
| 373 | } | |
| 374 | drop(pending); | |
| 375 | (members.events)(); | |
| 376 | } else if let Err(error) = members.grant(hello, line) { | |
| 377 | eprintln!("Live Share: could not admit a device: {error}"); | |
| 378 | let _ = line.send(kind::APPROVAL, &wire::Approval::Failed); | |
| 379 | } | |
| 380 | } | |
| 381 | Event::Left(hello) => { | |
| 382 | members.pending.lock().unwrap().remove(&hello.peer); | |
| 383 | (members.events)(); | |
| 384 | } | |
| 385 | Event::Changed => (members.events)(), | |
| 386 | _ => {} | |
| 387 | } | |
| 388 | }, | |
| 389 | ) | |
| 390 | } | |
| 391 | } |
crates/notebook/src/live/wire.rs+22-6| ... | ... | @@ -1,7 +1,6 @@ |
| 1 | 1 | //! What peers say to each other: an opening in the clear that meets through the secret both |
| 2 | 2 | //! hold (SPAKE2), then frames sealed under the keys it agreed, each a message kind and a CBOR |
| 3 | //! map. A later version adds kinds and fields; a reader skips the kinds and fields it doesn't | |
| 4 | //! know, so every version speaks to every other. | |
| 3 | //! map. Readers skip unknown kinds and fields within the same opening version. | |
| 5 | 4 | |
| 6 | 5 | use aes_gcm::{Aes256Gcm, KeyInit, aead::Aead}; |
| 7 | 6 | use hmac::{Hmac, Mac}; |
| ... | ... | @@ -10,8 +9,8 @@ use sha2::Sha256; |
| 10 | 9 | use spake2::{Ed25519Group, Identity, Password, Spake2}; |
| 11 | 10 | use std::io::{self, Read, Write}; |
| 12 | 11 | |
| 13 | /// The opening's version. Frames after it never change shape; they grow by kinds and fields. | |
| 14 | pub const VERSION: u16 = 2; | |
| 12 | /// The opening's version; peers must agree on its authentication and access rules. | |
| 13 | pub const VERSION: u16 = 3; | |
| 15 | 14 | |
| 16 | 15 | /// The version of the opening a peer of another version sent, as `open` fails with it. |
| 17 | 16 | #[derive(Clone, Copy, Debug, PartialEq, Eq)] |
| ... | ... | @@ -45,6 +44,7 @@ pub mod kind { |
| 45 | 44 | /// A section's bytes as a commit changed them: `Delta`. |
| 46 | 45 | pub const DELTA: u16 = 19; |
| 47 | 46 | pub const WELCOME: u16 = 32; |
| 47 | pub const APPROVAL: u16 = 33; | |
| 48 | 48 | /// Storage requests to a host, each a `Request` answered by a `Reply`. |
| 49 | 49 | pub const LIST: u16 = 257; |
| 50 | 50 | pub const STAMP: u16 = 258; |
| ... | ... | @@ -79,6 +79,7 @@ pub const KNOWN: &[u16] = &[ |
| 79 | 79 | kind::TOUCHED, |
| 80 | 80 | kind::DELTA, |
| 81 | 81 | kind::WELCOME, |
| 82 | kind::APPROVAL, | |
| 82 | 83 | kind::LIST, |
| 83 | 84 | kind::STAMP, |
| 84 | 85 | kind::READ, |
| ... | ... | @@ -122,6 +123,8 @@ pub struct Hello { |
| 122 | 123 | pub serves: Option<[u8; 16]>, |
| 123 | 124 | #[n(6)] |
| 124 | 125 | pub ops: Option<u16>, |
| 126 | #[n(7)] | |
| 127 | pub device: Option<String>, | |
| 125 | 128 | } |
| 126 | 129 | |
| 127 | 130 | impl Hello { |
| ... | ... | @@ -137,6 +140,7 @@ impl Hello { |
| 137 | 140 | kinds: KNOWN.to_vec(), |
| 138 | 141 | serves: None, |
| 139 | 142 | ops: Some(1), |
| 143 | device: None, | |
| 140 | 144 | }) |
| 141 | 145 | } |
| 142 | 146 | } |
| ... | ... | @@ -211,8 +215,7 @@ pub struct Bye { |
| 211 | 215 | pub reason: String, |
| 212 | 216 | } |
| 213 | 217 | |
| 214 | /// What the host of a share gives a peer that knew its code: the share's room, and names to | |
| 215 | /// show it by. | |
| 218 | /// A device's access credential, the current presence room, and the notebook's names. | |
| 216 | 219 | #[derive(Clone, Debug, PartialEq, Encode, Decode)] |
| 217 | 220 | #[cbor(map)] |
| 218 | 221 | pub struct Welcome { |
| ... | ... | @@ -225,6 +228,19 @@ pub struct Welcome { |
| 225 | 228 | /// The host's name for itself, as `Hello::name`. |
| 226 | 229 | #[n(3)] |
| 227 | 230 | pub host: String, |
| 231 | #[cbor(n(4), with = "minicbor::bytes")] | |
| 232 | pub room: [u8; 16], | |
| 233 | } | |
| 234 | ||
| 235 | #[derive(Clone, Debug, Encode, Decode)] | |
| 236 | #[cbor(index_only)] | |
| 237 | pub enum Approval { | |
| 238 | #[n(0)] | |
| 239 | Pending, | |
| 240 | #[n(1)] | |
| 241 | Declined, | |
| 242 | #[n(2)] | |
| 243 | Failed, | |
| 228 | 244 | } |
| 229 | 245 | |
| 230 | 246 | /// Paths a host's files changed at, by catalog path; `""` is the notebook's folder. |
crates/notebook/tests/live_membership.rs created+190| ... | ... | @@ -0,0 +1,190 @@ |
| 1 | #![cfg(feature = "live")] | |
| 2 | ||
| 3 | #[path = "support/live.rs"] | |
| 4 | mod live; | |
| 5 | use live::*; | |
| 6 | use notebook::live::share::{self, Guest, Host, Sharing}; | |
| 7 | use notebook::session::Notebook; | |
| 8 | use std::{ | |
| 9 | sync::{ | |
| 10 | Arc, Mutex, | |
| 11 | atomic::{AtomicBool, Ordering}, | |
| 12 | mpsc, | |
| 13 | }, | |
| 14 | time::Duration, | |
| 15 | }; | |
| 16 | ||
| 17 | #[test] | |
| 18 | fn approval_keeps_credentials_before_welcome_and_a_decline_grants_nothing() { | |
| 19 | let directory = tempfile::tempdir().unwrap(); | |
| 20 | let folder = notebook(directory.path()); | |
| 21 | let url = relay(Default::default()); | |
| 22 | let saved = Arc::new(Mutex::new(None)); | |
| 23 | let fail = Arc::new(AtomicBool::new(false)); | |
| 24 | let (kept, failing) = (Arc::clone(&saved), Arc::clone(&fail)); | |
| 25 | let mut sharing = Sharing::new("").unwrap(); | |
| 26 | sharing.approve = true; | |
| 27 | let host = Host::start( | |
| 28 | Notebook::open(&folder, directory.path().join("host")) | |
| 29 | .unwrap() | |
| 30 | .into_storage(), | |
| 31 | hello("Ada"), | |
| 32 | sharing, | |
| 33 | "Garden", | |
| 34 | None, | |
| 35 | Some(&url), | |
| 36 | || {}, | |
| 37 | move |sharing| { | |
| 38 | if failing.load(Ordering::Acquire) { | |
| 39 | return Err(std::io::ErrorKind::PermissionDenied.into()); | |
| 40 | } | |
| 41 | *kept.lock().unwrap() = Some(serde_json::to_vec(sharing).unwrap()); | |
| 42 | Ok(()) | |
| 43 | }, | |
| 44 | ) | |
| 45 | .unwrap(); | |
| 46 | let code = self::code(&host); | |
| 47 | let (reply, result) = mpsc::channel(); | |
| 48 | let joining = url.clone(); | |
| 49 | std::thread::spawn(move || { | |
| 50 | let _ = reply.send(share::join(hello("Grace"), &code, "", None, Some(&joining))); | |
| 51 | }); | |
| 52 | until("the request never arrived", || host.requests().len() == 1); | |
| 53 | assert!(result.try_recv().is_err()); | |
| 54 | assert!(host.devices().is_empty()); | |
| 55 | let peer = host.requests()[0].peer; | |
| 56 | fail.store(true, Ordering::Release); | |
| 57 | assert!(host.allow(&peer).is_err()); | |
| 58 | assert!(host.devices().is_empty()); | |
| 59 | assert_eq!(host.requests().len(), 1); | |
| 60 | fail.store(false, Ordering::Release); | |
| 61 | host.allow(&peer).unwrap(); | |
| 62 | let welcome = result | |
| 63 | .recv_timeout(Duration::from_secs(5)) | |
| 64 | .unwrap() | |
| 65 | .unwrap(); | |
| 66 | let persisted: Sharing = | |
| 67 | serde_json::from_slice(saved.lock().unwrap().as_ref().unwrap()).unwrap(); | |
| 68 | assert_eq!(persisted.members[0].secret, welcome.secret); | |
| 69 | assert_ne!(welcome.secret, welcome.room); | |
| 70 | fail.store(true, Ordering::Release); | |
| 71 | assert!(host.remove(&welcome.secret).is_err()); | |
| 72 | assert_eq!(host.sharing(), persisted); | |
| 73 | fail.store(false, Ordering::Release); | |
| 74 | ||
| 75 | let code = self::code(&host); | |
| 76 | let (reply, result) = mpsc::channel(); | |
| 77 | std::thread::spawn(move || { | |
| 78 | let _ = reply.send(share::join(hello("Alan"), &code, "", None, Some(&url))); | |
| 79 | }); | |
| 80 | until("the second request never arrived", || { | |
| 81 | host.requests().len() == 1 | |
| 82 | }); | |
| 83 | host.decline(&host.requests()[0].peer); | |
| 84 | assert_eq!( | |
| 85 | result.recv_timeout(Duration::from_secs(5)).unwrap(), | |
| 86 | Err(share::Refusal::Declined) | |
| 87 | ); | |
| 88 | assert_eq!(host.devices().len(), 1); | |
| 89 | } | |
| 90 | ||
| 91 | #[test] | |
| 92 | fn removal_retires_one_credential_and_other_devices_reconnect_after_a_restart() { | |
| 93 | let directory = tempfile::tempdir().unwrap(); | |
| 94 | let folder = notebook(directory.path()); | |
| 95 | let url = relay(Default::default()); | |
| 96 | let sharing = Sharing::new("").unwrap(); | |
| 97 | let host_cache = directory.path().join("host"); | |
| 98 | let host = host(&folder, &host_cache, &sharing, &url); | |
| 99 | let original_code = code(&host); | |
| 100 | let (alice, _) = guest( | |
| 101 | "Alice", | |
| 102 | &original_code, | |
| 103 | &url, | |
| 104 | &directory.path().join("alice"), | |
| 105 | ); | |
| 106 | let (bob, _) = guest("Bob", &original_code, &url, &directory.path().join("bob")); | |
| 107 | let before = host.sharing(); | |
| 108 | let alice_key = before | |
| 109 | .members | |
| 110 | .iter() | |
| 111 | .find(|device| device.name == "Alice") | |
| 112 | .unwrap() | |
| 113 | .secret; | |
| 114 | let bob_key = before | |
| 115 | .members | |
| 116 | .iter() | |
| 117 | .find(|device| device.name == "Bob") | |
| 118 | .unwrap() | |
| 119 | .secret; | |
| 120 | host.remove(&bob_key).unwrap(); | |
| 121 | until("Bob was not removed", || { | |
| 122 | bob.stopped() && bob.host().is_none() | |
| 123 | }); | |
| 124 | until("presence did not move to the new room", || { | |
| 125 | host.guests().len() == 1 | |
| 126 | }); | |
| 127 | let current = host.sharing(); | |
| 128 | assert_ne!(before.secret, current.secret); | |
| 129 | assert_eq!(current.members.len(), 1); | |
| 130 | assert_ne!(original_code, code(&host)); | |
| 131 | assert!(alice.host().is_some()); | |
| 132 | assert!(!alice.stopped()); | |
| 133 | let live = Notebook::open_hosted(Arc::clone(&alice), directory.path().join("alice")).unwrap(); | |
| 134 | assert_eq!(live.catalog().sections.len(), 2); | |
| 135 | let restarted: Sharing = | |
| 136 | serde_json::from_slice(&serde_json::to_vec(&current).unwrap()).unwrap(); | |
| 137 | drop(host); | |
| 138 | until("Alice's connection did not close", || { | |
| 139 | alice.host().is_none() | |
| 140 | }); | |
| 141 | let host = self::host(&folder, &host_cache, &restarted, &url); | |
| 142 | until("Alice did not reconnect", || alice.host().is_some()); | |
| 143 | assert_eq!(host.devices()[0].0.secret, alice_key); | |
| 144 | let forged = Guest::start( | |
| 145 | hello("Alice"), | |
| 146 | current.share, | |
| 147 | bob_key, | |
| 148 | None, | |
| 149 | Some(&url), | |
| 150 | || {}, | |
| 151 | ) | |
| 152 | .unwrap(); | |
| 153 | std::thread::sleep(Duration::from_millis(500)); | |
| 154 | assert!(forged.host().is_none()); | |
| 155 | assert!(Notebook::open_hosted(forged, directory.path().join("forged")).is_err()); | |
| 156 | } | |
| 157 | ||
| 158 | #[test] | |
| 159 | fn cancelling_a_join_retires_its_pending_request() { | |
| 160 | let directory = tempfile::tempdir().unwrap(); | |
| 161 | let folder = notebook(directory.path()); | |
| 162 | let url = relay(Default::default()); | |
| 163 | let mut sharing = Sharing::new("").unwrap(); | |
| 164 | sharing.approve = true; | |
| 165 | let host = host(&folder, &directory.path().join("host"), &sharing, &url); | |
| 166 | let code = code(&host); | |
| 167 | let alive = Arc::new(AtomicBool::new(true)); | |
| 168 | let continuing = Arc::clone(&alive); | |
| 169 | let (reply, result) = mpsc::channel(); | |
| 170 | std::thread::spawn(move || { | |
| 171 | let _ = reply.send(share::join_while( | |
| 172 | hello("Grace"), | |
| 173 | &code, | |
| 174 | "", | |
| 175 | None, | |
| 176 | Some(&url), | |
| 177 | |_| continuing.load(Ordering::Acquire), | |
| 178 | )); | |
| 179 | }); | |
| 180 | until("the request never arrived", || !host.requests().is_empty()); | |
| 181 | alive.store(false, Ordering::Release); | |
| 182 | assert_eq!( | |
| 183 | result.recv_timeout(Duration::from_secs(5)).unwrap(), | |
| 184 | Err(share::Refusal::Cancelled) | |
| 185 | ); | |
| 186 | until("the cancelled request stayed", || { | |
| 187 | host.requests().is_empty() | |
| 188 | }); | |
| 189 | assert!(host.devices().is_empty()); | |
| 190 | } |
crates/notebook/tests/live_share.rs+2| ... | ... | @@ -70,6 +70,7 @@ fn a_guest_queues_while_the_host_is_away() { |
| 70 | 70 | let image = std::fs::read(&file).unwrap(); |
| 71 | 71 | |
| 72 | 72 | // Ada's computer goes to sleep. |
| 73 | let sharing = host.sharing(); | |
| 73 | 74 | drop(host); |
| 74 | 75 | until("the host never left", || guest.host().is_none()); |
| 75 | 76 | let id = replace(&section, &image, 0..8, "Offline"); |
| ... | ... | @@ -110,6 +111,7 @@ fn two_guests_on_one_page_conflict_as_on_a_share() { |
| 110 | 111 | let file = folder.join("Garden.one"); |
| 111 | 112 | let image = std::fs::read(&file).unwrap(); |
| 112 | 113 | |
| 114 | let sharing = host.sharing(); | |
| 113 | 115 | drop(host); |
| 114 | 116 | until("the host never left", || { |
| 115 | 117 | grace.host().is_none() && alan.host().is_none() |
crates/notebook/tests/support/live.rs+1| ... | ... | @@ -83,6 +83,7 @@ pub fn host(folder: &Path, cache: &Path, sharing: &Sharing, url: &str) -> Host { |
| 83 | 83 | None, |
| 84 | 84 | Some(url), |
| 85 | 85 | || {}, |
| 86 | |_| Ok(()), | |
| 86 | 87 | ) |
| 87 | 88 | .unwrap() |
| 88 | 89 | } |
crates/relay/src/main.rs+2-2| ... | ... | @@ -12,14 +12,14 @@ underscores: SNOWBOUND_RELAY_LISTEN=127.0.0.1:23592. Flags win. |
| 12 | 12 | --listen ADDRESS where to listen (127.0.0.1:23592) |
| 13 | 13 | --trust-forwarded true|false count peers by X-Forwarded-For, behind a proxy (false) |
| 14 | 14 | --max-connections N (256) |
| 15 | --max-connections-per-address N per IPv4 address or IPv6 /64 (16) | |
| 15 | --max-connections-per-address N per IPv4 address or IPv6 /64 (128) | |
| 16 | 16 | --max-rooms N (128) |
| 17 | 17 | --max-room-peers N (64) |
| 18 | 18 | --max-message BYTES (262144) |
| 19 | 19 | --queue BYTES waiting to go to one peer before it is dropped (1048576) |
| 20 | 20 | --idle SECONDS silence before a connection is closed (600) |
| 21 | 21 | --room-bytes-per-second BYTES (4194304) |
| 22 | --joins-per-minute N per address (30) | |
| 22 | --joins-per-minute N per address (240) | |
| 23 | 23 | --room-joins-per-minute N (120) |
| 24 | 24 | --failures-per-minute N wrong codes per address before a lockout (10) |
| 25 | 25 | --burn-after N wrong codes before a code admits no one new (5) |
crates/relay/src/server.rs+2-2| ... | ... | @@ -48,14 +48,14 @@ impl Default for Config { |
| 48 | 48 | Self { |
| 49 | 49 | trust_forwarded: false, |
| 50 | 50 | max_connections: 256, |
| 51 | max_connections_per_address: 16, | |
| 51 | max_connections_per_address: 128, | |
| 52 | 52 | max_rooms: 128, |
| 53 | 53 | max_room_peers: 64, |
| 54 | 54 | max_message: 256 << 10, |
| 55 | 55 | queue: 1 << 20, |
| 56 | 56 | idle: Duration::from_secs(600), |
| 57 | 57 | room_bytes_per_second: 4 << 20, |
| 58 | joins_per_minute: 30, | |
| 58 | joins_per_minute: 240, | |
| 59 | 59 | room_joins_per_minute: 120, |
| 60 | 60 | failures_per_minute: 10, |
| 61 | 61 | burn_after: 5, |
crates/snowbound/src/live.rs+64-27| ... | ... | @@ -90,7 +90,17 @@ fn hello() -> io::Result<Hello> { |
| 90 | 90 | let picture = (picture && name == platform::user_name()) |
| 91 | 91 | .then(|| account_picture().clone()) |
| 92 | 92 | .flatten(); |
| 93 | Hello::new(name, picture) | |
| 93 | let mut hello = Hello::new(name, picture)?; | |
| 94 | hello.device = std::env::var("COMPUTERNAME").ok().or_else(|| { | |
| 95 | std::process::Command::new("hostname") | |
| 96 | .output() | |
| 97 | .ok() | |
| 98 | .filter(|output| output.status.success()) | |
| 99 | .and_then(|output| String::from_utf8(output.stdout).ok()) | |
| 100 | .map(|name| name.trim().to_owned()) | |
| 101 | .filter(|name| !name.is_empty()) | |
| 102 | }); | |
| 103 | Ok(hello) | |
| 94 | 104 | } |
| 95 | 105 | |
| 96 | 106 | /// 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<()> { |
| 141 | 151 | options.write(true).create(true).truncate(true); |
| 142 | 152 | #[cfg(unix)] |
| 143 | 153 | std::os::unix::fs::OpenOptionsExt::mode(&mut options, 0o600); |
| 144 | io::Write::write_all( | |
| 145 | &mut options.open(&partial)?, | |
| 146 | &serde_json::to_vec_pretty(value)?, | |
| 147 | )?; | |
| 148 | notebook::fs::rename(partial, file) | |
| 154 | let mut output = options.open(&partial)?; | |
| 155 | io::Write::write_all(&mut output, &serde_json::to_vec_pretty(value)?)?; | |
| 156 | output.sync_all()?; | |
| 157 | notebook::fs::rename(partial, file)?; | |
| 158 | #[cfg(unix)] | |
| 159 | std::fs::File::open(folder)?.sync_all()?; | |
| 160 | Ok(()) | |
| 149 | 161 | } |
| 150 | 162 | |
| 151 | 163 | /// A notebook another computer shares, as this one joined it. |
| ... | ... | @@ -157,6 +169,8 @@ struct Share { |
| 157 | 169 | host: String, |
| 158 | 170 | /// It has been listed here, so it opens without waiting for the host. |
| 159 | 171 | listed: bool, |
| 172 | #[serde(default)] | |
| 173 | protocol: u16, | |
| 160 | 174 | } |
| 161 | 175 | |
| 162 | 176 | const JOINED: &str = "joined.json"; |
| ... | ... | @@ -186,6 +200,11 @@ impl Joined { |
| 186 | 200 | )); |
| 187 | 201 | }; |
| 188 | 202 | let refused = |error: &dyn std::fmt::Display| (share.notebook.clone(), error.to_string()); |
| 203 | if share.protocol != live::wire::VERSION { | |
| 204 | return Err(refused( | |
| 205 | &"Open this notebook again with the person sharing’s current link or code.", | |
| 206 | )); | |
| 207 | } | |
| 189 | 208 | let guest = hello() |
| 190 | 209 | .and_then(|me| { |
| 191 | 210 | Guest::start( |
| ... | ... | @@ -241,9 +260,10 @@ impl Joined { |
| 241 | 260 | pub(crate) fn join( |
| 242 | 261 | code: &str, |
| 243 | 262 | password: &str, |
| 263 | waiting: impl Fn(bool) -> bool, | |
| 244 | 264 | ) -> Result<live::wire::Welcome, live::share::Refusal> { |
| 245 | 265 | let me = hello().map_err(|_| live::share::Refusal::Unreachable(live::Trouble::Other))?; |
| 246 | live::share::join(me, code, password, reach(), relay().as_deref()) | |
| 266 | live::share::join_while(me, code, password, reach(), relay().as_deref(), waiting) | |
| 247 | 267 | } |
| 248 | 268 | |
| 249 | 269 | /// 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<S |
| 259 | 279 | notebook: welcome.notebook, |
| 260 | 280 | host: welcome.host, |
| 261 | 281 | listed: false, |
| 282 | protocol: live::wire::VERSION, | |
| 262 | 283 | }, |
| 263 | 284 | ); |
| 264 | 285 | keep(&file, &shares)?; |
| ... | ... | @@ -289,7 +310,7 @@ pub(crate) struct Peers { |
| 289 | 310 | /// The notebooks this computer shares, by location. |
| 290 | 311 | pub(crate) hosts: BTreeMap<String, Arc<Host>>, |
| 291 | 312 | /// Shares as kept for the next launch. |
| 292 | sharing: Option<BTreeMap<String, Sharing>>, | |
| 313 | sharing: Option<Arc<Mutex<BTreeMap<String, Sharing>>>>, | |
| 293 | 314 | /// Shares starting on threads of their own, by location. |
| 294 | 315 | pub(crate) starting: HashMap<String, Option<String>>, |
| 295 | 316 | started: Option<Channel<(String, Result<Host, String>)>>, |
| ... | ... | @@ -410,7 +431,9 @@ impl State { |
| 410 | 431 | let sharing = self |
| 411 | 432 | .peers |
| 412 | 433 | .sharing |
| 413 | .get_or_insert_with(|| read_kept(&kept(&cache, HOSTING))) | |
| 434 | .get_or_insert_with(|| Arc::new(Mutex::new(read_kept(&kept(&cache, HOSTING))))) | |
| 435 | .lock() | |
| 436 | .unwrap() | |
| 414 | 437 | .clone(); |
| 415 | 438 | let open: Vec<Arc<Library>> = self.notebooks.clone(); |
| 416 | 439 | for library in &open { |
| ... | ... | @@ -443,25 +466,29 @@ impl State { |
| 443 | 466 | for location in closed { |
| 444 | 467 | self.stop_sharing(&location); |
| 445 | 468 | } |
| 446 | let now: BTreeMap<String, Sharing> = (self.peers.hosts.iter()) | |
| 447 | .map(|(location, host)| (location.clone(), host.sharing())) | |
| 448 | .collect(); | |
| 449 | let before: BTreeMap<String, Sharing> = (sharing.into_iter()) | |
| 450 | .filter(|(location, _)| { | |
| 451 | !self.peers.starting.contains_key(location) | |
| 452 | && open.iter().any(|library| library.location == *location) | |
| 469 | for host in self.peers.hosts.values() { | |
| 470 | host.code(); | |
| 471 | } | |
| 472 | if self.peers.share.is_none() | |
| 473 | && let Some(library) = open.iter().find(|library| { | |
| 474 | self.peers | |
| 475 | .hosts | |
| 476 | .get(&library.location) | |
| 477 | .is_some_and(|host| !host.requests().is_empty()) | |
| 453 | 478 | }) |
| 454 | .collect(); | |
| 455 | if now != before { | |
| 456 | if let Err(error) = keep(&kept(&cache, HOSTING), &now) { | |
| 457 | eprintln!("Keeping what this computer shares: {error}"); | |
| 458 | } | |
| 459 | self.peers.sharing = Some(now); | |
| 479 | { | |
| 480 | self.open_live_share(Arc::clone(library)); | |
| 460 | 481 | } |
| 461 | 482 | } |
| 462 | 483 | |
| 463 | 484 | /// Shares `library` as `sharing` says, on a thread of its own. |
| 464 | 485 | pub(crate) fn start_sharing(&mut self, library: &Arc<Library>, sharing: Sharing) { |
| 486 | let file = kept(&self.cache, HOSTING); | |
| 487 | let kept = Arc::clone( | |
| 488 | self.peers | |
| 489 | .sharing | |
| 490 | .get_or_insert_with(|| Arc::new(Mutex::new(read_kept(&file)))), | |
| 491 | ); | |
| 465 | 492 | self.peers.starting.insert(library.location.clone(), None); |
| 466 | 493 | let (started, _) = self.peers.started.get_or_insert_with(mpsc::channel); |
| 467 | 494 | let (started, library, redraw) = |
| ... | ... | @@ -470,6 +497,7 @@ impl State { |
| 470 | 497 | let host = (|| -> Result<Host, Box<dyn std::error::Error>> { |
| 471 | 498 | let storage = library.reopen()?.into_storage(); |
| 472 | 499 | let told = redraw.clone(); |
| 500 | let location = library.location.clone(); | |
| 473 | 501 | Ok(Host::start( |
| 474 | 502 | storage, |
| 475 | 503 | hello()?, |
| ... | ... | @@ -478,6 +506,14 @@ impl State { |
| 478 | 506 | reach(), |
| 479 | 507 | relay().as_deref(), |
| 480 | 508 | move || told.wake_by_ref(), |
| 509 | move |sharing| { | |
| 510 | let mut shares = kept.lock().unwrap(); | |
| 511 | let mut next = shares.clone(); | |
| 512 | next.insert(location.clone(), sharing.clone()); | |
| 513 | keep(&file, &next)?; | |
| 514 | *shares = next; | |
| 515 | Ok(()) | |
| 516 | }, | |
| 481 | 517 | )?) |
| 482 | 518 | })(); |
| 483 | 519 | let _ = started.send(( |
| ... | ... | @@ -521,11 +557,12 @@ impl State { |
| 521 | 557 | crate::spawn(move || host.stop()); |
| 522 | 558 | } |
| 523 | 559 | // Kept, the share would start again on the next frame, as after a relaunch. |
| 524 | if let Some(sharing) = &mut self.peers.sharing | |
| 525 | && sharing.remove(location).is_some() | |
| 526 | && let Err(error) = keep(&kept(&self.cache, HOSTING), sharing) | |
| 527 | { | |
| 528 | eprintln!("Keeping what this computer shares: {error}"); | |
| 560 | if let Some(sharing) = &self.peers.sharing { | |
| 561 | let mut sharing = sharing.lock().unwrap(); | |
| 562 | sharing.remove(location); | |
| 563 | if let Err(error) = keep(&kept(&self.cache, HOSTING), &*sharing) { | |
| 564 | eprintln!("Keeping what this computer shares: {error}"); | |
| 565 | } | |
| 529 | 566 | } |
| 530 | 567 | } |
| 531 | 568 |
crates/snowbound/src/share.rs+136-10| ... | ... | @@ -8,7 +8,11 @@ use notebook::live::{ |
| 8 | 8 | share::{self, Refusal, Sharing}, |
| 9 | 9 | wire::Welcome, |
| 10 | 10 | }; |
| 11 | use std::sync::{Arc, mpsc}; | |
| 11 | use std::sync::{ | |
| 12 | Arc, | |
| 13 | atomic::{AtomicBool, Ordering}, | |
| 14 | mpsc, | |
| 15 | }; | |
| 12 | 16 | use ui::{Anchor, Axis, Id, Spec, Ui, children, fill, fit, px}; |
| 13 | 17 | use winit::keyboard::NamedKey; |
| 14 | 18 | |
| ... | ... | @@ -41,6 +45,9 @@ pub(crate) struct ShareDialog { |
| 41 | 45 | library: Arc<Library>, |
| 42 | 46 | protect: bool, |
| 43 | 47 | password: String, |
| 48 | approve: bool, | |
| 49 | replies: crate::live::Channel<Result<(), String>>, | |
| 50 | error: Option<String>, | |
| 44 | 51 | /// What was copied last: the code, or the link. |
| 45 | 52 | copied: Option<&'static str>, |
| 46 | 53 | } |
| ... | ... | @@ -53,6 +60,14 @@ pub(crate) struct JoinDialog { |
| 53 | 60 | password: String, |
| 54 | 61 | status: Status, |
| 55 | 62 | replies: crate::live::Channel<Result<Welcome, Refusal>>, |
| 63 | alive: Arc<AtomicBool>, | |
| 64 | approval: Arc<AtomicBool>, | |
| 65 | } | |
| 66 | ||
| 67 | impl Drop for JoinDialog { | |
| 68 | fn drop(&mut self) { | |
| 69 | self.alive.store(false, Ordering::Release); | |
| 70 | } | |
| 56 | 71 | } |
| 57 | 72 | |
| 58 | 73 | enum Status { |
| ... | ... | @@ -92,6 +107,9 @@ fn refusal(refusal: &Refusal) -> String { |
| 92 | 107 | Refusal::Busy => "The Live Share relay is busy. Try again in a minute.".into(), |
| 93 | 108 | Refusal::Unreachable(trouble) => unreachable(*trouble), |
| 94 | 109 | Refusal::TimedOut => "The computer sharing didn’t answer. Try again.".into(), |
| 110 | Refusal::Declined => "The person sharing declined this request.".into(), | |
| 111 | Refusal::Cancelled => "Joining cancelled.".into(), | |
| 112 | Refusal::NotAdmitted => "The computer sharing couldn’t add this device. Ask the person sharing to check Live Share.".into(), | |
| 95 | 113 | } |
| 96 | 114 | } |
| 97 | 115 | |
| ... | ... | @@ -267,6 +285,9 @@ impl State { |
| 267 | 285 | library, |
| 268 | 286 | protect: false, |
| 269 | 287 | password: String::new(), |
| 288 | approve: true, | |
| 289 | replies: mpsc::channel(), | |
| 290 | error: None, | |
| 270 | 291 | copied: None, |
| 271 | 292 | }); |
| 272 | 293 | self.ui.open_popup(share_id()); |
| ... | ... | @@ -289,6 +310,8 @@ impl State { |
| 289 | 310 | password: String::new(), |
| 290 | 311 | status: Status::Idle, |
| 291 | 312 | replies: mpsc::channel(), |
| 313 | alive: Arc::new(AtomicBool::new(true)), | |
| 314 | approval: Arc::new(AtomicBool::new(false)), | |
| 292 | 315 | }); |
| 293 | 316 | self.ui.open_popup(join_id()); |
| 294 | 317 | self.ui.focus_all(code_field()); |
| ... | ... | @@ -307,6 +330,9 @@ impl State { |
| 307 | 330 | return; |
| 308 | 331 | } |
| 309 | 332 | let location = dialog.library.location.clone(); |
| 333 | for result in dialog.replies.1.try_iter() { | |
| 334 | dialog.error = result.err(); | |
| 335 | } | |
| 310 | 336 | let host = self.peers.hosts.get(&location).cloned(); |
| 311 | 337 | let starting = self.peers.starting.get(&location).cloned(); |
| 312 | 338 | let ui = &mut self.ui; |
| ... | ... | @@ -319,6 +345,7 @@ impl State { |
| 319 | 345 | frame(ui, share_id(), "Live Share"); |
| 320 | 346 | text(ui, "notebook", &dialog.library.name, true); |
| 321 | 347 | let (mut start, mut stop, mut copy) = (false, false, None); |
| 348 | let (mut approval, mut allow, mut decline, mut remove) = (None, None, None, None); | |
| 322 | 349 | match (&host, &starting) { |
| 323 | 350 | (Some(host), _) => { |
| 324 | 351 | let code = host.code(); |
| ... | ... | @@ -390,30 +417,85 @@ impl State { |
| 390 | 417 | ), |
| 391 | 418 | _ => {} |
| 392 | 419 | } |
| 420 | if ui::check_box(ui, "approve", "Ask before joining", host.sharing().approve) | |
| 421 | .clicked | |
| 422 | { | |
| 423 | approval = Some(!host.sharing().approve); | |
| 424 | } | |
| 425 | for request in host.requests() { | |
| 426 | ui.open( | |
| 427 | ("request", request.peer), | |
| 428 | Spec { | |
| 429 | size: [fill(), children()], | |
| 430 | gap: 8.0, | |
| 431 | ..Spec::default() | |
| 432 | }, | |
| 433 | ); | |
| 434 | ui.leaf( | |
| 435 | "name", | |
| 436 | Spec { | |
| 437 | size: [fill(), fit()], | |
| 438 | text: Some(&format!("{} wants to join", request.name)), | |
| 439 | overflow: ui::Overflow::Wrap, | |
| 440 | ..Spec::default() | |
| 441 | }, | |
| 442 | ); | |
| 443 | if ui::button(ui, "allow", "Allow").clicked { | |
| 444 | allow = Some(request.peer); | |
| 445 | } | |
| 446 | if ui::button(ui, "decline", "Decline").clicked { | |
| 447 | decline = Some(request.peer); | |
| 448 | } | |
| 449 | ui.close(); | |
| 450 | } | |
| 393 | 451 | ui.leaf( |
| 394 | 452 | "people", |
| 395 | 453 | Spec { |
| 396 | 454 | size: [fill(), px(theme.font_size * 2.0)], |
| 397 | text: Some("Connected"), | |
| 455 | text: Some("Devices"), | |
| 398 | 456 | bold: true, |
| 399 | 457 | role: Some(Role::Heading), |
| 400 | 458 | ..Spec::default() |
| 401 | 459 | }, |
| 402 | 460 | ); |
| 403 | let guests = host.guests(); | |
| 404 | if guests.is_empty() { | |
| 461 | let devices = host.devices(); | |
| 462 | if devices.is_empty() { | |
| 405 | 463 | text(ui, "nobody", "No one has joined yet.", true); |
| 406 | 464 | } |
| 407 | for guest in &guests { | |
| 465 | for (device, connected) in &devices { | |
| 466 | ui.open( | |
| 467 | ("device", device.secret), | |
| 468 | Spec { | |
| 469 | size: [fill(), children()], | |
| 470 | gap: 8.0, | |
| 471 | ..Spec::default() | |
| 472 | }, | |
| 473 | ); | |
| 408 | 474 | ui.leaf( |
| 409 | ("guest", guest.hello.peer), | |
| 475 | "name", | |
| 410 | 476 | Spec { |
| 411 | 477 | size: [fill(), px(theme.font_size * 1.6)], |
| 412 | text: Some(&guest.hello.name), | |
| 478 | text: Some(&match &device.device { | |
| 479 | Some(name) => format!("{} · {name}", device.name), | |
| 480 | None => device.name.clone(), | |
| 481 | }), | |
| 413 | 482 | role: Some(Role::ListItem), |
| 414 | 483 | ..Spec::default() |
| 415 | 484 | }, |
| 416 | 485 | ); |
| 486 | ui.leaf( | |
| 487 | "connection", | |
| 488 | Spec { | |
| 489 | size: [fit(), px(theme.font_size * 1.6)], | |
| 490 | text: Some(if *connected { "Connected" } else { "Offline" }), | |
| 491 | color: Some(theme.text_dim), | |
| 492 | ..Spec::default() | |
| 493 | }, | |
| 494 | ); | |
| 495 | if ui::button(ui, "remove", "Remove").clicked { | |
| 496 | remove = Some(device.secret); | |
| 497 | } | |
| 498 | ui.close(); | |
| 417 | 499 | } |
| 418 | 500 | buttons(ui); |
| 419 | 501 | stop = ui::button(ui, "stop", "Stop Sharing").clicked; |
| ... | ... | @@ -436,6 +518,9 @@ impl State { |
| 436 | 518 | ui.focus_all(password_field()); |
| 437 | 519 | } |
| 438 | 520 | } |
| 521 | if ui::check_box(ui, "approve", "Ask before joining", dialog.approve).clicked { | |
| 522 | dialog.approve = !dialog.approve; | |
| 523 | } | |
| 439 | 524 | if dialog.protect { |
| 440 | 525 | labelled(ui, "Password:", |ui| { |
| 441 | 526 | let spec = field_spec(ui); |
| ... | ... | @@ -455,6 +540,9 @@ impl State { |
| 455 | 540 | let offered = host.is_none() && !matches!(starting, Some(None)); |
| 456 | 541 | let done = ui::button(ui, "done", "Done").clicked || entered && !offered; |
| 457 | 542 | ui.close(); |
| 543 | if let Some(error) = &dialog.error { | |
| 544 | status(ui, error); | |
| 545 | } | |
| 458 | 546 | ui.close(); |
| 459 | 547 | if let Some((what, text)) = copy { |
| 460 | 548 | dialog.copied = self.clipboard.set_text(text).is_ok().then_some(what); |
| ... | ... | @@ -462,6 +550,27 @@ impl State { |
| 462 | 550 | if stop { |
| 463 | 551 | self.stop_sharing(&location); |
| 464 | 552 | } |
| 553 | if let Some(host) = host { | |
| 554 | if let Some(peer) = decline { | |
| 555 | host.decline(&peer); | |
| 556 | } | |
| 557 | if approval.is_some() || allow.is_some() || remove.is_some() { | |
| 558 | let replies = dialog.replies.0.clone(); | |
| 559 | let redraw = self.redraw.clone(); | |
| 560 | crate::spawn(move || { | |
| 561 | let result = if let Some(approve) = approval { | |
| 562 | host.approve(approve) | |
| 563 | } else if let Some(peer) = allow { | |
| 564 | host.allow(&peer) | |
| 565 | } else { | |
| 566 | host.remove(&remove.unwrap()) | |
| 567 | }; | |
| 568 | let _ = replies | |
| 569 | .send(result.map_err(|error| format!("Couldn’t change access: {error}"))); | |
| 570 | redraw.wake(); | |
| 571 | }); | |
| 572 | } | |
| 573 | } | |
| 465 | 574 | if start && !(dialog.protect && dialog.password.is_empty()) { |
| 466 | 575 | let password = if dialog.protect { |
| 467 | 576 | dialog.password.clone() |
| ... | ... | @@ -469,7 +578,10 @@ impl State { |
| 469 | 578 | String::new() |
| 470 | 579 | }; |
| 471 | 580 | match Sharing::new(&password) { |
| 472 | Ok(sharing) => self.start_sharing(&dialog.library, sharing), | |
| 581 | Ok(mut sharing) => { | |
| 582 | sharing.approve = dialog.approve; | |
| 583 | self.start_sharing(&dialog.library, sharing); | |
| 584 | } | |
| 473 | 585 | Err(error) => { |
| 474 | 586 | self.peers |
| 475 | 587 | .starting |
| ... | ... | @@ -539,7 +651,14 @@ impl State { |
| 539 | 651 | }); |
| 540 | 652 | match &dialog.status { |
| 541 | 653 | Status::Idle => {} |
| 542 | Status::Waiting => status(ui, "Connecting…"), | |
| 654 | Status::Waiting => status( | |
| 655 | ui, | |
| 656 | if dialog.approval.load(Ordering::Acquire) { | |
| 657 | "Needs approval" | |
| 658 | } else { | |
| 659 | "Connecting…" | |
| 660 | }, | |
| 661 | ), | |
| 543 | 662 | Status::Failed(message) => status(ui, message), |
| 544 | 663 | } |
| 545 | 664 | buttons(ui); |
| ... | ... | @@ -560,8 +679,15 @@ impl State { |
| 560 | 679 | dialog.status = Status::Waiting; |
| 561 | 680 | let (replies, redraw) = (dialog.replies.0.clone(), self.redraw.clone()); |
| 562 | 681 | let password = dialog.password.clone(); |
| 682 | let (alive, approval) = | |
| 683 | (Arc::clone(&dialog.alive), Arc::clone(&dialog.approval)); | |
| 563 | 684 | crate::spawn(move || { |
| 564 | let reply = crate::live::join(&code, &password); | |
| 685 | let reply = crate::live::join(&code, &password, |waiting| { | |
| 686 | if approval.swap(waiting, Ordering::AcqRel) != waiting { | |
| 687 | redraw.wake_by_ref(); | |
| 688 | } | |
| 689 | alive.load(Ordering::Acquire) | |
| 690 | }); | |
| 565 | 691 | let _ = replies.send(reply); |
| 566 | 692 | redraw.wake(); |
| 567 | 693 | }); |
crates/snowbound/tests/replay.rs+2-1| ... | ... | @@ -723,11 +723,12 @@ fn stop_sharing_ends_the_share() { |
| 723 | 723 | let mut steps = vec!["modifiers command shift", "key p", "modifiers", "settle"]; |
| 724 | 724 | steps.extend(["type Live Share", "settle", "key Enter", "settle"]); |
| 725 | 725 | steps.extend(["key Enter", "wait 1500", "accessibility shared"]); |
| 726 | // Copy, Copy Link, then Stop Sharing. | |
| 726 | // Copy, Copy Link, Ask before joining, then Stop Sharing. | |
| 727 | 727 | steps.extend([ |
| 728 | 728 | "key Tab", |
| 729 | 729 | "key Tab", |
| 730 | 730 | "key Tab", |
| 731 | "key Tab", | |
| 731 | 732 | "key Enter", |
| 732 | 733 | "wait 1500", |
| 733 | 734 | "accessibility stopped", |