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
9use 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};
13use crate::{Error, Result, background::Reports, discover, session::Storage};
14use onestore::{CommitError, CommitState, RevisionIndex, Stamp, Store, Transaction};
15use sha2::{Digest, Sha256};
16use 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};
28use web_time::Instant;
29
30#[path = "share/batch.rs"]
31mod batch;
32#[path = "share/membership.rs"]
33mod membership;
34#[cfg(target_arch = "wasm32")]
35#[path = "share/web.rs"]
36mod web;
37pub use membership::Host;
38
39/// The most bytes one message of a read or an upload carries.
40const CHUNK: usize = 128 << 10;
41/// The most bytes of chunks a guest has asked for and not yet been given.
42const WINDOW: usize = 512 << 10;
43/// How long a request waits for its reply.
44const TIMEOUT: Duration = Duration::from_secs(60);
45use crate::MAX_FILE_BYTES as LIMIT;
46/// Snapshots of files being read, per guest, and how long one is kept unread.
47const SNAPSHOTS: usize = 8;
48const 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.
51const 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.
54const VERSIONS: usize = 16;
55const 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.
58pub const SETTLE: Duration = Duration::from_millis(20);
59/// The longest a relay may ask a guest joining to wait before it is told so.
60const PATIENT: Duration = Duration::from_secs(10);
61/// What a host says leaving as it stops sharing.
62const STOPPED: &str = "stopped";
63/// What a host says hanging up on a guest that asks too much too fast.
64const FLOODED: &str = "flooded";
65/// A guest's requests a host holds at once, waiting and in hand, past which it hangs up.
66const QUEUED: usize = 32;
67/// The threads working through each guest's requests.
68const 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.
71const STARTS: f64 = 100.0;
72const 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`).
76pub 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)]
83pub 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)]
97pub struct Device {
98 pub secret: [u8; 16],
99 pub name: String,
100 pub device: Option<String>,
101}
102
103impl 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.
121pub fn location(share: &[u8; 16]) -> String {
122 format!("live://{}", super::hex(share))
123}
124
125/// Why joining failed.
126#[derive(Clone, Debug, PartialEq, Eq)]
127pub 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.
156pub 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.
167pub 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
186pub 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.
287struct 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.
306struct 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.
315type Image = (String, Stamp, Arc<Vec<u8>>);
316/// What a guest's requests left, by guest and request.
317type 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.
321fn 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`.
328fn folder(path: &str) -> String {
329 path.rsplit_once('/')
330 .map_or(String::new(), |(folder, _)| folder.to_owned())
331}
332
333fn refused(kind: io::ErrorKind, message: &str) -> Error {
334 io::Error::new(kind, message.to_owned()).into()
335}
336
337impl 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.
816fn 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 they are much less than `after` itself.
826fn 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).then_some(writes)
862}
863
864/// `image` with `writes`, `length` long: none where a write falls outside it.
865fn written(image: &[u8], length: u64, writes: &[Written]) -> Option<Vec<u8>> {
866 let length = usize::try_from(length)
867 .ok()
868 .filter(|length| *length >= image.len())?;
869 let mut next = image.to_vec();
870 next.resize(length, 0);
871 for write in writes {
872 let offset = usize::try_from(write.offset).ok()?;
873 next.get_mut(offset..offset.checked_add(write.bytes.len())?)?
874 .copy_from_slice(&write.bytes);
875 }
876 Some(next)
877}
878
879fn within(folder: &str, name: &str) -> String {
880 if folder.is_empty() {
881 name.to_owned()
882 } else {
883 format!("{folder}/{name}")
884 }
885}
886
887fn wire_entry(entry: &discover::Entry) -> WireEntry {
888 WireEntry {
889 name: entry.name.clone(),
890 kind: match entry.kind {
891 discover::EntryKind::File => 0,
892 discover::EntryKind::Directory => 1,
893 discover::EntryKind::Other => 2,
894 discover::EntryKind::Evicted => 3,
895 },
896 size: entry.listed.size,
897 modified: entry.listed.modified,
898 }
899}
900
901fn entry(entry: &WireEntry) -> discover::Entry {
902 discover::Entry {
903 name: entry.name.clone(),
904 kind: match entry.kind {
905 0 => discover::EntryKind::File,
906 1 => discover::EntryKind::Directory,
907 3 => discover::EntryKind::Evicted,
908 _ => discover::EntryKind::Other,
909 },
910 listed: discover::Listed {
911 size: entry.size,
912 modified: entry.modified,
913 },
914 }
915}
916
917fn state_number(state: CommitState) -> u8 {
918 match state {
919 CommitState::NotCommitted => 0,
920 CommitState::Unknown => 1,
921 CommitState::Committed => 2,
922 }
923}
924
925fn failed(id: u64, error: Error) -> Reply {
926 let (kind, state) = match &error {
927 Error::Io(error) | Error::RemoteIo(error) => (error.kind(), None),
928 Error::Remote(error) => (error.error.kind(), Some(state_number(error.state))),
929 Error::Document(_) => (io::ErrorKind::InvalidData, None),
930 _ => (io::ErrorKind::Other, None),
931 };
932 let message = match &error {
933 Error::Remote(error) => error.error.to_string(),
934 error => error.to_string(),
935 };
936 Reply {
937 id,
938 failure: Some(Failure {
939 kind: wire::error_number(kind),
940 message,
941 state,
942 }),
943 ..Reply::default()
944 }
945}
946
947/// How a request failed.
948enum Failed {
949 /// It never left this end.
950 Unsent(io::Error),
951 /// Its answer was lost: whatever it asked may have happened.
952 Lost(io::Error),
953 /// The host refused it.
954 Refused(Failure),
955}
956
957impl Failed {
958 fn io(self) -> io::Error {
959 match self {
960 Failed::Unsent(error) | Failed::Lost(error) => error,
961 Failed::Refused(failure) => {
962 io::Error::new(wire::error_kind(failure.kind), failure.message)
963 }
964 }
965 }
966
967 /// As a commit's failure: one never sent was not committed, one whose answer was lost may
968 /// have been.
969 fn commit(self) -> CommitError {
970 let state = match &self {
971 Failed::Unsent(_) => CommitState::NotCommitted,
972 Failed::Lost(_) => CommitState::Unknown,
973 Failed::Refused(failure) => match failure.state {
974 Some(1) => CommitState::Unknown,
975 Some(2) => CommitState::Committed,
976 _ => CommitState::NotCommitted,
977 },
978 };
979 CommitError {
980 state,
981 error: self.io(),
982 }
983 }
984}
985
986/// A guest's way to a share: the share's room, and the requests waiting on its host.
987pub struct Guest {
988 live: Live,
989 inner: Arc<Inner>,
990}
991
992struct Inner {
993 share: [u8; 16],
994 me: Hello,
995 reach: Option<Reach>,
996 relay: Option<String>,
997 events: Arc<dyn Fn() + Send + Sync>,
998 presence: Mutex<Option<([u8; 16], Live)>>,
999 here: Mutex<Presence>,
1000 /// The host and the line to it, while connected.
1001 host: Mutex<Option<(Arc<Hello>, Line)>>,
1002 pending: Mutex<HashMap<u64, crate::task::Answer<Reply>>>,
1003 next: AtomicU64,
1004 /// Where the host's reports of changed files go, while a background watches.
1005 watch: Mutex<Option<Reports>>,
1006 /// Why this device can no longer reach the share.
1007 ended: Mutex<Option<Ended>>,
1008 /// Bytes of chunks asked for and not yet given.
1009 asked: (Mutex<Chunks>, Condvar),
1010 /// The sections read lately, kept as the host's deltas change them, newest first.
1011 images: Mutex<Vec<(String, Arc<Vec<u8>>)>>,
1012 /// The stamps of sections whose image held is the host's now, as its last report of
1013 /// them said.
1014 current: Mutex<HashMap<String, Stamp>>,
1015 #[cfg(target_arch = "wasm32")]
1016 catalog: Mutex<Option<PathBuf>>,
1017}
1018
1019#[derive(Clone, Copy, Debug, PartialEq, Eq)]
1020pub enum Ended {
1021 Stopped,
1022 Removed,
1023}
1024
1025#[derive(Default)]
1026struct Chunks {
1027 bytes: usize,
1028 #[cfg(target_arch = "wasm32")]
1029 waiting: Vec<std::task::Waker>,
1030}
1031
1032struct Requested<'a> {
1033 inner: &'a Inner,
1034 id: u64,
1035 chunked: bool,
1036}
1037
1038impl Drop for Requested<'_> {
1039 fn drop(&mut self) {
1040 self.inner.pending.lock().unwrap().remove(&self.id);
1041 if self.chunked {
1042 let (asked, room) = &self.inner.asked;
1043 let mut asked = asked.lock().unwrap();
1044 asked.bytes -= CHUNK;
1045 #[cfg(target_arch = "wasm32")]
1046 let waiting = std::mem::take(&mut asked.waiting);
1047 drop(asked);
1048 room.notify_all();
1049 #[cfg(target_arch = "wasm32")]
1050 for waker in waiting {
1051 waker.wake();
1052 }
1053 }
1054 }
1055}
1056
1057impl Guest {
1058 /// Joins share `share` through its `secret` as `me`, where `reach` and `relay` say.
1059 /// `events` runs on a network thread whenever the host comes or goes, or the peers change.
1060 pub fn start(
1061 me: Hello,
1062 share: [u8; 16],
1063 secret: [u8; 16],
1064 reach: Option<Reach>,
1065 relay: Option<&str>,
1066 events: impl Fn() + Send + Sync + 'static,
1067 ) -> io::Result<Arc<Self>> {
1068 let inner = Arc::new(Inner {
1069 share,
1070 me: me.clone(),
1071 reach,
1072 relay: relay.map(str::to_owned),
1073 events: Arc::new(events),
1074 presence: Mutex::default(),
1075 here: Mutex::default(),
1076 host: Mutex::default(),
1077 pending: Mutex::default(),
1078 next: AtomicU64::new(1),
1079 watch: Mutex::default(),
1080 ended: Mutex::default(),
1081 asked: Default::default(),
1082 images: Mutex::default(),
1083 current: Mutex::default(),
1084 #[cfg(target_arch = "wasm32")]
1085 catalog: Mutex::default(),
1086 });
1087 let heard = Arc::clone(&inner);
1088 let live = Live::start(me, &Room::Notebook(secret), reach, relay, move |event| {
1089 heard.heard(event);
1090 (heard.events)();
1091 })?;
1092 Ok(Arc::new(Self { live, inner }))
1093 }
1094
1095 /// Where the notebook is kept: `location(share)`.
1096 pub fn location(&self) -> String {
1097 location(&self.inner.share)
1098 }
1099
1100 /// The host, while connected.
1101 pub fn host(&self) -> Option<Arc<Hello>> {
1102 let host = self.inner.host.lock().unwrap();
1103 host.as_ref().map(|(hello, _)| Arc::clone(hello))
1104 }
1105
1106 /// Why the host ended this device's access.
1107 pub fn ended(&self) -> Option<Ended> {
1108 *self.inner.ended.lock().unwrap()
1109 }
1110
1111 /// Everyone in the share's room: the host and the other guests.
1112 pub fn peers(&self) -> Vec<Peer> {
1113 self.inner
1114 .presence
1115 .lock()
1116 .unwrap()
1117 .as_ref()
1118 .map(|(_, live)| live.peers())
1119 .unwrap_or_default()
1120 }
1121
1122 pub fn set_presence(&self, presence: Presence) {
1123 *self.inner.here.lock().unwrap() = presence.clone();
1124 if let Some((_, live)) = &*self.inner.presence.lock().unwrap() {
1125 live.set_presence(presence);
1126 }
1127 }
1128
1129 /// Sends the host's reports of changed files to `reports` while it stays connected.
1130 pub(crate) fn watch(&self, reports: Reports) -> io::Result<()> {
1131 if self.host().is_none() {
1132 return Err(self.offline());
1133 }
1134 *self.inner.watch.lock().unwrap() = Some(reports);
1135 Ok(())
1136 }
1137
1138 fn offline(&self) -> io::Error {
1139 if let Some(version) = self.live.other_version() {
1140 return io::Error::new(io::ErrorKind::NotConnected, wire::Version(version));
1141 }
1142 let message = match &*self.inner.ended.lock().unwrap() {
1143 Some(Ended::Stopped) => "The host stopped sharing this notebook",
1144 Some(Ended::Removed) => {
1145 "This device was removed. Ask the person sharing for a new link or code."
1146 }
1147 None => "The computer sharing this notebook can’t be reached",
1148 };
1149 io::Error::new(io::ErrorKind::NotConnected, message)
1150 }
1151
1152 /// Asks the host `kind` of `request`, waiting for its reply.
1153 #[cfg(not(target_arch = "wasm32"))]
1154 fn request(&self, kind: u16, request: Request) -> std::result::Result<Reply, Failed> {
1155 crate::task::ready(self.request_async(kind, request)).map_err(Failed::Unsent)?
1156 }
1157
1158 async fn request_async(
1159 &self,
1160 kind: u16,
1161 mut request: Request,
1162 ) -> std::result::Result<Reply, Failed> {
1163 let Some((_, line)) = self.inner.host.lock().unwrap().clone() else {
1164 return Err(Failed::Unsent(self.offline()));
1165 };
1166 let chunked = matches!(kind, kind::READ | kind::READ_FILE | kind::PUT);
1167 if chunked {
1168 #[cfg(not(target_arch = "wasm32"))]
1169 {
1170 let (asked, _room) = &self.inner.asked;
1171 let mut asked = asked.lock().unwrap();
1172 while asked.bytes + CHUNK > WINDOW {
1173 asked = _room.wait(asked).unwrap();
1174 }
1175 asked.bytes += CHUNK;
1176 }
1177 #[cfg(target_arch = "wasm32")]
1178 std::future::poll_fn(|context| {
1179 let mut asked = self.inner.asked.0.lock().unwrap();
1180 if asked.bytes + CHUNK > WINDOW {
1181 if !asked
1182 .waiting
1183 .iter()
1184 .any(|waker| waker.will_wake(context.waker()))
1185 {
1186 asked.waiting.push(context.waker().clone());
1187 }
1188 std::task::Poll::Pending
1189 } else {
1190 asked.bytes += CHUNK;
1191 std::task::Poll::Ready(())
1192 }
1193 })
1194 .await;
1195 }
1196 let id = self.inner.next.fetch_add(1, Ordering::Relaxed);
1197 request.id = id;
1198 let (answer, answered) = crate::task::response();
1199 self.inner.pending.lock().unwrap().insert(id, answer);
1200 let _requested = Requested {
1201 inner: &self.inner,
1202 id,
1203 chunked,
1204 };
1205 let sent = line.send(kind, &request);
1206 let reply = match sent {
1207 Err(error) => Err(Failed::Unsent(error)),
1208 Ok(()) => match crate::task::answered(&answered, TIMEOUT).await {
1209 Ok(reply) => Ok(reply),
1210 Err(mpsc::RecvTimeoutError::Timeout) => Err(Failed::Lost(io::Error::new(
1211 io::ErrorKind::TimedOut,
1212 "The computer sharing this notebook didn’t answer",
1213 ))),
1214 Err(mpsc::RecvTimeoutError::Disconnected) => Err(Failed::Lost(self.offline())),
1215 },
1216 };
1217 let reply = reply?;
1218 match reply.failure {
1219 Some(failure) => Err(Failed::Refused(failure)),
1220 None => Ok(reply),
1221 }
1222 }
1223
1224 #[cfg(not(target_arch = "wasm32"))]
1225 fn ask(&self, kind: u16, request: Request) -> io::Result<Reply> {
1226 self.request(kind, request).map_err(Failed::io)
1227 }
1228
1229 async fn ask_async(&self, kind: u16, request: Request) -> io::Result<Reply> {
1230 self.request_async(kind, request).await.map_err(Failed::io)
1231 }
1232
1233 /// Puts `bytes` in `request`, or uploads them first where they are large.
1234 #[cfg(not(target_arch = "wasm32"))]
1235 fn carry(&self, request: &mut Request, bytes: Vec<u8>) -> std::result::Result<(), Failed> {
1236 crate::task::ready(self.carry_async(request, bytes)).map_err(Failed::Unsent)?
1237 }
1238
1239 async fn carry_async(
1240 &self,
1241 request: &mut Request,
1242 bytes: Vec<u8>,
1243 ) -> std::result::Result<(), Failed> {
1244 if bytes.len() > LIMIT {
1245 return Err(Failed::Unsent(io::ErrorKind::FileTooLarge.into()));
1246 }
1247 if bytes.len() <= CHUNK {
1248 request.bytes = Some(bytes);
1249 return Ok(());
1250 }
1251 let handle = self.inner.next.fetch_add(1, Ordering::Relaxed);
1252 for (index, chunk) in bytes.chunks(CHUNK).enumerate() {
1253 let put = Request {
1254 handle: Some(handle),
1255 offset: Some((index * CHUNK) as u64),
1256 bytes: Some(chunk.to_vec()),
1257 ..Request::default()
1258 };
1259 // Nothing the upload carries happens before the request that uses it.
1260 self.request_async(kind::PUT, put)
1261 .await
1262 .map_err(|failed| match failed {
1263 Failed::Lost(error) => Failed::Unsent(error),
1264 failed => failed,
1265 })?;
1266 }
1267 request.handle = Some(handle);
1268 Ok(())
1269 }
1270
1271 /// Reads the file at `path` a chunk at a time, as `kind` reads it; a section held, as
1272 /// the changes to it.
1273 #[cfg(not(target_arch = "wasm32"))]
1274 fn read(&self, kind: u16, path: &str, limit: usize) -> io::Result<Vec<u8>> {
1275 crate::task::ready(self.read_async(kind, path, limit))?
1276 }
1277
1278 async fn read_async(&self, kind: u16, path: &str, limit: usize) -> io::Result<Vec<u8>> {
1279 let held = (kind == kind::READ)
1280 .then(|| self.inner.image(path))
1281 .flatten();
1282 let first = self
1283 .ask_async(
1284 kind,
1285 Request {
1286 path: path.to_owned(),
1287 offset: Some(0),
1288 limit: Some(limit as u64),
1289 stamp: (held.as_deref())
1290 .and_then(|image| Stamp::of(image).ok())
1291 .map(|stamp| (&stamp).into()),
1292 ..Request::default()
1293 },
1294 )
1295 .await?;
1296 if let (Some(writes), Some(held)) = (&first.writes, held) {
1297 let stamp: Option<Stamp> = first.stamp.as_ref().and_then(|s| s.try_into().ok());
1298 let image = written(&held, first.length.unwrap_or_default(), writes)
1299 .filter(|image| stamp.is_some() && Stamp::of(image).ok() == stamp)
1300 .ok_or_else(|| io::Error::new(io::ErrorKind::InvalidData, "Changes that miss"))?;
1301 self.inner.hold(path, Arc::new(image.clone()));
1302 return Ok(image);
1303 }
1304 let length = first.length.unwrap_or_default() as usize;
1305 if length > limit {
1306 return Err(io::ErrorKind::FileTooLarge.into());
1307 }
1308 let mut image = first.bytes.unwrap_or_default();
1309 while image.len() < length {
1310 let chunk = self
1311 .ask_async(
1312 kind,
1313 Request {
1314 path: path.to_owned(),
1315 offset: Some(image.len() as u64),
1316 handle: first.handle,
1317 ..Request::default()
1318 },
1319 )
1320 .await?
1321 .bytes
1322 .unwrap_or_default();
1323 if chunk.is_empty() {
1324 return Err(io::ErrorKind::UnexpectedEof.into());
1325 }
1326 image.extend_from_slice(&chunk);
1327 }
1328 image.truncate(length);
1329 if kind == kind::READ {
1330 self.inner.hold(path, Arc::new(image.clone()));
1331 }
1332 Ok(image)
1333 }
1334
1335 /// The image of the section at `path` with stamp `stamp`, where this guest holds it.
1336 fn held(&self, path: &str, stamp: &Stamp) -> Option<Vec<u8>> {
1337 let image = self.inner.image(path)?;
1338 (Stamp::of(&image).ok().as_ref() == Some(stamp)).then(|| image.to_vec())
1339 }
1340
1341 #[cfg(not(target_arch = "wasm32"))]
1342 pub(crate) fn entries(&self, folder: &str) -> io::Result<Vec<discover::Entry>> {
1343 crate::task::ready(self.entries_async(folder))?
1344 }
1345
1346 pub(crate) async fn entries_async(&self, folder: &str) -> io::Result<Vec<discover::Entry>> {
1347 let reply = self
1348 .ask_async(
1349 kind::LIST,
1350 Request {
1351 path: folder.to_owned(),
1352 ..Request::default()
1353 },
1354 )
1355 .await?;
1356 let entries: Vec<_> = reply
1357 .entries
1358 .unwrap_or_default()
1359 .iter()
1360 .map(entry)
1361 .collect();
1362 #[cfg(target_arch = "wasm32")]
1363 self.cache_entries(folder, &entries).await?;
1364 Ok(entries)
1365 }
1366
1367 /// The stamp of the file at `path`: the host's last report of it where this guest holds
1368 /// that image, else the host's answer.
1369 #[cfg(not(target_arch = "wasm32"))]
1370 fn stamp(&self, path: &str) -> io::Result<Stamp> {
1371 crate::task::ready(self.stamp_async(path))?
1372 }
1373
1374 async fn stamp_async(&self, path: &str) -> io::Result<Stamp> {
1375 if let Some(stamp) = self.inner.current.lock().unwrap().get(path) {
1376 return Ok(stamp.clone());
1377 }
1378 let reply = self
1379 .ask_async(
1380 kind::STAMP,
1381 Request {
1382 path: path.to_owned(),
1383 ..Request::default()
1384 },
1385 )
1386 .await?;
1387 reply
1388 .stamp
1389 .as_ref()
1390 .ok_or_else(|| io::Error::from(io::ErrorKind::InvalidData))?
1391 .try_into()
1392 }
1393
1394 #[cfg(not(target_arch = "wasm32"))]
1395 fn commit(
1396 &self,
1397 path: &str,
1398 transaction: &Transaction,
1399 ) -> std::result::Result<(), CommitError> {
1400 crate::task::ready(self.commit_async(path, transaction)).map_err(|error| CommitError {
1401 state: CommitState::NotCommitted,
1402 error,
1403 })?
1404 }
1405
1406 async fn commit_async(
1407 &self,
1408 path: &str,
1409 transaction: &Transaction,
1410 ) -> std::result::Result<(), CommitError> {
1411 let mut request = Request {
1412 path: path.to_owned(),
1413 ..Request::default()
1414 };
1415 self.carry_async(&mut request, transaction.to_bytes())
1416 .await
1417 .map_err(Failed::commit)?;
1418 self.request_async(kind::COMMIT, request)
1419 .await
1420 .map(drop)
1421 .map_err(Failed::commit)
1422 }
1423
1424 #[cfg(not(target_arch = "wasm32"))]
1425 fn confirm(&self, path: &str, base: &Stamp) -> std::result::Result<(), CommitError> {
1426 crate::task::ready(self.confirm_async(path, base)).map_err(|error| CommitError {
1427 state: CommitState::NotCommitted,
1428 error,
1429 })?
1430 }
1431
1432 async fn confirm_async(
1433 &self,
1434 path: &str,
1435 base: &Stamp,
1436 ) -> std::result::Result<(), CommitError> {
1437 let request = Request {
1438 path: path.to_owned(),
1439 stamp: Some(base.into()),
1440 ..Request::default()
1441 };
1442 self.request_async(kind::CONFIRM, request)
1443 .await
1444 .map(drop)
1445 .map_err(Failed::commit)
1446 }
1447
1448 async fn edits_async(
1449 &self,
1450 path: &str,
1451 transaction: &Transaction,
1452 edits: &[crate::PendingEdit],
1453 revisions: &BTreeMap<onestore::ExGuid, onestore::ExGuid>,
1454 ) -> std::result::Result<(), CommitError> {
1455 let bytes = batch::encode(
1456 &batch::Edits {
1457 edits: edits
1458 .iter()
1459 .map(|edit| (edit.author.clone(), edit.edit.clone()))
1460 .collect(),
1461 revisions: revisions.clone(),
1462 },
1463 LIMIT,
1464 )
1465 .map_err(|error| CommitError {
1466 state: CommitState::NotCommitted,
1467 error,
1468 })?;
1469 let mut request = Request {
1470 path: path.to_owned(),
1471 stamp: Some(transaction.base().into()),
1472 ..Request::default()
1473 };
1474 self.carry_async(&mut request, bytes)
1475 .await
1476 .map_err(Failed::commit)?;
1477 match self.request_async(kind::EDITS, request).await {
1478 Ok(reply) => {
1479 if let Some(stamp) = reply.stamp {
1480 let stamp = Stamp::try_from(&stamp).map_err(|error| CommitError {
1481 state: CommitState::Unknown,
1482 error,
1483 })?;
1484 self.inner
1485 .current
1486 .lock()
1487 .unwrap()
1488 .insert(path.to_owned(), stamp);
1489 }
1490 Ok(())
1491 }
1492 Err(Failed::Refused(failure))
1493 if wire::error_kind(failure.kind) == io::ErrorKind::Unsupported =>
1494 {
1495 self.commit_async(path, transaction).await?;
1496 self.inner.published(path, transaction);
1497 Ok(())
1498 }
1499 Err(error) => Err(error.commit()),
1500 }
1501 }
1502
1503 /// A request on `path` that answers nothing but whether it happened.
1504 #[cfg(not(target_arch = "wasm32"))]
1505 fn verb(&self, kind: u16, path: &str, request: Request) -> Result<()> {
1506 self.ask(
1507 kind,
1508 Request {
1509 path: path.to_owned(),
1510 ..request
1511 },
1512 )?;
1513 Ok(())
1514 }
1515}
1516
1517impl Inner {
1518 fn image(&self, path: &str) -> Option<Arc<Vec<u8>>> {
1519 let images = self.images.lock().unwrap();
1520 images
1521 .iter()
1522 .find_map(|(held, image)| (held == path).then(|| Arc::clone(image)))
1523 }
1524
1525 fn hold(&self, path: &str, image: Arc<Vec<u8>>) {
1526 let mut images = self.images.lock().unwrap();
1527 images.retain(|(held, _)| held != path);
1528 images.insert(0, (path.to_owned(), image));
1529 images.truncate(IMAGES);
1530 }
1531
1532 /// Takes on a delta to an image held.
1533 fn apply(&self, delta: &Delta) {
1534 let held = self.image(&delta.path);
1535 let base: Option<Stamp> = (&delta.base).try_into().ok();
1536 if let Some(image) = held
1537 && Stamp::of(&image).ok() == base
1538 && let Some(next) = written(&image, delta.length, &delta.writes)
1539 {
1540 self.hold(&delta.path, Arc::new(next));
1541 }
1542 }
1543
1544 /// Takes on this guest's own commit to an image held; the host's report, not this, says
1545 /// whether the image is the host's, as another's commit may have followed.
1546 fn published(&self, path: &str, transaction: &Transaction) {
1547 if let Some(image) = self.image(path) {
1548 let mut next = (*image).clone();
1549 if transaction.apply(&mut next).is_ok() {
1550 self.hold(path, Arc::new(next));
1551 }
1552 }
1553 }
1554
1555 /// Takes the host's report of what the files at `paths` are now.
1556 fn touched(&self, touched: &Touched) {
1557 let mut current = self.current.lock().unwrap();
1558 for (at, path) in touched.paths.iter().enumerate() {
1559 let stamp = self.image(path).and_then(|image| Stamp::of(&image).ok());
1560 match stamp.filter(|stamp| touched.stamps.get(at) == Some(&Some(digest(stamp)))) {
1561 Some(stamp) => current.insert(path.clone(), stamp),
1562 None => current.remove(path),
1563 };
1564 }
1565 }
1566
1567 fn heard(self: &Arc<Self>, event: Event) {
1568 let serves = |hello: &Hello| {
1569 self.host
1570 .lock()
1571 .unwrap()
1572 .as_ref()
1573 .is_some_and(|(host, _)| host.peer == hello.peer)
1574 };
1575 match event {
1576 Event::Met(hello, line) if hello.serves == Some(self.share) => {
1577 *self.host.lock().unwrap() = Some((Arc::clone(hello), line.clone()));
1578 }
1579 Event::Left(hello) if serves(hello) => {
1580 *self.host.lock().unwrap() = None;
1581 drop(self.presence.lock().unwrap().take());
1582 self.current.lock().unwrap().clear();
1583 // Each request waiting hears its answer was lost.
1584 self.pending.lock().unwrap().clear();
1585 if let Some(reports) = self.watch.lock().unwrap().take() {
1586 reports.lost();
1587 }
1588 }
1589 Event::Frame {
1590 from, kind, body, ..
1591 } => match kind {
1592 kind::WELCOME if serves(from) => {
1593 if let Ok(welcome) = minicbor::decode::<Welcome>(body) {
1594 self.join_presence(welcome.room);
1595 }
1596 }
1597 kind::REPLY if serves(from) => {
1598 if let Ok(reply) = minicbor::decode::<Reply>(body)
1599 && let Some(waiting) = self.pending.lock().unwrap().remove(&reply.id)
1600 {
1601 let _ = waiting.send(reply);
1602 }
1603 }
1604 kind::DELTA if serves(from) => {
1605 if let Ok(delta) = minicbor::decode::<Delta>(body) {
1606 self.apply(&delta);
1607 }
1608 }
1609 kind::TOUCHED if serves(from) => {
1610 if let Ok(touched) = minicbor::decode::<Touched>(body) {
1611 self.touched(&touched);
1612 if let Some(reports) = &*self.watch.lock().unwrap() {
1613 reports.touched(&touched.paths);
1614 }
1615 }
1616 }
1617 kind::BYE if serves(from) => {
1618 if let Ok(bye) = minicbor::decode::<wire::Bye>(body) {
1619 let ended = match bye.reason.as_str() {
1620 STOPPED => Some(Ended::Stopped),
1621 "removed" => Some(Ended::Removed),
1622 _ => None,
1623 };
1624 if ended.is_some() {
1625 *self.ended.lock().unwrap() = ended;
1626 }
1627 }
1628 }
1629 _ => {}
1630 },
1631 _ => {}
1632 }
1633 }
1634
1635 fn join_presence(self: &Arc<Self>, secret: [u8; 16]) {
1636 let mut presence = self.presence.lock().unwrap();
1637 if presence.as_ref().is_some_and(|(held, _)| *held == secret) {
1638 return;
1639 }
1640 let inner = Arc::downgrade(self);
1641 let live = Live::start(
1642 self.me.clone(),
1643 &Room::Notebook(secret),
1644 self.reach,
1645 self.relay.as_deref(),
1646 move |event| {
1647 if let Some(inner) = inner.upgrade() {
1648 if matches!(
1649 event,
1650 Event::Frame {
1651 kind: kind::DELTA | kind::TOUCHED,
1652 ..
1653 }
1654 ) {
1655 inner.heard(event);
1656 }
1657 (inner.events)();
1658 }
1659 },
1660 );
1661 match live {
1662 Ok(live) => {
1663 live.set_presence(self.here.lock().unwrap().clone());
1664 *presence = Some((secret, live));
1665 let mut current = self.current.lock().unwrap();
1666 let paths: Vec<String> = current.keys().cloned().collect();
1667 current.clear();
1668 drop(current);
1669 if let Some(reports) = &*self.watch.lock().unwrap() {
1670 reports.touched(&paths);
1671 }
1672 }
1673 Err(error) => eprintln!("Live Share: could not join presence: {error}"),
1674 }
1675 }
1676}
1677
1678/// A section a Live Share host serves, as a guest's replica publishes to it.
1679pub struct HostedRemote {
1680 guest: Arc<Guest>,
1681 path: String,
1682 /// The stamp last asked for, which an image the guest holds may already have.
1683 seen: Option<Stamp>,
1684 rejected: bool,
1685 #[cfg(target_arch = "wasm32")]
1686 pending: Option<web::Pending>,
1687}
1688
1689impl HostedRemote {
1690 pub fn new(guest: &Arc<Guest>, path: &str) -> Self {
1691 Self {
1692 guest: Arc::clone(guest),
1693 path: path.to_owned(),
1694 seen: None,
1695 rejected: false,
1696 #[cfg(target_arch = "wasm32")]
1697 pending: None,
1698 }
1699 }
1700}
1701
1702#[cfg(not(target_arch = "wasm32"))]
1703impl crate::Remote for HostedRemote {
1704 fn accepts_edits(&self) -> bool {
1705 !self.rejected
1706 && self
1707 .guest
1708 .host()
1709 .is_some_and(|host| host.ops == Some(1) && host.kinds.contains(&kind::EDITS))
1710 }
1711
1712 fn publish_edits(
1713 &mut self,
1714 transaction: &Transaction,
1715 edits: &[crate::PendingEdit],
1716 revisions: &BTreeMap<onestore::ExGuid, onestore::ExGuid>,
1717 ) -> std::result::Result<(), CommitError> {
1718 let result =
1719 crate::task::ready(
1720 self.guest
1721 .edits_async(&self.path, transaction, edits, revisions),
1722 )
1723 .map_err(|error| CommitError {
1724 state: CommitState::NotCommitted,
1725 error,
1726 })?;
1727 if result.is_ok() {
1728 self.seen = None;
1729 }
1730 if result
1731 .as_ref()
1732 .is_err_and(|error| error.state == CommitState::NotCommitted)
1733 {
1734 self.rejected = true;
1735 self.guest.inner.current.lock().unwrap().remove(&self.path);
1736 }
1737 result
1738 }
1739
1740 fn read(&mut self) -> io::Result<Vec<u8>> {
1741 self.rejected = false;
1742 if let Some(image) = (self.seen.as_ref()).and_then(|seen| self.guest.held(&self.path, seen))
1743 {
1744 return Ok(image);
1745 }
1746 self.guest.read(kind::READ, &self.path, LIMIT)
1747 }
1748
1749 fn stamp(&mut self) -> io::Result<Stamp> {
1750 let stamp = self.guest.stamp(&self.path)?;
1751 self.seen = Some(stamp.clone());
1752 Ok(stamp)
1753 }
1754
1755 fn publish(&mut self, transaction: &Transaction) -> std::result::Result<(), CommitError> {
1756 if let Err(error) = self.guest.commit(&self.path, transaction) {
1757 // The host's image moved on, and its report may not have come yet.
1758 self.guest.inner.current.lock().unwrap().remove(&self.path);
1759 return Err(error);
1760 }
1761 self.guest.inner.published(&self.path, transaction);
1762 Ok(())
1763 }
1764
1765 fn confirm(&mut self, base: &Stamp) -> std::result::Result<(), CommitError> {
1766 self.guest.confirm(&self.path, base)
1767 }
1768}
1769
1770/// A notebook a Live Share host serves, as a guest's `session::Notebook` reaches it. Its
1771/// folders as last listed are kept at `listed`, so the notebook opens while the host can't
1772/// be reached.
1773pub struct Hosted {
1774 guest: Arc<Guest>,
1775 listed: PathBuf,
1776}
1777
1778/// Each folder's entries as last listed: name, `WireEntry::kind`, size and modified time.
1779type Listings = BTreeMap<String, Vec<(String, u8, u64, u64)>>;
1780
1781impl Hosted {
1782 pub(crate) fn new(guest: Arc<Guest>, listed: PathBuf) -> Self {
1783 #[cfg(target_arch = "wasm32")]
1784 {
1785 *guest.inner.catalog.lock().unwrap() = Some(listed.clone());
1786 }
1787 Self { guest, listed }
1788 }
1789}
1790
1791/// The host's folders, or as they were last listed while it can't be reached.
1792#[cfg(not(target_arch = "wasm32"))]
1793struct Source<'a> {
1794 guest: &'a Guest,
1795 kept: Listings,
1796 listed: Listings,
1797}
1798
1799#[cfg(not(target_arch = "wasm32"))]
1800impl discover::Source for Source<'_> {
1801 fn entries(&mut self, path: &str, limit: usize) -> io::Result<Vec<discover::Entry>> {
1802 let entries = match self.guest.entries(path) {
1803 Ok(entries) => entries,
1804 Err(error) if error.kind() == io::ErrorKind::NotConnected => self
1805 .kept
1806 .get(path)
1807 .ok_or(error)?
1808 .iter()
1809 .map(|(name, kind, size, modified)| {
1810 entry(&WireEntry {
1811 name: name.clone(),
1812 kind: *kind,
1813 size: *size,
1814 modified: *modified,
1815 })
1816 })
1817 .collect(),
1818 Err(error) => return Err(error),
1819 };
1820 if entries.len() > limit {
1821 return Err(io::ErrorKind::FileTooLarge.into());
1822 }
1823 self.listed.insert(
1824 path.to_owned(),
1825 entries
1826 .iter()
1827 .map(|listed| {
1828 let wire = wire_entry(listed);
1829 (wire.name, wire.kind, wire.size, wire.modified)
1830 })
1831 .collect(),
1832 );
1833 Ok(entries)
1834 }
1835
1836 fn read(&mut self, path: &str, limit: usize) -> io::Result<Vec<u8>> {
1837 self.guest.read(kind::READ, path, limit)
1838 }
1839
1840 fn read_asset(&mut self, path: &str, limit: usize) -> io::Result<Vec<u8>> {
1841 self.guest.read(kind::READ_FILE, path, limit)
1842 }
1843}
1844
1845#[cfg(not(target_arch = "wasm32"))]
1846impl Storage for Hosted {
1847 fn discover(
1848 &self,
1849 cache: &mut discover::Cache,
1850 limits: discover::Limits,
1851 ) -> Result<discover::Folder> {
1852 let kept = crate::fs::read(&self.listed)
1853 .ok()
1854 .and_then(|bytes| serde_json::from_slice(&bytes).ok())
1855 .unwrap_or_default();
1856 let mut source = Source {
1857 guest: &self.guest,
1858 kept,
1859 listed: Listings::new(),
1860 };
1861 let folder = cache.discover(&mut source, limits)?;
1862 if let Ok(json) = serde_json::to_vec(&source.listed) {
1863 let _ = crate::fs::write(&self.listed, json);
1864 }
1865 Ok(folder)
1866 }
1867
1868 fn location(&self) -> String {
1869 self.guest.location()
1870 }
1871
1872 fn entries(&self, folder: &str) -> io::Result<Vec<discover::Entry>> {
1873 self.guest.entries(folder)
1874 }
1875
1876 fn exists(&self, path: &str) -> bool {
1877 let request = Request {
1878 path: path.to_owned(),
1879 ..Request::default()
1880 };
1881 self.guest
1882 .ask(kind::EXISTS, request)
1883 .is_ok_and(|reply| reply.exists == Some(true))
1884 }
1885
1886 fn stamp(&self, path: &str) -> io::Result<Stamp> {
1887 self.guest.stamp(path)
1888 }
1889
1890 fn read(&self, path: &str) -> Result<Vec<u8>> {
1891 Ok(self.guest.read(kind::READ, path, LIMIT)?)
1892 }
1893
1894 fn read_file(&self, path: &str, limit: usize) -> Result<Vec<u8>> {
1895 Ok(self.guest.read(kind::READ_FILE, path, limit)?)
1896 }
1897
1898 fn create(&self, path: &str, bytes: &[u8]) -> Result<()> {
1899 if bytes.len() > LIMIT {
1900 return Err(io::Error::from(io::ErrorKind::FileTooLarge).into());
1901 }
1902 let mut request = Request {
1903 path: path.to_owned(),
1904 ..Request::default()
1905 };
1906 self.guest
1907 .carry(&mut request, bytes.to_vec())
1908 .map_err(Failed::io)?;
1909 self.guest.verb(kind::CREATE, path, request)
1910 }
1911
1912 fn create_directory(&self, path: &str) -> Result<()> {
1913 self.guest
1914 .verb(kind::CREATE_DIRECTORY, path, Request::default())
1915 }
1916
1917 fn hide(&self, path: &str) -> Result<()> {
1918 self.guest.verb(kind::HIDE, path, Request::default())
1919 }
1920
1921 fn rename(&self, from: &str, to: &str) -> Result<()> {
1922 let request = Request {
1923 to: Some(to.to_owned()),
1924 ..Request::default()
1925 };
1926 self.guest.verb(kind::RENAME, from, request)
1927 }
1928
1929 fn rename_root(&self, _: &str, _: &[String]) -> Result<String> {
1930 Err(refused(
1931 io::ErrorKind::Unsupported,
1932 "Only the computer sharing this notebook can rename its folder",
1933 ))
1934 }
1935
1936 fn replace(&self, from: &str, to: &str) -> Result<()> {
1937 let request = Request {
1938 to: Some(to.to_owned()),
1939 ..Request::default()
1940 };
1941 self.guest.verb(kind::REPLACE, from, request)
1942 }
1943
1944 fn delete(&self, path: &str) -> Result<()> {
1945 self.guest.verb(kind::DELETE, path, Request::default())
1946 }
1947
1948 fn place(&self, path: &str, ancestor: [u8; 16], name: &str) -> Result<()> {
1949 let request = Request {
1950 ancestor: Some(ancestor),
1951 name: Some(name.to_owned()),
1952 ..Request::default()
1953 };
1954 self.guest.verb(kind::PLACE, path, request)
1955 }
1956
1957 fn commit(&self, path: &str, transaction: &Transaction) -> Result<()> {
1958 Ok(self.guest.commit(path, transaction)?)
1959 }
1960
1961 fn confirm(&self, path: &str, base: &Stamp) -> std::result::Result<(), CommitError> {
1962 self.guest.confirm(path, base)
1963 }
1964
1965 fn supersede(&self, path: &str, base: &Stamp, with: &str) -> Result<()> {
1966 let request = Request {
1967 to: Some(with.to_owned()),
1968 stamp: Some(wire::WireStamp::from(base)),
1969 ..Request::default()
1970 };
1971 self.guest.verb(kind::SUPERSEDE, path, request)
1972 }
1973}
1974
1975#[cfg(test)]
1976mod tests {
1977 use super::*;
1978
1979 #[test]
1980 fn refused_uploads_release_their_bytes_and_share_one_guest_budget() {
1981 let directory = tempfile::tempdir().unwrap();
1982 std::fs::create_dir(directory.path().join("notebook")).unwrap();
1983 let served = Arc::new(Served {
1984 storage: crate::session::Notebook::open(
1985 directory.path().join("notebook"),
1986 directory.path().join("cache"),
1987 )
1988 .unwrap()
1989 .into_storage(),
1990 images: Mutex::default(),
1991 snapshots: Mutex::default(),
1992 puts: Mutex::default(),
1993 guests: Mutex::default(),
1994 writers: Mutex::default(),
1995 host: Mutex::default(),
1996 room: Mutex::default(),
1997 });
1998 let peer = [1; 16];
1999 let put = |handle, offset, bytes| {
2000 served.answer(
2001 &peer,
2002 kind::PUT,
2003 Request {
2004 handle: Some(handle),
2005 offset: Some(offset),
2006 bytes: Some(bytes),
2007 ..Request::default()
2008 },
2009 )
2010 };
2011 put(1, 0, vec![1; 3]).unwrap();
2012 put(2, 0, vec![2; 2]).unwrap();
2013 assert!(put(1, 2, vec![3]).is_err());
2014 assert_eq!(
2015 served.puts.lock().unwrap().get(&(peer, 2)).unwrap(),
2016 &[2; 2]
2017 );
2018 assert!(!served.puts.lock().unwrap().contains_key(&(peer, 1)));
2019 put(1, 0, vec![4]).unwrap();
2020 assert!(put(1, 1, Vec::new()).is_err());
2021 served
2022 .puts
2023 .lock()
2024 .unwrap()
2025 .insert((peer, 3), vec![0; LIMIT - 2]);
2026 assert!(put(2, 2, vec![5]).is_err());
2027 assert!(!served.puts.lock().unwrap().contains_key(&(peer, 2)));
2028 assert!(put(3, (LIMIT - 2) as u64, vec![6; 3]).is_err());
2029 assert!(served.puts.lock().unwrap().is_empty());
2030 put(4, 0, vec![7]).unwrap();
2031 }
2032
2033 /// A commit's writes, found by comparing images, rebuild the image after it from the one
2034 /// before, and a change touching most of the file is left to be read whole.
2035 #[test]
2036 fn deltas_rebuild_the_image_after_a_commit() {
2037 let before: Vec<u8> = (0..20_000u32).map(|at| (at % 251) as u8).collect();
2038 let mut after = before.clone();
2039 after[3] ^= 1;
2040 after[5000..5010].fill(9);
2041 after[5050] ^= 1;
2042 after[17_000] ^= 1;
2043 after.extend_from_slice(&[7; 3000]);
2044 let writes = delta(&before, &after).unwrap();
2045 assert_eq!(
2046 writes.len(),
2047 4,
2048 "the header, two runs merged as one, one more, the tail"
2049 );
2050 assert_eq!(
2051 written(&before, after.len() as u64, &writes).unwrap(),
2052 after
2053 );
2054 let rewritten: Vec<u8> = before.iter().map(|byte| byte ^ 1).collect();
2055 assert!(delta(&before, &rewritten).is_none());
2056 assert!(delta(&before, &before[..10_000]).is_none());
2057 let outside = [Written {
2058 offset: after.len() as u64,
2059 bytes: vec![1],
2060 }];
2061 assert!(written(&before, after.len() as u64, &outside).is_none());
2062 }
2063}