authorgravatar for git@paperclover.netclover caruso <git@paperclover.net> 2026-10-03 20:29:30-07:00
committergravatar for git@paperclover.netclover caruso <git@paperclover.net> 2026-10-03 21:12:52-07:00
log724f7aa7062d2a561f051778a85ee35deb71ebf7
tree4ed075c3aeb20dbf1f12ec36c9d2b4ecaf620bfa
parent725072f4398ea34fd120598c30e2c7c90714ffd3
signature Signed by SSH key SHA256:52mNGHRsVFBDED9IAX5pe+LRWUefqTbxEReunq21QvU

feat: edit Live Share notebooks in the browser

Join desktop hosts through sealed WebSockets using the shared Live Share protocol. Await remote operations without blocking the browser, keep ordered SQLite and OPFS writes durable before publication, and retain queued edits across reloads and reconnects. Refresh shared section catalogs, arm watches when all sections are held, and distinguish device removal from a stopped share. Hosting and section organization remain on desktop. Validated all CI lanes, browser approval/removal, bidirectional edits, offline reload/recovery, live section discovery, and a fresh OneNote 2010 cold-open of browser-written Unicode. Assisted-by: gpt-6.1-sol

33 files changed, 2242 insertions(+), 496 deletions(-)

arc/platforms.md+4
......@@ -185,6 +185,10 @@ keyboard, the toolbar and the macOS menu bar all run commands from it.
185185 `/Cache`. A storage worker keeps them in the origin's private file system, writing the
186186 byte ranges each burst changed through OPFS's synchronous handles, which only workers get.
187187 One tab at a time holds them.
188- Live Share joins desktop hosts through sealed WebSockets. Approval, device removal,
189 presence and edits use the desktop protocol; replicas await remote operations and an
190 ordered OPFS flush before publishing or returning a durable receipt. Hosting and section
191 organization stay on desktop.
188192- Open Notebook, where the browser has the File System Access API (Chromium), opens a folder
189193 of the user's: mirrored under `/Folders`, its handle kept in IndexedDB, other apps' writes
190194 read every few seconds. A browser takes no locks, so a commit there stands only once the
crates/notebook/Cargo.toml+7-5
......@@ -26,19 +26,21 @@ aes-gcm = "0.11.1"
2626base64 = { version = "0.23.1", default-features = false, features = ["std"] }
2727getrandom = "0.4.3"
2828hmac = "0.13.0"
29mdns-sd = { version = "0.21.4", default-features = false, optional = true }
3029minicbor = { version = "2.3.0", features = ["derive", "alloc"], optional = true }
3130spake2 = { version = "=0.5.0-pre.0", features = ["getrandom"], optional = true }
3231relay = { path = "../relay", optional = true }
33# A relay's TLS, trusting what the app's updates trust: the system's authorities, then Mozilla's.
34rustls = { version = "0.23", default-features = false, features = ["ring", "std", "tls12"], optional = true }
35rustls-native-certs = { version = "0.8", optional = true }
36webpki-root-certs = { version = "1.0", optional = true }
3732# OneNote packages (.onepkg) are cabinet files.
3833cab = "0.6"
3934# std's clock where there is one; the browser's on wasm32-unknown-unknown, whose std has none.
4035web-time = "1.1"
4136
37[target.'cfg(not(target_arch = "wasm32"))'.dependencies]
38mdns-sd = { version = "0.21.4", default-features = false, optional = true }
39# A relay's TLS, trusting the system's authorities, then Mozilla's.
40rustls = { version = "0.23", default-features = false, features = ["ring", "std", "tls12"], optional = true }
41rustls-native-certs = { version = "0.8", optional = true }
42webpki-root-certs = { version = "1.0", optional = true }
43
4244# The browser: its randomness, clock, timers and event loop, and SQLite's files kept with
4345# the notebooks' (`fs`).
4446[target.'cfg(target_arch = "wasm32")'.dependencies]
crates/notebook/src/background.rs+118-62
......@@ -79,7 +79,7 @@ pub(crate) struct Reports {
7979#[derive(Default)]
8080struct Watched {
8181 sections: Vec<Watch>,
82 /// Sections whose file changed since `changed` was last asked.
82 /// Sections or cached folders changed since `changed` was last asked.
8383 changed: Vec<String>,
8484 /// Folders a watch reported without naming the file, listed at `relist`.
8585 folders: BTreeSet<String>,
......@@ -127,7 +127,7 @@ impl Background {
127127 /// lists each folder, both by catalog path, arming any watch that reports the folder's
128128 /// changes through `Reports`, and whether one does; it runs again after a transport
129129 /// failure. `copies` keeps an offline copy of every section.
130 pub(crate) fn start<R, B, L>(
130 pub(crate) fn start<R, B, L, LF>(
131131 copies: bool,
132132 mut connect: impl FnMut(Reports) -> io::Result<((B, L), bool)> + Send + 'static,
133133 notify: impl Fn() + Send + 'static,
......@@ -135,7 +135,8 @@ impl Background {
135135 where
136136 R: Remote,
137137 B: FnMut(&str) -> R,
138 L: FnMut(&str) -> io::Result<Vec<Entry>>,
138 L: FnMut(String) -> LF,
139 LF: Future<Output = io::Result<Vec<Entry>>>,
139140 {
140141 let (signal, receiver) = Signal::new();
141142 let shared = Arc::new(Shared {
......@@ -192,33 +193,39 @@ impl Background {
192193 connection: connection + 1,
193194 };
194195 let connected = connect(reports);
195 let Ok(mut watched) = owner.watched.lock() else {
196 return;
197 };
198 // Whatever ended before now ended with the last connection.
199 watched.connection += 1;
200 watched.lost = false;
201 match connected {
202 Ok((mut files, reported)) => {
203 watched.watching(reported);
204 let folders = watched.rescanning();
205 drop(watched);
206 let listings = list(&mut files.1, folders);
207 let Ok(mut watched) = owner.watched.lock() else {
208 return;
209 };
210 watched.listed(&listings, Some(copies), Instant::now());
211 bound = Some(files);
212 }
213 Err(error) => {
214 let error = Error::RemoteIo(error);
215 let retry = now + RETRY;
216 for watch in &mut watched.sections {
217 news |= watch.fail(&error, None);
218 watch.due = watch.due.max(retry);
196 let folders = {
197 let Ok(mut watched) = owner.watched.lock() else {
198 return;
199 };
200 watched.connection += 1;
201 watched.lost = false;
202 match &connected {
203 Ok((_, reported)) => {
204 watched.watching(*reported);
205 Some(watched.rescanning())
206 }
207 Err(error) => {
208 let error = Error::RemoteIo(io::Error::new(
209 error.kind(),
210 error.to_string(),
211 ));
212 let retry = now + RETRY;
213 for watch in &mut watched.sections {
214 news |= watch.fail(&error, None);
215 watch.due = watch.due.max(retry);
216 }
217 watched.relist = watched.relist.map(|due| due.max(retry));
218 None
219219 }
220 watched.relist = watched.relist.map(|due| due.max(retry));
221220 }
221 };
222 if let (Ok((mut files, _)), Some(folders)) = (connected, folders) {
223 let listings = list(&mut files.1, folders).await;
224 let Ok(mut watched) = owner.watched.lock() else {
225 return;
226 };
227 watched.listed(&listings, Some(copies), Instant::now());
228 bound = Some(files);
222229 }
223230 continue;
224231 }
......@@ -226,7 +233,7 @@ impl Background {
226233 let Some((_, list_folder)) = &mut bound else {
227234 continue;
228235 };
229 let listings = list(list_folder, folders);
236 let listings = list(list_folder, folders).await;
230237 let Ok(mut watched) = owner.watched.lock() else {
231238 return;
232239 };
......@@ -244,14 +251,15 @@ impl Background {
244251 let Some((bind, _)) = &mut bound else {
245252 continue;
246253 };
247 let (queued, outcome) = step(
254 let (queued, outcome) = step_async(
248255 &mut bind(&path),
249256 replica.as_deref(),
250257 seen.as_deref(),
251258 current,
252259 image,
253260 copies,
254 );
261 )
262 .await;
255263 if outcome.as_ref().is_err_and(disconnected) {
256264 bound = None;
257265 }
......@@ -339,10 +347,12 @@ impl Background {
339347 crate::SmbRemote::new(Arc::clone(&bound), file, limit)
340348 };
341349 let root = root.clone();
342 let list = move |folder: &str| {
343 use crate::discover::Source;
344 crate::discover::Smb::new(&client, &root)?
345 .entries(folder, crate::session::LIMITS.entries)
350 let list = move |folder: String| {
351 std::future::ready((|| {
352 use crate::discover::Source;
353 crate::discover::Smb::new(&client, &root)?
354 .entries(&folder, crate::session::LIMITS.entries)
355 })())
346356 };
347357 Ok(((bind, list), reported))
348358 },
......@@ -364,7 +374,10 @@ impl Background {
364374 guest.watch(reports)?;
365375 let (bound, listing) = (Arc::clone(&guest), Arc::clone(&guest));
366376 let bind = move |path: &str| crate::live::share::HostedRemote::new(&bound, path);
367 let list = move |folder: &str| listing.entries(folder);
377 let list = move |folder: String| {
378 let listing = listing.clone();
379 async move { listing.entries_async(&folder).await }
380 };
368381 Ok(((bind, list), true))
369382 },
370383 notify,
......@@ -553,6 +566,16 @@ impl Shared {
553566
554567#[cfg_attr(not(any(feature = "smb", feature = "live")), allow(dead_code))]
555568impl Reports {
569 #[cfg(all(feature = "live", target_arch = "wasm32"))]
570 pub(crate) fn catalog(&self, folder: &str) {
571 if let Some(shared) = self.shared.upgrade()
572 && let Ok(mut watched) = shared.watched.lock()
573 && !watched.changed.iter().any(|path| path == folder)
574 {
575 watched.changed.push(folder.to_owned());
576 }
577 }
578
556579 /// The sections at or below these paths changed.
557580 pub(crate) fn touched(&self, paths: &[String]) {
558581 if let Some(shared) = self.shared.upgrade() {
......@@ -768,14 +791,14 @@ impl Watched {
768791 return Next::Wait(wait.map(|at| at.saturating_duration_since(now)));
769792 };
770793 let watch = &mut self.sections[index];
794 if !bound {
795 return Next::Connect(self.connection);
796 }
771797 if let Some(worker) = watch.held.as_ref().and_then(Weak::upgrade) {
772798 worker.wake();
773799 watch.due = now + interval;
774800 continue;
775801 }
776 if !bound {
777 return Next::Connect(self.connection);
778 }
779802 return Next::Check {
780803 path: watch.path.clone(),
781804 replica: watch.replica.clone(),
......@@ -789,17 +812,16 @@ impl Watched {
789812
790813/// Lists each of `folders` through `list`; a folder that cannot be listed lists nothing, so
791814/// that its sections are checked.
792fn list(
793 list: &mut impl FnMut(&str) -> io::Result<Vec<Entry>>,
815async fn list<F: Future<Output = io::Result<Vec<Entry>>>>(
816 list: &mut impl FnMut(String) -> F,
794817 folders: BTreeSet<String>,
795818) -> BTreeMap<String, Vec<Entry>> {
796 folders
797 .into_iter()
798 .map(|folder| {
799 let entries = list(&folder).unwrap_or_default();
800 (folder, entries)
801 })
802 .collect()
819 let mut listings = BTreeMap::new();
820 for folder in folders {
821 let entries = list(folder.clone()).await.unwrap_or_default();
822 listings.insert(folder, entries);
823 }
824 listings
803825}
804826
805827/// Deletes the replica at `replica` if it holds nothing unpublished and no one holds it.
......@@ -870,7 +892,7 @@ fn summary(status: &SyncStatus) -> (bool, Option<io::ErrorKind>, u64) {
870892/// with `current` is the stamp now, unread. With `copies`, a section without a replica gets
871893/// one from the file as it is now. `image`, the file as discovery read it, stands in for
872894/// reading it while the stamp is still its own.
873fn step<R: Remote>(
895async fn step_async<R: Remote>(
874896 remote: &mut R,
875897 replica: Option<&Path>,
876898 seen: Option<&Stamp>,
......@@ -880,7 +902,7 @@ fn step<R: Remote>(
880902) -> (Option<u64>, Result<(Stamp, bool)>) {
881903 let stamp = match seen.filter(|_| current) {
882904 Some(seen) => seen.clone(),
883 None => match remote.stamp() {
905 None => match crate::sync::awaited(remote, Remote::stamp).await {
884906 Ok(stamp) => stamp,
885907 Err(error) => return (None, Err(Error::RemoteIo(error))),
886908 },
......@@ -898,14 +920,18 @@ fn step<R: Remote>(
898920 if !copies {
899921 return (Some(0), Ok((stamp, moved)));
900922 }
901 let copied = (|| {
902 let image = remote.read().map_err(Error::RemoteIo)?;
923 let copied = (async {
924 let image = crate::sync::awaited(remote, Remote::read)
925 .await
926 .map_err(Error::RemoteIo)?;
903927 if let Some(folder) = replica.parent() {
904928 fs::create_dir_all(folder)?;
905929 }
906930 Replica::seed(replica, &image, None)?;
931 fs::durable().await?;
907932 Ok(Stamp::of(&image)?)
908 })();
933 })
934 .await;
909935 return match copied {
910936 Ok(stamp) => (Some(0), Ok((stamp, moved))),
911937 // A session made it first.
......@@ -916,7 +942,7 @@ fn step<R: Remote>(
916942 };
917943 }
918944 // A version kept beside the file leaves its stamp as it was.
919 let versions = match remote.versions() {
945 let versions = match crate::sync::awaited(remote, Remote::versions).await {
920946 Ok(versions) => versions,
921947 Err(error) => return (None, Err(Error::RemoteIo(error))),
922948 };
......@@ -932,16 +958,22 @@ fn step<R: Remote>(
932958 Err(error) if error.busy() => return (None, Ok((stamp, false))),
933959 Err(error) => return (None, Err(error)),
934960 };
935 let synced = (|| {
961 let synced = (async {
936962 let mut changed = moved;
937963 loop {
938 let synced = replica.sync_once(remote)?;
964 let synced = replica.sync_once_async(remote).await?;
939965 changed |= !synced.changed.is_empty();
940966 if !matches!(synced.edit, Some((_, EditStatus::Published { .. }))) {
941 return Ok((remote.stamp().map_err(Error::RemoteIo)?, changed));
967 return Ok((
968 crate::sync::awaited(remote, Remote::stamp)
969 .await
970 .map_err(Error::RemoteIo)?,
971 changed,
972 ));
942973 }
943974 }
944 })();
975 })
976 .await;
945977 let queued = replica
946978 .recovery_summary()
947979 .ok()
......@@ -958,6 +990,10 @@ struct Discovered<'a, R> {
958990}
959991
960992impl<R: Remote> Remote for Discovered<'_, R> {
993 fn pending(&mut self) -> Option<std::pin::Pin<Box<dyn Future<Output = ()> + '_>>> {
994 self.remote.pending()
995 }
996
961997 fn read(&mut self) -> io::Result<Vec<u8>> {
962998 match self.image.take() {
963999 Some((stamp, image)) if stamp == self.stamp => Ok(image),
......@@ -1036,6 +1072,20 @@ mod tests {
10361072 .collect()
10371073 }
10381074
1075 #[test]
1076 fn a_notebook_with_every_section_held_still_connects_its_watch() {
1077 let now = Instant::now();
1078 let mut watched = Watched::default();
1079 watched.watch(sections(1, true), now);
1080 let (held, woken) = Signal::new();
1081 watched.sections[0].held = Some(Arc::downgrade(&held));
1082 assert!(matches!(watched.next(now, false), Next::Connect(0)));
1083 watched.watching(true);
1084 assert!(held.watched.load(Ordering::Acquire));
1085 assert!(matches!(watched.next(now, true), Next::Wait(_)));
1086 assert!(woken.recv_timeout(Duration::from_millis(50)).is_ok());
1087 }
1088
10391089 /// The notebook's folders as a listing shows its first `count` sections.
10401090 fn listing(count: usize) -> BTreeMap<String, Vec<Entry>> {
10411091 let mut folders: BTreeMap<String, Vec<Entry>> = BTreeMap::new();
......@@ -1311,7 +1361,7 @@ mod tests {
13111361 armed.send(reports).unwrap();
13121362 let (files, listed) = (remote.clone(), remote.clone());
13131363 let bind = move |path: &str| File(files.clone(), path.to_owned());
1314 let list = move |folder: &str| Ok(listed.list(folder));
1364 let list = move |folder: String| std::future::ready(Ok(listed.list(&folder)));
13151365 Ok(((bind, list), true))
13161366 },
13171367 || {},
......@@ -1387,7 +1437,7 @@ mod tests {
13871437 armed.send(reports).unwrap();
13881438 let (files, listed) = (remote.clone(), remote.clone());
13891439 let bind = move |path: &str| File(files.clone(), path.to_owned());
1390 let list = move |folder: &str| Ok(listed.list(folder));
1440 let list = move |folder: String| std::future::ready(Ok(listed.list(&folder)));
13911441 Ok(((bind, list), true))
13921442 },
13931443 || {},
......@@ -1412,10 +1462,16 @@ mod tests {
14121462 quiet(&stamps, SETTLE).is_empty(),
14131463 "the worker reads, not this"
14141464 );
1415 // The watch ending hands polling back to the worker.
14161465 watch.lost();
14171466 assert!(woken.recv_timeout(SETTLE).is_ok());
1418 assert!(!signal.watched.load(Ordering::Acquire));
1467 let rearmed = watches.recv_timeout(SETTLE).unwrap();
1468 while !signal.watched.load(Ordering::Acquire) {
1469 woken.recv_timeout(SETTLE * 2).unwrap();
1470 }
1471 assert!(signal.watched.load(Ordering::Acquire));
1472 rearmed.touched(std::slice::from_ref(&held));
1473 assert!(woken.recv_timeout(SETTLE * 2).is_ok());
1474 assert!(quiet(&stamps, SETTLE).is_empty());
14191475 drop(signal);
14201476 drop(background);
14211477 }
crates/notebook/src/fs.rs+5
......@@ -19,3 +19,8 @@ pub fn commit_file(
1919mod web;
2020#[cfg(target_arch = "wasm32")]
2121pub use web::*;
22
23#[cfg(not(target_arch = "wasm32"))]
24pub async fn durable() -> std::io::Result<()> {
25 Ok(())
26}
crates/notebook/src/fs/web.rs+71-32
......@@ -32,6 +32,32 @@ struct Data {
3232const RANGES: usize = 64;
3333
3434impl Data {
35 fn flush(&mut self) {
36 if let Some(path) = self.path.clone() {
37 with(|files| {
38 if files.changed.remove(&path) {
39 files.durable.push((path, self.change()));
40 }
41 });
42 }
43 }
44
45 fn change(&mut self) -> Change {
46 let length = self.bytes.len();
47 let ranges = self
48 .unwritten
49 .replace(Vec::new())
50 .unwrap_or_else(|| std::iter::once(0..length).collect());
51 Change::File {
52 length: length as u64,
53 ranges: ranges
54 .into_iter()
55 .map(|range| range.start.min(length)..range.end.min(length))
56 .filter(|range| !range.is_empty())
57 .map(|range| (range.start as u64, self.bytes[range].to_vec()))
58 .collect(),
59 }
60 }
3561 fn new(
3662 bytes: Vec<u8>,
3763 modified: f64,
......@@ -79,6 +105,8 @@ struct Files {
79105 nodes: BTreeMap<PathBuf, Node>,
80106 /// Paths written, made or removed since the host last asked.
81107 changed: BTreeSet<PathBuf>,
108 /// Flushes and removals in order, including SQLite's WAL before its database checkpoint.
109 durable: Vec<(PathBuf, Change)>,
82110 /// Folders the host mirrors from elsewhere (`mount`).
83111 mounts: BTreeSet<PathBuf>,
84112 /// Images committed under a mount, waiting for the host to write them (`committed`).
......@@ -164,39 +192,24 @@ pub enum Change {
164192
165193/// Whether any path changed since `changes` was last called.
166194pub fn changed() -> bool {
167 FILES.with_borrow(|files| !files.changed.is_empty() || !files.committed.is_empty())
195 FILES.with_borrow(|files| {
196 !files.durable.is_empty() || !files.changed.is_empty() || !files.committed.is_empty()
197 })
168198}
169199
170/// The paths changed since the last call, each with how; a folder before what it holds.
200/// Ordered file flushes, followed by changes not yet flushed.
171201pub fn changes() -> Vec<(PathBuf, Change)> {
172202 FILES.with_borrow_mut(|files| {
173 std::mem::take(&mut files.changed)
174 .into_iter()
175 .map(|path| {
176 let change = match files.nodes.get(&path) {
177 None => Change::Removed,
178 Some(Node::Directory) => Change::Directory,
179 Some(Node::File(data)) => {
180 let mut data = data.borrow_mut();
181 let length = data.bytes.len();
182 let ranges = data
183 .unwritten
184 .replace(Vec::new())
185 .unwrap_or_else(|| std::iter::once(0..length).collect());
186 Change::File {
187 length: length as u64,
188 ranges: ranges
189 .into_iter()
190 .map(|range| range.start.min(length)..range.end.min(length))
191 .filter(|range| !range.is_empty())
192 .map(|range| (range.start as u64, data.bytes[range].to_vec()))
193 .collect(),
194 }
195 }
196 };
197 (path, change)
198 })
199 .collect()
203 let mut changes = std::mem::take(&mut files.durable);
204 changes.extend(std::mem::take(&mut files.changed).into_iter().map(|path| {
205 let change = match files.nodes.get(&path) {
206 None => Change::Removed,
207 Some(Node::Directory) => Change::Directory,
208 Some(Node::File(data)) => data.borrow_mut().change(),
209 };
210 (path, change)
211 }));
212 changes
200213 })
201214}
202215
......@@ -445,9 +458,13 @@ pub fn remove_file(path: impl AsRef<Path>) -> io::Result<()> {
445458 let path = normal(path.as_ref());
446459 with(|files| {
447460 let data = files.file(&path)?;
448 data.borrow_mut().path = None;
461 let mut data = data.borrow_mut();
462 if files.changed.remove(&path) {
463 files.durable.push((path.clone(), data.change()));
464 }
465 data.path = None;
449466 files.nodes.remove(&path);
450 files.changed.insert(path);
467 files.durable.push((path, Change::Removed));
451468 Ok(())
452469 })
453470}
......@@ -687,11 +704,12 @@ impl File {
687704 }
688705
689706 pub fn sync_all(&self) -> io::Result<()> {
707 self.data.borrow_mut().flush();
690708 Ok(())
691709 }
692710
693711 pub fn sync_data(&self) -> io::Result<()> {
694 Ok(())
712 self.sync_all()
695713 }
696714
697715 pub fn set_len(&self, size: u64) -> io::Result<()> {
......@@ -848,3 +866,24 @@ pub fn supersede_file(
848866 .map_err(not_committed)?;
849867 rename(with, path).map_err(not_committed)
850868}
869
870/// Waits until the browser host has durably written every ordered flush handed to it.
871pub async fn durable() -> io::Result<()> {
872 use wasm_bindgen::prelude::*;
873 #[wasm_bindgen]
874 extern "C" {
875 #[wasm_bindgen(catch, js_namespace = globalThis, js_name = snowboundFlushStorage)]
876 fn flush_storage() -> Result<js_sys::Promise, JsValue>;
877 }
878 let failure = |error: JsValue| {
879 io::Error::other(
880 error
881 .as_string()
882 .unwrap_or_else(|| "Browser storage failed".into()),
883 )
884 };
885 wasm_bindgen_futures::JsFuture::from(flush_storage().map_err(failure)?)
886 .await
887 .map_err(failure)?;
888 Ok(())
889}
crates/notebook/src/fs/web/sqlite.rs+1
......@@ -49,6 +49,7 @@ impl VfsFile for File {
4949 }
5050
5151 fn flush(&mut self) -> VfsResult<()> {
52 self.0.borrow_mut().flush();
5253 Ok(())
5354 }
5455
crates/notebook/src/lib.rs+1
......@@ -5,6 +5,7 @@
55pub mod discover;
66pub mod fs;
77#[cfg(feature = "live")]
8#[cfg_attr(target_arch = "wasm32", path = "live/web.rs")]
89pub mod live;
910#[cfg(feature = "smb")]
1011pub mod smb;
crates/notebook/src/live.rs+3-123
......@@ -19,7 +19,6 @@ pub use wire::{Caret, Guid, Hello, Presence, Spot};
1919
2020use mdns_sd::{IfKind, ServiceDaemon, ServiceEvent, ServiceInfo};
2121use minicbor::Encode;
22use sha2::{Digest, Sha256};
2322use std::{
2423 collections::BTreeMap,
2524 io::{self, Read, Write},
......@@ -46,124 +45,9 @@ const PATIENCE: Duration = Duration::from_secs(30);
4645/// Wrong tries of a code met off any relay before it admits no one new, as a relay burns one.
4746const TRIES: u32 = 5;
4847
49/// The secret peers meet through.
50#[derive(Clone, Debug, PartialEq, Eq)]
51pub enum Room {
52 /// A notebook's room, or a share's: a random secret only its members hold.
53 Notebook([u8; 16]),
54 /// A code typed on both ends (`code`), with a password where one is set: its room's
55 /// number names it on the network, and its secret and password only the two people know.
56 /// Its `owner`, the end sharing it, takes the number from a relay (the one the code has,
57 /// coming back; any free one for a secret alone) or picks one where it has no relay, and
58 /// judges every try, burning the code after too many wrong ones.
59 Code {
60 code: String,
61 password: String,
62 owner: bool,
63 },
64}
65
66impl Room {
67 /// The room a code typed on this end leads to.
68 pub fn join(code: &str, password: &str) -> Self {
69 Self::Code {
70 code: code.to_owned(),
71 password: password.to_owned(),
72 owner: false,
73 }
74 }
75
76 /// The room of a code this end shares: its secret, or a whole code to keep its number.
77 pub fn share(code: &str, password: &str) -> Self {
78 Self::Code {
79 code: code.to_owned(),
80 password: password.to_owned(),
81 owner: true,
82 }
83 }
84
85 /// What names the room in the clear: a hash of a room's secret, a code's number; none for
86 /// a code without one yet.
87 fn tag(&self) -> Option<String> {
88 match self {
89 Room::Notebook(id) => Some(hex(&Sha256::digest(
90 [&b"Snowbound room "[..], id].concat(),
91 )[..8])),
92 Room::Code { code, .. } => code_parts(code).0.map(|number| format!("code-{number}")),
93 }
94 }
95
96 fn secret(&self) -> Vec<u8> {
97 match self {
98 Room::Notebook(id) => id.to_vec(),
99 Room::Code { code, password, .. } => {
100 let mut secret = code_parts(code).1.into_bytes();
101 if !password.is_empty() {
102 secret.push(b'\n');
103 secret.extend_from_slice(password.as_bytes());
104 }
105 secret
106 }
107 }
108 }
109
110 fn owner(&self) -> bool {
111 matches!(self, Room::Code { owner: true, .. })
112 }
113}
114
115/// A code's room number, if it has one, and its secret: a whole code, or a secret alone.
116fn code_parts(code: &str) -> (Option<u32>, String) {
117 match code::parse(code) {
118 Some((number, secret)) => (Some(number), secret),
119 None => (None, code.trim().to_uppercase()),
120 }
121}
122
123/// Where peers are looked for.
124#[derive(Clone, Copy, Debug, PartialEq, Eq)]
125pub enum Reach {
126 /// Every network this computer is on.
127 Network,
128 /// This computer alone, for two copies of the app side by side.
129 Loopback,
130}
131
132#[derive(Clone, Debug)]
133pub struct Peer {
134 pub hello: Arc<Hello>,
135 /// None until the peer first says where it is.
136 pub presence: Option<Presence>,
137}
138
139/// What happened, as `Live::start`'s `events` hears it on a network thread.
140pub enum Event<'a> {
141 /// The peers, their presence, the code or the relay's answer changed.
142 Changed,
143 /// A stream to a peer opened, with the line to it.
144 Met(&'a Arc<Hello>, &'a Line),
145 /// A peer's stream ended.
146 Left(&'a Arc<Hello>),
147 /// A frame of a kind presence doesn't read itself, on a stream or to the group.
148 Frame {
149 from: &'a Arc<Hello>,
150 kind: u16,
151 body: &'a [u8],
152 },
153}
154
155/// How the relay last answered.
156#[derive(Clone, Debug, PartialEq, Eq)]
157pub enum Relayed {
158 /// Not asked yet, or no relay.
159 Unknown,
160 /// In the room.
161 Joined,
162 /// Answered with this HTTP status, and how long it asked to wait.
163 Refused(u16, Option<Duration>),
164 /// Not reached, and why.
165 Unreachable(Trouble),
166}
48mod model;
49pub use model::{Event, Peer, Reach, Relayed, Room};
50use model::{code_parts, hex};
16751
16852/// Presence on the network while it lives; dropping it leaves.
16953pub struct Live {
......@@ -909,9 +793,5 @@ fn decode(text: &str) -> String {
909793 String::from_utf8_lossy(&decoded).into_owned()
910794}
911795
912fn hex(bytes: &[u8]) -> String {
913 bytes.iter().map(|byte| format!("{byte:02x}")).collect()
914}
915
916796#[cfg(test)]
917797mod tests;
crates/notebook/src/live/model.rs created+149
......@@ -0,0 +1,149 @@
1use super::{Hello, Line, Presence, code};
2use sha2::{Digest, Sha256};
3use std::{sync::Arc, time::Duration};
4
5pub(super) fn hex(bytes: &[u8]) -> String {
6 bytes.iter().map(|byte| format!("{byte:02x}")).collect()
7}
8
9/// The secret peers meet through.
10#[derive(Clone, Debug, PartialEq, Eq)]
11pub enum Room {
12 /// A notebook's room, or a share's: a random secret only its members hold.
13 Notebook([u8; 16]),
14 /// A code typed on both ends (`code`), with a password where one is set: its room's
15 /// number names it on the network, and its secret and password only the two people know.
16 /// Its `owner`, the end sharing it, takes the number from a relay (the one the code has,
17 /// coming back; any free one for a secret alone) or picks one where it has no relay, and
18 /// judges every try, burning the code after too many wrong ones.
19 Code {
20 code: String,
21 password: String,
22 owner: bool,
23 },
24}
25
26impl Room {
27 /// The room a code typed on this end leads to.
28 pub fn join(code: &str, password: &str) -> Self {
29 Self::Code {
30 code: code.to_owned(),
31 password: password.to_owned(),
32 owner: false,
33 }
34 }
35
36 /// The room of a code this end shares: its secret, or a whole code to keep its number.
37 pub fn share(code: &str, password: &str) -> Self {
38 Self::Code {
39 code: code.to_owned(),
40 password: password.to_owned(),
41 owner: true,
42 }
43 }
44
45 /// What names the room in the clear: a hash of a room's secret, a code's number; none for
46 /// a code without one yet.
47 pub(super) fn tag(&self) -> Option<String> {
48 match self {
49 Room::Notebook(id) => Some(hex(&Sha256::digest(
50 [&b"Snowbound room "[..], id].concat(),
51 )[..8])),
52 Room::Code { code, .. } => code_parts(code).0.map(|number| format!("code-{number}")),
53 }
54 }
55
56 pub(super) fn secret(&self) -> Vec<u8> {
57 match self {
58 Room::Notebook(id) => id.to_vec(),
59 Room::Code { code, password, .. } => {
60 let mut secret = code_parts(code).1.into_bytes();
61 if !password.is_empty() {
62 secret.push(b'\n');
63 secret.extend_from_slice(password.as_bytes());
64 }
65 secret
66 }
67 }
68 }
69
70 #[cfg(not(target_arch = "wasm32"))]
71 pub(super) fn owner(&self) -> bool {
72 matches!(self, Room::Code { owner: true, .. })
73 }
74}
75
76/// A code's room number, if it has one, and its secret: a whole code, or a secret alone.
77pub(super) fn code_parts(code: &str) -> (Option<u32>, String) {
78 match code::parse(code) {
79 Some((number, secret)) => (Some(number), secret),
80 None => (None, code.trim().to_uppercase()),
81 }
82}
83
84/// Where peers are looked for.
85#[derive(Clone, Copy, Debug, PartialEq, Eq)]
86pub enum Reach {
87 /// Every network this computer is on.
88 Network,
89 /// This computer alone, for two copies of the app side by side.
90 Loopback,
91}
92
93#[derive(Clone, Debug)]
94pub struct Peer {
95 pub hello: Arc<Hello>,
96 /// None until the peer first says where it is.
97 pub presence: Option<Presence>,
98}
99
100/// What happened, as `Live::start`'s `events` hears it on a network thread.
101pub enum Event<'a> {
102 /// The peers, their presence, the code or the relay's answer changed.
103 Changed,
104 /// A stream to a peer opened, with the line to it.
105 Met(&'a Arc<Hello>, &'a Line),
106 /// A peer's stream ended.
107 Left(&'a Arc<Hello>),
108 /// A frame of a kind presence doesn't read itself, on a stream or to the group.
109 Frame {
110 from: &'a Arc<Hello>,
111 kind: u16,
112 body: &'a [u8],
113 },
114}
115
116/// How the relay last answered.
117#[derive(Clone, Debug, PartialEq, Eq)]
118pub enum Relayed {
119 /// Not asked yet, or no relay.
120 Unknown,
121 /// In the room.
122 Joined,
123 /// Answered with this HTTP status, and how long it asked to wait.
124 Refused(u16, Option<Duration>),
125 /// Not reached, and why.
126 Unreachable(Trouble),
127}
128
129/// Why a relay could not be reached, as a person can act on it.
130#[derive(Clone, Copy, Debug, PartialEq, Eq)]
131pub enum Trouble {
132 /// The relay's name, or the proxy's, didn't resolve.
133 Dns,
134 /// Nothing answered at the relay's address: a firewall, or the relay is down.
135 Unreachable,
136 TimedOut,
137 /// The proxy itself couldn't be reached.
138 ProxyUnreachable,
139 /// The proxy asks for a name and password (407).
140 ProxyAuthentication,
141 /// The proxy refused to connect to the relay, with this status.
142 ProxyRefused(u16),
143 /// The relay's certificate wasn't one the system trusts: something on the way presents its
144 /// own.
145 Certificate,
146 /// Something on the way refused both WebSockets and plain requests to the relay.
147 Blocked,
148 Other,
149}
crates/notebook/src/live/share.rs+316-122
......@@ -8,9 +8,7 @@
88
99use super::{
1010 Event, Hello, Line, Live, Peer, Presence, Reach, Relayed, Room, Sender,
11 wire::{
12 self, Delta, Failure, Reply, Request, Touched, Welcome, WireEntry, WireStamp, Written, kind,
13 },
11 wire::{self, Delta, Failure, Reply, Request, Touched, Welcome, WireEntry, Written, kind},
1412};
1513use crate::{Error, Result, background::Reports, discover, session::Storage};
1614use onestore::{CommitError, CommitState, RevisionIndex, Stamp, Store, Transaction};
......@@ -25,11 +23,17 @@ use std::{
2523 mpsc,
2624 },
2725 thread,
28 time::{Duration, Instant},
26 time::Duration,
2927};
28use web_time::Instant;
3029
30#[path = "share/batch.rs"]
3131mod batch;
32#[path = "share/membership.rs"]
3233mod membership;
34#[cfg(target_arch = "wasm32")]
35#[path = "share/web.rs"]
36mod web;
3337pub use membership::Host;
3438
3539/// The most bytes one message of a read or an upload carries.
......@@ -168,12 +172,31 @@ pub fn join_while(
168172 reach: Option<Reach>,
169173 relay: Option<&str>,
170174 continue_joining: impl Fn(bool) -> bool,
175) -> std::result::Result<Welcome, Refusal> {
176 crate::task::ready(join_while_async(
177 me,
178 code,
179 password,
180 reach,
181 relay,
182 continue_joining,
183 ))
184 .map_err(|_| Refusal::Busy)?
185}
186
187pub async fn join_while_async(
188 me: Hello,
189 code: &str,
190 password: &str,
191 reach: Option<Reach>,
192 relay: Option<&str>,
193 continue_joining: impl Fn(bool) -> bool,
171194) -> std::result::Result<Welcome, Refusal> {
172195 let code = self::code(code).ok_or(Refusal::Malformed)?;
173196 let (welcomed, welcome) = mpsc::channel();
174197 let approving = Arc::new(AtomicBool::new(false));
175198 let approval = Arc::clone(&approving);
176 let (changed, waiting) = mpsc::channel();
199 let (changed, waiting) = crate::task::channel();
177200 let live = Live::start(
178201 me,
179202 &Room::join(&code, password),
......@@ -207,7 +230,7 @@ pub fn join_while(
207230 }
208231 }
209232 _ => {
210 let _ = changed.send(());
233 let _ = changed.try_send(());
211234 }
212235 },
213236 )
......@@ -257,7 +280,7 @@ pub fn join_while(
257280 }
258281 _ => {}
259282 }
260 let _ = waiting.recv_timeout(Duration::from_millis(100));
283 crate::task::wait(&waiting, Some(Duration::from_millis(100))).await;
261284 }
262285}
263286
......@@ -963,26 +986,61 @@ struct Inner {
963986 here: Mutex<Presence>,
964987 /// The host and the line to it, while connected.
965988 host: Mutex<Option<(Arc<Hello>, Line)>>,
966 pending: Mutex<HashMap<u64, mpsc::Sender<Reply>>>,
989 pending: Mutex<HashMap<u64, crate::task::Answer<Reply>>>,
967990 next: AtomicU64,
968991 /// Where the host's reports of changed files go, while a background watches.
969992 watch: Mutex<Option<Reports>>,
970 /// The host stopped sharing.
993 /// Why this device can no longer reach the share.
971994 ended: Mutex<Option<Ended>>,
972995 /// Bytes of chunks asked for and not yet given.
973 asked: (Mutex<usize>, Condvar),
996 asked: (Mutex<Chunks>, Condvar),
974997 /// The sections read lately, kept as the host's deltas change them, newest first.
975998 images: Mutex<Vec<(String, Arc<Vec<u8>>)>>,
976999 /// The stamps of sections whose image held is the host's now, as its last report of
9771000 /// them said.
9781001 current: Mutex<HashMap<String, Stamp>>,
1002 #[cfg(target_arch = "wasm32")]
1003 catalog: Mutex<Option<PathBuf>>,
9791004}
9801005
981enum Ended {
1006#[derive(Clone, Copy, Debug, PartialEq, Eq)]
1007pub enum Ended {
9821008 Stopped,
9831009 Removed,
9841010}
9851011
1012#[derive(Default)]
1013struct Chunks {
1014 bytes: usize,
1015 #[cfg(target_arch = "wasm32")]
1016 waiting: Vec<std::task::Waker>,
1017}
1018
1019struct Requested<'a> {
1020 inner: &'a Inner,
1021 id: u64,
1022 chunked: bool,
1023}
1024
1025impl Drop for Requested<'_> {
1026 fn drop(&mut self) {
1027 self.inner.pending.lock().unwrap().remove(&self.id);
1028 if self.chunked {
1029 let (asked, room) = &self.inner.asked;
1030 let mut asked = asked.lock().unwrap();
1031 asked.bytes -= CHUNK;
1032 #[cfg(target_arch = "wasm32")]
1033 let waiting = std::mem::take(&mut asked.waiting);
1034 drop(asked);
1035 room.notify_all();
1036 #[cfg(target_arch = "wasm32")]
1037 for waker in waiting {
1038 waker.wake();
1039 }
1040 }
1041 }
1042}
1043
9861044impl Guest {
9871045 /// Joins share `share` through its `secret` as `me`, where `reach` and `relay` say.
9881046 /// `events` runs on a network thread whenever the host comes or goes, or the peers change.
......@@ -1010,6 +1068,8 @@ impl Guest {
10101068 asked: Default::default(),
10111069 images: Mutex::default(),
10121070 current: Mutex::default(),
1071 #[cfg(target_arch = "wasm32")]
1072 catalog: Mutex::default(),
10131073 });
10141074 let heard = Arc::clone(&inner);
10151075 let live = Live::start(me, &Room::Notebook(secret), reach, relay, move |event| {
......@@ -1030,9 +1090,9 @@ impl Guest {
10301090 host.as_ref().map(|(hello, _)| Arc::clone(hello))
10311091 }
10321092
1033 /// Whether the host said it stopped sharing.
1034 pub fn stopped(&self) -> bool {
1035 self.inner.ended.lock().unwrap().is_some()
1093 /// Why the host ended this device's access.
1094 pub fn ended(&self) -> Option<Ended> {
1095 *self.inner.ended.lock().unwrap()
10361096 }
10371097
10381098 /// Everyone in the share's room: the host and the other guests.
......@@ -1077,27 +1137,62 @@ impl Guest {
10771137 }
10781138
10791139 /// Asks the host `kind` of `request`, waiting for its reply.
1080 fn request(&self, kind: u16, mut request: Request) -> std::result::Result<Reply, Failed> {
1140 #[cfg(not(target_arch = "wasm32"))]
1141 fn request(&self, kind: u16, request: Request) -> std::result::Result<Reply, Failed> {
1142 crate::task::ready(self.request_async(kind, request)).map_err(Failed::Unsent)?
1143 }
1144
1145 async fn request_async(
1146 &self,
1147 kind: u16,
1148 mut request: Request,
1149 ) -> std::result::Result<Reply, Failed> {
10811150 let Some((_, line)) = self.inner.host.lock().unwrap().clone() else {
10821151 return Err(Failed::Unsent(self.offline()));
10831152 };
10841153 let chunked = matches!(kind, kind::READ | kind::READ_FILE | kind::PUT);
10851154 if chunked {
1086 let (asked, room) = &self.inner.asked;
1087 let mut asked = asked.lock().unwrap();
1088 while *asked + CHUNK > WINDOW {
1089 asked = room.wait(asked).unwrap();
1155 #[cfg(not(target_arch = "wasm32"))]
1156 {
1157 let (asked, _room) = &self.inner.asked;
1158 let mut asked = asked.lock().unwrap();
1159 while asked.bytes + CHUNK > WINDOW {
1160 asked = _room.wait(asked).unwrap();
1161 }
1162 asked.bytes += CHUNK;
10901163 }
1091 *asked += CHUNK;
1164 #[cfg(target_arch = "wasm32")]
1165 std::future::poll_fn(|context| {
1166 let mut asked = self.inner.asked.0.lock().unwrap();
1167 if asked.bytes + CHUNK > WINDOW {
1168 if !asked
1169 .waiting
1170 .iter()
1171 .any(|waker| waker.will_wake(context.waker()))
1172 {
1173 asked.waiting.push(context.waker().clone());
1174 }
1175 std::task::Poll::Pending
1176 } else {
1177 asked.bytes += CHUNK;
1178 std::task::Poll::Ready(())
1179 }
1180 })
1181 .await;
10921182 }
10931183 let id = self.inner.next.fetch_add(1, Ordering::Relaxed);
10941184 request.id = id;
1095 let (answer, answered) = mpsc::channel();
1185 let (answer, answered) = crate::task::response();
10961186 self.inner.pending.lock().unwrap().insert(id, answer);
1187 let _requested = Requested {
1188 inner: &self.inner,
1189 id,
1190 chunked,
1191 };
10971192 let sent = line.send(kind, &request);
10981193 let reply = match sent {
10991194 Err(error) => Err(Failed::Unsent(error)),
1100 Ok(()) => match answered.recv_timeout(TIMEOUT) {
1195 Ok(()) => match crate::task::answered(&answered, TIMEOUT).await {
11011196 Ok(reply) => Ok(reply),
11021197 Err(mpsc::RecvTimeoutError::Timeout) => Err(Failed::Lost(io::Error::new(
11031198 io::ErrorKind::TimedOut,
......@@ -1106,12 +1201,6 @@ impl Guest {
11061201 Err(mpsc::RecvTimeoutError::Disconnected) => Err(Failed::Lost(self.offline())),
11071202 },
11081203 };
1109 self.inner.pending.lock().unwrap().remove(&id);
1110 if chunked {
1111 let (asked, room) = &self.inner.asked;
1112 *asked.lock().unwrap() -= CHUNK;
1113 room.notify_all();
1114 }
11151204 let reply = reply?;
11161205 match reply.failure {
11171206 Some(failure) => Err(Failed::Refused(failure)),
......@@ -1119,12 +1208,26 @@ impl Guest {
11191208 }
11201209 }
11211210
1211 #[cfg(not(target_arch = "wasm32"))]
11221212 fn ask(&self, kind: u16, request: Request) -> io::Result<Reply> {
11231213 self.request(kind, request).map_err(Failed::io)
11241214 }
11251215
1216 async fn ask_async(&self, kind: u16, request: Request) -> io::Result<Reply> {
1217 self.request_async(kind, request).await.map_err(Failed::io)
1218 }
1219
11261220 /// Puts `bytes` in `request`, or uploads them first where they are large.
1221 #[cfg(not(target_arch = "wasm32"))]
11271222 fn carry(&self, request: &mut Request, bytes: Vec<u8>) -> std::result::Result<(), Failed> {
1223 crate::task::ready(self.carry_async(request, bytes)).map_err(Failed::Unsent)?
1224 }
1225
1226 async fn carry_async(
1227 &self,
1228 request: &mut Request,
1229 bytes: Vec<u8>,
1230 ) -> std::result::Result<(), Failed> {
11281231 if bytes.len() <= CHUNK {
11291232 request.bytes = Some(bytes);
11301233 return Ok(());
......@@ -1138,7 +1241,8 @@ impl Guest {
11381241 ..Request::default()
11391242 };
11401243 // Nothing the upload carries happens before the request that uses it.
1141 self.request(kind::PUT, put)
1244 self.request_async(kind::PUT, put)
1245 .await
11421246 .map_err(|failed| match failed {
11431247 Failed::Lost(error) => Failed::Unsent(error),
11441248 failed => failed,
......@@ -1150,22 +1254,29 @@ impl Guest {
11501254
11511255 /// Reads the file at `path` a chunk at a time, as `kind` reads it; a section held, as
11521256 /// the changes to it.
1257 #[cfg(not(target_arch = "wasm32"))]
11531258 fn read(&self, kind: u16, path: &str, limit: usize) -> io::Result<Vec<u8>> {
1259 crate::task::ready(self.read_async(kind, path, limit))?
1260 }
1261
1262 async fn read_async(&self, kind: u16, path: &str, limit: usize) -> io::Result<Vec<u8>> {
11541263 let held = (kind == kind::READ)
11551264 .then(|| self.inner.image(path))
11561265 .flatten();
1157 let first = self.ask(
1158 kind,
1159 Request {
1160 path: path.to_owned(),
1161 offset: Some(0),
1162 limit: Some(limit as u64),
1163 stamp: (held.as_deref())
1164 .and_then(|image| Stamp::of(image).ok())
1165 .map(|stamp| (&stamp).into()),
1166 ..Request::default()
1167 },
1168 )?;
1266 let first = self
1267 .ask_async(
1268 kind,
1269 Request {
1270 path: path.to_owned(),
1271 offset: Some(0),
1272 limit: Some(limit as u64),
1273 stamp: (held.as_deref())
1274 .and_then(|image| Stamp::of(image).ok())
1275 .map(|stamp| (&stamp).into()),
1276 ..Request::default()
1277 },
1278 )
1279 .await?;
11691280 if let (Some(writes), Some(held)) = (&first.writes, held) {
11701281 let stamp: Option<Stamp> = first.stamp.as_ref().and_then(|s| s.try_into().ok());
11711282 let image = written(&held, first.length.unwrap_or_default(), writes)
......@@ -1181,7 +1292,7 @@ impl Guest {
11811292 let mut image = first.bytes.unwrap_or_default();
11821293 while image.len() < length {
11831294 let chunk = self
1184 .ask(
1295 .ask_async(
11851296 kind,
11861297 Request {
11871298 path: path.to_owned(),
......@@ -1189,7 +1300,8 @@ impl Guest {
11891300 handle: first.handle,
11901301 ..Request::default()
11911302 },
1192 )?
1303 )
1304 .await?
11931305 .bytes
11941306 .unwrap_or_default();
11951307 if chunk.is_empty() {
......@@ -1210,35 +1322,52 @@ impl Guest {
12101322 (Stamp::of(&image).ok().as_ref() == Some(stamp)).then(|| image.to_vec())
12111323 }
12121324
1325 #[cfg(not(target_arch = "wasm32"))]
12131326 pub(crate) fn entries(&self, folder: &str) -> io::Result<Vec<discover::Entry>> {
1214 let reply = self.ask(
1215 kind::LIST,
1216 Request {
1217 path: folder.to_owned(),
1218 ..Request::default()
1219 },
1220 )?;
1221 Ok(reply
1327 crate::task::ready(self.entries_async(folder))?
1328 }
1329
1330 pub(crate) async fn entries_async(&self, folder: &str) -> io::Result<Vec<discover::Entry>> {
1331 let reply = self
1332 .ask_async(
1333 kind::LIST,
1334 Request {
1335 path: folder.to_owned(),
1336 ..Request::default()
1337 },
1338 )
1339 .await?;
1340 let entries: Vec<_> = reply
12221341 .entries
12231342 .unwrap_or_default()
12241343 .iter()
12251344 .map(entry)
1226 .collect())
1345 .collect();
1346 #[cfg(target_arch = "wasm32")]
1347 self.cache_entries(folder, &entries).await?;
1348 Ok(entries)
12271349 }
12281350
12291351 /// The stamp of the file at `path`: the host's last report of it where this guest holds
12301352 /// that image, else the host's answer.
1353 #[cfg(not(target_arch = "wasm32"))]
12311354 fn stamp(&self, path: &str) -> io::Result<Stamp> {
1355 crate::task::ready(self.stamp_async(path))?
1356 }
1357
1358 async fn stamp_async(&self, path: &str) -> io::Result<Stamp> {
12321359 if let Some(stamp) = self.inner.current.lock().unwrap().get(path) {
12331360 return Ok(stamp.clone());
12341361 }
1235 let reply = self.ask(
1236 kind::STAMP,
1237 Request {
1238 path: path.to_owned(),
1239 ..Request::default()
1240 },
1241 )?;
1362 let reply = self
1363 .ask_async(
1364 kind::STAMP,
1365 Request {
1366 path: path.to_owned(),
1367 ..Request::default()
1368 },
1369 )
1370 .await?;
12421371 reply
12431372 .stamp
12441373 .as_ref()
......@@ -1246,34 +1375,114 @@ impl Guest {
12461375 .try_into()
12471376 }
12481377
1378 #[cfg(not(target_arch = "wasm32"))]
12491379 fn commit(
12501380 &self,
12511381 path: &str,
12521382 transaction: &Transaction,
1383 ) -> std::result::Result<(), CommitError> {
1384 crate::task::ready(self.commit_async(path, transaction)).map_err(|error| CommitError {
1385 state: CommitState::NotCommitted,
1386 error,
1387 })?
1388 }
1389
1390 async fn commit_async(
1391 &self,
1392 path: &str,
1393 transaction: &Transaction,
12531394 ) -> std::result::Result<(), CommitError> {
12541395 let mut request = Request {
12551396 path: path.to_owned(),
12561397 ..Request::default()
12571398 };
1258 self.carry(&mut request, transaction.to_bytes())
1399 self.carry_async(&mut request, transaction.to_bytes())
1400 .await
12591401 .map_err(Failed::commit)?;
1260 self.request(kind::COMMIT, request)
1402 self.request_async(kind::COMMIT, request)
1403 .await
12611404 .map(drop)
12621405 .map_err(Failed::commit)
12631406 }
12641407
1408 #[cfg(not(target_arch = "wasm32"))]
12651409 fn confirm(&self, path: &str, base: &Stamp) -> std::result::Result<(), CommitError> {
1410 crate::task::ready(self.confirm_async(path, base)).map_err(|error| CommitError {
1411 state: CommitState::NotCommitted,
1412 error,
1413 })?
1414 }
1415
1416 async fn confirm_async(
1417 &self,
1418 path: &str,
1419 base: &Stamp,
1420 ) -> std::result::Result<(), CommitError> {
12661421 let request = Request {
12671422 path: path.to_owned(),
12681423 stamp: Some(base.into()),
12691424 ..Request::default()
12701425 };
1271 self.request(kind::CONFIRM, request)
1426 self.request_async(kind::CONFIRM, request)
1427 .await
12721428 .map(drop)
12731429 .map_err(Failed::commit)
12741430 }
12751431
1432 async fn edits_async(
1433 &self,
1434 path: &str,
1435 transaction: &Transaction,
1436 edits: &[crate::PendingEdit],
1437 revisions: &BTreeMap<onestore::ExGuid, onestore::ExGuid>,
1438 ) -> std::result::Result<(), CommitError> {
1439 let bytes = serde_json::to_vec(&batch::Edits {
1440 edits: edits
1441 .iter()
1442 .map(|edit| (edit.author.clone(), edit.edit.clone()))
1443 .collect(),
1444 revisions: revisions.clone(),
1445 })
1446 .map_err(|error| CommitError {
1447 state: CommitState::NotCommitted,
1448 error: io::Error::other(error),
1449 })?;
1450 let mut request = Request {
1451 path: path.to_owned(),
1452 stamp: Some(transaction.base().into()),
1453 ..Request::default()
1454 };
1455 self.carry_async(&mut request, bytes)
1456 .await
1457 .map_err(Failed::commit)?;
1458 match self.request_async(kind::EDITS, request).await {
1459 Ok(reply) => {
1460 if let Some(stamp) = reply.stamp {
1461 let stamp = Stamp::try_from(&stamp).map_err(|error| CommitError {
1462 state: CommitState::Unknown,
1463 error,
1464 })?;
1465 self.inner
1466 .current
1467 .lock()
1468 .unwrap()
1469 .insert(path.to_owned(), stamp);
1470 }
1471 Ok(())
1472 }
1473 Err(Failed::Refused(failure))
1474 if wire::error_kind(failure.kind) == io::ErrorKind::Unsupported =>
1475 {
1476 self.commit_async(path, transaction).await?;
1477 self.inner.published(path, transaction);
1478 Ok(())
1479 }
1480 Err(error) => Err(error.commit()),
1481 }
1482 }
1483
12761484 /// A request on `path` that answers nothing but whether it happened.
1485 #[cfg(not(target_arch = "wasm32"))]
12771486 fn verb(&self, kind: u16, path: &str, request: Request) -> Result<()> {
12781487 self.ask(
12791488 kind,
......@@ -1338,16 +1547,14 @@ impl Inner {
13381547
13391548 fn heard(self: &Arc<Self>, event: Event) {
13401549 let serves = |hello: &Hello| {
1341 hello.serves == Some(self.share)
1342 || self
1343 .host
1344 .lock()
1345 .unwrap()
1346 .as_ref()
1347 .is_some_and(|(host, _)| host.peer == hello.peer)
1550 self.host
1551 .lock()
1552 .unwrap()
1553 .as_ref()
1554 .is_some_and(|(host, _)| host.peer == hello.peer)
13481555 };
13491556 match event {
1350 Event::Met(hello, line) if serves(hello) => {
1557 Event::Met(hello, line) if hello.serves == Some(self.share) => {
13511558 *self.host.lock().unwrap() = Some((Arc::clone(hello), line.clone()));
13521559 }
13531560 Event::Left(hello) if serves(hello) => {
......@@ -1419,7 +1626,13 @@ impl Inner {
14191626 self.relay.as_deref(),
14201627 move |event| {
14211628 if let Some(inner) = inner.upgrade() {
1422 if matches!(event, Event::Frame { .. }) {
1629 if matches!(
1630 event,
1631 Event::Frame {
1632 kind: kind::DELTA | kind::TOUCHED,
1633 ..
1634 }
1635 ) {
14231636 inner.heard(event);
14241637 }
14251638 (inner.events)();
......@@ -1450,6 +1663,8 @@ pub struct HostedRemote {
14501663 /// The stamp last asked for, which an image the guest holds may already have.
14511664 seen: Option<Stamp>,
14521665 rejected: bool,
1666 #[cfg(target_arch = "wasm32")]
1667 pending: Option<web::Pending>,
14531668}
14541669
14551670impl HostedRemote {
......@@ -1459,10 +1674,13 @@ impl HostedRemote {
14591674 path: path.to_owned(),
14601675 seen: None,
14611676 rejected: false,
1677 #[cfg(target_arch = "wasm32")]
1678 pending: None,
14621679 }
14631680 }
14641681}
14651682
1683#[cfg(not(target_arch = "wasm32"))]
14661684impl crate::Remote for HostedRemote {
14671685 fn accepts_edits(&self) -> bool {
14681686 !self.rejected
......@@ -1478,57 +1696,26 @@ impl crate::Remote for HostedRemote {
14781696 edits: &[crate::PendingEdit],
14791697 revisions: &BTreeMap<onestore::ExGuid, onestore::ExGuid>,
14801698 ) -> std::result::Result<(), CommitError> {
1481 let bytes = serde_json::to_vec(&batch::Edits {
1482 edits: edits
1483 .iter()
1484 .map(|edit| (edit.author.clone(), edit.edit.clone()))
1485 .collect(),
1486 revisions: revisions.clone(),
1487 })
1488 .map_err(|error| CommitError {
1489 state: CommitState::NotCommitted,
1490 error: io::Error::other(error),
1491 })?;
1492 let mut request = Request {
1493 path: self.path.clone(),
1494 stamp: Some(transaction.base().into()),
1495 ..Request::default()
1496 };
1497 self.guest
1498 .carry(&mut request, bytes)
1499 .map_err(Failed::commit)?;
1500 let result = self.guest.request(kind::EDITS, request);
1501 match result {
1502 Ok(reply) => {
1503 if let Some(stamp) = reply.stamp {
1504 let stamp = Stamp::try_from(&stamp).map_err(|error| CommitError {
1505 state: CommitState::Unknown,
1506 error,
1507 })?;
1508 self.guest
1509 .inner
1510 .current
1511 .lock()
1512 .unwrap()
1513 .insert(self.path.clone(), stamp);
1514 }
1515 self.seen = None;
1516 Ok(())
1517 }
1518 Err(Failed::Refused(failure))
1519 if wire::error_kind(failure.kind) == io::ErrorKind::Unsupported =>
1520 {
1521 self.publish(transaction)
1522 }
1523 Err(error) => {
1524 let error = error.commit();
1525 if error.state == CommitState::NotCommitted {
1526 self.rejected = true;
1527 self.guest.inner.current.lock().unwrap().remove(&self.path);
1528 }
1529 Err(error)
1530 }
1699 let result =
1700 crate::task::ready(
1701 self.guest
1702 .edits_async(&self.path, transaction, edits, revisions),
1703 )
1704 .map_err(|error| CommitError {
1705 state: CommitState::NotCommitted,
1706 error,
1707 })?;
1708 if result.is_ok() {
1709 self.seen = None;
1710 }
1711 if result
1712 .as_ref()
1713 .is_err_and(|error| error.state == CommitState::NotCommitted)
1714 {
1715 self.rejected = true;
1716 self.guest.inner.current.lock().unwrap().remove(&self.path);
15311717 }
1718 result
15321719 }
15331720
15341721 fn read(&mut self) -> io::Result<Vec<u8>> {
......@@ -1574,17 +1761,23 @@ type Listings = BTreeMap<String, Vec<(String, u8, u64, u64)>>;
15741761
15751762impl Hosted {
15761763 pub(crate) fn new(guest: Arc<Guest>, listed: PathBuf) -> Self {
1764 #[cfg(target_arch = "wasm32")]
1765 {
1766 *guest.inner.catalog.lock().unwrap() = Some(listed.clone());
1767 }
15771768 Self { guest, listed }
15781769 }
15791770}
15801771
15811772/// The host's folders, or as they were last listed while it can't be reached.
1773#[cfg(not(target_arch = "wasm32"))]
15821774struct Source<'a> {
15831775 guest: &'a Guest,
15841776 kept: Listings,
15851777 listed: Listings,
15861778}
15871779
1780#[cfg(not(target_arch = "wasm32"))]
15881781impl discover::Source for Source<'_> {
15891782 fn entries(&mut self, path: &str, limit: usize) -> io::Result<Vec<discover::Entry>> {
15901783 let entries = match self.guest.entries(path) {
......@@ -1630,6 +1823,7 @@ impl discover::Source for Source<'_> {
16301823 }
16311824}
16321825
1826#[cfg(not(target_arch = "wasm32"))]
16331827impl Storage for Hosted {
16341828 fn discover(
16351829 &self,
......@@ -1749,7 +1943,7 @@ impl Storage for Hosted {
17491943 fn supersede(&self, path: &str, base: &Stamp, with: &str) -> Result<()> {
17501944 let request = Request {
17511945 to: Some(with.to_owned()),
1752 stamp: Some(WireStamp::from(base)),
1946 stamp: Some(wire::WireStamp::from(base)),
17531947 ..Request::default()
17541948 };
17551949 self.guest.verb(kind::SUPERSEDE, path, request)
crates/notebook/src/live/share/web.rs created+382
......@@ -0,0 +1,382 @@
1use super::*;
2use std::{future::Future, pin::Pin};
3
4pub(super) enum Outcome {
5 Read(Vec<u8>),
6 Stamp(Box<Stamp>),
7 Published,
8}
9
10pub(super) enum Pending {
11 Running(Pin<Box<dyn Future<Output = std::result::Result<Outcome, CommitError>>>>),
12 Ready(std::result::Result<Outcome, CommitError>),
13}
14
15fn uncommitted(error: io::Error) -> CommitError {
16 CommitError {
17 state: CommitState::NotCommitted,
18 error,
19 }
20}
21
22impl HostedRemote {
23 fn operation(
24 &mut self,
25 work: impl Future<Output = std::result::Result<Outcome, CommitError>> + 'static,
26 ) -> std::result::Result<Outcome, CommitError> {
27 match self.pending.take() {
28 Some(Pending::Ready(result)) => result,
29 pending => {
30 self.pending = Some(pending.unwrap_or_else(|| Pending::Running(Box::pin(work))));
31 Err(uncommitted(io::ErrorKind::WouldBlock.into()))
32 }
33 }
34 }
35
36 fn publication(
37 &mut self,
38 result: std::result::Result<Outcome, CommitError>,
39 ) -> std::result::Result<(), CommitError> {
40 match result {
41 Ok(Outcome::Published) => {
42 self.seen = None;
43 Ok(())
44 }
45 Ok(_) => Err(uncommitted(io::ErrorKind::InvalidInput.into())),
46 Err(error) => {
47 if error.error.kind() != io::ErrorKind::WouldBlock
48 && error.state == CommitState::NotCommitted
49 {
50 self.rejected = true;
51 self.guest.inner.current.lock().unwrap().remove(&self.path);
52 }
53 Err(error)
54 }
55 }
56 }
57}
58
59impl crate::Remote for HostedRemote {
60 fn pending(&mut self) -> Option<Pin<Box<dyn Future<Output = ()> + '_>>> {
61 let pending = self.pending.as_mut()?;
62 Some(Box::pin(async move {
63 if let Pending::Running(work) = pending {
64 *pending = Pending::Ready(work.await);
65 }
66 }))
67 }
68
69 fn read(&mut self) -> io::Result<Vec<u8>> {
70 let (guest, path, seen) = (self.guest.clone(), self.path.clone(), self.seen.clone());
71 let result = self
72 .operation(async move {
73 let image = match seen.as_ref().and_then(|stamp| guest.held(&path, stamp)) {
74 Some(image) => image,
75 None => guest
76 .read_async(kind::READ, &path, LIMIT)
77 .await
78 .map_err(uncommitted)?,
79 };
80 Ok(Outcome::Read(image))
81 })
82 .map_err(|error| error.error)?;
83 match result {
84 Outcome::Read(image) => {
85 self.rejected = false;
86 Ok(image)
87 }
88 _ => Err(io::ErrorKind::InvalidInput.into()),
89 }
90 }
91
92 fn stamp(&mut self) -> io::Result<Stamp> {
93 let (guest, path) = (self.guest.clone(), self.path.clone());
94 let result = self
95 .operation(async move {
96 guest
97 .stamp_async(&path)
98 .await
99 .map(|stamp| Outcome::Stamp(Box::new(stamp)))
100 .map_err(uncommitted)
101 })
102 .map_err(|error| error.error)?;
103 match result {
104 Outcome::Stamp(stamp) => {
105 self.seen = Some((*stamp).clone());
106 Ok(*stamp)
107 }
108 _ => Err(io::ErrorKind::InvalidInput.into()),
109 }
110 }
111
112 fn accepts_edits(&self) -> bool {
113 !self.rejected
114 && self
115 .guest
116 .host()
117 .is_some_and(|host| host.ops == Some(1) && host.kinds.contains(&kind::EDITS))
118 }
119
120 fn publish_edits(
121 &mut self,
122 transaction: &Transaction,
123 edits: &[crate::PendingEdit],
124 revisions: &BTreeMap<onestore::ExGuid, onestore::ExGuid>,
125 ) -> std::result::Result<(), CommitError> {
126 let (guest, path, transaction, edits, revisions) = (
127 self.guest.clone(),
128 self.path.clone(),
129 transaction.clone(),
130 edits.to_vec(),
131 revisions.clone(),
132 );
133 let result = self.operation(async move {
134 guest
135 .edits_async(&path, &transaction, &edits, &revisions)
136 .await?;
137 Ok(Outcome::Published)
138 });
139 self.publication(result)
140 }
141
142 fn publish(&mut self, transaction: &Transaction) -> std::result::Result<(), CommitError> {
143 let (guest, path, transaction) =
144 (self.guest.clone(), self.path.clone(), transaction.clone());
145 let result = self.operation(async move {
146 guest.commit_async(&path, &transaction).await?;
147 guest.inner.published(&path, &transaction);
148 Ok(Outcome::Published)
149 });
150 self.publication(result)
151 }
152
153 fn confirm(&mut self, base: &Stamp) -> std::result::Result<(), CommitError> {
154 let (guest, path, base) = (self.guest.clone(), self.path.clone(), base.clone());
155 let result = self.operation(async move {
156 guest.confirm_async(&path, &base).await?;
157 Ok(Outcome::Published)
158 });
159 self.publication(result)
160 }
161}
162
163struct Cached<'a>(&'a Hosted, Listings);
164
165impl Guest {
166 pub(super) async fn cache_entries(
167 &self,
168 folder: &str,
169 entries: &[discover::Entry],
170 ) -> io::Result<()> {
171 let Some(listed) = self.inner.catalog.lock().unwrap().clone() else {
172 return Ok(());
173 };
174 let files = listed.with_extension("files");
175 let mut listings: Listings = crate::fs::read(&listed)
176 .ok()
177 .and_then(|bytes| serde_json::from_slice(&bytes).ok())
178 .unwrap_or_default();
179 let before = listings.get(folder).cloned().unwrap_or_default();
180 let mut current = Vec::with_capacity(entries.len());
181 crate::fs::create_dir_all(files.join(folder))?;
182 for entry in entries {
183 if entry.name.contains(['/', '\\', '\0']) || matches!(entry.name.as_str(), "." | "..") {
184 return Err(io::ErrorKind::InvalidData.into());
185 }
186 let e = wire_entry(entry);
187 let listed = (e.name, e.kind, e.size, e.modified);
188 let path = if folder.is_empty() {
189 entry.name.clone()
190 } else {
191 format!("{folder}/{}", entry.name)
192 };
193 if entry.kind == discover::EntryKind::File
194 && (!before.contains(&listed) || crate::fs::metadata(files.join(&path)).is_err())
195 {
196 let kind = if path.ends_with(".one") || path.ends_with(".onetoc2") {
197 kind::READ
198 } else {
199 kind::READ_FILE
200 };
201 let bytes = self.read_async(kind, &path, LIMIT).await?;
202 crate::fs::write(files.join(path), bytes)?;
203 }
204 current.push(listed);
205 }
206 // Other folders may have finished listing while these files were being read.
207 if let Ok(bytes) = crate::fs::read(&listed) {
208 listings = serde_json::from_slice(&bytes).map_err(io::Error::other)?;
209 }
210 if listings.get(folder) != Some(&current) {
211 listings.insert(folder.to_owned(), current);
212 crate::fs::write(&listed, serde_json::to_vec(&listings)?)?;
213 crate::fs::durable().await?;
214 if let Some(reports) = &*self.inner.watch.lock().unwrap() {
215 reports.catalog(folder);
216 }
217 (self.inner.events)();
218 }
219 Ok(())
220 }
221}
222
223impl discover::Source for Cached<'_> {
224 fn entries(&mut self, path: &str, limit: usize) -> io::Result<Vec<discover::Entry>> {
225 let entries = self.1.get(path).ok_or(io::ErrorKind::NotConnected)?;
226 if entries.len() > limit {
227 return Err(io::ErrorKind::FileTooLarge.into());
228 }
229 Ok(entries
230 .iter()
231 .map(|(name, kind, size, modified)| {
232 entry(&WireEntry {
233 name: name.clone(),
234 kind: *kind,
235 size: *size,
236 modified: *modified,
237 })
238 })
239 .collect())
240 }
241 fn read(&mut self, path: &str, limit: usize) -> io::Result<Vec<u8>> {
242 let bytes = self
243 .0
244 .guest
245 .inner
246 .image(path)
247 .map(|image| image.to_vec())
248 .map(Ok)
249 .unwrap_or_else(|| crate::fs::read(self.0.files().join(path)))?;
250 if bytes.len() > limit {
251 return Err(io::ErrorKind::FileTooLarge.into());
252 }
253 Ok(bytes)
254 }
255 fn read_asset(&mut self, path: &str, limit: usize) -> io::Result<Vec<u8>> {
256 self.read(path, limit)
257 }
258}
259
260impl Hosted {
261 fn files(&self) -> PathBuf {
262 self.listed.with_extension("files")
263 }
264
265 /// Loads a bounded catalog and its files before the synchronous editor opens it.
266 pub async fn prepare(guest: Arc<Guest>, cache: &std::path::Path) -> io::Result<()> {
267 let listed =
268 crate::session::listing(cache, &guest.location()).with_extension("entries.json");
269 let hosted = Self::new(guest, listed);
270 let deadline = Instant::now() + Duration::from_secs(20);
271 let (_, waiting) = crate::task::channel();
272 while hosted.guest.host().is_none() {
273 if Instant::now() >= deadline {
274 return Err(hosted.guest.offline());
275 }
276 crate::task::wait(&waiting, Some(Duration::from_millis(50))).await;
277 }
278 let mut folders = vec![String::new()];
279 let mut count = 0;
280 while let Some(folder) = folders.pop() {
281 if folder.split('/').count() > 32 {
282 return Err(io::ErrorKind::FileTooLarge.into());
283 }
284 let entries = hosted.guest.entries_async(&folder).await?;
285 count += entries.len();
286 if count > 10000 {
287 return Err(io::ErrorKind::FileTooLarge.into());
288 }
289 crate::fs::create_dir_all(hosted.files().join(&folder))?;
290 for entry in &entries {
291 if entry.kind == discover::EntryKind::Directory {
292 folders.push(if folder.is_empty() {
293 entry.name.clone()
294 } else {
295 format!("{folder}/{}", entry.name)
296 });
297 }
298 }
299 }
300 crate::fs::durable().await
301 }
302}
303
304fn desktop() -> Error {
305 io::Error::new(
306 io::ErrorKind::Unsupported,
307 "Create and organize sections in desktop Snowbound.",
308 )
309 .into()
310}
311
312impl Storage for Hosted {
313 fn discover(
314 &self,
315 cache: &mut discover::Cache,
316 limits: discover::Limits,
317 ) -> Result<discover::Folder> {
318 let listings =
319 serde_json::from_slice(&crate::fs::read(&self.listed)?).map_err(io::Error::other)?;
320 Ok(cache.discover(&mut Cached(self, listings), limits)?)
321 }
322 fn location(&self) -> String {
323 self.guest.location()
324 }
325 fn entries(&self, path: &str) -> io::Result<Vec<discover::Entry>> {
326 let listings =
327 serde_json::from_slice(&crate::fs::read(&self.listed)?).map_err(io::Error::other)?;
328 discover::Source::entries(&mut Cached(self, listings), path, 10000)
329 }
330 fn exists(&self, path: &str) -> bool {
331 crate::fs::metadata(self.files().join(path)).is_ok()
332 }
333 fn stamp(&self, path: &str) -> io::Result<Stamp> {
334 Stamp::of(&self.read(path).map_err(io::Error::other)?).map_err(io::Error::other)
335 }
336 fn read(&self, path: &str) -> Result<Vec<u8>> {
337 self.read_file(path, LIMIT)
338 }
339 fn read_file(&self, path: &str, limit: usize) -> Result<Vec<u8>> {
340 Ok(discover::Source::read(
341 &mut Cached(self, Listings::new()),
342 path,
343 limit,
344 )?)
345 }
346 fn create(&self, _: &str, _: &[u8]) -> Result<()> {
347 Err(desktop())
348 }
349 fn create_directory(&self, _: &str) -> Result<()> {
350 Err(desktop())
351 }
352 fn hide(&self, _: &str) -> Result<()> {
353 Err(desktop())
354 }
355 fn rename(&self, _: &str, _: &str) -> Result<()> {
356 Err(desktop())
357 }
358 fn rename_root(&self, _: &str, _: &[String]) -> Result<String> {
359 Err(desktop())
360 }
361 fn replace(&self, _: &str, _: &str) -> Result<()> {
362 Err(desktop())
363 }
364 fn delete(&self, _: &str) -> Result<()> {
365 Err(desktop())
366 }
367 fn place(&self, _: &str, _: [u8; 16], _: &str) -> Result<()> {
368 Err(desktop())
369 }
370 fn commit(&self, _: &str, _: &Transaction) -> Result<()> {
371 Err(desktop())
372 }
373 fn confirm(&self, _: &str, _: &Stamp) -> std::result::Result<(), CommitError> {
374 Err(uncommitted(io::Error::new(
375 io::ErrorKind::Unsupported,
376 "Section management requires desktop Snowbound",
377 )))
378 }
379 fn supersede(&self, _: &str, _: &Stamp, _: &str) -> Result<()> {
380 Err(desktop())
381 }
382}
crates/notebook/src/live/transport.rs+1-21
......@@ -28,27 +28,7 @@ const MOST: usize = 4 << 20;
2828/// Set once WebSockets were refused, so later connections go straight to polling.
2929static POLLING: AtomicBool = AtomicBool::new(false);
3030
31/// Why a relay could not be reached, as a person can act on it.
32#[derive(Clone, Copy, Debug, PartialEq, Eq)]
33pub enum Trouble {
34 /// The relay's name, or the proxy's, didn't resolve.
35 Dns,
36 /// Nothing answered at the relay's address: a firewall, or the relay is down.
37 Unreachable,
38 TimedOut,
39 /// The proxy itself couldn't be reached.
40 ProxyUnreachable,
41 /// The proxy asks for a name and password (407).
42 ProxyAuthentication,
43 /// The proxy refused to connect to the relay, with this status.
44 ProxyRefused(u16),
45 /// The relay's certificate wasn't one the system trusts: something on the way presents its
46 /// own.
47 Certificate,
48 /// Something on the way refused both WebSockets and plain requests to the relay.
49 Blocked,
50 Other,
51}
31pub use super::model::Trouble;
5232
5333pub(super) enum Failure {
5434 Trouble(Trouble, io::Error),
crates/notebook/src/live/web.js created+30
......@@ -0,0 +1,30 @@
1const sockets = new Map();
2let next = 1;
3
4export function liveConnect(url, message, closed) {
5 const id = next++;
6 const socket = new WebSocket(url);
7 socket.binaryType = 'arraybuffer';
8 socket.onmessage = event => message(event.data);
9 socket.onclose = () => {
10 liveClose(id);
11 closed();
12 };
13 sockets.set(id, socket);
14 return id;
15}
16
17export function liveSend(id, bytes) {
18 const socket = sockets.get(id);
19 if (!socket || socket.readyState !== WebSocket.OPEN) throw new Error('Relay disconnected');
20 if (socket.bufferedAmount > 1048576) throw new Error('Relay cannot keep up');
21 socket.send(bytes);
22}
23
24export function liveClose(id) {
25 const socket = sockets.get(id);
26 sockets.delete(id);
27 if (!socket) return;
28 socket.onmessage = socket.onclose = null;
29 socket.close();
30}
crates/notebook/src/live/web.rs created+612
......@@ -0,0 +1,612 @@
1//! Browser Live Share: event-driven relay connections carrying the same sealed streams.
2
3pub use ::relay::code;
4#[path = "group.rs"]
5mod group;
6#[path = "model.rs"]
7mod model;
8#[path = "share.rs"]
9pub mod share;
10#[path = "wire.rs"]
11pub mod wire;
12use model::hex;
13pub use model::{Event, Peer, Reach, Relayed, Room, Trouble};
14pub use wire::{Caret, Guid, Hello, Presence, Spot};
15
16use minicbor::Encode;
17use std::{
18 collections::BTreeMap,
19 io,
20 sync::{
21 Arc, Mutex, Weak,
22 atomic::{AtomicBool, Ordering},
23 },
24 time::Duration,
25};
26use wasm_bindgen::prelude::*;
27use wire::{Side, kind};
28
29#[wasm_bindgen(module = "/src/live/web.js")]
30extern "C" {
31 #[wasm_bindgen(catch, js_name = liveConnect)]
32 fn connect(url: &str, message: &JsValue, closed: &JsValue) -> Result<u32, JsValue>;
33 #[wasm_bindgen(catch, js_name = liveSend)]
34 fn send(id: u32, bytes: &[u8]) -> Result<(), JsValue>;
35 #[wasm_bindgen(js_name = liveClose)]
36 fn close(id: u32);
37}
38
39pub struct Live(Arc<Shared>);
40
41struct Shared {
42 me: Hello,
43 room: Room,
44 url: String,
45 events: Box<dyn Fn(Event) + Send + Sync>,
46 stopped: AtomicBool,
47 state: Mutex<State>,
48}
49
50struct State {
51 socket: u32,
52 slot: u32,
53 relayed: Relayed,
54 failed: u32,
55 outdated: Option<u16>,
56 burned: bool,
57 streams: BTreeMap<u32, Stream>,
58 members: BTreeMap<u32, group::Member>,
59 sealer: Option<group::Sealer>,
60 presence: Presence,
61 scheduled: bool,
62}
63
64struct Stream {
65 bytes: Vec<u8>,
66 opening: Option<wire::Opening>,
67 receive: Option<wire::Sealer>,
68 line: Line,
69 peer: Option<Peer>,
70}
71
72#[derive(Clone)]
73pub struct Line {
74 shared: Weak<Shared>,
75 slot: u32,
76 sealer: Arc<Mutex<Option<wire::Sealer>>>,
77}
78
79impl Line {
80 pub fn send(&self, kind: u16, body: &impl Encode<()>) -> io::Result<()> {
81 let shared = self.shared.upgrade().ok_or(io::ErrorKind::NotConnected)?;
82 let mut bytes = Vec::new();
83 self.sealer
84 .lock()
85 .unwrap()
86 .as_mut()
87 .ok_or(io::ErrorKind::NotConnected)?
88 .send(&mut bytes, kind, body)?;
89 shared.stream(self.slot, &bytes)
90 }
91
92 pub fn hang_up(&self, reason: &str) {
93 let _ = self.send(
94 kind::BYE,
95 &wire::Bye {
96 reason: reason.into(),
97 },
98 );
99 }
100}
101
102#[derive(Clone)]
103pub struct Sender(Weak<Shared>);
104
105impl Sender {
106 pub fn send(&self, kind: u16, body: &impl Encode<()>, to: Option<&[[u8; 16]]>) {
107 if let Some(shared) = self.0.upgrade() {
108 match to {
109 None => {
110 let _ = shared.group(kind, body, None);
111 }
112 Some(to) => {
113 let slots: Vec<_> = shared
114 .state
115 .lock()
116 .unwrap()
117 .members
118 .iter()
119 .filter(|(_, m)| {
120 m.peer.as_ref().is_some_and(|p| to.contains(&p.hello.peer))
121 })
122 .map(|(slot, _)| *slot)
123 .collect();
124 for slot in slots {
125 let _ = shared.group(kind, body, Some(slot));
126 }
127 }
128 }
129 }
130 }
131}
132
133impl Live {
134 pub fn start(
135 me: Hello,
136 room: &Room,
137 _: Option<Reach>,
138 relay: Option<&str>,
139 events: impl Fn(Event) + Send + Sync + 'static,
140 ) -> io::Result<Self> {
141 let tag = room.tag().ok_or(io::ErrorKind::InvalidInput)?;
142 let relay = relay
143 .ok_or(io::ErrorKind::NotConnected)?
144 .trim_end_matches('/');
145 if !relay.starts_with("wss://") && !relay.starts_with("ws://") {
146 return Err(io::ErrorKind::InvalidInput.into());
147 }
148 let shared = Arc::new(Shared {
149 me,
150 room: room.clone(),
151 url: format!("{relay}/v1/room/{tag}"),
152 events: Box::new(events),
153 stopped: AtomicBool::new(false),
154 state: Mutex::new(State {
155 socket: 0,
156 slot: 0,
157 relayed: Relayed::Unknown,
158 failed: 0,
159 outdated: None,
160 burned: false,
161 streams: BTreeMap::new(),
162 members: BTreeMap::new(),
163 sealer: None,
164 presence: Presence::default(),
165 scheduled: false,
166 }),
167 });
168 shared.connect()?;
169 let weak = Arc::downgrade(&shared);
170 crate::task::spawn("live ping", move || async move {
171 let (_, wait) = crate::task::channel();
172 loop {
173 crate::task::wait(&wait, Some(Duration::from_secs(15))).await;
174 let Some(shared) = weak
175 .upgrade()
176 .filter(|s| !s.stopped.load(Ordering::Acquire))
177 else {
178 break;
179 };
180 let lines: Vec<_> = shared
181 .state
182 .lock()
183 .unwrap()
184 .streams
185 .values()
186 .map(|s| s.line.clone())
187 .collect();
188 for line in lines {
189 let _ = line.send(kind::PING, &());
190 }
191 if matches!(shared.room, Room::Notebook(_)) {
192 let _ = shared.group(kind::PING, &(), None);
193 }
194 }
195 })?;
196 Ok(Self(shared))
197 }
198
199 pub fn code(&self) -> Option<String> {
200 match &self.0.room {
201 Room::Code { code, .. } => Some(code.clone()),
202 _ => None,
203 }
204 }
205 pub fn relayed(&self) -> Relayed {
206 self.0.state.lock().unwrap().relayed.clone()
207 }
208 pub fn failed(&self) -> u32 {
209 self.0.state.lock().unwrap().failed
210 }
211 pub fn other_version(&self) -> Option<u16> {
212 self.0.state.lock().unwrap().outdated
213 }
214 pub fn burned(&self) -> bool {
215 self.0.state.lock().unwrap().burned
216 }
217 pub fn sender(&self) -> Sender {
218 Sender(Arc::downgrade(&self.0))
219 }
220 pub fn peers(&self) -> Vec<Peer> {
221 let state = self.0.state.lock().unwrap();
222 let mut peers = BTreeMap::new();
223 for peer in state
224 .members
225 .values()
226 .filter_map(|m| m.peer.as_ref())
227 .chain(state.streams.values().filter_map(|s| s.peer.as_ref()))
228 {
229 peers.insert(peer.hello.peer, peer.clone());
230 }
231 peers.into_values().collect()
232 }
233 pub fn line(&self, peer: &[u8; 16]) -> Option<Line> {
234 self.0
235 .state
236 .lock()
237 .unwrap()
238 .streams
239 .values()
240 .find(|s| s.peer.as_ref().is_some_and(|p| &p.hello.peer == peer))
241 .map(|s| s.line.clone())
242 }
243 pub fn set_presence(&self, presence: Presence) {
244 let mut state = self.0.state.lock().unwrap();
245 if state.presence == presence {
246 return;
247 }
248 state.presence = presence;
249 if std::mem::replace(&mut state.scheduled, true) {
250 return;
251 }
252 drop(state);
253 let weak = Arc::downgrade(&self.0);
254 let _ = crate::task::spawn("live presence", move || async move {
255 let (_, wait) = crate::task::channel();
256 crate::task::wait(&wait, Some(Duration::from_millis(100))).await;
257 if let Some(shared) = weak.upgrade() {
258 let presence = {
259 let mut state = shared.state.lock().unwrap();
260 state.scheduled = false;
261 state.presence.clone()
262 };
263 let _ = shared.group(kind::PRESENCE, &presence, None);
264 }
265 });
266 }
267 pub fn leave(self, reason: &str) {
268 let lines: Vec<_> = self
269 .0
270 .state
271 .lock()
272 .unwrap()
273 .streams
274 .values()
275 .map(|s| s.line.clone())
276 .collect();
277 for line in lines {
278 line.hang_up(reason);
279 }
280 }
281}
282
283impl Drop for Live {
284 fn drop(&mut self) {
285 self.0.stopped.store(true, Ordering::Release);
286 close(self.0.state.lock().unwrap().socket);
287 }
288}
289
290impl Shared {
291 fn connect(self: &Arc<Self>) -> io::Result<()> {
292 let weak = Arc::downgrade(self);
293 let message = Closure::<dyn FnMut(JsValue)>::new(move |data: JsValue| {
294 if let Some(shared) = weak.upgrade()
295 && let Err(error) = shared.heard(data)
296 {
297 if let Some(version) = error
298 .get_ref()
299 .and_then(|e| e.downcast_ref::<wire::Version>())
300 {
301 shared.state.lock().unwrap().outdated = Some(version.0);
302 } else if error.kind() == io::ErrorKind::InvalidData {
303 shared.state.lock().unwrap().failed += 1;
304 }
305 shared.disconnected();
306 }
307 })
308 .into_js_value();
309 let weak = Arc::downgrade(self);
310 let closed = Closure::<dyn FnMut()>::new(move || {
311 if let Some(shared) = weak.upgrade() {
312 shared.disconnected();
313 }
314 })
315 .into_js_value();
316 let socket = connect(&self.url, &message, &closed).map_err(js_error)?;
317 self.state.lock().unwrap().socket = socket;
318 Ok(())
319 }
320
321 fn disconnected(self: &Arc<Self>) {
322 let peers = {
323 let mut state = self.state.lock().unwrap();
324 close(state.socket);
325 let peers: Vec<_> = state
326 .streams
327 .values()
328 .filter_map(|s| s.peer.as_ref().map(|p| p.hello.clone()))
329 .collect();
330 state.streams.clear();
331 state.members.clear();
332 state.sealer = None;
333 state.relayed = Relayed::Unreachable(Trouble::Other);
334 peers
335 };
336 for hello in peers {
337 (self.events)(Event::Left(&hello));
338 }
339 (self.events)(Event::Changed);
340 let weak = Arc::downgrade(self);
341 let _ = crate::task::spawn("live reconnect", move || async move {
342 let (_, wait) = crate::task::channel();
343 crate::task::wait(&wait, Some(Duration::from_secs(2))).await;
344 if let Some(shared) = weak
345 .upgrade()
346 .filter(|s| !s.stopped.load(Ordering::Acquire))
347 {
348 let _ = shared.connect();
349 }
350 });
351 }
352
353 fn stream(&self, slot: u32, bytes: &[u8]) -> io::Result<()> {
354 let socket = self.state.lock().unwrap().socket;
355 for chunk in bytes.chunks(64 << 10) {
356 send(socket, &[&slot.to_be_bytes()[..], chunk].concat()).map_err(js_error)?;
357 }
358 Ok(())
359 }
360
361 fn group(&self, kind: u16, body: &impl Encode<()>, to: Option<u32>) -> io::Result<()> {
362 let body = minicbor::to_vec(body).map_err(io::Error::other)?;
363 let (socket, sealed) = {
364 let mut state = self.state.lock().unwrap();
365 let sealed = state
366 .sealer
367 .as_mut()
368 .ok_or(io::ErrorKind::NotConnected)?
369 .seal(kind, &body, to.is_none())?;
370 (state.socket, sealed)
371 };
372 let mut bytes = match to {
373 Some(slot) => [&(::relay::GROUP | 1).to_be_bytes()[..], &slot.to_be_bytes()].concat(),
374 None => ::relay::BROADCAST.to_be_bytes().to_vec(),
375 };
376 bytes.extend(sealed);
377 send(socket, &bytes).map_err(js_error)
378 }
379
380 fn meet(self: &Arc<Self>, slot: u32) -> io::Result<()> {
381 if self.state.lock().unwrap().streams.contains_key(&slot) {
382 return Ok(());
383 }
384 let tag = self.room.tag().ok_or(io::ErrorKind::InvalidInput)?;
385 let (opening, bytes) = wire::Opening::new(Side::Initiator, &tag, &self.room.secret())?;
386 let line = Line {
387 shared: Arc::downgrade(self),
388 slot,
389 sealer: Arc::new(Mutex::new(None)),
390 };
391 self.state.lock().unwrap().streams.insert(
392 slot,
393 Stream {
394 bytes: Vec::new(),
395 opening: Some(opening),
396 receive: None,
397 line,
398 peer: None,
399 },
400 );
401 self.stream(
402 slot,
403 &[&(bytes.len() as u32).to_be_bytes()[..], &bytes].concat(),
404 )
405 }
406
407 fn heard(self: &Arc<Self>, data: JsValue) -> io::Result<()> {
408 if let Some(text) = data.as_string() {
409 let notice = text
410 .parse::<::relay::Notice>()
411 .map_err(|_| io::ErrorKind::InvalidData)?;
412 match notice {
413 ::relay::Notice::Welcome { you, members } => {
414 {
415 let mut state = self.state.lock().unwrap();
416 state.slot = you;
417 state.relayed = Relayed::Joined;
418 if matches!(self.room, Room::Notebook(_)) {
419 state.sealer =
420 Some(group::Sealer::new(&group::Keys::new(&self.room.secret()))?);
421 }
422 }
423 if matches!(self.room, Room::Code { .. }) {
424 for slot in members.into_iter().filter(|slot| *slot < you) {
425 self.meet(slot)?;
426 }
427 } else {
428 self.group(kind::HELLO, &self.me, None)?;
429 let presence = self.state.lock().unwrap().presence.clone();
430 self.group(kind::PRESENCE, &presence, None)?;
431 }
432 }
433 ::relay::Notice::Left(slot) => {
434 let hello = {
435 let mut state = self.state.lock().unwrap();
436 state.members.remove(&slot);
437 state
438 .streams
439 .remove(&slot)
440 .and_then(|s| s.peer.map(|p| p.hello))
441 };
442 if let Some(hello) = hello {
443 (self.events)(Event::Left(&hello));
444 }
445 }
446 ::relay::Notice::Burned => self.state.lock().unwrap().burned = true,
447 _ => {}
448 }
449 (self.events)(Event::Changed);
450 return Ok(());
451 }
452 let bytes = js_sys::Uint8Array::new(&data).to_vec();
453 let (slot, bytes) = bytes
454 .split_first_chunk::<4>()
455 .ok_or(io::ErrorKind::InvalidData)?;
456 let slot = u32::from_be_bytes(*slot);
457 if slot & ::relay::GROUP != 0 {
458 return self.heard_group(slot & !::relay::GROUP, bytes);
459 }
460 self.heard_stream(slot, bytes)
461 }
462
463 fn heard_group(self: &Arc<Self>, slot: u32, frame: &[u8]) -> io::Result<()> {
464 let (kind, body, hello) = {
465 let mut state = self.state.lock().unwrap();
466 let member = match state.members.entry(slot) {
467 std::collections::btree_map::Entry::Occupied(e) => e.into_mut(),
468 std::collections::btree_map::Entry::Vacant(e) => e.insert(group::Member::new(
469 &group::Keys::new(&self.room.secret()),
470 frame,
471 )?),
472 };
473 let (kind, body) = member.open(frame)?;
474 (kind, body, member.hello().cloned())
475 };
476 match kind {
477 kind::HELLO | kind::HELLO_BACK => {
478 let hello: Hello =
479 minicbor::decode(&body).map_err(|_| io::ErrorKind::InvalidData)?;
480 if hello.peer == self.me.peer {
481 return Ok(());
482 }
483 let serves = hello.serves.is_some() && self.me.serves.is_none();
484 self.state
485 .lock()
486 .unwrap()
487 .members
488 .get_mut(&slot)
489 .unwrap()
490 .peer = Some(Peer {
491 hello: Arc::new(hello),
492 presence: None,
493 });
494 if kind == kind::HELLO {
495 self.group(kind::HELLO_BACK, &self.me, Some(slot))?;
496 let presence = self.state.lock().unwrap().presence.clone();
497 self.group(kind::PRESENCE, &presence, Some(slot))?;
498 }
499 if serves {
500 self.meet(slot)?;
501 }
502 (self.events)(Event::Changed);
503 }
504 kind::PRESENCE => {
505 let presence = minicbor::decode(&body).map_err(|_| io::ErrorKind::InvalidData)?;
506 if let Some(peer) = self
507 .state
508 .lock()
509 .unwrap()
510 .members
511 .get_mut(&slot)
512 .and_then(|m| m.peer.as_mut())
513 {
514 peer.presence = Some(presence);
515 }
516 (self.events)(Event::Changed);
517 }
518 kind if kind < 256 => {
519 if let Some(hello) = hello {
520 (self.events)(Event::Frame {
521 from: &hello,
522 kind,
523 body: &body,
524 });
525 }
526 }
527 _ => {}
528 }
529 Ok(())
530 }
531
532 fn heard_stream(self: &Arc<Self>, slot: u32, bytes: &[u8]) -> io::Result<()> {
533 {
534 let mut state = self.state.lock().unwrap();
535 let stream = state
536 .streams
537 .get_mut(&slot)
538 .ok_or(io::ErrorKind::InvalidData)?;
539 if stream.bytes.len() + bytes.len() > (16 << 20) + 4 {
540 return Err(io::ErrorKind::InvalidData.into());
541 }
542 stream.bytes.extend_from_slice(bytes);
543 }
544 loop {
545 let event = {
546 let mut state = self.state.lock().unwrap();
547 let stream = state.streams.get_mut(&slot).unwrap();
548 let Some(length) = stream.bytes.get(..4) else {
549 break;
550 };
551 let length = u32::from_be_bytes(length.try_into().unwrap()) as usize;
552 if length > 16 << 20 {
553 return Err(io::ErrorKind::InvalidData.into());
554 }
555 if stream.bytes.len() < length + 4 {
556 break;
557 }
558 let bytes: Vec<_> = stream.bytes.drain(..length + 4).collect();
559 if let Some(opening) = stream.opening.take() {
560 let (send, receive) = opening.finish(&bytes[4..])?;
561 *stream.line.sealer.lock().unwrap() = Some(send);
562 stream.receive = Some(receive);
563 Some((stream.line.clone(), None, kind::HELLO, Vec::new()))
564 } else {
565 let (kind, body) = stream
566 .receive
567 .as_mut()
568 .ok_or(io::ErrorKind::InvalidData)?
569 .receive(&mut io::Cursor::new(bytes))?;
570 if stream.peer.is_none() {
571 if kind != kind::HELLO {
572 return Err(io::ErrorKind::InvalidData.into());
573 }
574 let hello = minicbor::decode::<Hello>(&body)
575 .map_err(|_| io::ErrorKind::InvalidData)?;
576 stream.peer = Some(Peer {
577 hello: Arc::new(hello),
578 presence: None,
579 });
580 }
581 Some((
582 stream.line.clone(),
583 stream.peer.as_ref().map(|p| p.hello.clone()),
584 kind,
585 body,
586 ))
587 }
588 };
589 if let Some((line, hello, kind, body)) = event {
590 match hello {
591 None => line.send(kind::HELLO, &self.me)?,
592 Some(hello) if kind == kind::HELLO => (self.events)(Event::Met(&hello, &line)),
593 Some(hello) => (self.events)(Event::Frame {
594 from: &hello,
595 kind,
596 body: &body,
597 }),
598 }
599 }
600 }
601 Ok(())
602 }
603}
604
605fn js_error(error: JsValue) -> io::Error {
606 io::Error::new(
607 io::ErrorKind::NotConnected,
608 error
609 .as_string()
610 .unwrap_or_else(|| "Relay disconnected".into()),
611 )
612}
crates/notebook/src/live/wire.rs+60-36
......@@ -549,48 +549,72 @@ pub fn open(
549549 room: &str,
550550 secret: &[u8],
551551) -> io::Result<(Sealer, Sealer)> {
552 let password = Password::new(secret);
553 let [initiator, responder] = [b"initiator", b"responder"]
554 .map(|role| Identity::new(&[&role[..], room.as_bytes()].concat()));
555 let (pake, message) = match side {
556 Side::Initiator => Spake2::<Ed25519Group>::start_a(&password, &initiator, &responder),
557 Side::Responder => Spake2::<Ed25519Group>::start_b(&password, &initiator, &responder),
558 };
559 let ours = Open {
560 version: VERSION,
561 room: room.into(),
562 pake: message,
563 };
552 let (opening, ours) = Opening::new(side, room, secret)?;
564553 if side == Side::Initiator {
565 write_block(stream, &minicbor::to_vec(&ours).map_err(io::Error::other)?)?;
554 write_block(stream, &ours)?;
566555 }
567 let theirs: Open =
568 minicbor::decode(&read_block(stream)?).map_err(|_| invalid("A malformed opening"))?;
569 // A responder answers even a peer of another version, so that both ends can say which
570 // should update.
556 let theirs = read_block(stream)?;
571557 if side == Side::Responder {
572 write_block(stream, &minicbor::to_vec(&ours).map_err(io::Error::other)?)?;
558 write_block(stream, &ours)?;
573559 }
574 if theirs.version != VERSION {
575 return Err(io::Error::new(
576 io::ErrorKind::Unsupported,
577 Version(theirs.version),
578 ));
560 opening.finish(&theirs)
561}
562
563pub(super) struct Opening {
564 pake: Spake2<Ed25519Group>,
565 side: Side,
566 room: String,
567}
568
569impl Opening {
570 pub(super) fn new(side: Side, room: &str, secret: &[u8]) -> io::Result<(Self, Vec<u8>)> {
571 let password = Password::new(secret);
572 let [initiator, responder] = [b"initiator", b"responder"]
573 .map(|role| Identity::new(&[&role[..], room.as_bytes()].concat()));
574 let (pake, message) = match side {
575 Side::Initiator => Spake2::<Ed25519Group>::start_a(&password, &initiator, &responder),
576 Side::Responder => Spake2::<Ed25519Group>::start_b(&password, &initiator, &responder),
577 };
578 let ours = Open {
579 version: VERSION,
580 room: room.into(),
581 pake: message,
582 };
583 let message = minicbor::to_vec(&ours).map_err(io::Error::other)?;
584 Ok((
585 Self {
586 pake,
587 side,
588 room: room.to_owned(),
589 },
590 message,
591 ))
579592 }
580 if theirs.room != room {
581 return Err(invalid("The peer means another room"));
593
594 pub(super) fn finish(self, bytes: &[u8]) -> io::Result<(Sealer, Sealer)> {
595 let theirs: Open = minicbor::decode(bytes).map_err(|_| invalid("A malformed opening"))?;
596 if theirs.version != VERSION {
597 return Err(io::Error::new(
598 io::ErrorKind::Unsupported,
599 Version(theirs.version),
600 ));
601 }
602 if theirs.room != self.room {
603 return Err(invalid("The peer means another room"));
604 }
605 let key = self
606 .pake
607 .finish(&theirs.pake)
608 .map_err(|_| invalid("A malformed key exchange"))?;
609 let [from_initiator, from_responder] = [
610 Sealer::new(&key, b"Snowbound live v1 initiator"),
611 Sealer::new(&key, b"Snowbound live v1 responder"),
612 ];
613 Ok(match self.side {
614 Side::Initiator => (from_initiator, from_responder),
615 Side::Responder => (from_responder, from_initiator),
616 })
582617 }
583 let key = pake
584 .finish(&theirs.pake)
585 .map_err(|_| invalid("A malformed key exchange"))?;
586 let [from_initiator, from_responder] = [
587 Sealer::new(&key, b"Snowbound live v1 initiator"),
588 Sealer::new(&key, b"Snowbound live v1 responder"),
589 ];
590 Ok(match side {
591 Side::Initiator => (from_initiator, from_responder),
592 Side::Responder => (from_responder, from_initiator),
593 })
594618}
595619
596620fn write_block(to: &mut impl Write, bytes: &[u8]) -> io::Result<()> {
crates/notebook/src/session.rs+2-2
......@@ -1588,9 +1588,9 @@ impl Notebook {
15881588 let (folder, remote) = (root.clone(), remote.clone());
15891589 let bind = move |path: &str| remote(&folder.join(path));
15901590 let mut local = discover::Local::open(&root)?;
1591 let list = move |folder: &str| {
1591 let list = move |folder: String| {
15921592 use discover::Source;
1593 local.entries(folder, LIMITS.entries)
1593 std::future::ready(local.entries(&folder, LIMITS.entries))
15941594 };
15951595 Ok(((bind, list), watched))
15961596 },
crates/notebook/src/sync.rs+109-23
......@@ -14,6 +14,11 @@ use std::{
1414/// Errors retain publication state; confirmation checks the stamp, flushes, and notifies
1515/// cached readers.
1616pub trait Remote {
17 /// Completes a non-blocking operation that returned `WouldBlock`; retrying that
18 /// operation then takes its result. Blocking providers leave this absent.
19 fn pending(&mut self) -> Option<std::pin::Pin<Box<dyn Future<Output = ()> + '_>>> {
20 None
21 }
1722 fn read(&mut self) -> io::Result<Vec<u8>>;
1823 /// The file's stamp without reading its body or coordinating with writers. While it is
1924 /// the last observed image's, synchronization neither reads nor revalidates the file.
......@@ -54,6 +59,38 @@ pub trait Remote {
5459 }
5560}
5661
62pub(crate) trait Waiting {
63 fn waiting(&self) -> bool;
64}
65
66impl Waiting for io::Error {
67 fn waiting(&self) -> bool {
68 self.kind() == io::ErrorKind::WouldBlock
69 }
70}
71
72impl Waiting for CommitError {
73 fn waiting(&self) -> bool {
74 self.error.kind() == io::ErrorKind::WouldBlock
75 }
76}
77
78pub(crate) async fn awaited<R: Remote, T, E: Waiting>(
79 remote: &mut R,
80 mut operation: impl FnMut(&mut R) -> std::result::Result<T, E>,
81) -> std::result::Result<T, E> {
82 loop {
83 let result = operation(remote);
84 if result.as_ref().is_err_and(Waiting::waiting)
85 && let Some(pending) = remote.pending()
86 {
87 pending.await;
88 continue;
89 }
90 return result;
91 }
92}
93
5794#[derive(Debug, Clone, PartialEq, Eq)]
5895pub enum EditStatus {
5996 Pending,
......@@ -189,29 +226,53 @@ impl Replica {
189226 /// Versions the remote keeps beside the file merge into it first, each published as one
190227 /// more revision and then retired (`resolve.rs`).
191228 pub fn sync_once(&self, remote: &mut impl Remote) -> Result<Synced> {
229 crate::task::ready(self.sync_once_async(remote))?
230 }
231
232 /// The same guarded synchronization step, awaiting non-blocking remote operations.
233 pub async fn sync_once_async(&self, remote: &mut impl Remote) -> Result<Synced> {
234 let result = self.sync_once_inner(remote).await;
235 crate::fs::durable().await?;
236 result
237 }
238
239 /// Ownership uses `try_lock`; the cache mutex is released before every await.
240 #[allow(clippy::await_holding_lock)]
241 async fn sync_once_inner(&self, remote: &mut impl Remote) -> Result<Synced> {
192242 let _owner = self.sync_owner()?;
193 for version in remote.versions().map_err(Error::RemoteIo)? {
194 let image = remote.version(&version.id).map_err(Error::RemoteIo)?;
195 let current = remote.read().map_err(Error::RemoteIo)?;
243 for version in awaited(remote, Remote::versions)
244 .await
245 .map_err(Error::RemoteIo)?
246 {
247 let image = awaited(remote, |remote| remote.version(&version.id))
248 .await
249 .map_err(Error::RemoteIo)?;
250 let current = awaited(remote, Remote::read)
251 .await
252 .map_err(Error::RemoteIo)?;
196253 let device = version.device.as_deref().unwrap_or("Another device");
197254 let keep = match crate::resolve::merge(&current, &image, device) {
198255 Ok(Merged::Held) => false,
199256 Ok(Merged::Publish(transaction)) => {
200 remote.publish(&transaction)?;
257 awaited(remote, |remote| remote.publish(&transaction)).await?;
201258 false
202259 }
203260 // A version this cannot merge is kept whole rather than lost.
204261 Ok(Merged::Foreign) | Err(Error::Document(_) | Error::Rejected(_)) => true,
205262 Err(error) => return Err(error),
206263 };
207 remote.retire(&version.id, keep).map_err(Error::RemoteIo)?;
264 awaited(remote, |remote| remote.retire(&version.id, keep))
265 .await
266 .map_err(Error::RemoteIo)?;
208267 }
209268 let state = state(&*self.lock()?)?;
210269 let batched = self.section.key.is_none()
211270 && state.queued
212271 && state.blocked.is_none()
213272 && remote.accepts_edits();
214 let observed = remote.stamp().map_err(Error::RemoteIo)?;
273 let observed = awaited(remote, Remote::stamp)
274 .await
275 .map_err(Error::RemoteIo)?;
215276 if let Some(blocked) = &state.blocked
216277 && observed == *state.remote.as_ref().unwrap_or(&state.base)
217278 {
......@@ -229,7 +290,9 @@ impl Replica {
229290 let image = match observed {
230291 observed if observed == state.base || batched => None,
231292 _ => {
232 let image = remote.read().map_err(Error::RemoteIo)?;
293 let image = awaited(remote, Remote::read)
294 .await
295 .map_err(Error::RemoteIo)?;
233296 // Protected elsewhere: a section written anew, which only its key reads.
234297 if self.section.key.is_none() && crate::discover::locked(&Store::parse(&image)?) {
235298 return Err(Error::RemoteIo(io::Error::new(
......@@ -297,7 +360,8 @@ impl Replica {
297360 .collect(),
298361 )
299362 };
300 if let Err(error) = remote.confirm(&Stamp::of(&image)?) {
363 let stamp = Stamp::of(&image)?;
364 if let Err(error) = awaited(remote, |remote| remote.confirm(&stamp)).await {
301365 if error.state == CommitState::Committed {
302366 self.acknowledge(*batch, sealed.as_ref(), receipts.as_ref())?;
303367 }
......@@ -351,7 +415,7 @@ impl Replica {
351415 let Some(transaction) = transaction else {
352416 // Edits that changed nothing are published once the remote's image is durable.
353417 let base = base::base_stamp(&*self.lock()?)?;
354 if let Err(error) = remote.confirm(&base) {
418 if let Err(error) = awaited(remote, |remote| remote.confirm(&base)).await {
355419 if error.state == CommitState::Committed {
356420 self.acknowledge(batch, None, None)?;
357421 }
......@@ -364,18 +428,28 @@ impl Replica {
364428 changed,
365429 });
366430 };
431 if let Err(error) = crate::fs::durable().await {
432 self.lock()?
433 .execute("UPDATE batches SET attempted=0 WHERE id=?1", [batch])?;
434 return Err(error.into());
435 }
367436 let published = if self.section.key.is_none() && remote.accepts_edits() {
368 let connection = self.lock()?;
369 let edits = queue::load(&connection, None, Some(batch))?;
370 let revisions = decode_revisions(&connection.query_row(
371 "SELECT revisions FROM batches WHERE id=?1",
372 [batch],
373 |row| row.get::<_, String>(0),
374 )?)?;
375 drop(connection);
376 remote.publish_edits(&transaction, &edits, &revisions)
437 let (edits, revisions) = {
438 let connection = self.lock()?;
439 let edits = queue::load(&connection, None, Some(batch))?;
440 let revisions = decode_revisions(&connection.query_row(
441 "SELECT revisions FROM batches WHERE id=?1",
442 [batch],
443 |row| row.get::<_, String>(0),
444 )?)?;
445 (edits, revisions)
446 };
447 awaited(remote, |remote| {
448 remote.publish_edits(&transaction, &edits, &revisions)
449 })
450 .await
377451 } else {
378 remote.publish(&transaction)
452 awaited(remote, |remote| remote.publish(&transaction)).await
379453 };
380454 match published {
381455 Ok(()) => {}
......@@ -393,7 +467,13 @@ impl Replica {
393467 self.acknowledge(batch, Some(&transaction), None)?;
394468 let revision = self.receipt(id)?;
395469 if batched {
396 changed.extend(self.rebase(Some(remote.read().map_err(Error::RemoteIo)?))?);
470 changed.extend(
471 self.rebase(Some(
472 awaited(remote, Remote::read)
473 .await
474 .map_err(Error::RemoteIo)?,
475 ))?,
476 );
397477 changed.sort();
398478 changed.dedup();
399479 }
......@@ -550,15 +630,21 @@ impl Replica {
550630 /// Whether nothing is queued and the remote still has the base's stamp, or the queue
551631 /// is blocked on a remote that has not changed since; reads neither image and does not
552632 /// lock the cache during remote I/O.
553 pub(crate) fn settled(&self, remote: &mut impl Remote) -> Result<bool> {
633 pub(crate) async fn settled(&self, remote: &mut impl Remote) -> Result<bool> {
554634 let state = state(&*self.lock()?)?;
555635 let expected = match &state.blocked {
556636 Some(_) => state.remote.unwrap_or(state.base),
557637 None if !state.queued => state.base,
558638 None => return Ok(false),
559639 };
560 Ok(remote.stamp().map_err(Error::RemoteIo)? == expected
561 && remote.versions().map_err(Error::RemoteIo)?.is_empty())
640 Ok(awaited(remote, Remote::stamp)
641 .await
642 .map_err(Error::RemoteIo)?
643 == expected
644 && awaited(remote, Remote::versions)
645 .await
646 .map_err(Error::RemoteIo)?
647 .is_empty())
562648 }
563649}
564650
crates/notebook/src/task.rs+101-3
......@@ -1,7 +1,6 @@
11//! The sync worker's and background's loops: threads natively, and in the browser, which
2//! gives a module one thread, tasks on its event loop. A loop is an `async` block whose only
3//! waits are `wait`; natively each wait blocks its thread, so `complete` runs the block to the
4//! end in one poll.
2//! gives a module one thread, tasks on its event loop. Native waits block their thread, so
3//! `complete` runs the block to the end in one poll.
54
65use std::{io, time::Duration};
76
......@@ -11,6 +10,22 @@ pub(crate) use std::{
1110 thread::JoinHandle,
1211};
1312
13#[cfg(all(feature = "live", not(target_arch = "wasm32")))]
14pub(crate) type Answer<T> = std::sync::mpsc::Sender<T>;
15
16#[cfg(all(feature = "live", not(target_arch = "wasm32")))]
17pub(crate) fn response<T>() -> (Answer<T>, std::sync::mpsc::Receiver<T>) {
18 std::sync::mpsc::channel()
19}
20
21#[cfg(all(feature = "live", not(target_arch = "wasm32")))]
22pub(crate) async fn answered<T>(
23 receiver: &std::sync::mpsc::Receiver<T>,
24 timeout: Duration,
25) -> Result<T, std::sync::mpsc::RecvTimeoutError> {
26 receiver.recv_timeout(timeout)
27}
28
1429/// A wake that waits, at most one, as `mpsc::sync_channel(1)` keeps one.
1530#[cfg(not(target_arch = "wasm32"))]
1631pub(crate) fn channel() -> (Sender<()>, Receiver<()>) {
......@@ -50,6 +65,18 @@ pub(crate) fn complete<T>(work: impl Future<Output = T>) -> T {
5065 }
5166}
5267
68/// Polls a synchronous entry point once; browser I/O must use the async entry point.
69pub(crate) fn ready<T>(work: impl Future<Output = T>) -> io::Result<T> {
70 let mut work = std::pin::pin!(work);
71 match work
72 .as_mut()
73 .poll(&mut std::task::Context::from_waker(std::task::Waker::noop()))
74 {
75 std::task::Poll::Ready(output) => Ok(output),
76 std::task::Poll::Pending => Err(io::ErrorKind::WouldBlock.into()),
77 }
78}
79
5380#[cfg(target_arch = "wasm32")]
5481pub(crate) use web::*;
5582
......@@ -62,6 +89,77 @@ mod web {
6289 };
6390 use wasm_bindgen::{JsCast, prelude::*};
6491
92 struct Response<T> {
93 value: Option<T>,
94 closed: bool,
95 waiting: Option<Waker>,
96 }
97
98 pub(crate) struct Answer<T>(Arc<Mutex<Response<T>>>);
99 pub(crate) struct Answered<T>(Arc<Mutex<Response<T>>>);
100
101 pub(crate) fn response<T>() -> (Answer<T>, Answered<T>) {
102 let response = Arc::new(Mutex::new(Response {
103 value: None,
104 closed: false,
105 waiting: None,
106 }));
107 (Answer(response.clone()), Answered(response))
108 }
109
110 impl<T> Answer<T> {
111 pub(crate) fn send(self, value: T) -> Result<(), ()> {
112 let wake = {
113 let mut response = self.0.lock().map_err(|_| ())?;
114 response.value = Some(value);
115 response.waiting.take()
116 };
117 if let Some(wake) = wake {
118 wake.wake();
119 }
120 Ok(())
121 }
122 }
123
124 impl<T> Drop for Answer<T> {
125 fn drop(&mut self) {
126 let wake = {
127 let mut response = self.0.lock().unwrap();
128 response.closed = true;
129 response.waiting.take()
130 };
131 if let Some(wake) = wake {
132 wake.wake();
133 }
134 }
135 }
136
137 pub(crate) async fn answered<T>(
138 receiver: &Answered<T>,
139 timeout: Duration,
140 ) -> Result<T, std::sync::mpsc::RecvTimeoutError> {
141 let deadline = web_time::Instant::now() + timeout;
142 let mut armed = false;
143 std::future::poll_fn(|context| {
144 let mut response = receiver.0.lock().unwrap();
145 if let Some(value) = response.value.take() {
146 return Poll::Ready(Ok(value));
147 }
148 if response.closed {
149 return Poll::Ready(Err(std::sync::mpsc::RecvTimeoutError::Disconnected));
150 }
151 if web_time::Instant::now() >= deadline {
152 return Poll::Ready(Err(std::sync::mpsc::RecvTimeoutError::Timeout));
153 }
154 response.waiting = Some(context.waker().clone());
155 if !std::mem::replace(&mut armed, true) {
156 after(timeout, context.waker().clone());
157 }
158 Poll::Pending
159 })
160 .await
161 }
162
65163 #[derive(Default)]
66164 struct Bell {
67165 rung: bool,
crates/notebook/src/worker.rs+2-2
......@@ -214,13 +214,13 @@ impl Replica {
214214 }
215215 let result = match remote.as_mut() {
216216 Some(remote) => {
217 if reported && replica.settled(remote).unwrap_or(false) {
217 if reported && replica.settled(remote).await.unwrap_or(false) {
218218 worker_signal.requested.store(false, Ordering::Release);
219219 worker_signal.synced.store(crate::now(), Ordering::Release);
220220 rest().await;
221221 continue;
222222 }
223 let result = replica.sync_once(remote);
223 let result = replica.sync_once_async(remote).await;
224224 reported = result.is_ok();
225225 result
226226 }
crates/notebook/tests/live_membership.rs+2-2
......@@ -119,7 +119,7 @@ fn removal_retires_one_credential_and_other_devices_reconnect_after_a_restart()
119119 .secret;
120120 host.remove(&bob_key).unwrap();
121121 until("Bob was not removed", || {
122 bob.stopped() && bob.host().is_none()
122 bob.ended() == Some(share::Ended::Removed) && bob.host().is_none()
123123 });
124124 until("presence did not move to the new room", || {
125125 host.guests().len() == 1
......@@ -129,7 +129,7 @@ fn removal_retires_one_credential_and_other_devices_reconnect_after_a_restart()
129129 assert_eq!(current.members.len(), 1);
130130 assert_ne!(original_code, code(&host));
131131 assert!(alice.host().is_some());
132 assert!(!alice.stopped());
132 assert_eq!(alice.ended(), None);
133133 let live = Notebook::open_hosted(Arc::clone(&alice), directory.path().join("alice")).unwrap();
134134 assert_eq!(live.catalog().sections.len(), 2);
135135 let restarted: Sharing =
crates/notebook/tests/live_share.rs+1-1
......@@ -398,7 +398,7 @@ fn stopping_lets_every_guest_go_and_retires_the_code() {
398398 let (guest, _notebook) = guest("Grace", &code, &url, &directory.path().join("grace"));
399399 host.stop();
400400 until("the guest never heard the host stop", || {
401 guest.stopped() && guest.host().is_none()
401 guest.ended() == Some(share::Ended::Stopped) && guest.host().is_none()
402402 });
403403 assert!(host.code().is_none() && host.guests().is_empty());
404404 assert!(matches!(
crates/notebook/tests/sync.rs+100
......@@ -20,6 +20,106 @@ fn save(cache: &Replica, text: ExGuid, range: Range<u32>, replacement: &str) ->
2020 .unwrap()
2121}
2222
23#[test]
24fn asynchronous_io_keeps_the_attempt_and_recovers_a_lost_reply() {
25 use std::{
26 future::Future,
27 pin::Pin,
28 task::{Context, Poll, Waker},
29 };
30
31 struct Deferred {
32 server: Server,
33 ready: bool,
34 }
35 impl Remote for Deferred {
36 fn pending(&mut self) -> Option<Pin<Box<dyn Future<Output = ()> + '_>>> {
37 let mut yielded = false;
38 Some(Box::pin(std::future::poll_fn(move |context| {
39 if yielded {
40 self.ready = true;
41 Poll::Ready(())
42 } else {
43 yielded = true;
44 context.waker().wake_by_ref();
45 Poll::Pending
46 }
47 })))
48 }
49 fn read(&mut self) -> io::Result<Vec<u8>> {
50 if !std::mem::take(&mut self.ready) {
51 return Err(io::ErrorKind::WouldBlock.into());
52 }
53 self.server.read()
54 }
55 fn stamp(&mut self) -> io::Result<onestore::Stamp> {
56 if !std::mem::take(&mut self.ready) {
57 return Err(io::ErrorKind::WouldBlock.into());
58 }
59 self.server.stamp()
60 }
61 fn publish(&mut self, transaction: &Transaction) -> Result<(), CommitError> {
62 if !std::mem::take(&mut self.ready) {
63 return Err(CommitError {
64 state: CommitState::NotCommitted,
65 error: io::ErrorKind::WouldBlock.into(),
66 });
67 }
68 self.server.publish(transaction)
69 }
70 fn confirm(&mut self, base: &onestore::Stamp) -> Result<(), CommitError> {
71 if !std::mem::take(&mut self.ready) {
72 return Err(CommitError {
73 state: CommitState::NotCommitted,
74 error: io::ErrorKind::WouldBlock.into(),
75 });
76 }
77 self.server.confirm(base)
78 }
79 }
80 fn run<T>(work: impl Future<Output = T>) -> T {
81 let mut work = std::pin::pin!(work);
82 let mut polls = 0;
83 loop {
84 polls += 1;
85 assert!(polls < 100, "An async operation stopped making progress");
86 if let Poll::Ready(result) = work.as_mut().poll(&mut Context::from_waker(Waker::noop()))
87 {
88 assert!(polls > 1);
89 return result;
90 }
91 }
92 }
93 let directory = tempfile::tempdir().unwrap();
94 let path = directory.path().join("cache.sqlite");
95 let source = onestore::create_section("async.one", "Original", "Fixture").unwrap();
96 let (_, object, _) = text(&source);
97 let cache = Replica::create(&path, &source).unwrap();
98 let id = save(&cache, object, 0..0, "Browser ").unwrap();
99 let mut remote = Deferred {
100 server: Server::new(&source),
101 ready: false,
102 };
103 remote.server.fault = Fault::UnknownAfter;
104 assert!(matches!(
105 run(cache.sync_once_async(&mut remote)),
106 Err(Error::Remote(CommitError {
107 state: CommitState::Unknown,
108 ..
109 }))
110 ));
111 assert_eq!(remote.server.publications, 1);
112 drop(cache);
113 let cache = Replica::open(&path).unwrap();
114 assert!(
115 matches!(run(cache.sync_once_async(&mut remote)).unwrap().edit, Some((published, EditStatus::Published { .. })) if published == id)
116 );
117 assert_eq!(remote.server.publications, 1);
118 assert_eq!(text(&remote.server.visible).2, "Browser Original");
119 assert_eq!(remote.server.visible, remote.server.durable);
120 assert!(cache.pending().unwrap().is_empty());
121}
122
23123#[test]
24124fn an_unchanged_stamp_publishes_and_settles_without_reading_the_remote() {
25125 struct Counted {
crates/snowbound/src/commands.rs+1
......@@ -1163,6 +1163,7 @@ impl State {
11631163 #[cfg(feature = "live")]
11641164 Id::LiveShare => enabled(
11651165 !modal
1166 && cfg!(not(target_arch = "wasm32"))
11661167 && self.notebook().is_some_and(|library| {
11671168 library.catalog().is_some() && library.joined.is_none()
11681169 }),
crates/snowbound/src/live.rs+25-12
......@@ -147,7 +147,7 @@ fn keep(file: &Path, value: &impl serde::Serialize) -> io::Result<()> {
147147 let folder = file.parent().unwrap_or(Path::new("."));
148148 notebook::fs::create_dir_all(folder)?;
149149 let partial = file.with_extension("partial");
150 let mut options = std::fs::OpenOptions::new();
150 let mut options = notebook::fs::OpenOptions::new();
151151 options.write(true).create(true).truncate(true);
152152 #[cfg(unix)]
153153 std::os::unix::fs::OpenOptionsExt::mode(&mut options, 0o600);
......@@ -217,8 +217,9 @@ impl Joined {
217217 )
218218 })
219219 .map_err(|error| refused(&error))?;
220 #[cfg(not(target_arch = "wasm32"))]
220221 if !share.listed {
221 let since = std::time::Instant::now();
222 let since = web_time::Instant::now();
222223 while guest.host().is_none() && since.elapsed() < FIRST_LISTING {
223224 std::thread::sleep(std::time::Duration::from_millis(50));
224225 }
......@@ -257,17 +258,17 @@ impl Joined {
257258}
258259
259260/// Meets whoever shares `code` (and `password`): the share it welcomes this computer to.
260pub(crate) fn join(
261pub(crate) async fn join(
261262 code: &str,
262263 password: &str,
263264 waiting: impl Fn(bool) -> bool,
264265) -> Result<live::wire::Welcome, live::share::Refusal> {
265266 let me = hello().map_err(|_| live::share::Refusal::Unreachable(live::Trouble::Other))?;
266 live::share::join_while(me, code, password, reach(), relay().as_deref(), waiting)
267 live::share::join_while_async(me, code, password, reach(), relay().as_deref(), waiting).await
267268}
268269
269270/// Keeps `welcome` as the notebook this computer joined: its location.
270pub(crate) fn joined(cache: &Path, welcome: live::wire::Welcome) -> io::Result<String> {
271pub(crate) async fn joined(cache: &Path, welcome: live::wire::Welcome) -> io::Result<String> {
271272 let location = live::share::location(&welcome.share);
272273 let file = kept(cache, JOINED);
273274 let mut shares: BTreeMap<String, Share> = read_kept(&file);
......@@ -276,20 +277,32 @@ pub(crate) fn joined(cache: &Path, welcome: live::wire::Welcome) -> io::Result<S
276277 Share {
277278 share: welcome.share,
278279 secret: welcome.secret,
279 notebook: welcome.notebook,
280 host: welcome.host,
280 notebook: welcome.notebook.clone(),
281 host: welcome.host.clone(),
281282 listed: false,
282283 protocol: live::wire::VERSION,
283284 },
284285 );
285286 keep(&file, &shares)?;
287 #[cfg(target_arch = "wasm32")]
288 {
289 let guest = Guest::start(
290 hello()?,
291 welcome.share,
292 welcome.secret,
293 None,
294 relay().as_deref(),
295 crate::library::notify_background,
296 )?;
297 live::share::Hosted::prepare(guest, cache).await?;
298 }
286299 Ok(location)
287300}
288301
289302/// A peer to go to, since when, and the section and page already asked to open.
290303type Going = (
291304 [u8; 16],
292 std::time::Instant,
305 web_time::Instant,
293306 Option<(Option<[u8; 16]>, Option<live::Guid>)>,
294307);
295308
......@@ -319,9 +332,9 @@ pub(crate) struct Peers {
319332 /// Each peer's picture as drawn, cut to a circle, decoded once.
320333 pictures: HashMap<[u8; 16], Option<draw::RasterImage>>,
321334 /// Where each peer's caret was last seen, and when it got there.
322 moved: HashMap<[u8; 16], (Caret, std::time::Instant)>,
335 moved: HashMap<[u8; 16], (Caret, web_time::Instant)>,
323336 /// Where each peer was last seen, and when they last moved, for ordering the avatars.
324 active: HashMap<[u8; 16], (Option<live::Presence>, std::time::Instant)>,
337 active: HashMap<[u8; 16], (Option<live::Presence>, web_time::Instant)>,
325338 /// The peer an avatar's click goes to, since when, and the place already asked to open.
326339 going: Option<Going>,
327340 pub(crate) share: Option<crate::share::ShareDialog>,
......@@ -650,7 +663,7 @@ impl State {
650663 if peers.is_empty() {
651664 return;
652665 }
653 let now = std::time::Instant::now();
666 let now = web_time::Instant::now();
654667 for peer in &peers {
655668 let id = peer.hello.peer;
656669 let seen = self
......@@ -921,7 +934,7 @@ impl State {
921934 (y * viewport.scale + viewport.origin[1]) / scale,
922935 ]
923936 };
924 let now = std::time::Instant::now();
937 let now = web_time::Instant::now();
925938 let (mut placed, mut above, mut below) = (Vec::new(), Vec::new(), Vec::new());
926939 for peer in self.connected() {
927940 let Some(caret) = (peer.presence.as_ref())
crates/snowbound/src/main.rs+33
......@@ -4247,6 +4247,23 @@ impl State {
42474247 .filter(|(_, paths)| !paths.is_empty())
42484248 .collect();
42494249 for (library, paths) in &reported {
4250 #[cfg(all(feature = "live", target_arch = "wasm32"))]
4251 if library.joined.is_some() && paths.iter().any(|path| !path.ends_with(".one")) {
4252 let refreshed = Arc::new(library.with(library.reopen()?));
4253 let shown = self
4254 .session
4255 .as_ref()
4256 .filter(|session| session.library.location == library.location)
4257 .and_then(|session| {
4258 let path = &session.tabs[session.tab].path;
4259 if refreshed.contains(path) {
4260 Some(path.clone())
4261 } else {
4262 refreshed.first_section()
4263 }
4264 });
4265 self.adopt(refreshed, shown.as_deref())?;
4266 }
42504267 self.sections_changed(library, paths.clone());
42514268 }
42524269 let changed: Vec<String> = (reported.iter())
......@@ -6184,6 +6201,22 @@ fn spawn(work: impl FnOnce() + Send + 'static) {
61846201 platform::defer(work);
61856202}
61866203
6204/// Runs asynchronous work beside the frame, creating the future on its executor.
6205#[cfg(feature = "live")]
6206fn spawn_async<F: Future<Output = ()> + 'static>(work: impl FnOnce() -> F + Send + 'static) {
6207 RUNNING.fetch_add(1, Ordering::Relaxed);
6208 #[cfg(not(target_arch = "wasm32"))]
6209 std::thread::spawn(move || {
6210 pollster::block_on(work());
6211 RUNNING.fetch_sub(1, Ordering::Release);
6212 });
6213 #[cfg(target_arch = "wasm32")]
6214 wasm_bindgen_futures::spawn_local(async move {
6215 work().await;
6216 RUNNING.fetch_sub(1, Ordering::Release);
6217 });
6218}
6219
61876220/// FILETIME now: when an edit happened, which its modification times record.
61886221fn filetime() -> u64 {
61896222 let unix = web_time::SystemTime::now()
crates/snowbound/src/menus.rs+1-1
......@@ -601,7 +601,7 @@ impl State {
601601 }
602602 actions.extend([
603603 item(Action::CopyLink, "Copy Link to Notebook", false, true),
604 #[cfg(feature = "live")]
604 #[cfg(all(feature = "live", not(target_arch = "wasm32")))]
605605 item(
606606 Action::LiveShare,
607607 commands::command(commands::Id::LiveShare).title,
crates/snowbound/src/share.rs+21-17
......@@ -6,7 +6,6 @@ use accesskit::Role;
66use notebook::live::{
77 Relayed, Trouble,
88 share::{self, Refusal, Sharing},
9 wire::Welcome,
109};
1110use std::sync::{
1211 Arc,
......@@ -59,7 +58,7 @@ pub(crate) struct JoinDialog {
5958 code: String,
6059 password: String,
6160 status: Status,
62 replies: crate::live::Channel<Result<Welcome, Refusal>>,
61 replies: crate::live::Channel<Result<String, String>>,
6362 alive: Arc<AtomicBool>,
6463 approval: Arc<AtomicBool>,
6564}
......@@ -607,21 +606,14 @@ impl State {
607606 let mut welcomed = None;
608607 for reply in dialog.replies.1.try_iter() {
609608 match reply {
610 Ok(welcome) => welcomed = Some(welcome),
611 Err(refused) => dialog.status = Status::Failed(refusal(&refused)),
609 Ok(location) => welcomed = Some(location),
610 Err(message) => dialog.status = Status::Failed(message),
612611 }
613612 }
614 if let Some(welcome) = welcomed {
615 match crate::live::joined(&self.cache, welcome) {
616 Ok(location) => {
617 self.ui.close_popup(join_id());
618 self.open_notebook(location, None);
619 return;
620 }
621 Err(error) => {
622 dialog.status = Status::Failed(format!("Couldn’t keep the code: {error}"));
623 }
624 }
613 if let Some(location) = welcomed {
614 self.ui.close_popup(join_id());
615 self.open_notebook(location, None);
616 return;
625617 }
626618 let ui = &mut self.ui;
627619 let owners = [join_id(), code_field(), join_password()];
......@@ -679,15 +671,27 @@ impl State {
679671 dialog.status = Status::Waiting;
680672 let (replies, redraw) = (dialog.replies.0.clone(), self.redraw.clone());
681673 let password = dialog.password.clone();
674 let cache = self.cache.clone();
682675 let (alive, approval) =
683676 (Arc::clone(&dialog.alive), Arc::clone(&dialog.approval));
684 crate::spawn(move || {
677 crate::spawn_async(move || async move {
685678 let reply = crate::live::join(&code, &password, |waiting| {
686679 if approval.swap(waiting, Ordering::AcqRel) != waiting {
687680 redraw.wake_by_ref();
688681 }
689682 alive.load(Ordering::Acquire)
690 });
683 })
684 .await;
685 approval.store(false, Ordering::Release);
686 redraw.wake_by_ref();
687 let reply = match reply {
688 Ok(welcome) => {
689 crate::live::joined(&cache, welcome).await.map_err(|error| {
690 format!("Couldn’t open the shared notebook: {error}")
691 })
692 }
693 Err(error) => Err(refusal(&error)),
694 };
691695 let _ = replies.send(reply);
692696 redraw.wake();
693697 });
crates/snowbound/src/sync.rs+15-9
......@@ -265,8 +265,7 @@ struct Facts<'a> {
265265 changes_listed: bool,
266266 /// Whether the open section's file can be shown in the file manager.
267267 local: bool,
268 /// The computer sharing this notebook by Live Share stopped sharing it.
269 stopped: bool,
268 ended: Option<&'a str>,
270269}
271270
272271/// What the reader picked in the popup.
......@@ -541,9 +540,9 @@ fn build(ui: &mut Ui, facts: &Facts) -> Picked {
541540 let conflict;
542541 // A problem's advice stays while working offline, so turning it on moves nothing.
543542 let advice = match state {
544 SyncState::NotConnected if facts.stopped => Some(format!(
545 "The person sharing this notebook stopped sharing it. Your changes stay on {THIS}."
546 )),
543 SyncState::NotConnected if facts.ended.is_some() => facts
544 .ended
545 .map(|ended| format!("{ended} Your changes stay on {THIS}.")),
547546 SyncState::NotConnected => Some(match sync.queued {
548547 0 => format!("Can’t reach {host}. Sync continues when it’s back."),
549548 _ => format!("Can’t reach {host}. {waiting} will sync when it’s back."),
......@@ -1089,12 +1088,19 @@ impl State {
10891088 location: &library.location,
10901089 place: library.place(),
10911090 #[cfg(feature = "live")]
1092 stopped: library
1091 ended: library
10931092 .joined
10941093 .as_ref()
1095 .is_some_and(|joined| joined.guest.stopped()),
1094 .and_then(|joined| match joined.guest.ended()? {
1095 notebook::live::share::Ended::Stopped => {
1096 Some("The person sharing this notebook stopped sharing it.")
1097 }
1098 notebook::live::share::Ended::Removed => Some(
1099 "This device was removed. Ask the person sharing for a new link or code.",
1100 ),
1101 }),
10961102 #[cfg(not(feature = "live"))]
1097 stopped: false,
1103 ended: None,
10981104 notice: library.notice.as_deref(),
10991105 sections: sections(&library, session),
11001106 conflicts: session
......@@ -1212,7 +1218,7 @@ mod tests {
12121218 update: &IDLE,
12131219 changes_listed: false,
12141220 local: true,
1215 stopped: false,
1221 ended: None,
12161222 }
12171223 }
12181224
crates/snowbound/src/web.rs+12
......@@ -1185,6 +1185,18 @@ async fn open(
11851185 None => state.create_notebook(own.clone())?,
11861186 }
11871187 }
1188 #[cfg(feature = "live")]
1189 if let Some(code) = web_sys::window()
1190 .and_then(|window| window.location().search().ok())
1191 .and_then(|query| {
1192 query
1193 .trim_start_matches('?')
1194 .split('&')
1195 .find_map(|pair| pair.strip_prefix("join=").map(str::to_owned))
1196 })
1197 {
1198 state.join_link(code);
1199 }
11881200 Ok(state)
11891201}
11901202
crates/snowbound/web/glue.js+50-21
......@@ -23,6 +23,7 @@ function storageWorker() {
2323 const handles = new Map();
2424 let root;
2525 let queue = Promise.resolve();
26 const pending = [];
2627
2728 const names = (path) => path.split("/").filter(Boolean);
2829 const folder = async (parts) => {
......@@ -47,7 +48,14 @@ function storageWorker() {
4748 };
4849 const write = async (path, length, ranges) => {
4950 const open = await handle(path);
50 for (const [offset, bytes] of ranges) open.write(bytes, { at: offset });
51 for (const [offset, bytes] of ranges) {
52 let written = 0;
53 while (written < bytes.length) {
54 const count = open.write(bytes.subarray(written), { at: offset + written });
55 if (!count) throw new Error("Browser storage could not finish writing");
56 written += count;
57 }
58 }
5159 open.truncate(length);
5260 open.flush();
5361 };
......@@ -103,25 +111,34 @@ function storageWorker() {
103111 });
104112 };
105113
114 const flush = async () => {
115 while (pending.length) {
116 for (const [path, length, ranges] of pending[0])
117 if (length === undefined) await remove(path);
118 else if (length === null) await folder(names(path));
119 else await write(path, length, ranges);
120 pending.shift();
121 }
122 };
123
106124 onmessage = ({ data }) => {
107 queue = queue
108 .then(async () => {
125 queue = queue.then(async () => {
126 try {
109127 if (data.kind === "load") {
110128 root = await navigator.storage.getDirectory();
111129 let files = await list(root, "", []);
112 if (!files.length) {
113 await migrate();
114 files = await list(root, "", []);
115 }
116 postMessage(files, files.flatMap(([, bytes]) => (bytes ? [bytes.buffer] : [])));
117 } else if (data.kind === "settle") postMessage(null);
118 else
119 for (const [path, length, ranges] of data.changes)
120 if (length === undefined) await remove(path);
121 else if (length === null) await folder(names(path));
122 else await write(path, length, ranges);
123 })
124 .catch((error) => console.error("Keeping files", error));
130 if (!files.length) { await migrate(); files = await list(root, "", []); }
131 postMessage(files, files.flatMap(([, bytes]) => bytes ? [bytes.buffer] : []));
132 } else {
133 if (data.kind === "store") pending.push(data.changes);
134 await flush();
135 if (data.kind === "settle") postMessage({ error: null });
136 }
137 } catch (error) {
138 if (data.kind === "settle" || data.kind === "load") postMessage({ error: String(error) });
139 else console.error("Keeping files", error);
140 }
141 });
125142 };
126143}
127144
......@@ -132,15 +149,25 @@ export function loadFiles() {
132149 storage.postMessage({ kind: "load" });
133150 return new Promise((resolve, reject) => {
134151 storage.onmessage = ({ data }) => {
135 storage.onmessage = () => settling.shift()?.();
136 resolve(data);
152 storage.onmessage = ({ data }) => {
153 const waiting = settling.shift();
154 if (data.error) waiting?.reject(new Error(data.error));
155 else waiting?.resolve();
156 };
157 if (data.error) reject(new Error(data.error));
158 else resolve(data);
159 };
160 storage.onerror = (event) => {
161 storageError = new Error(`The storage worker failed: ${event.message}`);
162 reject(storageError);
163 for (const waiting of settling.splice(0)) waiting.reject(storageError);
137164 };
138 storage.onerror = (error) => reject(new Error(`The storage worker failed: ${error.message}`));
139165 });
140166}
141167
142168// Callers waiting for the storage worker to finish what it was given.
143169const settling = [];
170let storageError;
144171
145172/** Writes `[path]` (removed), `[path, null]` (a folder) and `[path, length, ranges]` entries;
146173 * those under a folder of the user's go there, a committed section only where nothing else
......@@ -322,8 +349,9 @@ export function fetchDictionary(name) {
322349
323350/** Resolves once every file handed over so far is written. */
324351function settled() {
325 const kept = new Promise((resolve) => {
326 settling.push(resolve);
352 if (storageError) return Promise.reject(storageError);
353 const kept = new Promise((resolve, reject) => {
354 settling.push({ resolve, reject });
327355 storage.postMessage({ kind: "settle" });
328356 });
329357 return Promise.all([kept, writing]);
......@@ -597,6 +625,7 @@ export function adoptCanvas(fresh) {
597625/** Wires the page's canvas, text area and file input to `module`'s exports. */
598626export function attach(module) {
599627 wasm = module;
628 globalThis.snowboundFlushStorage = () => { wasm.flush(); return settled(); };
600629 canvas = document.getElementById("page");
601630 input = document.getElementById("input");
602631 picker = document.getElementById("files");
crates/snowbound/web/index.html+5
......@@ -205,6 +205,11 @@
205205 ]);
206206 status.textContent = "Starting Snowbound";
207207 await snowbound.start(snowbound, fonts, dictionaries);
208 const opened = new URL(location.href);
209 if (opened.searchParams.has("join")) {
210 opened.searchParams.delete("join");
211 history.replaceState(null, "", opened);
212 }
208213 sessionStorage.removeItem("snowbound.refetched");
209214 } catch (error) {
210215 // A module and JavaScript from different builds can't link: an old cached copy. Once,
tools/ci.py+1-1
......@@ -80,7 +80,7 @@ def lanes():
8080 result.append(dict(
8181 name='web', minutes=20, packages=['snowbound'], paths=('crates/snowbound/web/', 'tools/release_web.py'),
8282 # Clippy builds it; the module itself is linked by release_web.py.
83 commands=[cargo('clippy', '-p', 'snowbound', '--target', 'wasm32-unknown-unknown', '--no-default-features', '--features', 'wgpu', '--', '-D', 'warnings')],
83 commands=[cargo('clippy', '-p', 'snowbound', '--target', 'wasm32-unknown-unknown', '--no-default-features', '--features', 'wgpu,live', '--', '-D', 'warnings')],
8484 environment={**(wasm or {}), 'CARGO_TARGET_DIR': str(TARGET / 'wasm')},
8585 missing='needs the wasm32-unknown-unknown target' if 'wasm32-unknown-unknown' not in targets.split()
8686 else None if wasm else 'needs nix for a clang that builds for wasm32'))
tools/release_web.py+1-1
......@@ -100,7 +100,7 @@ def build(out):
100100 """Writes index.html, the module, its JavaScript and the fonts to `out`, the module knowing
101101 when it was built."""
102102 built = str(int(time.time()))
103 run(['cargo', 'build', '--locked', '-p', 'snowbound', '--release', '--target', 'wasm32-unknown-unknown', '--no-default-features', '--features', 'wgpu',
103 run(['cargo', 'build', '--locked', '-p', 'snowbound', '--release', '--target', 'wasm32-unknown-unknown', '--no-default-features', '--features', 'wgpu,live',
104104 *PROFILE], env={**environment(), 'SNOWBOUND_WEB_BUILD': built})
105105 if out.exists():
106106 shutil.rmtree(out)