| 1 | //! Live Share: a notebook one Snowbound holds, opened on others through a short code. The |
| 2 | //! host serves its notebook's storage verbs (`session::Storage`) and batches guests' ops into |
| 3 | //! guarded publications. A guest runs the same replica and durable queue as on an SMB share; |
| 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 | //! the next, so a relay never holds much for a slow peer. |
| 8 | |
| 9 | use super::{ |
| 10 | Event, Hello, Line, Live, Peer, Presence, Reach, Relayed, Room, Sender, |
| 11 | wire::{self, Delta, Failure, Reply, Request, Touched, Welcome, WireEntry, Written, kind}, |
| 12 | }; |
| 13 | use crate::{Error, Result, background::Reports, discover, session::Storage}; |
| 14 | use onestore::{CommitError, CommitState, RevisionIndex, Stamp, Store, Transaction}; |
| 15 | use sha2::{Digest, Sha256}; |
| 16 | use std::{ |
| 17 | collections::{BTreeMap, HashMap}, |
| 18 | io, |
| 19 | path::PathBuf, |
| 20 | sync::{ |
| 21 | Arc, Condvar, Mutex, |
| 22 | atomic::{AtomicBool, AtomicU64, Ordering}, |
| 23 | mpsc, |
| 24 | }, |
| 25 | thread, |
| 26 | time::Duration, |
| 27 | }; |
| 28 | use web_time::Instant; |
| 29 | |
| 30 | #[path = "share/batch.rs"] |
| 31 | mod batch; |
| 32 | #[path = "share/membership.rs"] |
| 33 | mod membership; |
| 34 | #[cfg(target_arch = "wasm32")] |
| 35 | #[path = "share/web.rs"] |
| 36 | mod web; |
| 37 | pub use membership::Host; |
| 38 | |
| 39 | /// The most bytes one message of a read or an upload carries. |
| 40 | const CHUNK: usize = 128 << 10; |
| 41 | /// The most bytes of chunks a guest has asked for and not yet been given. |
| 42 | const WINDOW: usize = 512 << 10; |
| 43 | /// How long a request waits for its reply. |
| 44 | const TIMEOUT: Duration = Duration::from_secs(60); |
| 45 | use crate::MAX_FILE_BYTES as LIMIT; |
| 46 | /// Snapshots of files being read, per guest, and how long one is kept unread. |
| 47 | const SNAPSHOTS: usize = 8; |
| 48 | const SNAPSHOT_AGE: Duration = Duration::from_secs(120); |
| 49 | /// Sections whose image a host keeps to check guests' commits on and tell them what changed, |
| 50 | /// and a guest keeps to take those changes on without reading. |
| 51 | const IMAGES: usize = 4; |
| 52 | /// Earlier images of those a host keeps besides, to tell a guest that missed a change what |
| 53 | /// it was, and the most bytes they hold. |
| 54 | const VERSIONS: usize = 16; |
| 55 | const VERSIONS_BYTES: usize = 64 << 20; |
| 56 | /// How long a report of a change settles in a guest's background: its host reports each |
| 57 | /// commit once, at once. |
| 58 | pub const SETTLE: Duration = Duration::from_millis(20); |
| 59 | /// The longest a relay may ask a guest joining to wait before it is told so. |
| 60 | const PATIENT: Duration = Duration::from_secs(10); |
| 61 | /// What a host says leaving as it stops sharing. |
| 62 | const STOPPED: &str = "stopped"; |
| 63 | /// What a host says hanging up on a guest that asks too much too fast. |
| 64 | const FLOODED: &str = "flooded"; |
| 65 | /// A guest's requests a host holds at once, waiting and in hand, past which it hangs up. |
| 66 | const QUEUED: usize = 32; |
| 67 | /// The threads working through each guest's requests. |
| 68 | const WORKERS: usize = 2; |
| 69 | /// Requests a guest may start each second, and at once: every request but a read's later |
| 70 | /// chunks and an upload's, which come from memory or go to it. |
| 71 | const STARTS: f64 = 100.0; |
| 72 | const BURST: f64 = 200.0; |
| 73 | |
| 74 | /// A code as typed, `7kq 4mz 9xr`, as it is shown, `7KQ-4MZ-9XR`; none where it is not a |
| 75 | /// code (`super::code`). |
| 76 | pub fn code(typed: &str) -> Option<String> { |
| 77 | let (number, secret) = super::code::parse(typed)?; |
| 78 | super::code::format(number, &secret) |
| 79 | } |
| 80 | |
| 81 | /// What a host keeps of a share to take it up again after a relaunch. |
| 82 | #[derive(Clone, Debug, PartialEq, Eq, serde::Serialize, serde::Deserialize)] |
| 83 | pub struct Sharing { |
| 84 | pub share: [u8; 16], |
| 85 | /// The presence key, rotated when a device is removed. |
| 86 | pub secret: [u8; 16], |
| 87 | /// The code, or its secret alone until it has a number (`super::code`). |
| 88 | pub code: String, |
| 89 | pub password: String, |
| 90 | #[serde(default)] |
| 91 | pub approve: bool, |
| 92 | #[serde(default)] |
| 93 | pub members: Vec<Device>, |
| 94 | } |
| 95 | |
| 96 | #[derive(Clone, Debug, PartialEq, Eq, serde::Serialize, serde::Deserialize)] |
| 97 | pub struct Device { |
| 98 | pub secret: [u8; 16], |
| 99 | pub name: String, |
| 100 | pub device: Option<String>, |
| 101 | } |
| 102 | |
| 103 | impl Sharing { |
| 104 | /// A new share: its own id, room secret and code. |
| 105 | pub fn new(password: &str) -> io::Result<Self> { |
| 106 | let mut random = [0; 32]; |
| 107 | getrandom::fill(&mut random) |
| 108 | .map_err(|_| io::Error::other("System random source failed"))?; |
| 109 | Ok(Self { |
| 110 | share: random[..16].try_into().expect("16 bytes"), |
| 111 | secret: random[16..].try_into().expect("16 bytes"), |
| 112 | code: super::code::secret()?, |
| 113 | password: password.to_owned(), |
| 114 | approve: false, |
| 115 | members: Vec::new(), |
| 116 | }) |
| 117 | } |
| 118 | } |
| 119 | |
| 120 | /// Where a guest keeps a share's notebook: `live://` and the share's id. |
| 121 | pub fn location(share: &[u8; 16]) -> String { |
| 122 | format!("live://{}", super::hex(share)) |
| 123 | } |
| 124 | |
| 125 | /// Why joining failed. |
| 126 | #[derive(Clone, Debug, PartialEq, Eq)] |
| 127 | pub enum Refusal { |
| 128 | /// Not a code, or one mistyped, which spends none of the relay's tries. |
| 129 | Malformed, |
| 130 | /// The person sharing runs a Snowbound of another Live Share version: the newer one's |
| 131 | /// `true` where it is theirs, so this one should update. |
| 132 | Version { |
| 133 | theirs_newer: bool, |
| 134 | }, |
| 135 | /// The code's secret or password is wrong. |
| 136 | Wrong, |
| 137 | /// No one shares with the code's number now. |
| 138 | NoOne, |
| 139 | /// The code had too many wrong tries. |
| 140 | Expired, |
| 141 | /// Too many wrong codes from this network; try again after the wait, where known. |
| 142 | TooMany(Option<Duration>), |
| 143 | /// The relay is full. |
| 144 | Busy, |
| 145 | /// The relay couldn't be reached, and why, and no one answered on this network. |
| 146 | Unreachable(super::Trouble), |
| 147 | /// The relay let this end in, but no one answered. |
| 148 | TimedOut, |
| 149 | Declined, |
| 150 | Cancelled, |
| 151 | NotAdmitted, |
| 152 | } |
| 153 | |
| 154 | /// Meets the host sharing `code` (and `password`) as `me`, on the networks `reach` names and |
| 155 | /// through `relay`: the share it welcomes this end to. |
| 156 | pub fn join( |
| 157 | me: Hello, |
| 158 | code: &str, |
| 159 | password: &str, |
| 160 | reach: Option<Reach>, |
| 161 | relay: Option<&str>, |
| 162 | ) -> std::result::Result<Welcome, Refusal> { |
| 163 | join_while(me, code, password, reach, relay, |_| true) |
| 164 | } |
| 165 | |
| 166 | /// Joins while `waiting` returns true, reporting whether the host is deciding approval. |
| 167 | pub fn join_while( |
| 168 | me: Hello, |
| 169 | code: &str, |
| 170 | password: &str, |
| 171 | reach: Option<Reach>, |
| 172 | relay: Option<&str>, |
| 173 | continue_joining: impl Fn(bool) -> bool, |
| 174 | ) -> std::result::Result<Welcome, Refusal> { |
| 175 | crate::task::ready(join_while_async( |
| 176 | me, |
| 177 | code, |
| 178 | password, |
| 179 | reach, |
| 180 | relay, |
| 181 | continue_joining, |
| 182 | )) |
| 183 | .map_err(|_| Refusal::Busy)? |
| 184 | } |
| 185 | |
| 186 | pub async fn join_while_async( |
| 187 | me: Hello, |
| 188 | code: &str, |
| 189 | password: &str, |
| 190 | reach: Option<Reach>, |
| 191 | relay: Option<&str>, |
| 192 | continue_joining: impl Fn(bool) -> bool, |
| 193 | ) -> std::result::Result<Welcome, Refusal> { |
| 194 | let code = self::code(code).ok_or(Refusal::Malformed)?; |
| 195 | let (welcomed, welcome) = mpsc::channel(); |
| 196 | let approving = Arc::new(AtomicBool::new(false)); |
| 197 | let approval = Arc::clone(&approving); |
| 198 | let (changed, waiting) = crate::task::channel(); |
| 199 | let live = Live::start( |
| 200 | me, |
| 201 | &Room::join(&code, password), |
| 202 | reach, |
| 203 | relay, |
| 204 | move |event| match event { |
| 205 | Event::Frame { |
| 206 | kind: kind::WELCOME, |
| 207 | body, |
| 208 | .. |
| 209 | } => { |
| 210 | if let Ok(body) = minicbor::decode::<Welcome>(body) { |
| 211 | let _ = welcomed.send(Ok(body)); |
| 212 | } |
| 213 | } |
| 214 | Event::Frame { |
| 215 | kind: kind::APPROVAL, |
| 216 | body, |
| 217 | .. |
| 218 | } => { |
| 219 | if let Ok(body) = minicbor::decode::<wire::Approval>(body) { |
| 220 | match body { |
| 221 | wire::Approval::Pending => approval.store(true, Ordering::Release), |
| 222 | wire::Approval::Declined => { |
| 223 | let _ = welcomed.send(Err(Refusal::Declined)); |
| 224 | } |
| 225 | wire::Approval::Failed => { |
| 226 | let _ = welcomed.send(Err(Refusal::NotAdmitted)); |
| 227 | } |
| 228 | } |
| 229 | } |
| 230 | } |
| 231 | _ => { |
| 232 | let _ = changed.try_send(()); |
| 233 | } |
| 234 | }, |
| 235 | ) |
| 236 | .map_err(|_| Refusal::Unreachable(super::Trouble::Other))?; |
| 237 | let start = Instant::now(); |
| 238 | loop { |
| 239 | if !continue_joining(approving.load(Ordering::Acquire)) { |
| 240 | return Err(Refusal::Cancelled); |
| 241 | } |
| 242 | if let Ok(welcome) = welcome.try_recv() { |
| 243 | return welcome; |
| 244 | } |
| 245 | if let Some(version) = live.other_version() { |
| 246 | return Err(Refusal::Version { |
| 247 | theirs_newer: version > wire::VERSION, |
| 248 | }); |
| 249 | } |
| 250 | if live.failed() > 0 { |
| 251 | return Err(Refusal::Wrong); |
| 252 | } |
| 253 | // A peer on this network may yet answer what the relay refused. |
| 254 | let waited = start.elapsed(); |
| 255 | let settled = waited > Duration::from_secs(3) || reach.is_none(); |
| 256 | match live.relayed() { |
| 257 | Relayed::Refused(404, _) if settled => return Err(Refusal::NoOne), |
| 258 | Relayed::Refused(410, _) => return Err(Refusal::Expired), |
| 259 | // A short wait, as a crowd joining at once meets, passes as the relay is asked |
| 260 | // again. |
| 261 | Relayed::Refused(429, wait) if wait.is_none_or(|wait| wait > PATIENT) => { |
| 262 | return Err(Refusal::TooMany(wait)); |
| 263 | } |
| 264 | Relayed::Refused(503, _) if settled => return Err(Refusal::Busy), |
| 265 | Relayed::Unreachable(trouble) if waited > Duration::from_secs(10) => { |
| 266 | return Err(Refusal::Unreachable(trouble)); |
| 267 | } |
| 268 | Relayed::Unknown if relay.is_none() && waited > Duration::from_secs(10) => { |
| 269 | return Err(Refusal::NoOne); |
| 270 | } |
| 271 | _ if waited |
| 272 | > Duration::from_secs(if approving.load(Ordering::Acquire) { |
| 273 | 300 |
| 274 | } else { |
| 275 | 20 |
| 276 | }) => |
| 277 | { |
| 278 | return Err(Refusal::TimedOut); |
| 279 | } |
| 280 | _ => {} |
| 281 | } |
| 282 | crate::task::wait(&waiting, Some(Duration::from_millis(100))).await; |
| 283 | } |
| 284 | } |
| 285 | |
| 286 | /// A host's side of its guests' storage requests. |
| 287 | struct Served { |
| 288 | storage: Box<dyn Storage>, |
| 289 | /// Sections' images by path, with their stamps, to check commits on. |
| 290 | images: Mutex<Vec<Image>>, |
| 291 | /// Files being read a chunk at a time, by guest and the read's first request. |
| 292 | snapshots: Mutex<ByGuest<(Arc<Vec<u8>>, Instant)>>, |
| 293 | /// Bytes a later request carries, by guest and upload. |
| 294 | puts: Mutex<ByGuest<Vec<u8>>>, |
| 295 | guests: Mutex<BTreeMap<[u8; 16], Admitted>>, |
| 296 | writers: Mutex<HashMap<String, mpsc::SyncSender<batch::Waiting>>>, |
| 297 | /// Hears the paths guests changed, as the host's own notebook should. |
| 298 | host: Mutex<Option<crate::session::Listener>>, |
| 299 | /// The share's room, to tell guests what changed. |
| 300 | room: Mutex<Option<Sender>>, |
| 301 | } |
| 302 | |
| 303 | /// A guest as its host serves it: the line to it, its requests waiting for its workers, how |
| 304 | /// many more it may start now, and the sections it read lately, newest first, as it keeps |
| 305 | /// their images. |
| 306 | struct Admitted { |
| 307 | line: Line, |
| 308 | queue: mpsc::SyncSender<(u16, Vec<u8>)>, |
| 309 | starts: f64, |
| 310 | counted: Instant, |
| 311 | held: Vec<String>, |
| 312 | } |
| 313 | |
| 314 | /// A section's path, stamp and image. |
| 315 | type Image = (String, Stamp, Arc<Vec<u8>>); |
| 316 | /// What a guest's requests left, by guest and request. |
| 317 | type ByGuest<T> = HashMap<([u8; 16], u64), T>; |
| 318 | |
| 319 | /// Whether a guest may name `path`: a catalog path inside the notebook, and not presence's |
| 320 | /// own secret, which guests meet in the share's room instead. |
| 321 | fn allowed(path: &str) -> bool { |
| 322 | path.split('/').all(|part| { |
| 323 | !part.is_empty() && part != "." && part != ".." && !part.contains(['\\', '\0', ':']) |
| 324 | }) && !path.to_ascii_lowercase().starts_with(".snowbound/live") |
| 325 | } |
| 326 | |
| 327 | /// The folder holding catalog path `path`. |
| 328 | fn folder(path: &str) -> String { |
| 329 | path.rsplit_once('/') |
| 330 | .map_or(String::new(), |(folder, _)| folder.to_owned()) |
| 331 | } |
| 332 | |
| 333 | fn refused(kind: io::ErrorKind, message: &str) -> Error { |
| 334 | io::Error::new(kind, message.to_owned()).into() |
| 335 | } |
| 336 | |
| 337 | impl Served { |
| 338 | /// Serves `peer` on `line`, through workers of its own that end as it leaves. |
| 339 | fn admit(self: &Arc<Self>, peer: [u8; 16], line: &Line) { |
| 340 | let (queue, waiting) = mpsc::sync_channel::<(u16, Vec<u8>)>(QUEUED - WORKERS); |
| 341 | let waiting = Arc::new(Mutex::new(waiting)); |
| 342 | for _ in 0..WORKERS { |
| 343 | let (served, waiting, line) = |
| 344 | (Arc::downgrade(self), Arc::clone(&waiting), line.clone()); |
| 345 | thread::spawn(move || { |
| 346 | loop { |
| 347 | let next = waiting.lock().unwrap().recv(); |
| 348 | let (Ok((kind, body)), Some(served)) = (next, served.upgrade()) else { |
| 349 | return; |
| 350 | }; |
| 351 | let _ = line.send(kind::REPLY, &served.handle(&peer, kind, &body)); |
| 352 | } |
| 353 | }); |
| 354 | } |
| 355 | let guest = Admitted { |
| 356 | line: line.clone(), |
| 357 | queue, |
| 358 | starts: BURST, |
| 359 | counted: Instant::now(), |
| 360 | held: Vec::new(), |
| 361 | }; |
| 362 | self.guests.lock().unwrap().insert(peer, guest); |
| 363 | } |
| 364 | |
| 365 | /// Hands a request from `peer` to its workers, or hangs up on a guest that has too many |
| 366 | /// waiting or starts them too fast. |
| 367 | fn queue(&self, peer: &[u8; 16], kind: u16, body: &[u8]) { |
| 368 | let mut guests = self.guests.lock().unwrap(); |
| 369 | let Some(guest) = guests.get_mut(peer) else { |
| 370 | return; |
| 371 | }; |
| 372 | let continued = kind == kind::PUT |
| 373 | || matches!(kind, kind::READ | kind::READ_FILE) |
| 374 | && minicbor::decode::<Request>(body).is_ok_and(|request| request.handle.is_some()); |
| 375 | let now = Instant::now(); |
| 376 | guest.starts = |
| 377 | (guest.starts + now.duration_since(guest.counted).as_secs_f64() * STARTS).min(BURST); |
| 378 | guest.counted = now; |
| 379 | if !continued { |
| 380 | guest.starts -= 1.0; |
| 381 | } |
| 382 | if guest.starts < 0.0 || guest.queue.try_send((kind, body.to_vec())).is_err() { |
| 383 | guest.line.hang_up(FLOODED); |
| 384 | guests.remove(peer); |
| 385 | } |
| 386 | } |
| 387 | |
| 388 | fn forget(&self, peer: &[u8; 16]) { |
| 389 | self.guests.lock().unwrap().remove(peer); |
| 390 | self.snapshots |
| 391 | .lock() |
| 392 | .unwrap() |
| 393 | .retain(|(guest, _), _| guest != peer); |
| 394 | self.puts |
| 395 | .lock() |
| 396 | .unwrap() |
| 397 | .retain(|(guest, _), _| guest != peer); |
| 398 | } |
| 399 | |
| 400 | /// Tells every guest, and the host, that a guest changed the files at `paths`. |
| 401 | fn changed(&self, paths: &[String]) { |
| 402 | self.tell(paths); |
| 403 | if let Some(listener) = &*self.host.lock().unwrap() { |
| 404 | listener(paths); |
| 405 | } |
| 406 | } |
| 407 | |
| 408 | /// Tells every guest the files at `paths` changed, and what they are now. |
| 409 | fn tell(&self, paths: &[String]) { |
| 410 | let touched = Touched { |
| 411 | paths: paths.to_vec(), |
| 412 | stamps: (paths.iter()) |
| 413 | .map(|path| self.storage.stamp(path).ok().map(|stamp| digest(&stamp))) |
| 414 | .collect(), |
| 415 | }; |
| 416 | if let Some(room) = self.room.lock().unwrap().as_ref() { |
| 417 | room.send(kind::TOUCHED, &touched, None); |
| 418 | } |
| 419 | } |
| 420 | |
| 421 | /// Notes that `guest` holds the section at `path`, as it now keeps its image. |
| 422 | fn holds(guest: &mut Admitted, path: &str) { |
| 423 | guest.held.retain(|held| held != path); |
| 424 | guest.held.insert(0, path.to_owned()); |
| 425 | guest.held.truncate(IMAGES); |
| 426 | } |
| 427 | |
| 428 | fn handle(self: &Arc<Self>, peer: &[u8; 16], kind: u16, body: &[u8]) -> Reply { |
| 429 | let request = match minicbor::decode::<Request>(body) { |
| 430 | Ok(request) => request, |
| 431 | Err(_) => { |
| 432 | return failed( |
| 433 | 0, |
| 434 | refused(io::ErrorKind::InvalidData, "A malformed request"), |
| 435 | ); |
| 436 | } |
| 437 | }; |
| 438 | let id = request.id; |
| 439 | if !self.guests.lock().unwrap().contains_key(peer) { |
| 440 | return failed( |
| 441 | id, |
| 442 | refused( |
| 443 | io::ErrorKind::PermissionDenied, |
| 444 | "The device is no longer connected", |
| 445 | ), |
| 446 | ); |
| 447 | } |
| 448 | match self.answer(peer, kind, request) { |
| 449 | Ok(reply) => Reply { id, ..reply }, |
| 450 | Err(error) => failed(id, error), |
| 451 | } |
| 452 | } |
| 453 | |
| 454 | fn answer(self: &Arc<Self>, peer: &[u8; 16], kind: u16, request: Request) -> Result<Reply> { |
| 455 | let path = request.path.as_str(); |
| 456 | if !(path.is_empty() && matches!(kind, kind::LIST | kind::PUT) || allowed(path)) |
| 457 | || request.to.as_deref().is_some_and(|to| !allowed(to)) |
| 458 | { |
| 459 | return Err(refused( |
| 460 | io::ErrorKind::PermissionDenied, |
| 461 | "Outside the notebook", |
| 462 | )); |
| 463 | } |
| 464 | let to = || { |
| 465 | request |
| 466 | .to |
| 467 | .as_deref() |
| 468 | .ok_or_else(|| refused(io::ErrorKind::InvalidInput, "No target")) |
| 469 | }; |
| 470 | let stamp = || -> Result<Stamp> { |
| 471 | Ok(request |
| 472 | .stamp |
| 473 | .as_ref() |
| 474 | .ok_or_else(|| refused(io::ErrorKind::InvalidInput, "No stamp"))? |
| 475 | .try_into()?) |
| 476 | }; |
| 477 | let done = Reply::default(); |
| 478 | Ok(match kind { |
| 479 | kind::LIST => Reply { |
| 480 | entries: Some( |
| 481 | self.storage |
| 482 | .entries(path)? |
| 483 | .into_iter() |
| 484 | .filter(|entry| allowed(&within(path, &entry.name))) |
| 485 | .map(|entry| wire_entry(&entry)) |
| 486 | .collect(), |
| 487 | ), |
| 488 | ..done |
| 489 | }, |
| 490 | kind::STAMP => Reply { |
| 491 | stamp: Some((&self.storage.stamp(path)?).into()), |
| 492 | ..done |
| 493 | }, |
| 494 | kind::EXISTS => Reply { |
| 495 | exists: Some(self.storage.exists(path)), |
| 496 | ..done |
| 497 | }, |
| 498 | kind::READ | kind::READ_FILE => self.read(peer, kind, &request)?, |
| 499 | kind::PUT => { |
| 500 | let handle = request.handle.unwrap_or_default(); |
| 501 | let bytes = request.bytes.unwrap_or_default(); |
| 502 | let mut puts = self.puts.lock().unwrap(); |
| 503 | let key = (*peer, handle); |
| 504 | let length = puts.get(&key).map_or(0, Vec::len); |
| 505 | if request.offset != Some(length as u64) || bytes.is_empty() { |
| 506 | puts.remove(&key); |
| 507 | return Err(refused( |
| 508 | io::ErrorKind::InvalidInput, |
| 509 | "An upload out of order", |
| 510 | )); |
| 511 | } |
| 512 | let held: usize = puts |
| 513 | .iter() |
| 514 | .filter(|((guest, _), _)| guest == peer) |
| 515 | .map(|(_, bytes)| bytes.len()) |
| 516 | .sum(); |
| 517 | if bytes.len() > LIMIT.saturating_sub(held) { |
| 518 | puts.remove(&key); |
| 519 | return Err(refused( |
| 520 | io::ErrorKind::FileTooLarge, |
| 521 | "An upload is too large", |
| 522 | )); |
| 523 | } |
| 524 | puts.entry(key).or_default().extend_from_slice(&bytes); |
| 525 | done |
| 526 | } |
| 527 | kind::COMMIT => { |
| 528 | let transaction = Transaction::from_bytes(&self.carried(peer, &request)?)?; |
| 529 | self.commit(peer, path, &transaction)?; |
| 530 | self.changed(&[path.to_owned()]); |
| 531 | done |
| 532 | } |
| 533 | kind::EDITS => self.batch(peer, request)?, |
| 534 | kind::CONFIRM => { |
| 535 | self.storage.confirm(path, &stamp()?)?; |
| 536 | done |
| 537 | } |
| 538 | kind::CREATE => { |
| 539 | self.storage.create(path, &self.carried(peer, &request)?)?; |
| 540 | self.changed(&[folder(path)]); |
| 541 | done |
| 542 | } |
| 543 | kind::CREATE_DIRECTORY => { |
| 544 | self.storage.create_directory(path)?; |
| 545 | self.changed(&[folder(path)]); |
| 546 | done |
| 547 | } |
| 548 | kind::HIDE => { |
| 549 | self.storage.hide(path)?; |
| 550 | done |
| 551 | } |
| 552 | kind::RENAME | kind::REPLACE => { |
| 553 | let to = to()?; |
| 554 | if kind == kind::RENAME { |
| 555 | self.storage.rename(path, to)?; |
| 556 | } else { |
| 557 | self.storage.replace(path, to)?; |
| 558 | } |
| 559 | self.changed(&[folder(path), folder(to)]); |
| 560 | done |
| 561 | } |
| 562 | kind::DELETE => { |
| 563 | self.storage.delete(path)?; |
| 564 | self.changed(&[folder(path)]); |
| 565 | done |
| 566 | } |
| 567 | kind::PLACE => { |
| 568 | let ancestor = request |
| 569 | .ancestor |
| 570 | .ok_or_else(|| refused(io::ErrorKind::InvalidInput, "No ancestor"))?; |
| 571 | let name = request.name.as_deref().unwrap_or_default(); |
| 572 | self.storage.place(path, ancestor, name)?; |
| 573 | self.changed(&[path.to_owned()]); |
| 574 | done |
| 575 | } |
| 576 | kind::SUPERSEDE => { |
| 577 | self.storage.supersede(path, &stamp()?, to()?)?; |
| 578 | self.changed(&[folder(path)]); |
| 579 | done |
| 580 | } |
| 581 | _ => return Err(refused(io::ErrorKind::Unsupported, "An unknown request")), |
| 582 | }) |
| 583 | } |
| 584 | |
| 585 | /// The bytes a request carries, itself or in its uploads. |
| 586 | fn carried(&self, peer: &[u8; 16], request: &Request) -> Result<Vec<u8>> { |
| 587 | match (&request.bytes, request.handle) { |
| 588 | (Some(bytes), _) => Ok(bytes.clone()), |
| 589 | (None, Some(handle)) => self |
| 590 | .puts |
| 591 | .lock() |
| 592 | .unwrap() |
| 593 | .remove(&(*peer, handle)) |
| 594 | .ok_or_else(|| refused(io::ErrorKind::InvalidInput, "No such upload")), |
| 595 | (None, None) => Ok(Vec::new()), |
| 596 | } |
| 597 | } |
| 598 | |
| 599 | /// A chunk of a file as it stood when its read began. |
| 600 | fn read(&self, peer: &[u8; 16], kind: u16, request: &Request) -> Result<Reply> { |
| 601 | let offset = request.offset.unwrap_or_default() as usize; |
| 602 | let mut snapshots = self.snapshots.lock().unwrap(); |
| 603 | snapshots.retain(|_, (_, read)| read.elapsed() < SNAPSHOT_AGE); |
| 604 | if let (kind::READ, None, Some(base)) = (kind, request.handle, &request.stamp) |
| 605 | && let Ok(base) = Stamp::try_from(base) |
| 606 | && let Some(changes) = self.changes(&request.path, &base)? |
| 607 | { |
| 608 | if let Some(guest) = self.guests.lock().unwrap().get_mut(peer) { |
| 609 | Self::holds(guest, &request.path); |
| 610 | } |
| 611 | return Ok(changes); |
| 612 | } |
| 613 | let (image, handle) = match request.handle { |
| 614 | Some(handle) => { |
| 615 | let (image, read) = snapshots |
| 616 | .get_mut(&(*peer, handle)) |
| 617 | .ok_or_else(|| refused(io::ErrorKind::TimedOut, "The read went stale"))?; |
| 618 | *read = Instant::now(); |
| 619 | (Arc::clone(image), handle) |
| 620 | } |
| 621 | None => { |
| 622 | drop(snapshots); |
| 623 | if kind == kind::READ |
| 624 | && let Some(guest) = self.guests.lock().unwrap().get_mut(peer) |
| 625 | { |
| 626 | Self::holds(guest, &request.path); |
| 627 | } |
| 628 | let limit = (request.limit.unwrap_or(LIMIT as u64) as usize).min(LIMIT); |
| 629 | let image = match kind { |
| 630 | kind::READ => self.image(&request.path)?, |
| 631 | _ => Arc::new(self.storage.read_file(&request.path, limit)?), |
| 632 | }; |
| 633 | if image.len() > limit { |
| 634 | return Err(io::Error::from(io::ErrorKind::FileTooLarge).into()); |
| 635 | } |
| 636 | snapshots = self.snapshots.lock().unwrap(); |
| 637 | if snapshots.keys().filter(|(guest, _)| guest == peer).count() >= SNAPSHOTS { |
| 638 | snapshots.retain(|(guest, _), _| guest != peer); |
| 639 | } |
| 640 | (image, request.id) |
| 641 | } |
| 642 | }; |
| 643 | let end = image.len().min(offset.saturating_add(CHUNK)); |
| 644 | let bytes = image.get(offset..end).unwrap_or_default().to_vec(); |
| 645 | if end < image.len() { |
| 646 | snapshots.insert((*peer, handle), (Arc::clone(&image), Instant::now())); |
| 647 | } else { |
| 648 | snapshots.remove(&(*peer, handle)); |
| 649 | } |
| 650 | Ok(Reply { |
| 651 | bytes: Some(bytes), |
| 652 | length: Some(image.len() as u64), |
| 653 | handle: Some(handle), |
| 654 | stamp: (kind == kind::READ) |
| 655 | .then(|| Stamp::of(&image).ok()) |
| 656 | .flatten() |
| 657 | .map(|stamp| (&stamp).into()), |
| 658 | ..Reply::default() |
| 659 | }) |
| 660 | } |
| 661 | |
| 662 | /// What changed in the section at `path` since the image with stamp `base`, where that |
| 663 | /// image is kept and the changes are much smaller than the section. |
| 664 | fn changes(&self, path: &str, base: &Stamp) -> Result<Option<Reply>> { |
| 665 | let before = self |
| 666 | .images |
| 667 | .lock() |
| 668 | .unwrap() |
| 669 | .iter() |
| 670 | .find_map(|(held, at, image)| (held == path && at == base).then(|| Arc::clone(image))); |
| 671 | let Some(before) = before else { |
| 672 | return Ok(None); |
| 673 | }; |
| 674 | let after = self.image(path)?; |
| 675 | let writes = match delta(&before, &after) { |
| 676 | Some(writes) => writes, |
| 677 | None if before == after => Vec::new(), |
| 678 | None => return Ok(None), |
| 679 | }; |
| 680 | Ok(Some(Reply { |
| 681 | writes: Some(writes), |
| 682 | length: Some(after.len() as u64), |
| 683 | stamp: Some((&Stamp::of(&after)?).into()), |
| 684 | ..Reply::default() |
| 685 | })) |
| 686 | } |
| 687 | |
| 688 | /// The section or TOC at `path` as it stands: the image kept for it while its stamp holds. |
| 689 | fn image(&self, path: &str) -> Result<Arc<Vec<u8>>> { |
| 690 | let stamp = self.storage.stamp(path)?; |
| 691 | let kept = self |
| 692 | .images |
| 693 | .lock() |
| 694 | .unwrap() |
| 695 | .iter() |
| 696 | .find_map(|(held, at, image)| { |
| 697 | (held == path && *at == stamp).then(|| Arc::clone(image)) |
| 698 | }); |
| 699 | if let Some(image) = kept { |
| 700 | return Ok(image); |
| 701 | } |
| 702 | let image = Arc::new(self.storage.read(path)?); |
| 703 | self.keep(path, Arc::clone(&image)); |
| 704 | Ok(image) |
| 705 | } |
| 706 | |
| 707 | /// Keeps `image` as the section at `path` now, with the newest of each of `IMAGES` |
| 708 | /// sections and earlier images up to `VERSIONS` and `VERSIONS_BYTES`. |
| 709 | fn keep(&self, path: &str, image: Arc<Vec<u8>>) { |
| 710 | let Ok(stamp) = Stamp::of(&image) else { |
| 711 | return; |
| 712 | }; |
| 713 | let mut images = self.images.lock().unwrap(); |
| 714 | images.retain(|(held, at, _)| held != path || *at != stamp); |
| 715 | images.insert(0, (path.to_owned(), stamp, image)); |
| 716 | let (mut newest, mut versions, mut bytes) = (Vec::new(), 0, 0); |
| 717 | images.retain(|(held, _, image)| { |
| 718 | if !newest.contains(held) { |
| 719 | newest.push(held.clone()); |
| 720 | return newest.len() <= IMAGES; |
| 721 | } |
| 722 | versions += 1; |
| 723 | bytes += image.len(); |
| 724 | newest.iter().take(IMAGES).any(|kept| kept == held) |
| 725 | && versions <= VERSIONS |
| 726 | && bytes <= VERSIONS_BYTES |
| 727 | }); |
| 728 | } |
| 729 | |
| 730 | /// Commits `guest`'s transaction once the section it makes parses. |
| 731 | fn commit(&self, guest: &[u8; 16], path: &str, transaction: &Transaction) -> Result<()> { |
| 732 | let not_committed = |error: io::Error| { |
| 733 | Error::Remote(CommitError { |
| 734 | state: CommitState::NotCommitted, |
| 735 | error, |
| 736 | }) |
| 737 | }; |
| 738 | let image = self.image(path).map_err(|error| match error { |
| 739 | Error::Io(error) => not_committed(error), |
| 740 | error => error, |
| 741 | })?; |
| 742 | if Stamp::of(&image).ok().as_ref() != Some(transaction.base()) { |
| 743 | return Err(not_committed(io::Error::new( |
| 744 | io::ErrorKind::ResourceBusy, |
| 745 | "The section changed since", |
| 746 | ))); |
| 747 | } |
| 748 | let mut next = (*image).clone(); |
| 749 | let checked = transaction.apply(&mut next).and_then(|()| { |
| 750 | let store = Store::parse(&next)?; |
| 751 | RevisionIndex::parse(&store).map(drop) |
| 752 | }); |
| 753 | if let Err(error) = checked { |
| 754 | return Err(not_committed(io::Error::new( |
| 755 | io::ErrorKind::InvalidData, |
| 756 | error.to_string(), |
| 757 | ))); |
| 758 | } |
| 759 | self.storage.commit(path, transaction)?; |
| 760 | self.tell_delta(path, &image, &next, Some(guest)); |
| 761 | self.keep(path, Arc::new(next)); |
| 762 | Ok(()) |
| 763 | } |
| 764 | |
| 765 | /// Tells the guests holding the section at `path` what a commit changed in it, from |
| 766 | /// `before` to `after`, where that is small enough to send: all but the guest whose |
| 767 | /// commit it was, which has it already. |
| 768 | fn tell_delta(&self, path: &str, before: &[u8], after: &[u8], committed: Option<&[u8; 16]>) { |
| 769 | let (Ok(base), Some(writes)) = (Stamp::of(before), delta(before, after)) else { |
| 770 | return; |
| 771 | }; |
| 772 | let delta = Delta { |
| 773 | path: path.to_owned(), |
| 774 | base: (&base).into(), |
| 775 | length: after.len() as u64, |
| 776 | writes, |
| 777 | }; |
| 778 | let mut guests = self.guests.lock().unwrap(); |
| 779 | let holding: Vec<[u8; 16]> = guests |
| 780 | .iter_mut() |
| 781 | .filter(|(id, guest)| Some(*id) != committed && guest.held.iter().any(|h| h == path)) |
| 782 | .map(|(id, guest)| { |
| 783 | Self::holds(guest, path); |
| 784 | *id |
| 785 | }) |
| 786 | .collect(); |
| 787 | drop(guests); |
| 788 | if let (Some(room), false) = (self.room.lock().unwrap().as_ref(), holding.is_empty()) { |
| 789 | room.send(kind::DELTA, &delta, Some(&holding)); |
| 790 | } |
| 791 | } |
| 792 | |
| 793 | /// The host's own change to the section at `path`, which guests hear as a delta where |
| 794 | /// the host kept the image before it. |
| 795 | fn changed_here(&self, path: &str) { |
| 796 | let kept = self |
| 797 | .images |
| 798 | .lock() |
| 799 | .unwrap() |
| 800 | .iter() |
| 801 | .find_map(|(held, at, image)| (held == path).then(|| (at.clone(), Arc::clone(image)))); |
| 802 | let Some((at, before)) = kept else { |
| 803 | return; |
| 804 | }; |
| 805 | if self.storage.stamp(path).is_ok_and(|now| now == at) { |
| 806 | return; |
| 807 | } |
| 808 | if let Ok(after) = self.storage.read(path) { |
| 809 | self.tell_delta(path, &before, &after, None); |
| 810 | self.keep(path, Arc::new(after)); |
| 811 | } |
| 812 | } |
| 813 | } |
| 814 | |
| 815 | /// A stamp as `Touched` names it. |
| 816 | fn digest(stamp: &Stamp) -> u64 { |
| 817 | let hash = Sha256::new() |
| 818 | .chain_update(stamp.header) |
| 819 | .chain_update(stamp.length.to_le_bytes()) |
| 820 | .finalize(); |
| 821 | u64::from_le_bytes(hash[..8].try_into().expect("eight bytes")) |
| 822 | } |
| 823 | |
| 824 | /// The writes that make `after` of `before`, a commit's appended bytes, patches and header, |
| 825 | /// where their encoding fits a chunk and they are much less than `after` itself. |
| 826 | fn delta(before: &[u8], after: &[u8]) -> Option<Vec<Written>> { |
| 827 | const BLOCK: usize = 4096; |
| 828 | if before.len() < 1024 || after.len() < before.len() { |
| 829 | return None; |
| 830 | } |
| 831 | let mut writes = vec![Written { |
| 832 | offset: 0, |
| 833 | bytes: after[..1024].to_vec(), |
| 834 | }]; |
| 835 | let mut at = 1024; |
| 836 | while at < before.len() { |
| 837 | let end = (at + BLOCK).min(before.len()); |
| 838 | if before[at..end] != after[at..end] { |
| 839 | let first = (at..end).find(|&i| before[i] != after[i]).unwrap_or(at); |
| 840 | let last = (at..end).rfind(|&i| before[i] != after[i]).unwrap_or(first) + 1; |
| 841 | match writes.last_mut() { |
| 842 | Some(write) if write.offset as usize + write.bytes.len() + 64 >= first => { |
| 843 | let from = write.offset as usize; |
| 844 | write.bytes = after[from..last].to_vec(); |
| 845 | } |
| 846 | _ => writes.push(Written { |
| 847 | offset: first as u64, |
| 848 | bytes: after[first..last].to_vec(), |
| 849 | }), |
| 850 | } |
| 851 | } |
| 852 | at = end; |
| 853 | } |
| 854 | if after.len() > before.len() { |
| 855 | writes.push(Written { |
| 856 | offset: before.len() as u64, |
| 857 | bytes: after[before.len()..].to_vec(), |
| 858 | }); |
| 859 | } |
| 860 | let sent: usize = writes.iter().map(|write| write.bytes.len()).sum(); |
| 861 | (sent <= after.len() / 2 && sent <= CHUNK && minicbor::to_vec(&writes).ok()?.len() <= CHUNK) |
| 862 | .then_some(writes) |
| 863 | } |
| 864 | |
| 865 | /// `image` with `writes`, `length` long: none where a write falls outside it. |
| 866 | fn written(image: &[u8], length: u64, writes: &[Written]) -> Option<Vec<u8>> { |
| 867 | let length = usize::try_from(length) |
| 868 | .ok() |
| 869 | .filter(|length| *length >= image.len())?; |
| 870 | let mut next = image.to_vec(); |
| 871 | next.resize(length, 0); |
| 872 | for write in writes { |
| 873 | let offset = usize::try_from(write.offset).ok()?; |
| 874 | next.get_mut(offset..offset.checked_add(write.bytes.len())?)? |
| 875 | .copy_from_slice(&write.bytes); |
| 876 | } |
| 877 | Some(next) |
| 878 | } |
| 879 | |
| 880 | fn within(folder: &str, name: &str) -> String { |
| 881 | if folder.is_empty() { |
| 882 | name.to_owned() |
| 883 | } else { |
| 884 | format!("{folder}/{name}") |
| 885 | } |
| 886 | } |
| 887 | |
| 888 | fn wire_entry(entry: &discover::Entry) -> WireEntry { |
| 889 | WireEntry { |
| 890 | name: entry.name.clone(), |
| 891 | kind: match entry.kind { |
| 892 | discover::EntryKind::File => 0, |
| 893 | discover::EntryKind::Directory => 1, |
| 894 | discover::EntryKind::Other => 2, |
| 895 | discover::EntryKind::Evicted => 3, |
| 896 | }, |
| 897 | size: entry.listed.size, |
| 898 | modified: entry.listed.modified, |
| 899 | } |
| 900 | } |
| 901 | |
| 902 | fn entry(entry: &WireEntry) -> discover::Entry { |
| 903 | discover::Entry { |
| 904 | name: entry.name.clone(), |
| 905 | kind: match entry.kind { |
| 906 | 0 => discover::EntryKind::File, |
| 907 | 1 => discover::EntryKind::Directory, |
| 908 | 3 => discover::EntryKind::Evicted, |
| 909 | _ => discover::EntryKind::Other, |
| 910 | }, |
| 911 | listed: discover::Listed { |
| 912 | size: entry.size, |
| 913 | modified: entry.modified, |
| 914 | }, |
| 915 | } |
| 916 | } |
| 917 | |
| 918 | fn state_number(state: CommitState) -> u8 { |
| 919 | match state { |
| 920 | CommitState::NotCommitted => 0, |
| 921 | CommitState::Unknown => 1, |
| 922 | CommitState::Committed => 2, |
| 923 | } |
| 924 | } |
| 925 | |
| 926 | fn failed(id: u64, error: Error) -> Reply { |
| 927 | let (kind, state) = match &error { |
| 928 | Error::Io(error) | Error::RemoteIo(error) => (error.kind(), None), |
| 929 | Error::Remote(error) => (error.error.kind(), Some(state_number(error.state))), |
| 930 | Error::Document(_) => (io::ErrorKind::InvalidData, None), |
| 931 | _ => (io::ErrorKind::Other, None), |
| 932 | }; |
| 933 | let message = match &error { |
| 934 | Error::Remote(error) => error.error.to_string(), |
| 935 | error => error.to_string(), |
| 936 | }; |
| 937 | Reply { |
| 938 | id, |
| 939 | failure: Some(Failure { |
| 940 | kind: wire::error_number(kind), |
| 941 | message, |
| 942 | state, |
| 943 | }), |
| 944 | ..Reply::default() |
| 945 | } |
| 946 | } |
| 947 | |
| 948 | /// How a request failed. |
| 949 | enum Failed { |
| 950 | /// It never left this end. |
| 951 | Unsent(io::Error), |
| 952 | /// Its answer was lost: whatever it asked may have happened. |
| 953 | Lost(io::Error), |
| 954 | /// The host refused it. |
| 955 | Refused(Failure), |
| 956 | } |
| 957 | |
| 958 | impl Failed { |
| 959 | fn io(self) -> io::Error { |
| 960 | match self { |
| 961 | Failed::Unsent(error) | Failed::Lost(error) => error, |
| 962 | Failed::Refused(failure) => { |
| 963 | io::Error::new(wire::error_kind(failure.kind), failure.message) |
| 964 | } |
| 965 | } |
| 966 | } |
| 967 | |
| 968 | /// As a commit's failure: one never sent was not committed, one whose answer was lost may |
| 969 | /// have been. |
| 970 | fn commit(self) -> CommitError { |
| 971 | let state = match &self { |
| 972 | Failed::Unsent(_) => CommitState::NotCommitted, |
| 973 | Failed::Lost(_) => CommitState::Unknown, |
| 974 | Failed::Refused(failure) => match failure.state { |
| 975 | Some(1) => CommitState::Unknown, |
| 976 | Some(2) => CommitState::Committed, |
| 977 | _ => CommitState::NotCommitted, |
| 978 | }, |
| 979 | }; |
| 980 | CommitError { |
| 981 | state, |
| 982 | error: self.io(), |
| 983 | } |
| 984 | } |
| 985 | } |
| 986 | |
| 987 | /// A guest's way to a share: the share's room, and the requests waiting on its host. |
| 988 | pub struct Guest { |
| 989 | live: Live, |
| 990 | inner: Arc<Inner>, |
| 991 | } |
| 992 | |
| 993 | struct Inner { |
| 994 | share: [u8; 16], |
| 995 | me: Hello, |
| 996 | reach: Option<Reach>, |
| 997 | relay: Option<String>, |
| 998 | events: Arc<dyn Fn() + Send + Sync>, |
| 999 | presence: Mutex<Option<([u8; 16], Live)>>, |
| 1000 | here: Mutex<Presence>, |
| 1001 | /// The host and the line to it, while connected. |
| 1002 | host: Mutex<Option<(Arc<Hello>, Line)>>, |
| 1003 | pending: Mutex<HashMap<u64, crate::task::Answer<Reply>>>, |
| 1004 | next: AtomicU64, |
| 1005 | /// Where the host's reports of changed files go, while a background watches. |
| 1006 | watch: Mutex<Option<Arc<Reports>>>, |
| 1007 | /// Why this device can no longer reach the share. |
| 1008 | ended: Mutex<Option<Ended>>, |
| 1009 | /// Bytes of chunks asked for and not yet given. |
| 1010 | asked: (Mutex<Chunks>, Condvar), |
| 1011 | /// The sections read lately, kept as the host's deltas change them, newest first. |
| 1012 | images: Mutex<Vec<(String, Arc<Vec<u8>>)>>, |
| 1013 | /// The stamps of sections whose image held is the host's now, as its last report of |
| 1014 | /// them said. |
| 1015 | current: Mutex<HashMap<String, Stamp>>, |
| 1016 | #[cfg(target_arch = "wasm32")] |
| 1017 | catalog: Mutex<Option<PathBuf>>, |
| 1018 | } |
| 1019 | |
| 1020 | #[derive(Clone, Copy, Debug, PartialEq, Eq)] |
| 1021 | pub enum Ended { |
| 1022 | Stopped, |
| 1023 | Removed, |
| 1024 | } |
| 1025 | |
| 1026 | #[derive(Default)] |
| 1027 | struct Chunks { |
| 1028 | bytes: usize, |
| 1029 | #[cfg(target_arch = "wasm32")] |
| 1030 | waiting: Vec<std::task::Waker>, |
| 1031 | } |
| 1032 | |
| 1033 | struct Requested<'a> { |
| 1034 | inner: &'a Inner, |
| 1035 | id: u64, |
| 1036 | chunked: bool, |
| 1037 | } |
| 1038 | |
| 1039 | impl Drop for Requested<'_> { |
| 1040 | fn drop(&mut self) { |
| 1041 | self.inner.pending.lock().unwrap().remove(&self.id); |
| 1042 | if self.chunked { |
| 1043 | let (asked, room) = &self.inner.asked; |
| 1044 | let mut asked = asked.lock().unwrap(); |
| 1045 | asked.bytes -= CHUNK; |
| 1046 | #[cfg(target_arch = "wasm32")] |
| 1047 | let waiting = std::mem::take(&mut asked.waiting); |
| 1048 | drop(asked); |
| 1049 | room.notify_all(); |
| 1050 | #[cfg(target_arch = "wasm32")] |
| 1051 | for waker in waiting { |
| 1052 | waker.wake(); |
| 1053 | } |
| 1054 | } |
| 1055 | } |
| 1056 | } |
| 1057 | |
| 1058 | impl Guest { |
| 1059 | /// Joins share `share` through its `secret` as `me`, where `reach` and `relay` say. |
| 1060 | /// `events` runs on a network thread whenever the host comes or goes, or the peers change. |
| 1061 | pub fn start( |
| 1062 | me: Hello, |
| 1063 | share: [u8; 16], |
| 1064 | secret: [u8; 16], |
| 1065 | reach: Option<Reach>, |
| 1066 | relay: Option<&str>, |
| 1067 | events: impl Fn() + Send + Sync + 'static, |
| 1068 | ) -> io::Result<Arc<Self>> { |
| 1069 | let inner = Arc::new(Inner { |
| 1070 | share, |
| 1071 | me: me.clone(), |
| 1072 | reach, |
| 1073 | relay: relay.map(str::to_owned), |
| 1074 | events: Arc::new(events), |
| 1075 | presence: Mutex::default(), |
| 1076 | here: Mutex::default(), |
| 1077 | host: Mutex::default(), |
| 1078 | pending: Mutex::default(), |
| 1079 | next: AtomicU64::new(1), |
| 1080 | watch: Mutex::default(), |
| 1081 | ended: Mutex::default(), |
| 1082 | asked: Default::default(), |
| 1083 | images: Mutex::default(), |
| 1084 | current: Mutex::default(), |
| 1085 | #[cfg(target_arch = "wasm32")] |
| 1086 | catalog: Mutex::default(), |
| 1087 | }); |
| 1088 | let heard = Arc::clone(&inner); |
| 1089 | let live = Live::start(me, &Room::Notebook(secret), reach, relay, move |event| { |
| 1090 | heard.heard(event); |
| 1091 | (heard.events)(); |
| 1092 | })?; |
| 1093 | Ok(Arc::new(Self { live, inner })) |
| 1094 | } |
| 1095 | |
| 1096 | /// Where the notebook is kept: `location(share)`. |
| 1097 | pub fn location(&self) -> String { |
| 1098 | location(&self.inner.share) |
| 1099 | } |
| 1100 | |
| 1101 | /// The host, while connected. |
| 1102 | pub fn host(&self) -> Option<Arc<Hello>> { |
| 1103 | let host = self.inner.host.lock().unwrap(); |
| 1104 | host.as_ref().map(|(hello, _)| Arc::clone(hello)) |
| 1105 | } |
| 1106 | |
| 1107 | /// Why the host ended this device's access. |
| 1108 | pub fn ended(&self) -> Option<Ended> { |
| 1109 | *self.inner.ended.lock().unwrap() |
| 1110 | } |
| 1111 | |
| 1112 | /// Everyone in the share's room: the host and the other guests. |
| 1113 | pub fn peers(&self) -> Vec<Peer> { |
| 1114 | self.inner |
| 1115 | .presence |
| 1116 | .lock() |
| 1117 | .unwrap() |
| 1118 | .as_ref() |
| 1119 | .map(|(_, live)| live.peers()) |
| 1120 | .unwrap_or_default() |
| 1121 | } |
| 1122 | |
| 1123 | pub fn set_presence(&self, presence: Presence) { |
| 1124 | *self.inner.here.lock().unwrap() = presence.clone(); |
| 1125 | if let Some((_, live)) = &*self.inner.presence.lock().unwrap() { |
| 1126 | live.set_presence(presence); |
| 1127 | } |
| 1128 | } |
| 1129 | |
| 1130 | /// Sends the host's reports of changed files to `reports` while it stays connected. |
| 1131 | pub(crate) fn watch(&self, reports: Reports) -> io::Result<()> { |
| 1132 | if self.host().is_none() { |
| 1133 | return Err(self.offline()); |
| 1134 | } |
| 1135 | *self.inner.watch.lock().unwrap() = Some(Arc::new(reports)); |
| 1136 | Ok(()) |
| 1137 | } |
| 1138 | |
| 1139 | fn offline(&self) -> io::Error { |
| 1140 | if let Some(version) = self.live.other_version() { |
| 1141 | return io::Error::new(io::ErrorKind::NotConnected, wire::Version(version)); |
| 1142 | } |
| 1143 | let message = match &*self.inner.ended.lock().unwrap() { |
| 1144 | Some(Ended::Stopped) => "The host stopped sharing this notebook", |
| 1145 | Some(Ended::Removed) => { |
| 1146 | "This device was removed. Ask the person sharing for a new link or code." |
| 1147 | } |
| 1148 | None => "The computer sharing this notebook can’t be reached", |
| 1149 | }; |
| 1150 | io::Error::new(io::ErrorKind::NotConnected, message) |
| 1151 | } |
| 1152 | |
| 1153 | /// Asks the host `kind` of `request`, waiting for its reply. |
| 1154 | #[cfg(not(target_arch = "wasm32"))] |
| 1155 | fn request(&self, kind: u16, request: Request) -> std::result::Result<Reply, Failed> { |
| 1156 | crate::task::ready(self.request_async(kind, request)).map_err(Failed::Unsent)? |
| 1157 | } |
| 1158 | |
| 1159 | async fn request_async( |
| 1160 | &self, |
| 1161 | kind: u16, |
| 1162 | mut request: Request, |
| 1163 | ) -> std::result::Result<Reply, Failed> { |
| 1164 | let Some((_, line)) = self.inner.host.lock().unwrap().clone() else { |
| 1165 | return Err(Failed::Unsent(self.offline())); |
| 1166 | }; |
| 1167 | let chunked = matches!(kind, kind::READ | kind::READ_FILE | kind::PUT); |
| 1168 | if chunked { |
| 1169 | #[cfg(not(target_arch = "wasm32"))] |
| 1170 | { |
| 1171 | let (asked, _room) = &self.inner.asked; |
| 1172 | let mut asked = asked.lock().unwrap(); |
| 1173 | while asked.bytes + CHUNK > WINDOW { |
| 1174 | asked = _room.wait(asked).unwrap(); |
| 1175 | } |
| 1176 | asked.bytes += CHUNK; |
| 1177 | } |
| 1178 | #[cfg(target_arch = "wasm32")] |
| 1179 | std::future::poll_fn(|context| { |
| 1180 | let mut asked = self.inner.asked.0.lock().unwrap(); |
| 1181 | if asked.bytes + CHUNK > WINDOW { |
| 1182 | if !asked |
| 1183 | .waiting |
| 1184 | .iter() |
| 1185 | .any(|waker| waker.will_wake(context.waker())) |
| 1186 | { |
| 1187 | asked.waiting.push(context.waker().clone()); |
| 1188 | } |
| 1189 | std::task::Poll::Pending |
| 1190 | } else { |
| 1191 | asked.bytes += CHUNK; |
| 1192 | std::task::Poll::Ready(()) |
| 1193 | } |
| 1194 | }) |
| 1195 | .await; |
| 1196 | } |
| 1197 | let id = self.inner.next.fetch_add(1, Ordering::Relaxed); |
| 1198 | request.id = id; |
| 1199 | let (answer, answered) = crate::task::response(); |
| 1200 | self.inner.pending.lock().unwrap().insert(id, answer); |
| 1201 | let _requested = Requested { |
| 1202 | inner: &self.inner, |
| 1203 | id, |
| 1204 | chunked, |
| 1205 | }; |
| 1206 | let sent = line.send(kind, &request); |
| 1207 | let reply = match sent { |
| 1208 | Err(error) => Err(Failed::Unsent(error)), |
| 1209 | Ok(()) => match crate::task::answered(&answered, TIMEOUT).await { |
| 1210 | Ok(reply) => Ok(reply), |
| 1211 | Err(mpsc::RecvTimeoutError::Timeout) => Err(Failed::Lost(io::Error::new( |
| 1212 | io::ErrorKind::TimedOut, |
| 1213 | "The computer sharing this notebook didn’t answer", |
| 1214 | ))), |
| 1215 | Err(mpsc::RecvTimeoutError::Disconnected) => Err(Failed::Lost(self.offline())), |
| 1216 | }, |
| 1217 | }; |
| 1218 | let reply = reply?; |
| 1219 | match reply.failure { |
| 1220 | Some(failure) => Err(Failed::Refused(failure)), |
| 1221 | None => Ok(reply), |
| 1222 | } |
| 1223 | } |
| 1224 | |
| 1225 | #[cfg(not(target_arch = "wasm32"))] |
| 1226 | fn ask(&self, kind: u16, request: Request) -> io::Result<Reply> { |
| 1227 | self.request(kind, request).map_err(Failed::io) |
| 1228 | } |
| 1229 | |
| 1230 | async fn ask_async(&self, kind: u16, request: Request) -> io::Result<Reply> { |
| 1231 | self.request_async(kind, request).await.map_err(Failed::io) |
| 1232 | } |
| 1233 | |
| 1234 | /// Puts `bytes` in `request`, or uploads them first where they are large. |
| 1235 | #[cfg(not(target_arch = "wasm32"))] |
| 1236 | fn carry(&self, request: &mut Request, bytes: Vec<u8>) -> std::result::Result<(), Failed> { |
| 1237 | crate::task::ready(self.carry_async(request, bytes)).map_err(Failed::Unsent)? |
| 1238 | } |
| 1239 | |
| 1240 | async fn carry_async( |
| 1241 | &self, |
| 1242 | request: &mut Request, |
| 1243 | bytes: Vec<u8>, |
| 1244 | ) -> std::result::Result<(), Failed> { |
| 1245 | if bytes.len() > LIMIT { |
| 1246 | return Err(Failed::Unsent(io::ErrorKind::FileTooLarge.into())); |
| 1247 | } |
| 1248 | if bytes.len() <= CHUNK { |
| 1249 | request.bytes = Some(bytes); |
| 1250 | return Ok(()); |
| 1251 | } |
| 1252 | let handle = self.inner.next.fetch_add(1, Ordering::Relaxed); |
| 1253 | for (index, chunk) in bytes.chunks(CHUNK).enumerate() { |
| 1254 | let put = Request { |
| 1255 | handle: Some(handle), |
| 1256 | offset: Some((index * CHUNK) as u64), |
| 1257 | bytes: Some(chunk.to_vec()), |
| 1258 | ..Request::default() |
| 1259 | }; |
| 1260 | // Nothing the upload carries happens before the request that uses it. |
| 1261 | self.request_async(kind::PUT, put) |
| 1262 | .await |
| 1263 | .map_err(|failed| match failed { |
| 1264 | Failed::Lost(error) => Failed::Unsent(error), |
| 1265 | failed => failed, |
| 1266 | })?; |
| 1267 | } |
| 1268 | request.handle = Some(handle); |
| 1269 | Ok(()) |
| 1270 | } |
| 1271 | |
| 1272 | /// Reads the file at `path` a chunk at a time, as `kind` reads it; a section held, as |
| 1273 | /// the changes to it. |
| 1274 | #[cfg(not(target_arch = "wasm32"))] |
| 1275 | fn read(&self, kind: u16, path: &str, limit: usize) -> io::Result<Vec<u8>> { |
| 1276 | crate::task::ready(self.read_async(kind, path, limit))? |
| 1277 | } |
| 1278 | |
| 1279 | async fn read_async(&self, kind: u16, path: &str, limit: usize) -> io::Result<Vec<u8>> { |
| 1280 | let held = (kind == kind::READ) |
| 1281 | .then(|| self.inner.image(path)) |
| 1282 | .flatten(); |
| 1283 | let first = self |
| 1284 | .ask_async( |
| 1285 | kind, |
| 1286 | Request { |
| 1287 | path: path.to_owned(), |
| 1288 | offset: Some(0), |
| 1289 | limit: Some(limit as u64), |
| 1290 | stamp: (held.as_deref()) |
| 1291 | .and_then(|image| Stamp::of(image).ok()) |
| 1292 | .map(|stamp| (&stamp).into()), |
| 1293 | ..Request::default() |
| 1294 | }, |
| 1295 | ) |
| 1296 | .await?; |
| 1297 | if let (Some(writes), Some(held)) = (&first.writes, held) { |
| 1298 | let stamp: Option<Stamp> = first.stamp.as_ref().and_then(|s| s.try_into().ok()); |
| 1299 | let image = written(&held, first.length.unwrap_or_default(), writes) |
| 1300 | .filter(|image| stamp.is_some() && Stamp::of(image).ok() == stamp) |
| 1301 | .ok_or_else(|| io::Error::new(io::ErrorKind::InvalidData, "Changes that miss"))?; |
| 1302 | self.inner.hold(path, Arc::new(image.clone())); |
| 1303 | return Ok(image); |
| 1304 | } |
| 1305 | let length = first.length.unwrap_or_default() as usize; |
| 1306 | if length > limit { |
| 1307 | return Err(io::ErrorKind::FileTooLarge.into()); |
| 1308 | } |
| 1309 | let mut image = first.bytes.unwrap_or_default(); |
| 1310 | while image.len() < length { |
| 1311 | let chunk = self |
| 1312 | .ask_async( |
| 1313 | kind, |
| 1314 | Request { |
| 1315 | path: path.to_owned(), |
| 1316 | offset: Some(image.len() as u64), |
| 1317 | handle: first.handle, |
| 1318 | ..Request::default() |
| 1319 | }, |
| 1320 | ) |
| 1321 | .await? |
| 1322 | .bytes |
| 1323 | .unwrap_or_default(); |
| 1324 | if chunk.is_empty() { |
| 1325 | return Err(io::ErrorKind::UnexpectedEof.into()); |
| 1326 | } |
| 1327 | image.extend_from_slice(&chunk); |
| 1328 | } |
| 1329 | image.truncate(length); |
| 1330 | if kind == kind::READ { |
| 1331 | self.inner.hold(path, Arc::new(image.clone())); |
| 1332 | } |
| 1333 | Ok(image) |
| 1334 | } |
| 1335 | |
| 1336 | /// The image of the section at `path` with stamp `stamp`, where this guest holds it. |
| 1337 | fn held(&self, path: &str, stamp: &Stamp) -> Option<Vec<u8>> { |
| 1338 | let image = self.inner.image(path)?; |
| 1339 | (Stamp::of(&image).ok().as_ref() == Some(stamp)).then(|| image.to_vec()) |
| 1340 | } |
| 1341 | |
| 1342 | #[cfg(not(target_arch = "wasm32"))] |
| 1343 | pub(crate) fn entries(&self, folder: &str) -> io::Result<Vec<discover::Entry>> { |
| 1344 | crate::task::ready(self.entries_async(folder))? |
| 1345 | } |
| 1346 | |
| 1347 | pub(crate) async fn entries_async(&self, folder: &str) -> io::Result<Vec<discover::Entry>> { |
| 1348 | let reply = self |
| 1349 | .ask_async( |
| 1350 | kind::LIST, |
| 1351 | Request { |
| 1352 | path: folder.to_owned(), |
| 1353 | ..Request::default() |
| 1354 | }, |
| 1355 | ) |
| 1356 | .await?; |
| 1357 | let entries: Vec<_> = reply |
| 1358 | .entries |
| 1359 | .unwrap_or_default() |
| 1360 | .iter() |
| 1361 | .map(entry) |
| 1362 | .collect(); |
| 1363 | #[cfg(target_arch = "wasm32")] |
| 1364 | self.cache_entries(folder, &entries).await?; |
| 1365 | Ok(entries) |
| 1366 | } |
| 1367 | |
| 1368 | /// The stamp of the file at `path`: the host's last report of it where this guest holds |
| 1369 | /// that image, else the host's answer. |
| 1370 | #[cfg(not(target_arch = "wasm32"))] |
| 1371 | fn stamp(&self, path: &str) -> io::Result<Stamp> { |
| 1372 | crate::task::ready(self.stamp_async(path))? |
| 1373 | } |
| 1374 | |
| 1375 | async fn stamp_async(&self, path: &str) -> io::Result<Stamp> { |
| 1376 | if let Some(stamp) = self.inner.current.lock().unwrap().get(path) { |
| 1377 | return Ok(stamp.clone()); |
| 1378 | } |
| 1379 | let reply = self |
| 1380 | .ask_async( |
| 1381 | kind::STAMP, |
| 1382 | Request { |
| 1383 | path: path.to_owned(), |
| 1384 | ..Request::default() |
| 1385 | }, |
| 1386 | ) |
| 1387 | .await?; |
| 1388 | reply |
| 1389 | .stamp |
| 1390 | .as_ref() |
| 1391 | .ok_or_else(|| io::Error::from(io::ErrorKind::InvalidData))? |
| 1392 | .try_into() |
| 1393 | } |
| 1394 | |
| 1395 | #[cfg(not(target_arch = "wasm32"))] |
| 1396 | fn commit( |
| 1397 | &self, |
| 1398 | path: &str, |
| 1399 | transaction: &Transaction, |
| 1400 | ) -> std::result::Result<(), CommitError> { |
| 1401 | crate::task::ready(self.commit_async(path, transaction)).map_err(|error| CommitError { |
| 1402 | state: CommitState::NotCommitted, |
| 1403 | error, |
| 1404 | })? |
| 1405 | } |
| 1406 | |
| 1407 | async fn commit_async( |
| 1408 | &self, |
| 1409 | path: &str, |
| 1410 | transaction: &Transaction, |
| 1411 | ) -> std::result::Result<(), CommitError> { |
| 1412 | let mut request = Request { |
| 1413 | path: path.to_owned(), |
| 1414 | ..Request::default() |
| 1415 | }; |
| 1416 | self.carry_async(&mut request, transaction.to_bytes()) |
| 1417 | .await |
| 1418 | .map_err(Failed::commit)?; |
| 1419 | self.request_async(kind::COMMIT, request) |
| 1420 | .await |
| 1421 | .map(drop) |
| 1422 | .map_err(Failed::commit) |
| 1423 | } |
| 1424 | |
| 1425 | #[cfg(not(target_arch = "wasm32"))] |
| 1426 | fn confirm(&self, path: &str, base: &Stamp) -> std::result::Result<(), CommitError> { |
| 1427 | crate::task::ready(self.confirm_async(path, base)).map_err(|error| CommitError { |
| 1428 | state: CommitState::NotCommitted, |
| 1429 | error, |
| 1430 | })? |
| 1431 | } |
| 1432 | |
| 1433 | async fn confirm_async( |
| 1434 | &self, |
| 1435 | path: &str, |
| 1436 | base: &Stamp, |
| 1437 | ) -> std::result::Result<(), CommitError> { |
| 1438 | let request = Request { |
| 1439 | path: path.to_owned(), |
| 1440 | stamp: Some(base.into()), |
| 1441 | ..Request::default() |
| 1442 | }; |
| 1443 | self.request_async(kind::CONFIRM, request) |
| 1444 | .await |
| 1445 | .map(drop) |
| 1446 | .map_err(Failed::commit) |
| 1447 | } |
| 1448 | |
| 1449 | async fn edits_async( |
| 1450 | &self, |
| 1451 | path: &str, |
| 1452 | transaction: &Transaction, |
| 1453 | edits: &[crate::PendingEdit], |
| 1454 | revisions: &BTreeMap<onestore::ExGuid, onestore::ExGuid>, |
| 1455 | ) -> std::result::Result<(), CommitError> { |
| 1456 | let bytes = batch::encode( |
| 1457 | &batch::Edits { |
| 1458 | edits: edits |
| 1459 | .iter() |
| 1460 | .map(|edit| (edit.author.clone(), edit.edit.clone())) |
| 1461 | .collect(), |
| 1462 | revisions: revisions.clone(), |
| 1463 | }, |
| 1464 | LIMIT, |
| 1465 | ) |
| 1466 | .map_err(|error| CommitError { |
| 1467 | state: CommitState::NotCommitted, |
| 1468 | error, |
| 1469 | })?; |
| 1470 | let mut request = Request { |
| 1471 | path: path.to_owned(), |
| 1472 | stamp: Some(transaction.base().into()), |
| 1473 | ..Request::default() |
| 1474 | }; |
| 1475 | self.carry_async(&mut request, bytes) |
| 1476 | .await |
| 1477 | .map_err(Failed::commit)?; |
| 1478 | match self.request_async(kind::EDITS, request).await { |
| 1479 | Ok(reply) => { |
| 1480 | if let Some(stamp) = reply.stamp { |
| 1481 | let stamp = Stamp::try_from(&stamp).map_err(|error| CommitError { |
| 1482 | state: CommitState::Unknown, |
| 1483 | error, |
| 1484 | })?; |
| 1485 | self.inner |
| 1486 | .current |
| 1487 | .lock() |
| 1488 | .unwrap() |
| 1489 | .insert(path.to_owned(), stamp); |
| 1490 | } |
| 1491 | Ok(()) |
| 1492 | } |
| 1493 | Err(Failed::Refused(failure)) |
| 1494 | if wire::error_kind(failure.kind) == io::ErrorKind::Unsupported => |
| 1495 | { |
| 1496 | self.commit_async(path, transaction).await?; |
| 1497 | self.inner.published(path, transaction); |
| 1498 | Ok(()) |
| 1499 | } |
| 1500 | Err(error) => Err(error.commit()), |
| 1501 | } |
| 1502 | } |
| 1503 | |
| 1504 | /// A request on `path` that answers nothing but whether it happened. |
| 1505 | #[cfg(not(target_arch = "wasm32"))] |
| 1506 | fn verb(&self, kind: u16, path: &str, request: Request) -> Result<()> { |
| 1507 | self.ask( |
| 1508 | kind, |
| 1509 | Request { |
| 1510 | path: path.to_owned(), |
| 1511 | ..request |
| 1512 | }, |
| 1513 | )?; |
| 1514 | Ok(()) |
| 1515 | } |
| 1516 | } |
| 1517 | |
| 1518 | impl Inner { |
| 1519 | fn image(&self, path: &str) -> Option<Arc<Vec<u8>>> { |
| 1520 | let images = self.images.lock().unwrap(); |
| 1521 | images |
| 1522 | .iter() |
| 1523 | .find_map(|(held, image)| (held == path).then(|| Arc::clone(image))) |
| 1524 | } |
| 1525 | |
| 1526 | fn hold(&self, path: &str, image: Arc<Vec<u8>>) { |
| 1527 | let mut images = self.images.lock().unwrap(); |
| 1528 | images.retain(|(held, _)| held != path); |
| 1529 | images.insert(0, (path.to_owned(), image)); |
| 1530 | images.truncate(IMAGES); |
| 1531 | } |
| 1532 | |
| 1533 | /// Takes on a delta to an image held. |
| 1534 | fn apply(&self, delta: &Delta) { |
| 1535 | let held = self.image(&delta.path); |
| 1536 | let base: Option<Stamp> = (&delta.base).try_into().ok(); |
| 1537 | if let Some(image) = held |
| 1538 | && Stamp::of(&image).ok() == base |
| 1539 | && let Some(next) = written(&image, delta.length, &delta.writes) |
| 1540 | { |
| 1541 | self.hold(&delta.path, Arc::new(next)); |
| 1542 | } |
| 1543 | } |
| 1544 | |
| 1545 | /// Takes on this guest's own commit to an image held; the host's report, not this, says |
| 1546 | /// whether the image is the host's, as another's commit may have followed. |
| 1547 | fn published(&self, path: &str, transaction: &Transaction) { |
| 1548 | if let Some(image) = self.image(path) { |
| 1549 | let mut next = (*image).clone(); |
| 1550 | if transaction.apply(&mut next).is_ok() { |
| 1551 | self.hold(path, Arc::new(next)); |
| 1552 | } |
| 1553 | } |
| 1554 | } |
| 1555 | |
| 1556 | /// Takes the host's report of what the files at `paths` are now. |
| 1557 | fn touched(&self, touched: &Touched) { |
| 1558 | let mut current = self.current.lock().unwrap(); |
| 1559 | for (at, path) in touched.paths.iter().enumerate() { |
| 1560 | let stamp = self.image(path).and_then(|image| Stamp::of(&image).ok()); |
| 1561 | match stamp.filter(|stamp| touched.stamps.get(at) == Some(&Some(digest(stamp)))) { |
| 1562 | Some(stamp) => current.insert(path.clone(), stamp), |
| 1563 | None => current.remove(path), |
| 1564 | }; |
| 1565 | } |
| 1566 | } |
| 1567 | |
| 1568 | fn heard(self: &Arc<Self>, event: Event) { |
| 1569 | let serves = |hello: &Hello| { |
| 1570 | self.host |
| 1571 | .lock() |
| 1572 | .unwrap() |
| 1573 | .as_ref() |
| 1574 | .is_some_and(|(host, _)| host.peer == hello.peer) |
| 1575 | }; |
| 1576 | match event { |
| 1577 | Event::Met(hello, line) if hello.serves == Some(self.share) => { |
| 1578 | *self.host.lock().unwrap() = Some((Arc::clone(hello), line.clone())); |
| 1579 | } |
| 1580 | Event::Left(hello) if serves(hello) => { |
| 1581 | *self.host.lock().unwrap() = None; |
| 1582 | drop(self.presence.lock().unwrap().take()); |
| 1583 | self.current.lock().unwrap().clear(); |
| 1584 | // Each request waiting hears its answer was lost. |
| 1585 | self.pending.lock().unwrap().clear(); |
| 1586 | if let Some(reports) = self.watch.lock().unwrap().take() { |
| 1587 | reports.lost(); |
| 1588 | } |
| 1589 | } |
| 1590 | Event::Frame { |
| 1591 | from, kind, body, .. |
| 1592 | } => match kind { |
| 1593 | kind::WELCOME if serves(from) => { |
| 1594 | if let Ok(welcome) = minicbor::decode::<Welcome>(body) { |
| 1595 | self.join_presence(welcome.room); |
| 1596 | } |
| 1597 | } |
| 1598 | kind::REPLY if serves(from) => { |
| 1599 | if let Ok(reply) = minicbor::decode::<Reply>(body) |
| 1600 | && let Some(waiting) = self.pending.lock().unwrap().remove(&reply.id) |
| 1601 | { |
| 1602 | let _ = waiting.send(reply); |
| 1603 | } |
| 1604 | } |
| 1605 | kind::DELTA if serves(from) => { |
| 1606 | if let Ok(delta) = minicbor::decode::<Delta>(body) { |
| 1607 | self.apply(&delta); |
| 1608 | } |
| 1609 | } |
| 1610 | kind::TOUCHED if serves(from) => { |
| 1611 | if let Ok(touched) = minicbor::decode::<Touched>(body) { |
| 1612 | self.touched(&touched); |
| 1613 | let reports = self.watch.lock().unwrap().clone(); |
| 1614 | if let Some(reports) = reports { |
| 1615 | reports.touched(&touched.paths); |
| 1616 | } |
| 1617 | } |
| 1618 | } |
| 1619 | kind::BYE if serves(from) => { |
| 1620 | if let Ok(bye) = minicbor::decode::<wire::Bye>(body) { |
| 1621 | let ended = match bye.reason.as_str() { |
| 1622 | STOPPED => Some(Ended::Stopped), |
| 1623 | "removed" => Some(Ended::Removed), |
| 1624 | _ => None, |
| 1625 | }; |
| 1626 | if ended.is_some() { |
| 1627 | *self.ended.lock().unwrap() = ended; |
| 1628 | } |
| 1629 | } |
| 1630 | } |
| 1631 | _ => {} |
| 1632 | }, |
| 1633 | _ => {} |
| 1634 | } |
| 1635 | } |
| 1636 | |
| 1637 | fn join_presence(self: &Arc<Self>, secret: [u8; 16]) { |
| 1638 | let mut presence = self.presence.lock().unwrap(); |
| 1639 | if presence.as_ref().is_some_and(|(held, _)| *held == secret) { |
| 1640 | return; |
| 1641 | } |
| 1642 | let inner = Arc::downgrade(self); |
| 1643 | let connected = Mutex::new(None); |
| 1644 | let live = Live::start( |
| 1645 | self.me.clone(), |
| 1646 | &Room::Notebook(secret), |
| 1647 | self.reach, |
| 1648 | self.relay.as_deref(), |
| 1649 | move |event| { |
| 1650 | if let Some(inner) = inner.upgrade() { |
| 1651 | if matches!(event, Event::Changed) { |
| 1652 | let mut connected = connected.lock().unwrap(); |
| 1653 | let host = inner |
| 1654 | .host |
| 1655 | .lock() |
| 1656 | .unwrap() |
| 1657 | .as_ref() |
| 1658 | .map(|(host, _)| host.peer); |
| 1659 | let peers = inner |
| 1660 | .presence |
| 1661 | .lock() |
| 1662 | .unwrap() |
| 1663 | .as_ref() |
| 1664 | .map(|(_, live)| live.peers()) |
| 1665 | .unwrap_or_default(); |
| 1666 | let present = |
| 1667 | host.filter(|host| peers.iter().any(|peer| peer.hello.peer == *host)); |
| 1668 | if *connected != present { |
| 1669 | *connected = present; |
| 1670 | inner.current.lock().unwrap().clear(); |
| 1671 | drop(connected); |
| 1672 | let reports = inner.watch.lock().unwrap().clone(); |
| 1673 | if present.is_some() |
| 1674 | && let Some(reports) = reports |
| 1675 | { |
| 1676 | reports.touched(&[String::new()]); |
| 1677 | } |
| 1678 | } |
| 1679 | } |
| 1680 | if matches!( |
| 1681 | event, |
| 1682 | Event::Frame { |
| 1683 | kind: kind::DELTA | kind::TOUCHED, |
| 1684 | .. |
| 1685 | } |
| 1686 | ) { |
| 1687 | inner.heard(event); |
| 1688 | } |
| 1689 | (inner.events)(); |
| 1690 | } |
| 1691 | }, |
| 1692 | ); |
| 1693 | match live { |
| 1694 | Ok(live) => { |
| 1695 | live.set_presence(self.here.lock().unwrap().clone()); |
| 1696 | *presence = Some((secret, live)); |
| 1697 | drop(presence); |
| 1698 | let mut current = self.current.lock().unwrap(); |
| 1699 | let paths: Vec<String> = current.keys().cloned().collect(); |
| 1700 | current.clear(); |
| 1701 | drop(current); |
| 1702 | let reports = self.watch.lock().unwrap().clone(); |
| 1703 | if let Some(reports) = reports { |
| 1704 | reports.touched(&paths); |
| 1705 | } |
| 1706 | } |
| 1707 | Err(error) => eprintln!("Live Share: could not join presence: {error}"), |
| 1708 | } |
| 1709 | } |
| 1710 | } |
| 1711 | |
| 1712 | /// A section a Live Share host serves, as a guest's replica publishes to it. |
| 1713 | pub struct HostedRemote { |
| 1714 | guest: Arc<Guest>, |
| 1715 | path: String, |
| 1716 | /// The stamp last asked for, which an image the guest holds may already have. |
| 1717 | seen: Option<Stamp>, |
| 1718 | rejected: bool, |
| 1719 | #[cfg(target_arch = "wasm32")] |
| 1720 | pending: Option<web::Pending>, |
| 1721 | } |
| 1722 | |
| 1723 | impl HostedRemote { |
| 1724 | pub fn new(guest: &Arc<Guest>, path: &str) -> Self { |
| 1725 | Self { |
| 1726 | guest: Arc::clone(guest), |
| 1727 | path: path.to_owned(), |
| 1728 | seen: None, |
| 1729 | rejected: false, |
| 1730 | #[cfg(target_arch = "wasm32")] |
| 1731 | pending: None, |
| 1732 | } |
| 1733 | } |
| 1734 | } |
| 1735 | |
| 1736 | #[cfg(not(target_arch = "wasm32"))] |
| 1737 | impl crate::Remote for HostedRemote { |
| 1738 | fn accepts_edits(&self) -> bool { |
| 1739 | !self.rejected |
| 1740 | && self |
| 1741 | .guest |
| 1742 | .host() |
| 1743 | .is_some_and(|host| host.ops == Some(1) && host.kinds.contains(&kind::EDITS)) |
| 1744 | } |
| 1745 | |
| 1746 | fn publish_edits( |
| 1747 | &mut self, |
| 1748 | transaction: &Transaction, |
| 1749 | edits: &[crate::PendingEdit], |
| 1750 | revisions: &BTreeMap<onestore::ExGuid, onestore::ExGuid>, |
| 1751 | ) -> std::result::Result<(), CommitError> { |
| 1752 | let result = |
| 1753 | crate::task::ready( |
| 1754 | self.guest |
| 1755 | .edits_async(&self.path, transaction, edits, revisions), |
| 1756 | ) |
| 1757 | .map_err(|error| CommitError { |
| 1758 | state: CommitState::NotCommitted, |
| 1759 | error, |
| 1760 | })?; |
| 1761 | if result.is_ok() { |
| 1762 | self.seen = None; |
| 1763 | } |
| 1764 | if result |
| 1765 | .as_ref() |
| 1766 | .is_err_and(|error| error.state == CommitState::NotCommitted) |
| 1767 | { |
| 1768 | self.rejected = true; |
| 1769 | self.guest.inner.current.lock().unwrap().remove(&self.path); |
| 1770 | } |
| 1771 | result |
| 1772 | } |
| 1773 | |
| 1774 | fn read(&mut self) -> io::Result<Vec<u8>> { |
| 1775 | self.rejected = false; |
| 1776 | if let Some(image) = (self.seen.as_ref()).and_then(|seen| self.guest.held(&self.path, seen)) |
| 1777 | { |
| 1778 | return Ok(image); |
| 1779 | } |
| 1780 | self.guest.read(kind::READ, &self.path, LIMIT) |
| 1781 | } |
| 1782 | |
| 1783 | fn stamp(&mut self) -> io::Result<Stamp> { |
| 1784 | let stamp = self.guest.stamp(&self.path)?; |
| 1785 | self.seen = Some(stamp.clone()); |
| 1786 | Ok(stamp) |
| 1787 | } |
| 1788 | |
| 1789 | fn publish(&mut self, transaction: &Transaction) -> std::result::Result<(), CommitError> { |
| 1790 | if let Err(error) = self.guest.commit(&self.path, transaction) { |
| 1791 | // The host's image moved on, and its report may not have come yet. |
| 1792 | self.guest.inner.current.lock().unwrap().remove(&self.path); |
| 1793 | return Err(error); |
| 1794 | } |
| 1795 | self.guest.inner.published(&self.path, transaction); |
| 1796 | Ok(()) |
| 1797 | } |
| 1798 | |
| 1799 | fn confirm(&mut self, base: &Stamp) -> std::result::Result<(), CommitError> { |
| 1800 | self.guest.confirm(&self.path, base) |
| 1801 | } |
| 1802 | } |
| 1803 | |
| 1804 | /// A notebook a Live Share host serves, as a guest's `session::Notebook` reaches it. Its |
| 1805 | /// folders as last listed are kept at `listed`, so the notebook opens while the host can't |
| 1806 | /// be reached. |
| 1807 | pub struct Hosted { |
| 1808 | guest: Arc<Guest>, |
| 1809 | listed: PathBuf, |
| 1810 | } |
| 1811 | |
| 1812 | /// Each folder's entries as last listed: name, `WireEntry::kind`, size and modified time. |
| 1813 | type Listings = BTreeMap<String, Vec<(String, u8, u64, u64)>>; |
| 1814 | |
| 1815 | impl Hosted { |
| 1816 | pub(crate) fn new(guest: Arc<Guest>, listed: PathBuf) -> Self { |
| 1817 | #[cfg(target_arch = "wasm32")] |
| 1818 | { |
| 1819 | *guest.inner.catalog.lock().unwrap() = Some(listed.clone()); |
| 1820 | } |
| 1821 | Self { guest, listed } |
| 1822 | } |
| 1823 | } |
| 1824 | |
| 1825 | /// The host's folders, or as they were last listed while it can't be reached. |
| 1826 | #[cfg(not(target_arch = "wasm32"))] |
| 1827 | struct Source<'a> { |
| 1828 | guest: &'a Guest, |
| 1829 | kept: Listings, |
| 1830 | listed: Listings, |
| 1831 | } |
| 1832 | |
| 1833 | #[cfg(not(target_arch = "wasm32"))] |
| 1834 | impl discover::Source for Source<'_> { |
| 1835 | fn entries(&mut self, path: &str, limit: usize) -> io::Result<Vec<discover::Entry>> { |
| 1836 | let entries = match self.guest.entries(path) { |
| 1837 | Ok(entries) => entries, |
| 1838 | Err(error) if error.kind() == io::ErrorKind::NotConnected => self |
| 1839 | .kept |
| 1840 | .get(path) |
| 1841 | .ok_or(error)? |
| 1842 | .iter() |
| 1843 | .map(|(name, kind, size, modified)| { |
| 1844 | entry(&WireEntry { |
| 1845 | name: name.clone(), |
| 1846 | kind: *kind, |
| 1847 | size: *size, |
| 1848 | modified: *modified, |
| 1849 | }) |
| 1850 | }) |
| 1851 | .collect(), |
| 1852 | Err(error) => return Err(error), |
| 1853 | }; |
| 1854 | if entries.len() > limit { |
| 1855 | return Err(io::ErrorKind::FileTooLarge.into()); |
| 1856 | } |
| 1857 | self.listed.insert( |
| 1858 | path.to_owned(), |
| 1859 | entries |
| 1860 | .iter() |
| 1861 | .map(|listed| { |
| 1862 | let wire = wire_entry(listed); |
| 1863 | (wire.name, wire.kind, wire.size, wire.modified) |
| 1864 | }) |
| 1865 | .collect(), |
| 1866 | ); |
| 1867 | Ok(entries) |
| 1868 | } |
| 1869 | |
| 1870 | fn read(&mut self, path: &str, limit: usize) -> io::Result<Vec<u8>> { |
| 1871 | self.guest.read(kind::READ, path, limit) |
| 1872 | } |
| 1873 | |
| 1874 | fn read_asset(&mut self, path: &str, limit: usize) -> io::Result<Vec<u8>> { |
| 1875 | self.guest.read(kind::READ_FILE, path, limit) |
| 1876 | } |
| 1877 | } |
| 1878 | |
| 1879 | #[cfg(not(target_arch = "wasm32"))] |
| 1880 | impl Storage for Hosted { |
| 1881 | fn discover( |
| 1882 | &self, |
| 1883 | cache: &mut discover::Cache, |
| 1884 | limits: discover::Limits, |
| 1885 | ) -> Result<discover::Folder> { |
| 1886 | let kept = crate::fs::read(&self.listed) |
| 1887 | .ok() |
| 1888 | .and_then(|bytes| serde_json::from_slice(&bytes).ok()) |
| 1889 | .unwrap_or_default(); |
| 1890 | let mut source = Source { |
| 1891 | guest: &self.guest, |
| 1892 | kept, |
| 1893 | listed: Listings::new(), |
| 1894 | }; |
| 1895 | let folder = cache.discover(&mut source, limits)?; |
| 1896 | if let Ok(json) = serde_json::to_vec(&source.listed) { |
| 1897 | let _ = crate::fs::write(&self.listed, json); |
| 1898 | } |
| 1899 | Ok(folder) |
| 1900 | } |
| 1901 | |
| 1902 | fn location(&self) -> String { |
| 1903 | self.guest.location() |
| 1904 | } |
| 1905 | |
| 1906 | fn entries(&self, folder: &str) -> io::Result<Vec<discover::Entry>> { |
| 1907 | self.guest.entries(folder) |
| 1908 | } |
| 1909 | |
| 1910 | fn exists(&self, path: &str) -> bool { |
| 1911 | let request = Request { |
| 1912 | path: path.to_owned(), |
| 1913 | ..Request::default() |
| 1914 | }; |
| 1915 | self.guest |
| 1916 | .ask(kind::EXISTS, request) |
| 1917 | .is_ok_and(|reply| reply.exists == Some(true)) |
| 1918 | } |
| 1919 | |
| 1920 | fn stamp(&self, path: &str) -> io::Result<Stamp> { |
| 1921 | self.guest.stamp(path) |
| 1922 | } |
| 1923 | |
| 1924 | fn read(&self, path: &str) -> Result<Vec<u8>> { |
| 1925 | Ok(self.guest.read(kind::READ, path, LIMIT)?) |
| 1926 | } |
| 1927 | |
| 1928 | fn read_file(&self, path: &str, limit: usize) -> Result<Vec<u8>> { |
| 1929 | Ok(self.guest.read(kind::READ_FILE, path, limit)?) |
| 1930 | } |
| 1931 | |
| 1932 | fn create(&self, path: &str, bytes: &[u8]) -> Result<()> { |
| 1933 | if bytes.len() > LIMIT { |
| 1934 | return Err(io::Error::from(io::ErrorKind::FileTooLarge).into()); |
| 1935 | } |
| 1936 | let mut request = Request { |
| 1937 | path: path.to_owned(), |
| 1938 | ..Request::default() |
| 1939 | }; |
| 1940 | self.guest |
| 1941 | .carry(&mut request, bytes.to_vec()) |
| 1942 | .map_err(Failed::io)?; |
| 1943 | self.guest.verb(kind::CREATE, path, request) |
| 1944 | } |
| 1945 | |
| 1946 | fn create_directory(&self, path: &str) -> Result<()> { |
| 1947 | self.guest |
| 1948 | .verb(kind::CREATE_DIRECTORY, path, Request::default()) |
| 1949 | } |
| 1950 | |
| 1951 | fn hide(&self, path: &str) -> Result<()> { |
| 1952 | self.guest.verb(kind::HIDE, path, Request::default()) |
| 1953 | } |
| 1954 | |
| 1955 | fn rename(&self, from: &str, to: &str) -> Result<()> { |
| 1956 | let request = Request { |
| 1957 | to: Some(to.to_owned()), |
| 1958 | ..Request::default() |
| 1959 | }; |
| 1960 | self.guest.verb(kind::RENAME, from, request) |
| 1961 | } |
| 1962 | |
| 1963 | fn rename_root(&self, _: &str, _: &[String]) -> Result<String> { |
| 1964 | Err(refused( |
| 1965 | io::ErrorKind::Unsupported, |
| 1966 | "Only the computer sharing this notebook can rename its folder", |
| 1967 | )) |
| 1968 | } |
| 1969 | |
| 1970 | fn replace(&self, from: &str, to: &str) -> Result<()> { |
| 1971 | let request = Request { |
| 1972 | to: Some(to.to_owned()), |
| 1973 | ..Request::default() |
| 1974 | }; |
| 1975 | self.guest.verb(kind::REPLACE, from, request) |
| 1976 | } |
| 1977 | |
| 1978 | fn delete(&self, path: &str) -> Result<()> { |
| 1979 | self.guest.verb(kind::DELETE, path, Request::default()) |
| 1980 | } |
| 1981 | |
| 1982 | fn place(&self, path: &str, ancestor: [u8; 16], name: &str) -> Result<()> { |
| 1983 | let request = Request { |
| 1984 | ancestor: Some(ancestor), |
| 1985 | name: Some(name.to_owned()), |
| 1986 | ..Request::default() |
| 1987 | }; |
| 1988 | self.guest.verb(kind::PLACE, path, request) |
| 1989 | } |
| 1990 | |
| 1991 | fn commit(&self, path: &str, transaction: &Transaction) -> Result<()> { |
| 1992 | Ok(self.guest.commit(path, transaction)?) |
| 1993 | } |
| 1994 | |
| 1995 | fn confirm(&self, path: &str, base: &Stamp) -> std::result::Result<(), CommitError> { |
| 1996 | self.guest.confirm(path, base) |
| 1997 | } |
| 1998 | |
| 1999 | fn supersede(&self, path: &str, base: &Stamp, with: &str) -> Result<()> { |
| 2000 | let request = Request { |
| 2001 | to: Some(with.to_owned()), |
| 2002 | stamp: Some(wire::WireStamp::from(base)), |
| 2003 | ..Request::default() |
| 2004 | }; |
| 2005 | self.guest.verb(kind::SUPERSEDE, path, request) |
| 2006 | } |
| 2007 | } |
| 2008 | |
| 2009 | #[cfg(test)] |
| 2010 | mod tests { |
| 2011 | use super::*; |
| 2012 | |
| 2013 | #[test] |
| 2014 | fn refused_uploads_release_their_bytes_and_share_one_guest_budget() { |
| 2015 | let directory = tempfile::tempdir().unwrap(); |
| 2016 | std::fs::create_dir(directory.path().join("notebook")).unwrap(); |
| 2017 | let served = Arc::new(Served { |
| 2018 | storage: crate::session::Notebook::open( |
| 2019 | directory.path().join("notebook"), |
| 2020 | directory.path().join("cache"), |
| 2021 | ) |
| 2022 | .unwrap() |
| 2023 | .into_storage(), |
| 2024 | images: Mutex::default(), |
| 2025 | snapshots: Mutex::default(), |
| 2026 | puts: Mutex::default(), |
| 2027 | guests: Mutex::default(), |
| 2028 | writers: Mutex::default(), |
| 2029 | host: Mutex::default(), |
| 2030 | room: Mutex::default(), |
| 2031 | }); |
| 2032 | let peer = [1; 16]; |
| 2033 | let put = |handle, offset, bytes| { |
| 2034 | served.answer( |
| 2035 | &peer, |
| 2036 | kind::PUT, |
| 2037 | Request { |
| 2038 | handle: Some(handle), |
| 2039 | offset: Some(offset), |
| 2040 | bytes: Some(bytes), |
| 2041 | ..Request::default() |
| 2042 | }, |
| 2043 | ) |
| 2044 | }; |
| 2045 | put(1, 0, vec![1; 3]).unwrap(); |
| 2046 | put(2, 0, vec![2; 2]).unwrap(); |
| 2047 | assert!(put(1, 2, vec![3]).is_err()); |
| 2048 | assert_eq!( |
| 2049 | served.puts.lock().unwrap().get(&(peer, 2)).unwrap(), |
| 2050 | &[2; 2] |
| 2051 | ); |
| 2052 | assert!(!served.puts.lock().unwrap().contains_key(&(peer, 1))); |
| 2053 | put(1, 0, vec![4]).unwrap(); |
| 2054 | assert!(put(1, 1, Vec::new()).is_err()); |
| 2055 | served |
| 2056 | .puts |
| 2057 | .lock() |
| 2058 | .unwrap() |
| 2059 | .insert((peer, 3), vec![0; LIMIT - 2]); |
| 2060 | assert!(put(2, 2, vec![5]).is_err()); |
| 2061 | assert!(!served.puts.lock().unwrap().contains_key(&(peer, 2))); |
| 2062 | assert!(put(3, (LIMIT - 2) as u64, vec![6; 3]).is_err()); |
| 2063 | assert!(served.puts.lock().unwrap().is_empty()); |
| 2064 | put(4, 0, vec![7]).unwrap(); |
| 2065 | } |
| 2066 | |
| 2067 | /// A commit's writes, found by comparing images, rebuild the image after it from the one |
| 2068 | /// before, and a change touching most of the file is left to be read whole. |
| 2069 | #[test] |
| 2070 | fn deltas_rebuild_the_image_after_a_commit() { |
| 2071 | let before: Vec<u8> = (0..20_000u32).map(|at| (at % 251) as u8).collect(); |
| 2072 | let mut after = before.clone(); |
| 2073 | after[3] ^= 1; |
| 2074 | after[5000..5010].fill(9); |
| 2075 | after[5050] ^= 1; |
| 2076 | after[17_000] ^= 1; |
| 2077 | after.extend_from_slice(&[7; 3000]); |
| 2078 | let writes = delta(&before, &after).unwrap(); |
| 2079 | assert_eq!( |
| 2080 | writes.len(), |
| 2081 | 4, |
| 2082 | "the header, two runs merged as one, one more, the tail" |
| 2083 | ); |
| 2084 | assert_eq!( |
| 2085 | written(&before, after.len() as u64, &writes).unwrap(), |
| 2086 | after |
| 2087 | ); |
| 2088 | let rewritten: Vec<u8> = before.iter().map(|byte| byte ^ 1).collect(); |
| 2089 | assert!(delta(&before, &rewritten).is_none()); |
| 2090 | assert!(delta(&before, &before[..10_000]).is_none()); |
| 2091 | let outside = [Written { |
| 2092 | offset: after.len() as u64, |
| 2093 | bytes: vec![1], |
| 2094 | }]; |
| 2095 | assert!(written(&before, after.len() as u64, &outside).is_none()); |
| 2096 | } |
| 2097 | |
| 2098 | #[test] |
| 2099 | fn differential_reads_include_metadata_in_the_chunk_budget() { |
| 2100 | let before = vec![0; CHUNK * 4]; |
| 2101 | let mut after = before.clone(); |
| 2102 | after.resize(before.len() + CHUNK - 1024, 1); |
| 2103 | assert!(delta(&before, &after).is_none()); |
| 2104 | after.truncate(before.len() + CHUNK / 2); |
| 2105 | let writes = delta(&before, &after).unwrap(); |
| 2106 | assert!(minicbor::to_vec(&writes).unwrap().len() <= CHUNK); |
| 2107 | assert_eq!( |
| 2108 | written(&before, after.len() as u64, &writes).unwrap(), |
| 2109 | after |
| 2110 | ); |
| 2111 | } |
| 2112 | } |