diff --git a/arc/platforms.md b/arc/platforms.md index 0afb491628c40898af1d56e5cb61ced2b7702c9c..b47241d9db47a74539e9c9111d0cdec50441e018 100644 --- a/arc/platforms.md +++ b/arc/platforms.md @@ -185,6 +185,10 @@ keyboard, the toolbar and the macOS menu bar all run commands from it. `/Cache`. A storage worker keeps them in the origin's private file system, writing the byte ranges each burst changed through OPFS's synchronous handles, which only workers get. One tab at a time holds them. +- Live Share joins desktop hosts through sealed WebSockets. Approval, device removal, + presence and edits use the desktop protocol; replicas await remote operations and an + ordered OPFS flush before publishing or returning a durable receipt. Hosting and section + organization stay on desktop. - Open Notebook, where the browser has the File System Access API (Chromium), opens a folder of the user's: mirrored under `/Folders`, its handle kept in IndexedDB, other apps' writes read every few seconds. A browser takes no locks, so a commit there stands only once the diff --git a/crates/notebook/Cargo.toml b/crates/notebook/Cargo.toml index 89071756ff0d4478dcb0700aca8cd864cd152d80..c6528b3e83a9102b1830ce86347ec4e289d63c3b 100644 --- a/crates/notebook/Cargo.toml +++ b/crates/notebook/Cargo.toml @@ -26,19 +26,21 @@ aes-gcm = "0.11.1" base64 = { version = "0.23.1", default-features = false, features = ["std"] } getrandom = "0.4.3" hmac = "0.13.0" -mdns-sd = { version = "0.21.4", default-features = false, optional = true } minicbor = { version = "2.3.0", features = ["derive", "alloc"], optional = true } spake2 = { version = "=0.5.0-pre.0", features = ["getrandom"], optional = true } relay = { path = "../relay", optional = true } -# A relay's TLS, trusting what the app's updates trust: the system's authorities, then Mozilla's. -rustls = { version = "0.23", default-features = false, features = ["ring", "std", "tls12"], optional = true } -rustls-native-certs = { version = "0.8", optional = true } -webpki-root-certs = { version = "1.0", optional = true } # OneNote packages (.onepkg) are cabinet files. cab = "0.6" # std's clock where there is one; the browser's on wasm32-unknown-unknown, whose std has none. web-time = "1.1" +[target.'cfg(not(target_arch = "wasm32"))'.dependencies] +mdns-sd = { version = "0.21.4", default-features = false, optional = true } +# A relay's TLS, trusting the system's authorities, then Mozilla's. +rustls = { version = "0.23", default-features = false, features = ["ring", "std", "tls12"], optional = true } +rustls-native-certs = { version = "0.8", optional = true } +webpki-root-certs = { version = "1.0", optional = true } + # The browser: its randomness, clock, timers and event loop, and SQLite's files kept with # the notebooks' (`fs`). [target.'cfg(target_arch = "wasm32")'.dependencies] diff --git a/crates/notebook/src/background.rs b/crates/notebook/src/background.rs index 3f5a6a0cc1e94d296931553f5cc2fa392466cce0..6802c83a324026be0e39500ce01b10801ee95b77 100644 --- a/crates/notebook/src/background.rs +++ b/crates/notebook/src/background.rs @@ -79,7 +79,7 @@ pub(crate) struct Reports { #[derive(Default)] struct Watched { sections: Vec, - /// Sections whose file changed since `changed` was last asked. + /// Sections or cached folders changed since `changed` was last asked. changed: Vec, /// Folders a watch reported without naming the file, listed at `relist`. folders: BTreeSet, @@ -127,7 +127,7 @@ impl Background { /// lists each folder, both by catalog path, arming any watch that reports the folder's /// changes through `Reports`, and whether one does; it runs again after a transport /// failure. `copies` keeps an offline copy of every section. - pub(crate) fn start( + pub(crate) fn start( copies: bool, mut connect: impl FnMut(Reports) -> io::Result<((B, L), bool)> + Send + 'static, notify: impl Fn() + Send + 'static, @@ -135,7 +135,8 @@ impl Background { where R: Remote, B: FnMut(&str) -> R, - L: FnMut(&str) -> io::Result>, + L: FnMut(String) -> LF, + LF: Future>>, { let (signal, receiver) = Signal::new(); let shared = Arc::new(Shared { @@ -192,33 +193,39 @@ impl Background { connection: connection + 1, }; let connected = connect(reports); - let Ok(mut watched) = owner.watched.lock() else { - return; + let folders = { + let Ok(mut watched) = owner.watched.lock() else { + return; + }; + watched.connection += 1; + watched.lost = false; + match &connected { + Ok((_, reported)) => { + watched.watching(*reported); + Some(watched.rescanning()) + } + Err(error) => { + let error = Error::RemoteIo(io::Error::new( + error.kind(), + error.to_string(), + )); + let retry = now + RETRY; + for watch in &mut watched.sections { + news |= watch.fail(&error, None); + watch.due = watch.due.max(retry); + } + watched.relist = watched.relist.map(|due| due.max(retry)); + None + } + } }; - // Whatever ended before now ended with the last connection. - watched.connection += 1; - watched.lost = false; - match connected { - Ok((mut files, reported)) => { - watched.watching(reported); - let folders = watched.rescanning(); - drop(watched); - let listings = list(&mut files.1, folders); - let Ok(mut watched) = owner.watched.lock() else { - return; - }; - watched.listed(&listings, Some(copies), Instant::now()); - bound = Some(files); - } - Err(error) => { - let error = Error::RemoteIo(error); - let retry = now + RETRY; - for watch in &mut watched.sections { - news |= watch.fail(&error, None); - watch.due = watch.due.max(retry); - } - watched.relist = watched.relist.map(|due| due.max(retry)); - } + if let (Ok((mut files, _)), Some(folders)) = (connected, folders) { + let listings = list(&mut files.1, folders).await; + let Ok(mut watched) = owner.watched.lock() else { + return; + }; + watched.listed(&listings, Some(copies), Instant::now()); + bound = Some(files); } continue; } @@ -226,7 +233,7 @@ impl Background { let Some((_, list_folder)) = &mut bound else { continue; }; - let listings = list(list_folder, folders); + let listings = list(list_folder, folders).await; let Ok(mut watched) = owner.watched.lock() else { return; }; @@ -244,14 +251,15 @@ impl Background { let Some((bind, _)) = &mut bound else { continue; }; - let (queued, outcome) = step( + let (queued, outcome) = step_async( &mut bind(&path), replica.as_deref(), seen.as_deref(), current, image, copies, - ); + ) + .await; if outcome.as_ref().is_err_and(disconnected) { bound = None; } @@ -339,10 +347,12 @@ impl Background { crate::SmbRemote::new(Arc::clone(&bound), file, limit) }; let root = root.clone(); - let list = move |folder: &str| { - use crate::discover::Source; - crate::discover::Smb::new(&client, &root)? - .entries(folder, crate::session::LIMITS.entries) + let list = move |folder: String| { + std::future::ready((|| { + use crate::discover::Source; + crate::discover::Smb::new(&client, &root)? + .entries(&folder, crate::session::LIMITS.entries) + })()) }; Ok(((bind, list), reported)) }, @@ -364,7 +374,10 @@ impl Background { guest.watch(reports)?; let (bound, listing) = (Arc::clone(&guest), Arc::clone(&guest)); let bind = move |path: &str| crate::live::share::HostedRemote::new(&bound, path); - let list = move |folder: &str| listing.entries(folder); + let list = move |folder: String| { + let listing = listing.clone(); + async move { listing.entries_async(&folder).await } + }; Ok(((bind, list), true)) }, notify, @@ -553,6 +566,16 @@ impl Shared { #[cfg_attr(not(any(feature = "smb", feature = "live")), allow(dead_code))] impl Reports { + #[cfg(all(feature = "live", target_arch = "wasm32"))] + pub(crate) fn catalog(&self, folder: &str) { + if let Some(shared) = self.shared.upgrade() + && let Ok(mut watched) = shared.watched.lock() + && !watched.changed.iter().any(|path| path == folder) + { + watched.changed.push(folder.to_owned()); + } + } + /// The sections at or below these paths changed. pub(crate) fn touched(&self, paths: &[String]) { if let Some(shared) = self.shared.upgrade() { @@ -768,14 +791,14 @@ impl Watched { return Next::Wait(wait.map(|at| at.saturating_duration_since(now))); }; let watch = &mut self.sections[index]; + if !bound { + return Next::Connect(self.connection); + } if let Some(worker) = watch.held.as_ref().and_then(Weak::upgrade) { worker.wake(); watch.due = now + interval; continue; } - if !bound { - return Next::Connect(self.connection); - } return Next::Check { path: watch.path.clone(), replica: watch.replica.clone(), @@ -789,17 +812,16 @@ impl Watched { /// Lists each of `folders` through `list`; a folder that cannot be listed lists nothing, so /// that its sections are checked. -fn list( - list: &mut impl FnMut(&str) -> io::Result>, +async fn list>>>( + list: &mut impl FnMut(String) -> F, folders: BTreeSet, ) -> BTreeMap> { - folders - .into_iter() - .map(|folder| { - let entries = list(&folder).unwrap_or_default(); - (folder, entries) - }) - .collect() + let mut listings = BTreeMap::new(); + for folder in folders { + let entries = list(folder.clone()).await.unwrap_or_default(); + listings.insert(folder, entries); + } + listings } /// 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, u64) { /// with `current` is the stamp now, unread. With `copies`, a section without a replica gets /// one from the file as it is now. `image`, the file as discovery read it, stands in for /// reading it while the stamp is still its own. -fn step( +async fn step_async( remote: &mut R, replica: Option<&Path>, seen: Option<&Stamp>, @@ -880,7 +902,7 @@ fn step( ) -> (Option, Result<(Stamp, bool)>) { let stamp = match seen.filter(|_| current) { Some(seen) => seen.clone(), - None => match remote.stamp() { + None => match crate::sync::awaited(remote, Remote::stamp).await { Ok(stamp) => stamp, Err(error) => return (None, Err(Error::RemoteIo(error))), }, @@ -898,14 +920,18 @@ fn step( if !copies { return (Some(0), Ok((stamp, moved))); } - let copied = (|| { - let image = remote.read().map_err(Error::RemoteIo)?; + let copied = (async { + let image = crate::sync::awaited(remote, Remote::read) + .await + .map_err(Error::RemoteIo)?; if let Some(folder) = replica.parent() { fs::create_dir_all(folder)?; } Replica::seed(replica, &image, None)?; + fs::durable().await?; Ok(Stamp::of(&image)?) - })(); + }) + .await; return match copied { Ok(stamp) => (Some(0), Ok((stamp, moved))), // A session made it first. @@ -916,7 +942,7 @@ fn step( }; } // A version kept beside the file leaves its stamp as it was. - let versions = match remote.versions() { + let versions = match crate::sync::awaited(remote, Remote::versions).await { Ok(versions) => versions, Err(error) => return (None, Err(Error::RemoteIo(error))), }; @@ -932,16 +958,22 @@ fn step( Err(error) if error.busy() => return (None, Ok((stamp, false))), Err(error) => return (None, Err(error)), }; - let synced = (|| { + let synced = (async { let mut changed = moved; loop { - let synced = replica.sync_once(remote)?; + let synced = replica.sync_once_async(remote).await?; changed |= !synced.changed.is_empty(); if !matches!(synced.edit, Some((_, EditStatus::Published { .. }))) { - return Ok((remote.stamp().map_err(Error::RemoteIo)?, changed)); + return Ok(( + crate::sync::awaited(remote, Remote::stamp) + .await + .map_err(Error::RemoteIo)?, + changed, + )); } } - })(); + }) + .await; let queued = replica .recovery_summary() .ok() @@ -958,6 +990,10 @@ struct Discovered<'a, R> { } impl Remote for Discovered<'_, R> { + fn pending(&mut self) -> Option + '_>>> { + self.remote.pending() + } + fn read(&mut self) -> io::Result> { match self.image.take() { Some((stamp, image)) if stamp == self.stamp => Ok(image), @@ -1036,6 +1072,20 @@ mod tests { .collect() } + #[test] + fn a_notebook_with_every_section_held_still_connects_its_watch() { + let now = Instant::now(); + let mut watched = Watched::default(); + watched.watch(sections(1, true), now); + let (held, woken) = Signal::new(); + watched.sections[0].held = Some(Arc::downgrade(&held)); + assert!(matches!(watched.next(now, false), Next::Connect(0))); + watched.watching(true); + assert!(held.watched.load(Ordering::Acquire)); + assert!(matches!(watched.next(now, true), Next::Wait(_))); + assert!(woken.recv_timeout(Duration::from_millis(50)).is_ok()); + } + /// The notebook's folders as a listing shows its first `count` sections. fn listing(count: usize) -> BTreeMap> { let mut folders: BTreeMap> = BTreeMap::new(); @@ -1311,7 +1361,7 @@ mod tests { armed.send(reports).unwrap(); let (files, listed) = (remote.clone(), remote.clone()); let bind = move |path: &str| File(files.clone(), path.to_owned()); - let list = move |folder: &str| Ok(listed.list(folder)); + let list = move |folder: String| std::future::ready(Ok(listed.list(&folder))); Ok(((bind, list), true)) }, || {}, @@ -1387,7 +1437,7 @@ mod tests { armed.send(reports).unwrap(); let (files, listed) = (remote.clone(), remote.clone()); let bind = move |path: &str| File(files.clone(), path.to_owned()); - let list = move |folder: &str| Ok(listed.list(folder)); + let list = move |folder: String| std::future::ready(Ok(listed.list(&folder))); Ok(((bind, list), true)) }, || {}, @@ -1412,10 +1462,16 @@ mod tests { quiet(&stamps, SETTLE).is_empty(), "the worker reads, not this" ); - // The watch ending hands polling back to the worker. watch.lost(); assert!(woken.recv_timeout(SETTLE).is_ok()); - assert!(!signal.watched.load(Ordering::Acquire)); + let rearmed = watches.recv_timeout(SETTLE).unwrap(); + while !signal.watched.load(Ordering::Acquire) { + woken.recv_timeout(SETTLE * 2).unwrap(); + } + assert!(signal.watched.load(Ordering::Acquire)); + rearmed.touched(std::slice::from_ref(&held)); + assert!(woken.recv_timeout(SETTLE * 2).is_ok()); + assert!(quiet(&stamps, SETTLE).is_empty()); drop(signal); drop(background); } diff --git a/crates/notebook/src/fs.rs b/crates/notebook/src/fs.rs index d991f487d3ae80fc3ef593929d5917868f83ccb8..8ec0a8e5491aeb10063c1a1916898ef93f7103fc 100644 --- a/crates/notebook/src/fs.rs +++ b/crates/notebook/src/fs.rs @@ -19,3 +19,8 @@ pub fn commit_file( mod web; #[cfg(target_arch = "wasm32")] pub use web::*; + +#[cfg(not(target_arch = "wasm32"))] +pub async fn durable() -> std::io::Result<()> { + Ok(()) +} diff --git a/crates/notebook/src/fs/web.rs b/crates/notebook/src/fs/web.rs index d4979c80ceb5f62e01c23e3f17309570e585de9f..11f0e22bc7c515eca9b433f1ebc59421d7e49571 100644 --- a/crates/notebook/src/fs/web.rs +++ b/crates/notebook/src/fs/web.rs @@ -32,6 +32,32 @@ struct Data { const RANGES: usize = 64; impl Data { + fn flush(&mut self) { + if let Some(path) = self.path.clone() { + with(|files| { + if files.changed.remove(&path) { + files.durable.push((path, self.change())); + } + }); + } + } + + fn change(&mut self) -> Change { + let length = self.bytes.len(); + let ranges = self + .unwritten + .replace(Vec::new()) + .unwrap_or_else(|| std::iter::once(0..length).collect()); + Change::File { + length: length as u64, + ranges: ranges + .into_iter() + .map(|range| range.start.min(length)..range.end.min(length)) + .filter(|range| !range.is_empty()) + .map(|range| (range.start as u64, self.bytes[range].to_vec())) + .collect(), + } + } fn new( bytes: Vec, modified: f64, @@ -79,6 +105,8 @@ struct Files { nodes: BTreeMap, /// Paths written, made or removed since the host last asked. changed: BTreeSet, + /// Flushes and removals in order, including SQLite's WAL before its database checkpoint. + durable: Vec<(PathBuf, Change)>, /// Folders the host mirrors from elsewhere (`mount`). mounts: BTreeSet, /// Images committed under a mount, waiting for the host to write them (`committed`). @@ -164,39 +192,24 @@ pub enum Change { /// Whether any path changed since `changes` was last called. pub fn changed() -> bool { - FILES.with_borrow(|files| !files.changed.is_empty() || !files.committed.is_empty()) + FILES.with_borrow(|files| { + !files.durable.is_empty() || !files.changed.is_empty() || !files.committed.is_empty() + }) } -/// The paths changed since the last call, each with how; a folder before what it holds. +/// Ordered file flushes, followed by changes not yet flushed. pub fn changes() -> Vec<(PathBuf, Change)> { FILES.with_borrow_mut(|files| { - std::mem::take(&mut files.changed) - .into_iter() - .map(|path| { - let change = match files.nodes.get(&path) { - None => Change::Removed, - Some(Node::Directory) => Change::Directory, - Some(Node::File(data)) => { - let mut data = data.borrow_mut(); - let length = data.bytes.len(); - let ranges = data - .unwritten - .replace(Vec::new()) - .unwrap_or_else(|| std::iter::once(0..length).collect()); - Change::File { - length: length as u64, - ranges: ranges - .into_iter() - .map(|range| range.start.min(length)..range.end.min(length)) - .filter(|range| !range.is_empty()) - .map(|range| (range.start as u64, data.bytes[range].to_vec())) - .collect(), - } - } - }; - (path, change) - }) - .collect() + let mut changes = std::mem::take(&mut files.durable); + changes.extend(std::mem::take(&mut files.changed).into_iter().map(|path| { + let change = match files.nodes.get(&path) { + None => Change::Removed, + Some(Node::Directory) => Change::Directory, + Some(Node::File(data)) => data.borrow_mut().change(), + }; + (path, change) + })); + changes }) } @@ -445,9 +458,13 @@ pub fn remove_file(path: impl AsRef) -> io::Result<()> { let path = normal(path.as_ref()); with(|files| { let data = files.file(&path)?; - data.borrow_mut().path = None; + let mut data = data.borrow_mut(); + if files.changed.remove(&path) { + files.durable.push((path.clone(), data.change())); + } + data.path = None; files.nodes.remove(&path); - files.changed.insert(path); + files.durable.push((path, Change::Removed)); Ok(()) }) } @@ -687,11 +704,12 @@ impl File { } pub fn sync_all(&self) -> io::Result<()> { + self.data.borrow_mut().flush(); Ok(()) } pub fn sync_data(&self) -> io::Result<()> { - Ok(()) + self.sync_all() } pub fn set_len(&self, size: u64) -> io::Result<()> { @@ -848,3 +866,24 @@ pub fn supersede_file( .map_err(not_committed)?; rename(with, path).map_err(not_committed) } + +/// Waits until the browser host has durably written every ordered flush handed to it. +pub async fn durable() -> io::Result<()> { + use wasm_bindgen::prelude::*; + #[wasm_bindgen] + extern "C" { + #[wasm_bindgen(catch, js_namespace = globalThis, js_name = snowboundFlushStorage)] + fn flush_storage() -> Result; + } + let failure = |error: JsValue| { + io::Error::other( + error + .as_string() + .unwrap_or_else(|| "Browser storage failed".into()), + ) + }; + wasm_bindgen_futures::JsFuture::from(flush_storage().map_err(failure)?) + .await + .map_err(failure)?; + Ok(()) +} diff --git a/crates/notebook/src/fs/web/sqlite.rs b/crates/notebook/src/fs/web/sqlite.rs index f9915a8b4192dfcba54f0ba6d03c45bcd92792eb..c619435b54070a00cc857f38f02fa29d9c2e0433 100644 --- a/crates/notebook/src/fs/web/sqlite.rs +++ b/crates/notebook/src/fs/web/sqlite.rs @@ -49,6 +49,7 @@ impl VfsFile for File { } fn flush(&mut self) -> VfsResult<()> { + self.0.borrow_mut().flush(); Ok(()) } diff --git a/crates/notebook/src/lib.rs b/crates/notebook/src/lib.rs index 7160fc053dfeb6b7ba143a017ab8111957fa9b15..9410e8cc99f3563e6fd50bf6c7c278c2a9c7ba26 100644 --- a/crates/notebook/src/lib.rs +++ b/crates/notebook/src/lib.rs @@ -5,6 +5,7 @@ pub mod discover; pub mod fs; #[cfg(feature = "live")] +#[cfg_attr(target_arch = "wasm32", path = "live/web.rs")] pub mod live; #[cfg(feature = "smb")] pub mod smb; diff --git a/crates/notebook/src/live.rs b/crates/notebook/src/live.rs index c47e7ed5a054c3f9efcefcefc23c19a68c41583f..75bddc30729038f2993d08e39cc932ddf14d5b2e 100644 --- a/crates/notebook/src/live.rs +++ b/crates/notebook/src/live.rs @@ -19,7 +19,6 @@ pub use wire::{Caret, Guid, Hello, Presence, Spot}; use mdns_sd::{IfKind, ServiceDaemon, ServiceEvent, ServiceInfo}; use minicbor::Encode; -use sha2::{Digest, Sha256}; use std::{ collections::BTreeMap, io::{self, Read, Write}, @@ -46,124 +45,9 @@ const PATIENCE: Duration = Duration::from_secs(30); /// Wrong tries of a code met off any relay before it admits no one new, as a relay burns one. const TRIES: u32 = 5; -/// The secret peers meet through. -#[derive(Clone, Debug, PartialEq, Eq)] -pub enum Room { - /// A notebook's room, or a share's: a random secret only its members hold. - Notebook([u8; 16]), - /// A code typed on both ends (`code`), with a password where one is set: its room's - /// number names it on the network, and its secret and password only the two people know. - /// Its `owner`, the end sharing it, takes the number from a relay (the one the code has, - /// coming back; any free one for a secret alone) or picks one where it has no relay, and - /// judges every try, burning the code after too many wrong ones. - Code { - code: String, - password: String, - owner: bool, - }, -} - -impl Room { - /// The room a code typed on this end leads to. - pub fn join(code: &str, password: &str) -> Self { - Self::Code { - code: code.to_owned(), - password: password.to_owned(), - owner: false, - } - } - - /// The room of a code this end shares: its secret, or a whole code to keep its number. - pub fn share(code: &str, password: &str) -> Self { - Self::Code { - code: code.to_owned(), - password: password.to_owned(), - owner: true, - } - } - - /// What names the room in the clear: a hash of a room's secret, a code's number; none for - /// a code without one yet. - fn tag(&self) -> Option { - match self { - Room::Notebook(id) => Some(hex(&Sha256::digest( - [&b"Snowbound room "[..], id].concat(), - )[..8])), - Room::Code { code, .. } => code_parts(code).0.map(|number| format!("code-{number}")), - } - } - - fn secret(&self) -> Vec { - match self { - Room::Notebook(id) => id.to_vec(), - Room::Code { code, password, .. } => { - let mut secret = code_parts(code).1.into_bytes(); - if !password.is_empty() { - secret.push(b'\n'); - secret.extend_from_slice(password.as_bytes()); - } - secret - } - } - } - - fn owner(&self) -> bool { - matches!(self, Room::Code { owner: true, .. }) - } -} - -/// A code's room number, if it has one, and its secret: a whole code, or a secret alone. -fn code_parts(code: &str) -> (Option, String) { - match code::parse(code) { - Some((number, secret)) => (Some(number), secret), - None => (None, code.trim().to_uppercase()), - } -} - -/// Where peers are looked for. -#[derive(Clone, Copy, Debug, PartialEq, Eq)] -pub enum Reach { - /// Every network this computer is on. - Network, - /// This computer alone, for two copies of the app side by side. - Loopback, -} - -#[derive(Clone, Debug)] -pub struct Peer { - pub hello: Arc, - /// None until the peer first says where it is. - pub presence: Option, -} - -/// What happened, as `Live::start`'s `events` hears it on a network thread. -pub enum Event<'a> { - /// The peers, their presence, the code or the relay's answer changed. - Changed, - /// A stream to a peer opened, with the line to it. - Met(&'a Arc, &'a Line), - /// A peer's stream ended. - Left(&'a Arc), - /// A frame of a kind presence doesn't read itself, on a stream or to the group. - Frame { - from: &'a Arc, - kind: u16, - body: &'a [u8], - }, -} - -/// How the relay last answered. -#[derive(Clone, Debug, PartialEq, Eq)] -pub enum Relayed { - /// Not asked yet, or no relay. - Unknown, - /// In the room. - Joined, - /// Answered with this HTTP status, and how long it asked to wait. - Refused(u16, Option), - /// Not reached, and why. - Unreachable(Trouble), -} +mod model; +pub use model::{Event, Peer, Reach, Relayed, Room}; +use model::{code_parts, hex}; /// Presence on the network while it lives; dropping it leaves. pub struct Live { @@ -909,9 +793,5 @@ fn decode(text: &str) -> String { String::from_utf8_lossy(&decoded).into_owned() } -fn hex(bytes: &[u8]) -> String { - bytes.iter().map(|byte| format!("{byte:02x}")).collect() -} - #[cfg(test)] mod tests; diff --git a/crates/notebook/src/live/model.rs b/crates/notebook/src/live/model.rs new file mode 100644 index 0000000000000000000000000000000000000000..cf07432cafd8ef78df5b0b9d1127940894a5e45b --- /dev/null +++ b/crates/notebook/src/live/model.rs @@ -0,0 +1,149 @@ +use super::{Hello, Line, Presence, code}; +use sha2::{Digest, Sha256}; +use std::{sync::Arc, time::Duration}; + +pub(super) fn hex(bytes: &[u8]) -> String { + bytes.iter().map(|byte| format!("{byte:02x}")).collect() +} + +/// The secret peers meet through. +#[derive(Clone, Debug, PartialEq, Eq)] +pub enum Room { + /// A notebook's room, or a share's: a random secret only its members hold. + Notebook([u8; 16]), + /// A code typed on both ends (`code`), with a password where one is set: its room's + /// number names it on the network, and its secret and password only the two people know. + /// Its `owner`, the end sharing it, takes the number from a relay (the one the code has, + /// coming back; any free one for a secret alone) or picks one where it has no relay, and + /// judges every try, burning the code after too many wrong ones. + Code { + code: String, + password: String, + owner: bool, + }, +} + +impl Room { + /// The room a code typed on this end leads to. + pub fn join(code: &str, password: &str) -> Self { + Self::Code { + code: code.to_owned(), + password: password.to_owned(), + owner: false, + } + } + + /// The room of a code this end shares: its secret, or a whole code to keep its number. + pub fn share(code: &str, password: &str) -> Self { + Self::Code { + code: code.to_owned(), + password: password.to_owned(), + owner: true, + } + } + + /// What names the room in the clear: a hash of a room's secret, a code's number; none for + /// a code without one yet. + pub(super) fn tag(&self) -> Option { + match self { + Room::Notebook(id) => Some(hex(&Sha256::digest( + [&b"Snowbound room "[..], id].concat(), + )[..8])), + Room::Code { code, .. } => code_parts(code).0.map(|number| format!("code-{number}")), + } + } + + pub(super) fn secret(&self) -> Vec { + match self { + Room::Notebook(id) => id.to_vec(), + Room::Code { code, password, .. } => { + let mut secret = code_parts(code).1.into_bytes(); + if !password.is_empty() { + secret.push(b'\n'); + secret.extend_from_slice(password.as_bytes()); + } + secret + } + } + } + + #[cfg(not(target_arch = "wasm32"))] + pub(super) fn owner(&self) -> bool { + matches!(self, Room::Code { owner: true, .. }) + } +} + +/// A code's room number, if it has one, and its secret: a whole code, or a secret alone. +pub(super) fn code_parts(code: &str) -> (Option, String) { + match code::parse(code) { + Some((number, secret)) => (Some(number), secret), + None => (None, code.trim().to_uppercase()), + } +} + +/// Where peers are looked for. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum Reach { + /// Every network this computer is on. + Network, + /// This computer alone, for two copies of the app side by side. + Loopback, +} + +#[derive(Clone, Debug)] +pub struct Peer { + pub hello: Arc, + /// None until the peer first says where it is. + pub presence: Option, +} + +/// What happened, as `Live::start`'s `events` hears it on a network thread. +pub enum Event<'a> { + /// The peers, their presence, the code or the relay's answer changed. + Changed, + /// A stream to a peer opened, with the line to it. + Met(&'a Arc, &'a Line), + /// A peer's stream ended. + Left(&'a Arc), + /// A frame of a kind presence doesn't read itself, on a stream or to the group. + Frame { + from: &'a Arc, + kind: u16, + body: &'a [u8], + }, +} + +/// How the relay last answered. +#[derive(Clone, Debug, PartialEq, Eq)] +pub enum Relayed { + /// Not asked yet, or no relay. + Unknown, + /// In the room. + Joined, + /// Answered with this HTTP status, and how long it asked to wait. + Refused(u16, Option), + /// Not reached, and why. + Unreachable(Trouble), +} + +/// Why a relay could not be reached, as a person can act on it. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum Trouble { + /// The relay's name, or the proxy's, didn't resolve. + Dns, + /// Nothing answered at the relay's address: a firewall, or the relay is down. + Unreachable, + TimedOut, + /// The proxy itself couldn't be reached. + ProxyUnreachable, + /// The proxy asks for a name and password (407). + ProxyAuthentication, + /// The proxy refused to connect to the relay, with this status. + ProxyRefused(u16), + /// The relay's certificate wasn't one the system trusts: something on the way presents its + /// own. + Certificate, + /// Something on the way refused both WebSockets and plain requests to the relay. + Blocked, + Other, +} diff --git a/crates/notebook/src/live/share.rs b/crates/notebook/src/live/share.rs index a8f5e15f1b6b0ee78bc54d5778879a73fe8dff32..84ac703c5f098247a4afdcb2f9a102be78dd49df 100644 --- a/crates/notebook/src/live/share.rs +++ b/crates/notebook/src/live/share.rs @@ -8,9 +8,7 @@ use super::{ Event, Hello, Line, Live, Peer, Presence, Reach, Relayed, Room, Sender, - wire::{ - self, Delta, Failure, Reply, Request, Touched, Welcome, WireEntry, WireStamp, Written, kind, - }, + wire::{self, Delta, Failure, Reply, Request, Touched, Welcome, WireEntry, Written, kind}, }; use crate::{Error, Result, background::Reports, discover, session::Storage}; use onestore::{CommitError, CommitState, RevisionIndex, Stamp, Store, Transaction}; @@ -25,11 +23,17 @@ use std::{ mpsc, }, thread, - time::{Duration, Instant}, + time::Duration, }; +use web_time::Instant; +#[path = "share/batch.rs"] mod batch; +#[path = "share/membership.rs"] mod membership; +#[cfg(target_arch = "wasm32")] +#[path = "share/web.rs"] +mod web; pub use membership::Host; /// The most bytes one message of a read or an upload carries. @@ -168,12 +172,31 @@ pub fn join_while( reach: Option, relay: Option<&str>, continue_joining: impl Fn(bool) -> bool, +) -> std::result::Result { + crate::task::ready(join_while_async( + me, + code, + password, + reach, + relay, + continue_joining, + )) + .map_err(|_| Refusal::Busy)? +} + +pub async fn join_while_async( + me: Hello, + code: &str, + password: &str, + reach: Option, + relay: Option<&str>, + continue_joining: impl Fn(bool) -> bool, ) -> std::result::Result { let code = self::code(code).ok_or(Refusal::Malformed)?; let (welcomed, welcome) = mpsc::channel(); let approving = Arc::new(AtomicBool::new(false)); let approval = Arc::clone(&approving); - let (changed, waiting) = mpsc::channel(); + let (changed, waiting) = crate::task::channel(); let live = Live::start( me, &Room::join(&code, password), @@ -207,7 +230,7 @@ pub fn join_while( } } _ => { - let _ = changed.send(()); + let _ = changed.try_send(()); } }, ) @@ -257,7 +280,7 @@ pub fn join_while( } _ => {} } - let _ = waiting.recv_timeout(Duration::from_millis(100)); + crate::task::wait(&waiting, Some(Duration::from_millis(100))).await; } } @@ -963,26 +986,61 @@ struct Inner { here: Mutex, /// The host and the line to it, while connected. host: Mutex, Line)>>, - pending: Mutex>>, + pending: Mutex>>, next: AtomicU64, /// Where the host's reports of changed files go, while a background watches. watch: Mutex>, - /// The host stopped sharing. + /// Why this device can no longer reach the share. ended: Mutex>, /// Bytes of chunks asked for and not yet given. - asked: (Mutex, Condvar), + asked: (Mutex, Condvar), /// The sections read lately, kept as the host's deltas change them, newest first. images: Mutex>)>>, /// The stamps of sections whose image held is the host's now, as its last report of /// them said. current: Mutex>, + #[cfg(target_arch = "wasm32")] + catalog: Mutex>, } -enum Ended { +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum Ended { Stopped, Removed, } +#[derive(Default)] +struct Chunks { + bytes: usize, + #[cfg(target_arch = "wasm32")] + waiting: Vec, +} + +struct Requested<'a> { + inner: &'a Inner, + id: u64, + chunked: bool, +} + +impl Drop for Requested<'_> { + fn drop(&mut self) { + self.inner.pending.lock().unwrap().remove(&self.id); + if self.chunked { + let (asked, room) = &self.inner.asked; + let mut asked = asked.lock().unwrap(); + asked.bytes -= CHUNK; + #[cfg(target_arch = "wasm32")] + let waiting = std::mem::take(&mut asked.waiting); + drop(asked); + room.notify_all(); + #[cfg(target_arch = "wasm32")] + for waker in waiting { + waker.wake(); + } + } + } +} + impl Guest { /// Joins share `share` through its `secret` as `me`, where `reach` and `relay` say. /// `events` runs on a network thread whenever the host comes or goes, or the peers change. @@ -1010,6 +1068,8 @@ impl Guest { asked: Default::default(), images: Mutex::default(), current: Mutex::default(), + #[cfg(target_arch = "wasm32")] + catalog: Mutex::default(), }); let heard = Arc::clone(&inner); let live = Live::start(me, &Room::Notebook(secret), reach, relay, move |event| { @@ -1030,9 +1090,9 @@ impl Guest { host.as_ref().map(|(hello, _)| Arc::clone(hello)) } - /// Whether the host said it stopped sharing. - pub fn stopped(&self) -> bool { - self.inner.ended.lock().unwrap().is_some() + /// Why the host ended this device's access. + pub fn ended(&self) -> Option { + *self.inner.ended.lock().unwrap() } /// Everyone in the share's room: the host and the other guests. @@ -1077,27 +1137,62 @@ impl Guest { } /// Asks the host `kind` of `request`, waiting for its reply. - fn request(&self, kind: u16, mut request: Request) -> std::result::Result { + #[cfg(not(target_arch = "wasm32"))] + fn request(&self, kind: u16, request: Request) -> std::result::Result { + crate::task::ready(self.request_async(kind, request)).map_err(Failed::Unsent)? + } + + async fn request_async( + &self, + kind: u16, + mut request: Request, + ) -> std::result::Result { let Some((_, line)) = self.inner.host.lock().unwrap().clone() else { return Err(Failed::Unsent(self.offline())); }; let chunked = matches!(kind, kind::READ | kind::READ_FILE | kind::PUT); if chunked { - let (asked, room) = &self.inner.asked; - let mut asked = asked.lock().unwrap(); - while *asked + CHUNK > WINDOW { - asked = room.wait(asked).unwrap(); + #[cfg(not(target_arch = "wasm32"))] + { + let (asked, _room) = &self.inner.asked; + let mut asked = asked.lock().unwrap(); + while asked.bytes + CHUNK > WINDOW { + asked = _room.wait(asked).unwrap(); + } + asked.bytes += CHUNK; } - *asked += CHUNK; + #[cfg(target_arch = "wasm32")] + std::future::poll_fn(|context| { + let mut asked = self.inner.asked.0.lock().unwrap(); + if asked.bytes + CHUNK > WINDOW { + if !asked + .waiting + .iter() + .any(|waker| waker.will_wake(context.waker())) + { + asked.waiting.push(context.waker().clone()); + } + std::task::Poll::Pending + } else { + asked.bytes += CHUNK; + std::task::Poll::Ready(()) + } + }) + .await; } let id = self.inner.next.fetch_add(1, Ordering::Relaxed); request.id = id; - let (answer, answered) = mpsc::channel(); + let (answer, answered) = crate::task::response(); self.inner.pending.lock().unwrap().insert(id, answer); + let _requested = Requested { + inner: &self.inner, + id, + chunked, + }; let sent = line.send(kind, &request); let reply = match sent { Err(error) => Err(Failed::Unsent(error)), - Ok(()) => match answered.recv_timeout(TIMEOUT) { + Ok(()) => match crate::task::answered(&answered, TIMEOUT).await { Ok(reply) => Ok(reply), Err(mpsc::RecvTimeoutError::Timeout) => Err(Failed::Lost(io::Error::new( io::ErrorKind::TimedOut, @@ -1106,12 +1201,6 @@ impl Guest { Err(mpsc::RecvTimeoutError::Disconnected) => Err(Failed::Lost(self.offline())), }, }; - self.inner.pending.lock().unwrap().remove(&id); - if chunked { - let (asked, room) = &self.inner.asked; - *asked.lock().unwrap() -= CHUNK; - room.notify_all(); - } let reply = reply?; match reply.failure { Some(failure) => Err(Failed::Refused(failure)), @@ -1119,12 +1208,26 @@ impl Guest { } } + #[cfg(not(target_arch = "wasm32"))] fn ask(&self, kind: u16, request: Request) -> io::Result { self.request(kind, request).map_err(Failed::io) } + async fn ask_async(&self, kind: u16, request: Request) -> io::Result { + self.request_async(kind, request).await.map_err(Failed::io) + } + /// Puts `bytes` in `request`, or uploads them first where they are large. + #[cfg(not(target_arch = "wasm32"))] fn carry(&self, request: &mut Request, bytes: Vec) -> std::result::Result<(), Failed> { + crate::task::ready(self.carry_async(request, bytes)).map_err(Failed::Unsent)? + } + + async fn carry_async( + &self, + request: &mut Request, + bytes: Vec, + ) -> std::result::Result<(), Failed> { if bytes.len() <= CHUNK { request.bytes = Some(bytes); return Ok(()); @@ -1138,7 +1241,8 @@ impl Guest { ..Request::default() }; // Nothing the upload carries happens before the request that uses it. - self.request(kind::PUT, put) + self.request_async(kind::PUT, put) + .await .map_err(|failed| match failed { Failed::Lost(error) => Failed::Unsent(error), failed => failed, @@ -1150,22 +1254,29 @@ impl Guest { /// Reads the file at `path` a chunk at a time, as `kind` reads it; a section held, as /// the changes to it. + #[cfg(not(target_arch = "wasm32"))] fn read(&self, kind: u16, path: &str, limit: usize) -> io::Result> { + crate::task::ready(self.read_async(kind, path, limit))? + } + + async fn read_async(&self, kind: u16, path: &str, limit: usize) -> io::Result> { let held = (kind == kind::READ) .then(|| self.inner.image(path)) .flatten(); - let first = self.ask( - kind, - Request { - path: path.to_owned(), - offset: Some(0), - limit: Some(limit as u64), - stamp: (held.as_deref()) - .and_then(|image| Stamp::of(image).ok()) - .map(|stamp| (&stamp).into()), - ..Request::default() - }, - )?; + let first = self + .ask_async( + kind, + Request { + path: path.to_owned(), + offset: Some(0), + limit: Some(limit as u64), + stamp: (held.as_deref()) + .and_then(|image| Stamp::of(image).ok()) + .map(|stamp| (&stamp).into()), + ..Request::default() + }, + ) + .await?; if let (Some(writes), Some(held)) = (&first.writes, held) { let stamp: Option = first.stamp.as_ref().and_then(|s| s.try_into().ok()); let image = written(&held, first.length.unwrap_or_default(), writes) @@ -1181,7 +1292,7 @@ impl Guest { let mut image = first.bytes.unwrap_or_default(); while image.len() < length { let chunk = self - .ask( + .ask_async( kind, Request { path: path.to_owned(), @@ -1189,7 +1300,8 @@ impl Guest { handle: first.handle, ..Request::default() }, - )? + ) + .await? .bytes .unwrap_or_default(); if chunk.is_empty() { @@ -1210,35 +1322,52 @@ impl Guest { (Stamp::of(&image).ok().as_ref() == Some(stamp)).then(|| image.to_vec()) } + #[cfg(not(target_arch = "wasm32"))] pub(crate) fn entries(&self, folder: &str) -> io::Result> { - let reply = self.ask( - kind::LIST, - Request { - path: folder.to_owned(), - ..Request::default() - }, - )?; - Ok(reply + crate::task::ready(self.entries_async(folder))? + } + + pub(crate) async fn entries_async(&self, folder: &str) -> io::Result> { + let reply = self + .ask_async( + kind::LIST, + Request { + path: folder.to_owned(), + ..Request::default() + }, + ) + .await?; + let entries: Vec<_> = reply .entries .unwrap_or_default() .iter() .map(entry) - .collect()) + .collect(); + #[cfg(target_arch = "wasm32")] + self.cache_entries(folder, &entries).await?; + Ok(entries) } /// The stamp of the file at `path`: the host's last report of it where this guest holds /// that image, else the host's answer. + #[cfg(not(target_arch = "wasm32"))] fn stamp(&self, path: &str) -> io::Result { + crate::task::ready(self.stamp_async(path))? + } + + async fn stamp_async(&self, path: &str) -> io::Result { if let Some(stamp) = self.inner.current.lock().unwrap().get(path) { return Ok(stamp.clone()); } - let reply = self.ask( - kind::STAMP, - Request { - path: path.to_owned(), - ..Request::default() - }, - )?; + let reply = self + .ask_async( + kind::STAMP, + Request { + path: path.to_owned(), + ..Request::default() + }, + ) + .await?; reply .stamp .as_ref() @@ -1246,34 +1375,114 @@ impl Guest { .try_into() } + #[cfg(not(target_arch = "wasm32"))] fn commit( &self, path: &str, transaction: &Transaction, + ) -> std::result::Result<(), CommitError> { + crate::task::ready(self.commit_async(path, transaction)).map_err(|error| CommitError { + state: CommitState::NotCommitted, + error, + })? + } + + async fn commit_async( + &self, + path: &str, + transaction: &Transaction, ) -> std::result::Result<(), CommitError> { let mut request = Request { path: path.to_owned(), ..Request::default() }; - self.carry(&mut request, transaction.to_bytes()) + self.carry_async(&mut request, transaction.to_bytes()) + .await .map_err(Failed::commit)?; - self.request(kind::COMMIT, request) + self.request_async(kind::COMMIT, request) + .await .map(drop) .map_err(Failed::commit) } + #[cfg(not(target_arch = "wasm32"))] fn confirm(&self, path: &str, base: &Stamp) -> std::result::Result<(), CommitError> { + crate::task::ready(self.confirm_async(path, base)).map_err(|error| CommitError { + state: CommitState::NotCommitted, + error, + })? + } + + async fn confirm_async( + &self, + path: &str, + base: &Stamp, + ) -> std::result::Result<(), CommitError> { let request = Request { path: path.to_owned(), stamp: Some(base.into()), ..Request::default() }; - self.request(kind::CONFIRM, request) + self.request_async(kind::CONFIRM, request) + .await .map(drop) .map_err(Failed::commit) } + async fn edits_async( + &self, + path: &str, + transaction: &Transaction, + edits: &[crate::PendingEdit], + revisions: &BTreeMap, + ) -> std::result::Result<(), CommitError> { + let bytes = serde_json::to_vec(&batch::Edits { + edits: edits + .iter() + .map(|edit| (edit.author.clone(), edit.edit.clone())) + .collect(), + revisions: revisions.clone(), + }) + .map_err(|error| CommitError { + state: CommitState::NotCommitted, + error: io::Error::other(error), + })?; + let mut request = Request { + path: path.to_owned(), + stamp: Some(transaction.base().into()), + ..Request::default() + }; + self.carry_async(&mut request, bytes) + .await + .map_err(Failed::commit)?; + match self.request_async(kind::EDITS, request).await { + Ok(reply) => { + if let Some(stamp) = reply.stamp { + let stamp = Stamp::try_from(&stamp).map_err(|error| CommitError { + state: CommitState::Unknown, + error, + })?; + self.inner + .current + .lock() + .unwrap() + .insert(path.to_owned(), stamp); + } + Ok(()) + } + Err(Failed::Refused(failure)) + if wire::error_kind(failure.kind) == io::ErrorKind::Unsupported => + { + self.commit_async(path, transaction).await?; + self.inner.published(path, transaction); + Ok(()) + } + Err(error) => Err(error.commit()), + } + } + /// A request on `path` that answers nothing but whether it happened. + #[cfg(not(target_arch = "wasm32"))] fn verb(&self, kind: u16, path: &str, request: Request) -> Result<()> { self.ask( kind, @@ -1338,16 +1547,14 @@ impl Inner { fn heard(self: &Arc, event: Event) { let serves = |hello: &Hello| { - hello.serves == Some(self.share) - || self - .host - .lock() - .unwrap() - .as_ref() - .is_some_and(|(host, _)| host.peer == hello.peer) + self.host + .lock() + .unwrap() + .as_ref() + .is_some_and(|(host, _)| host.peer == hello.peer) }; match event { - Event::Met(hello, line) if serves(hello) => { + Event::Met(hello, line) if hello.serves == Some(self.share) => { *self.host.lock().unwrap() = Some((Arc::clone(hello), line.clone())); } Event::Left(hello) if serves(hello) => { @@ -1419,7 +1626,13 @@ impl Inner { self.relay.as_deref(), move |event| { if let Some(inner) = inner.upgrade() { - if matches!(event, Event::Frame { .. }) { + if matches!( + event, + Event::Frame { + kind: kind::DELTA | kind::TOUCHED, + .. + } + ) { inner.heard(event); } (inner.events)(); @@ -1450,6 +1663,8 @@ pub struct HostedRemote { /// The stamp last asked for, which an image the guest holds may already have. seen: Option, rejected: bool, + #[cfg(target_arch = "wasm32")] + pending: Option, } impl HostedRemote { @@ -1459,10 +1674,13 @@ impl HostedRemote { path: path.to_owned(), seen: None, rejected: false, + #[cfg(target_arch = "wasm32")] + pending: None, } } } +#[cfg(not(target_arch = "wasm32"))] impl crate::Remote for HostedRemote { fn accepts_edits(&self) -> bool { !self.rejected @@ -1478,57 +1696,26 @@ impl crate::Remote for HostedRemote { edits: &[crate::PendingEdit], revisions: &BTreeMap, ) -> std::result::Result<(), CommitError> { - let bytes = serde_json::to_vec(&batch::Edits { - edits: edits - .iter() - .map(|edit| (edit.author.clone(), edit.edit.clone())) - .collect(), - revisions: revisions.clone(), - }) - .map_err(|error| CommitError { - state: CommitState::NotCommitted, - error: io::Error::other(error), - })?; - let mut request = Request { - path: self.path.clone(), - stamp: Some(transaction.base().into()), - ..Request::default() - }; - self.guest - .carry(&mut request, bytes) - .map_err(Failed::commit)?; - let result = self.guest.request(kind::EDITS, request); - match result { - Ok(reply) => { - if let Some(stamp) = reply.stamp { - let stamp = Stamp::try_from(&stamp).map_err(|error| CommitError { - state: CommitState::Unknown, - error, - })?; - self.guest - .inner - .current - .lock() - .unwrap() - .insert(self.path.clone(), stamp); - } - self.seen = None; - Ok(()) - } - Err(Failed::Refused(failure)) - if wire::error_kind(failure.kind) == io::ErrorKind::Unsupported => - { - self.publish(transaction) - } - Err(error) => { - let error = error.commit(); - if error.state == CommitState::NotCommitted { - self.rejected = true; - self.guest.inner.current.lock().unwrap().remove(&self.path); - } - Err(error) - } + let result = + crate::task::ready( + self.guest + .edits_async(&self.path, transaction, edits, revisions), + ) + .map_err(|error| CommitError { + state: CommitState::NotCommitted, + error, + })?; + if result.is_ok() { + self.seen = None; } + if result + .as_ref() + .is_err_and(|error| error.state == CommitState::NotCommitted) + { + self.rejected = true; + self.guest.inner.current.lock().unwrap().remove(&self.path); + } + result } fn read(&mut self) -> io::Result> { @@ -1574,17 +1761,23 @@ type Listings = BTreeMap>; impl Hosted { pub(crate) fn new(guest: Arc, listed: PathBuf) -> Self { + #[cfg(target_arch = "wasm32")] + { + *guest.inner.catalog.lock().unwrap() = Some(listed.clone()); + } Self { guest, listed } } } /// The host's folders, or as they were last listed while it can't be reached. +#[cfg(not(target_arch = "wasm32"))] struct Source<'a> { guest: &'a Guest, kept: Listings, listed: Listings, } +#[cfg(not(target_arch = "wasm32"))] impl discover::Source for Source<'_> { fn entries(&mut self, path: &str, limit: usize) -> io::Result> { let entries = match self.guest.entries(path) { @@ -1630,6 +1823,7 @@ impl discover::Source for Source<'_> { } } +#[cfg(not(target_arch = "wasm32"))] impl Storage for Hosted { fn discover( &self, @@ -1749,7 +1943,7 @@ impl Storage for Hosted { fn supersede(&self, path: &str, base: &Stamp, with: &str) -> Result<()> { let request = Request { to: Some(with.to_owned()), - stamp: Some(WireStamp::from(base)), + stamp: Some(wire::WireStamp::from(base)), ..Request::default() }; self.guest.verb(kind::SUPERSEDE, path, request) diff --git a/crates/notebook/src/live/share/web.rs b/crates/notebook/src/live/share/web.rs new file mode 100644 index 0000000000000000000000000000000000000000..cd5fcc65644f462a18d4666f0120958100f20471 --- /dev/null +++ b/crates/notebook/src/live/share/web.rs @@ -0,0 +1,382 @@ +use super::*; +use std::{future::Future, pin::Pin}; + +pub(super) enum Outcome { + Read(Vec), + Stamp(Box), + Published, +} + +pub(super) enum Pending { + Running(Pin>>>), + Ready(std::result::Result), +} + +fn uncommitted(error: io::Error) -> CommitError { + CommitError { + state: CommitState::NotCommitted, + error, + } +} + +impl HostedRemote { + fn operation( + &mut self, + work: impl Future> + 'static, + ) -> std::result::Result { + match self.pending.take() { + Some(Pending::Ready(result)) => result, + pending => { + self.pending = Some(pending.unwrap_or_else(|| Pending::Running(Box::pin(work)))); + Err(uncommitted(io::ErrorKind::WouldBlock.into())) + } + } + } + + fn publication( + &mut self, + result: std::result::Result, + ) -> std::result::Result<(), CommitError> { + match result { + Ok(Outcome::Published) => { + self.seen = None; + Ok(()) + } + Ok(_) => Err(uncommitted(io::ErrorKind::InvalidInput.into())), + Err(error) => { + if error.error.kind() != io::ErrorKind::WouldBlock + && error.state == CommitState::NotCommitted + { + self.rejected = true; + self.guest.inner.current.lock().unwrap().remove(&self.path); + } + Err(error) + } + } + } +} + +impl crate::Remote for HostedRemote { + fn pending(&mut self) -> Option + '_>>> { + let pending = self.pending.as_mut()?; + Some(Box::pin(async move { + if let Pending::Running(work) = pending { + *pending = Pending::Ready(work.await); + } + })) + } + + fn read(&mut self) -> io::Result> { + let (guest, path, seen) = (self.guest.clone(), self.path.clone(), self.seen.clone()); + let result = self + .operation(async move { + let image = match seen.as_ref().and_then(|stamp| guest.held(&path, stamp)) { + Some(image) => image, + None => guest + .read_async(kind::READ, &path, LIMIT) + .await + .map_err(uncommitted)?, + }; + Ok(Outcome::Read(image)) + }) + .map_err(|error| error.error)?; + match result { + Outcome::Read(image) => { + self.rejected = false; + Ok(image) + } + _ => Err(io::ErrorKind::InvalidInput.into()), + } + } + + fn stamp(&mut self) -> io::Result { + let (guest, path) = (self.guest.clone(), self.path.clone()); + let result = self + .operation(async move { + guest + .stamp_async(&path) + .await + .map(|stamp| Outcome::Stamp(Box::new(stamp))) + .map_err(uncommitted) + }) + .map_err(|error| error.error)?; + match result { + Outcome::Stamp(stamp) => { + self.seen = Some((*stamp).clone()); + Ok(*stamp) + } + _ => Err(io::ErrorKind::InvalidInput.into()), + } + } + + fn accepts_edits(&self) -> bool { + !self.rejected + && self + .guest + .host() + .is_some_and(|host| host.ops == Some(1) && host.kinds.contains(&kind::EDITS)) + } + + fn publish_edits( + &mut self, + transaction: &Transaction, + edits: &[crate::PendingEdit], + revisions: &BTreeMap, + ) -> std::result::Result<(), CommitError> { + let (guest, path, transaction, edits, revisions) = ( + self.guest.clone(), + self.path.clone(), + transaction.clone(), + edits.to_vec(), + revisions.clone(), + ); + let result = self.operation(async move { + guest + .edits_async(&path, &transaction, &edits, &revisions) + .await?; + Ok(Outcome::Published) + }); + self.publication(result) + } + + fn publish(&mut self, transaction: &Transaction) -> std::result::Result<(), CommitError> { + let (guest, path, transaction) = + (self.guest.clone(), self.path.clone(), transaction.clone()); + let result = self.operation(async move { + guest.commit_async(&path, &transaction).await?; + guest.inner.published(&path, &transaction); + Ok(Outcome::Published) + }); + self.publication(result) + } + + fn confirm(&mut self, base: &Stamp) -> std::result::Result<(), CommitError> { + let (guest, path, base) = (self.guest.clone(), self.path.clone(), base.clone()); + let result = self.operation(async move { + guest.confirm_async(&path, &base).await?; + Ok(Outcome::Published) + }); + self.publication(result) + } +} + +struct Cached<'a>(&'a Hosted, Listings); + +impl Guest { + pub(super) async fn cache_entries( + &self, + folder: &str, + entries: &[discover::Entry], + ) -> io::Result<()> { + let Some(listed) = self.inner.catalog.lock().unwrap().clone() else { + return Ok(()); + }; + let files = listed.with_extension("files"); + let mut listings: Listings = crate::fs::read(&listed) + .ok() + .and_then(|bytes| serde_json::from_slice(&bytes).ok()) + .unwrap_or_default(); + let before = listings.get(folder).cloned().unwrap_or_default(); + let mut current = Vec::with_capacity(entries.len()); + crate::fs::create_dir_all(files.join(folder))?; + for entry in entries { + if entry.name.contains(['/', '\\', '\0']) || matches!(entry.name.as_str(), "." | "..") { + return Err(io::ErrorKind::InvalidData.into()); + } + let e = wire_entry(entry); + let listed = (e.name, e.kind, e.size, e.modified); + let path = if folder.is_empty() { + entry.name.clone() + } else { + format!("{folder}/{}", entry.name) + }; + if entry.kind == discover::EntryKind::File + && (!before.contains(&listed) || crate::fs::metadata(files.join(&path)).is_err()) + { + let kind = if path.ends_with(".one") || path.ends_with(".onetoc2") { + kind::READ + } else { + kind::READ_FILE + }; + let bytes = self.read_async(kind, &path, LIMIT).await?; + crate::fs::write(files.join(path), bytes)?; + } + current.push(listed); + } + // Other folders may have finished listing while these files were being read. + if let Ok(bytes) = crate::fs::read(&listed) { + listings = serde_json::from_slice(&bytes).map_err(io::Error::other)?; + } + if listings.get(folder) != Some(¤t) { + listings.insert(folder.to_owned(), current); + crate::fs::write(&listed, serde_json::to_vec(&listings)?)?; + crate::fs::durable().await?; + if let Some(reports) = &*self.inner.watch.lock().unwrap() { + reports.catalog(folder); + } + (self.inner.events)(); + } + Ok(()) + } +} + +impl discover::Source for Cached<'_> { + fn entries(&mut self, path: &str, limit: usize) -> io::Result> { + let entries = self.1.get(path).ok_or(io::ErrorKind::NotConnected)?; + if entries.len() > limit { + return Err(io::ErrorKind::FileTooLarge.into()); + } + Ok(entries + .iter() + .map(|(name, kind, size, modified)| { + entry(&WireEntry { + name: name.clone(), + kind: *kind, + size: *size, + modified: *modified, + }) + }) + .collect()) + } + fn read(&mut self, path: &str, limit: usize) -> io::Result> { + let bytes = self + .0 + .guest + .inner + .image(path) + .map(|image| image.to_vec()) + .map(Ok) + .unwrap_or_else(|| crate::fs::read(self.0.files().join(path)))?; + if bytes.len() > limit { + return Err(io::ErrorKind::FileTooLarge.into()); + } + Ok(bytes) + } + fn read_asset(&mut self, path: &str, limit: usize) -> io::Result> { + self.read(path, limit) + } +} + +impl Hosted { + fn files(&self) -> PathBuf { + self.listed.with_extension("files") + } + + /// Loads a bounded catalog and its files before the synchronous editor opens it. + pub async fn prepare(guest: Arc, cache: &std::path::Path) -> io::Result<()> { + let listed = + crate::session::listing(cache, &guest.location()).with_extension("entries.json"); + let hosted = Self::new(guest, listed); + let deadline = Instant::now() + Duration::from_secs(20); + let (_, waiting) = crate::task::channel(); + while hosted.guest.host().is_none() { + if Instant::now() >= deadline { + return Err(hosted.guest.offline()); + } + crate::task::wait(&waiting, Some(Duration::from_millis(50))).await; + } + let mut folders = vec![String::new()]; + let mut count = 0; + while let Some(folder) = folders.pop() { + if folder.split('/').count() > 32 { + return Err(io::ErrorKind::FileTooLarge.into()); + } + let entries = hosted.guest.entries_async(&folder).await?; + count += entries.len(); + if count > 10000 { + return Err(io::ErrorKind::FileTooLarge.into()); + } + crate::fs::create_dir_all(hosted.files().join(&folder))?; + for entry in &entries { + if entry.kind == discover::EntryKind::Directory { + folders.push(if folder.is_empty() { + entry.name.clone() + } else { + format!("{folder}/{}", entry.name) + }); + } + } + } + crate::fs::durable().await + } +} + +fn desktop() -> Error { + io::Error::new( + io::ErrorKind::Unsupported, + "Create and organize sections in desktop Snowbound.", + ) + .into() +} + +impl Storage for Hosted { + fn discover( + &self, + cache: &mut discover::Cache, + limits: discover::Limits, + ) -> Result { + let listings = + serde_json::from_slice(&crate::fs::read(&self.listed)?).map_err(io::Error::other)?; + Ok(cache.discover(&mut Cached(self, listings), limits)?) + } + fn location(&self) -> String { + self.guest.location() + } + fn entries(&self, path: &str) -> io::Result> { + let listings = + serde_json::from_slice(&crate::fs::read(&self.listed)?).map_err(io::Error::other)?; + discover::Source::entries(&mut Cached(self, listings), path, 10000) + } + fn exists(&self, path: &str) -> bool { + crate::fs::metadata(self.files().join(path)).is_ok() + } + fn stamp(&self, path: &str) -> io::Result { + Stamp::of(&self.read(path).map_err(io::Error::other)?).map_err(io::Error::other) + } + fn read(&self, path: &str) -> Result> { + self.read_file(path, LIMIT) + } + fn read_file(&self, path: &str, limit: usize) -> Result> { + Ok(discover::Source::read( + &mut Cached(self, Listings::new()), + path, + limit, + )?) + } + fn create(&self, _: &str, _: &[u8]) -> Result<()> { + Err(desktop()) + } + fn create_directory(&self, _: &str) -> Result<()> { + Err(desktop()) + } + fn hide(&self, _: &str) -> Result<()> { + Err(desktop()) + } + fn rename(&self, _: &str, _: &str) -> Result<()> { + Err(desktop()) + } + fn rename_root(&self, _: &str, _: &[String]) -> Result { + Err(desktop()) + } + fn replace(&self, _: &str, _: &str) -> Result<()> { + Err(desktop()) + } + fn delete(&self, _: &str) -> Result<()> { + Err(desktop()) + } + fn place(&self, _: &str, _: [u8; 16], _: &str) -> Result<()> { + Err(desktop()) + } + fn commit(&self, _: &str, _: &Transaction) -> Result<()> { + Err(desktop()) + } + fn confirm(&self, _: &str, _: &Stamp) -> std::result::Result<(), CommitError> { + Err(uncommitted(io::Error::new( + io::ErrorKind::Unsupported, + "Section management requires desktop Snowbound", + ))) + } + fn supersede(&self, _: &str, _: &Stamp, _: &str) -> Result<()> { + Err(desktop()) + } +} diff --git a/crates/notebook/src/live/transport.rs b/crates/notebook/src/live/transport.rs index 5435b9d6f9fb1f5bc86b49e6f72cad8e86ed1505..30088c82cd871762bdb75e9d203922845459dae0 100644 --- a/crates/notebook/src/live/transport.rs +++ b/crates/notebook/src/live/transport.rs @@ -28,27 +28,7 @@ const MOST: usize = 4 << 20; /// Set once WebSockets were refused, so later connections go straight to polling. static POLLING: AtomicBool = AtomicBool::new(false); -/// Why a relay could not be reached, as a person can act on it. -#[derive(Clone, Copy, Debug, PartialEq, Eq)] -pub enum Trouble { - /// The relay's name, or the proxy's, didn't resolve. - Dns, - /// Nothing answered at the relay's address: a firewall, or the relay is down. - Unreachable, - TimedOut, - /// The proxy itself couldn't be reached. - ProxyUnreachable, - /// The proxy asks for a name and password (407). - ProxyAuthentication, - /// The proxy refused to connect to the relay, with this status. - ProxyRefused(u16), - /// The relay's certificate wasn't one the system trusts: something on the way presents its - /// own. - Certificate, - /// Something on the way refused both WebSockets and plain requests to the relay. - Blocked, - Other, -} +pub use super::model::Trouble; pub(super) enum Failure { Trouble(Trouble, io::Error), diff --git a/crates/notebook/src/live/web.js b/crates/notebook/src/live/web.js new file mode 100644 index 0000000000000000000000000000000000000000..603b99508a9fa378e1901557de4d30434601b17e --- /dev/null +++ b/crates/notebook/src/live/web.js @@ -0,0 +1,30 @@ +const sockets = new Map(); +let next = 1; + +export function liveConnect(url, message, closed) { + const id = next++; + const socket = new WebSocket(url); + socket.binaryType = 'arraybuffer'; + socket.onmessage = event => message(event.data); + socket.onclose = () => { + liveClose(id); + closed(); + }; + sockets.set(id, socket); + return id; +} + +export function liveSend(id, bytes) { + const socket = sockets.get(id); + if (!socket || socket.readyState !== WebSocket.OPEN) throw new Error('Relay disconnected'); + if (socket.bufferedAmount > 1048576) throw new Error('Relay cannot keep up'); + socket.send(bytes); +} + +export function liveClose(id) { + const socket = sockets.get(id); + sockets.delete(id); + if (!socket) return; + socket.onmessage = socket.onclose = null; + socket.close(); +} diff --git a/crates/notebook/src/live/web.rs b/crates/notebook/src/live/web.rs new file mode 100644 index 0000000000000000000000000000000000000000..f2c1c6927c7d943745a398195ccb3d0e49fb5b98 --- /dev/null +++ b/crates/notebook/src/live/web.rs @@ -0,0 +1,612 @@ +//! Browser Live Share: event-driven relay connections carrying the same sealed streams. + +pub use ::relay::code; +#[path = "group.rs"] +mod group; +#[path = "model.rs"] +mod model; +#[path = "share.rs"] +pub mod share; +#[path = "wire.rs"] +pub mod wire; +use model::hex; +pub use model::{Event, Peer, Reach, Relayed, Room, Trouble}; +pub use wire::{Caret, Guid, Hello, Presence, Spot}; + +use minicbor::Encode; +use std::{ + collections::BTreeMap, + io, + sync::{ + Arc, Mutex, Weak, + atomic::{AtomicBool, Ordering}, + }, + time::Duration, +}; +use wasm_bindgen::prelude::*; +use wire::{Side, kind}; + +#[wasm_bindgen(module = "/src/live/web.js")] +extern "C" { + #[wasm_bindgen(catch, js_name = liveConnect)] + fn connect(url: &str, message: &JsValue, closed: &JsValue) -> Result; + #[wasm_bindgen(catch, js_name = liveSend)] + fn send(id: u32, bytes: &[u8]) -> Result<(), JsValue>; + #[wasm_bindgen(js_name = liveClose)] + fn close(id: u32); +} + +pub struct Live(Arc); + +struct Shared { + me: Hello, + room: Room, + url: String, + events: Box, + stopped: AtomicBool, + state: Mutex, +} + +struct State { + socket: u32, + slot: u32, + relayed: Relayed, + failed: u32, + outdated: Option, + burned: bool, + streams: BTreeMap, + members: BTreeMap, + sealer: Option, + presence: Presence, + scheduled: bool, +} + +struct Stream { + bytes: Vec, + opening: Option, + receive: Option, + line: Line, + peer: Option, +} + +#[derive(Clone)] +pub struct Line { + shared: Weak, + slot: u32, + sealer: Arc>>, +} + +impl Line { + pub fn send(&self, kind: u16, body: &impl Encode<()>) -> io::Result<()> { + let shared = self.shared.upgrade().ok_or(io::ErrorKind::NotConnected)?; + let mut bytes = Vec::new(); + self.sealer + .lock() + .unwrap() + .as_mut() + .ok_or(io::ErrorKind::NotConnected)? + .send(&mut bytes, kind, body)?; + shared.stream(self.slot, &bytes) + } + + pub fn hang_up(&self, reason: &str) { + let _ = self.send( + kind::BYE, + &wire::Bye { + reason: reason.into(), + }, + ); + } +} + +#[derive(Clone)] +pub struct Sender(Weak); + +impl Sender { + pub fn send(&self, kind: u16, body: &impl Encode<()>, to: Option<&[[u8; 16]]>) { + if let Some(shared) = self.0.upgrade() { + match to { + None => { + let _ = shared.group(kind, body, None); + } + Some(to) => { + let slots: Vec<_> = shared + .state + .lock() + .unwrap() + .members + .iter() + .filter(|(_, m)| { + m.peer.as_ref().is_some_and(|p| to.contains(&p.hello.peer)) + }) + .map(|(slot, _)| *slot) + .collect(); + for slot in slots { + let _ = shared.group(kind, body, Some(slot)); + } + } + } + } + } +} + +impl Live { + pub fn start( + me: Hello, + room: &Room, + _: Option, + relay: Option<&str>, + events: impl Fn(Event) + Send + Sync + 'static, + ) -> io::Result { + let tag = room.tag().ok_or(io::ErrorKind::InvalidInput)?; + let relay = relay + .ok_or(io::ErrorKind::NotConnected)? + .trim_end_matches('/'); + if !relay.starts_with("wss://") && !relay.starts_with("ws://") { + return Err(io::ErrorKind::InvalidInput.into()); + } + let shared = Arc::new(Shared { + me, + room: room.clone(), + url: format!("{relay}/v1/room/{tag}"), + events: Box::new(events), + stopped: AtomicBool::new(false), + state: Mutex::new(State { + socket: 0, + slot: 0, + relayed: Relayed::Unknown, + failed: 0, + outdated: None, + burned: false, + streams: BTreeMap::new(), + members: BTreeMap::new(), + sealer: None, + presence: Presence::default(), + scheduled: false, + }), + }); + shared.connect()?; + let weak = Arc::downgrade(&shared); + crate::task::spawn("live ping", move || async move { + let (_, wait) = crate::task::channel(); + loop { + crate::task::wait(&wait, Some(Duration::from_secs(15))).await; + let Some(shared) = weak + .upgrade() + .filter(|s| !s.stopped.load(Ordering::Acquire)) + else { + break; + }; + let lines: Vec<_> = shared + .state + .lock() + .unwrap() + .streams + .values() + .map(|s| s.line.clone()) + .collect(); + for line in lines { + let _ = line.send(kind::PING, &()); + } + if matches!(shared.room, Room::Notebook(_)) { + let _ = shared.group(kind::PING, &(), None); + } + } + })?; + Ok(Self(shared)) + } + + pub fn code(&self) -> Option { + match &self.0.room { + Room::Code { code, .. } => Some(code.clone()), + _ => None, + } + } + pub fn relayed(&self) -> Relayed { + self.0.state.lock().unwrap().relayed.clone() + } + pub fn failed(&self) -> u32 { + self.0.state.lock().unwrap().failed + } + pub fn other_version(&self) -> Option { + self.0.state.lock().unwrap().outdated + } + pub fn burned(&self) -> bool { + self.0.state.lock().unwrap().burned + } + pub fn sender(&self) -> Sender { + Sender(Arc::downgrade(&self.0)) + } + pub fn peers(&self) -> Vec { + let state = self.0.state.lock().unwrap(); + let mut peers = BTreeMap::new(); + for peer in state + .members + .values() + .filter_map(|m| m.peer.as_ref()) + .chain(state.streams.values().filter_map(|s| s.peer.as_ref())) + { + peers.insert(peer.hello.peer, peer.clone()); + } + peers.into_values().collect() + } + pub fn line(&self, peer: &[u8; 16]) -> Option { + self.0 + .state + .lock() + .unwrap() + .streams + .values() + .find(|s| s.peer.as_ref().is_some_and(|p| &p.hello.peer == peer)) + .map(|s| s.line.clone()) + } + pub fn set_presence(&self, presence: Presence) { + let mut state = self.0.state.lock().unwrap(); + if state.presence == presence { + return; + } + state.presence = presence; + if std::mem::replace(&mut state.scheduled, true) { + return; + } + drop(state); + let weak = Arc::downgrade(&self.0); + let _ = crate::task::spawn("live presence", move || async move { + let (_, wait) = crate::task::channel(); + crate::task::wait(&wait, Some(Duration::from_millis(100))).await; + if let Some(shared) = weak.upgrade() { + let presence = { + let mut state = shared.state.lock().unwrap(); + state.scheduled = false; + state.presence.clone() + }; + let _ = shared.group(kind::PRESENCE, &presence, None); + } + }); + } + pub fn leave(self, reason: &str) { + let lines: Vec<_> = self + .0 + .state + .lock() + .unwrap() + .streams + .values() + .map(|s| s.line.clone()) + .collect(); + for line in lines { + line.hang_up(reason); + } + } +} + +impl Drop for Live { + fn drop(&mut self) { + self.0.stopped.store(true, Ordering::Release); + close(self.0.state.lock().unwrap().socket); + } +} + +impl Shared { + fn connect(self: &Arc) -> io::Result<()> { + let weak = Arc::downgrade(self); + let message = Closure::::new(move |data: JsValue| { + if let Some(shared) = weak.upgrade() + && let Err(error) = shared.heard(data) + { + if let Some(version) = error + .get_ref() + .and_then(|e| e.downcast_ref::()) + { + shared.state.lock().unwrap().outdated = Some(version.0); + } else if error.kind() == io::ErrorKind::InvalidData { + shared.state.lock().unwrap().failed += 1; + } + shared.disconnected(); + } + }) + .into_js_value(); + let weak = Arc::downgrade(self); + let closed = Closure::::new(move || { + if let Some(shared) = weak.upgrade() { + shared.disconnected(); + } + }) + .into_js_value(); + let socket = connect(&self.url, &message, &closed).map_err(js_error)?; + self.state.lock().unwrap().socket = socket; + Ok(()) + } + + fn disconnected(self: &Arc) { + let peers = { + let mut state = self.state.lock().unwrap(); + close(state.socket); + let peers: Vec<_> = state + .streams + .values() + .filter_map(|s| s.peer.as_ref().map(|p| p.hello.clone())) + .collect(); + state.streams.clear(); + state.members.clear(); + state.sealer = None; + state.relayed = Relayed::Unreachable(Trouble::Other); + peers + }; + for hello in peers { + (self.events)(Event::Left(&hello)); + } + (self.events)(Event::Changed); + let weak = Arc::downgrade(self); + let _ = crate::task::spawn("live reconnect", move || async move { + let (_, wait) = crate::task::channel(); + crate::task::wait(&wait, Some(Duration::from_secs(2))).await; + if let Some(shared) = weak + .upgrade() + .filter(|s| !s.stopped.load(Ordering::Acquire)) + { + let _ = shared.connect(); + } + }); + } + + fn stream(&self, slot: u32, bytes: &[u8]) -> io::Result<()> { + let socket = self.state.lock().unwrap().socket; + for chunk in bytes.chunks(64 << 10) { + send(socket, &[&slot.to_be_bytes()[..], chunk].concat()).map_err(js_error)?; + } + Ok(()) + } + + fn group(&self, kind: u16, body: &impl Encode<()>, to: Option) -> io::Result<()> { + let body = minicbor::to_vec(body).map_err(io::Error::other)?; + let (socket, sealed) = { + let mut state = self.state.lock().unwrap(); + let sealed = state + .sealer + .as_mut() + .ok_or(io::ErrorKind::NotConnected)? + .seal(kind, &body, to.is_none())?; + (state.socket, sealed) + }; + let mut bytes = match to { + Some(slot) => [&(::relay::GROUP | 1).to_be_bytes()[..], &slot.to_be_bytes()].concat(), + None => ::relay::BROADCAST.to_be_bytes().to_vec(), + }; + bytes.extend(sealed); + send(socket, &bytes).map_err(js_error) + } + + fn meet(self: &Arc, slot: u32) -> io::Result<()> { + if self.state.lock().unwrap().streams.contains_key(&slot) { + return Ok(()); + } + let tag = self.room.tag().ok_or(io::ErrorKind::InvalidInput)?; + let (opening, bytes) = wire::Opening::new(Side::Initiator, &tag, &self.room.secret())?; + let line = Line { + shared: Arc::downgrade(self), + slot, + sealer: Arc::new(Mutex::new(None)), + }; + self.state.lock().unwrap().streams.insert( + slot, + Stream { + bytes: Vec::new(), + opening: Some(opening), + receive: None, + line, + peer: None, + }, + ); + self.stream( + slot, + &[&(bytes.len() as u32).to_be_bytes()[..], &bytes].concat(), + ) + } + + fn heard(self: &Arc, data: JsValue) -> io::Result<()> { + if let Some(text) = data.as_string() { + let notice = text + .parse::<::relay::Notice>() + .map_err(|_| io::ErrorKind::InvalidData)?; + match notice { + ::relay::Notice::Welcome { you, members } => { + { + let mut state = self.state.lock().unwrap(); + state.slot = you; + state.relayed = Relayed::Joined; + if matches!(self.room, Room::Notebook(_)) { + state.sealer = + Some(group::Sealer::new(&group::Keys::new(&self.room.secret()))?); + } + } + if matches!(self.room, Room::Code { .. }) { + for slot in members.into_iter().filter(|slot| *slot < you) { + self.meet(slot)?; + } + } else { + self.group(kind::HELLO, &self.me, None)?; + let presence = self.state.lock().unwrap().presence.clone(); + self.group(kind::PRESENCE, &presence, None)?; + } + } + ::relay::Notice::Left(slot) => { + let hello = { + let mut state = self.state.lock().unwrap(); + state.members.remove(&slot); + state + .streams + .remove(&slot) + .and_then(|s| s.peer.map(|p| p.hello)) + }; + if let Some(hello) = hello { + (self.events)(Event::Left(&hello)); + } + } + ::relay::Notice::Burned => self.state.lock().unwrap().burned = true, + _ => {} + } + (self.events)(Event::Changed); + return Ok(()); + } + let bytes = js_sys::Uint8Array::new(&data).to_vec(); + let (slot, bytes) = bytes + .split_first_chunk::<4>() + .ok_or(io::ErrorKind::InvalidData)?; + let slot = u32::from_be_bytes(*slot); + if slot & ::relay::GROUP != 0 { + return self.heard_group(slot & !::relay::GROUP, bytes); + } + self.heard_stream(slot, bytes) + } + + fn heard_group(self: &Arc, slot: u32, frame: &[u8]) -> io::Result<()> { + let (kind, body, hello) = { + let mut state = self.state.lock().unwrap(); + let member = match state.members.entry(slot) { + std::collections::btree_map::Entry::Occupied(e) => e.into_mut(), + std::collections::btree_map::Entry::Vacant(e) => e.insert(group::Member::new( + &group::Keys::new(&self.room.secret()), + frame, + )?), + }; + let (kind, body) = member.open(frame)?; + (kind, body, member.hello().cloned()) + }; + match kind { + kind::HELLO | kind::HELLO_BACK => { + let hello: Hello = + minicbor::decode(&body).map_err(|_| io::ErrorKind::InvalidData)?; + if hello.peer == self.me.peer { + return Ok(()); + } + let serves = hello.serves.is_some() && self.me.serves.is_none(); + self.state + .lock() + .unwrap() + .members + .get_mut(&slot) + .unwrap() + .peer = Some(Peer { + hello: Arc::new(hello), + presence: None, + }); + if kind == kind::HELLO { + self.group(kind::HELLO_BACK, &self.me, Some(slot))?; + let presence = self.state.lock().unwrap().presence.clone(); + self.group(kind::PRESENCE, &presence, Some(slot))?; + } + if serves { + self.meet(slot)?; + } + (self.events)(Event::Changed); + } + kind::PRESENCE => { + let presence = minicbor::decode(&body).map_err(|_| io::ErrorKind::InvalidData)?; + if let Some(peer) = self + .state + .lock() + .unwrap() + .members + .get_mut(&slot) + .and_then(|m| m.peer.as_mut()) + { + peer.presence = Some(presence); + } + (self.events)(Event::Changed); + } + kind if kind < 256 => { + if let Some(hello) = hello { + (self.events)(Event::Frame { + from: &hello, + kind, + body: &body, + }); + } + } + _ => {} + } + Ok(()) + } + + fn heard_stream(self: &Arc, slot: u32, bytes: &[u8]) -> io::Result<()> { + { + let mut state = self.state.lock().unwrap(); + let stream = state + .streams + .get_mut(&slot) + .ok_or(io::ErrorKind::InvalidData)?; + if stream.bytes.len() + bytes.len() > (16 << 20) + 4 { + return Err(io::ErrorKind::InvalidData.into()); + } + stream.bytes.extend_from_slice(bytes); + } + loop { + let event = { + let mut state = self.state.lock().unwrap(); + let stream = state.streams.get_mut(&slot).unwrap(); + let Some(length) = stream.bytes.get(..4) else { + break; + }; + let length = u32::from_be_bytes(length.try_into().unwrap()) as usize; + if length > 16 << 20 { + return Err(io::ErrorKind::InvalidData.into()); + } + if stream.bytes.len() < length + 4 { + break; + } + let bytes: Vec<_> = stream.bytes.drain(..length + 4).collect(); + if let Some(opening) = stream.opening.take() { + let (send, receive) = opening.finish(&bytes[4..])?; + *stream.line.sealer.lock().unwrap() = Some(send); + stream.receive = Some(receive); + Some((stream.line.clone(), None, kind::HELLO, Vec::new())) + } else { + let (kind, body) = stream + .receive + .as_mut() + .ok_or(io::ErrorKind::InvalidData)? + .receive(&mut io::Cursor::new(bytes))?; + if stream.peer.is_none() { + if kind != kind::HELLO { + return Err(io::ErrorKind::InvalidData.into()); + } + let hello = minicbor::decode::(&body) + .map_err(|_| io::ErrorKind::InvalidData)?; + stream.peer = Some(Peer { + hello: Arc::new(hello), + presence: None, + }); + } + Some(( + stream.line.clone(), + stream.peer.as_ref().map(|p| p.hello.clone()), + kind, + body, + )) + } + }; + if let Some((line, hello, kind, body)) = event { + match hello { + None => line.send(kind::HELLO, &self.me)?, + Some(hello) if kind == kind::HELLO => (self.events)(Event::Met(&hello, &line)), + Some(hello) => (self.events)(Event::Frame { + from: &hello, + kind, + body: &body, + }), + } + } + } + Ok(()) + } +} + +fn js_error(error: JsValue) -> io::Error { + io::Error::new( + io::ErrorKind::NotConnected, + error + .as_string() + .unwrap_or_else(|| "Relay disconnected".into()), + ) +} diff --git a/crates/notebook/src/live/wire.rs b/crates/notebook/src/live/wire.rs index 892135de87a079f6173267e61349916a4d8c8cee..8cbb9070a8773a71ab170e41977d5aa3cdaf4dcb 100644 --- a/crates/notebook/src/live/wire.rs +++ b/crates/notebook/src/live/wire.rs @@ -549,48 +549,72 @@ pub fn open( room: &str, secret: &[u8], ) -> io::Result<(Sealer, Sealer)> { - let password = Password::new(secret); - let [initiator, responder] = [b"initiator", b"responder"] - .map(|role| Identity::new(&[&role[..], room.as_bytes()].concat())); - let (pake, message) = match side { - Side::Initiator => Spake2::::start_a(&password, &initiator, &responder), - Side::Responder => Spake2::::start_b(&password, &initiator, &responder), - }; - let ours = Open { - version: VERSION, - room: room.into(), - pake: message, - }; + let (opening, ours) = Opening::new(side, room, secret)?; if side == Side::Initiator { - write_block(stream, &minicbor::to_vec(&ours).map_err(io::Error::other)?)?; + write_block(stream, &ours)?; } - let theirs: Open = - minicbor::decode(&read_block(stream)?).map_err(|_| invalid("A malformed opening"))?; - // A responder answers even a peer of another version, so that both ends can say which - // should update. + let theirs = read_block(stream)?; if side == Side::Responder { - write_block(stream, &minicbor::to_vec(&ours).map_err(io::Error::other)?)?; + write_block(stream, &ours)?; } - if theirs.version != VERSION { - return Err(io::Error::new( - io::ErrorKind::Unsupported, - Version(theirs.version), - )); + opening.finish(&theirs) +} + +pub(super) struct Opening { + pake: Spake2, + side: Side, + room: String, +} + +impl Opening { + pub(super) fn new(side: Side, room: &str, secret: &[u8]) -> io::Result<(Self, Vec)> { + let password = Password::new(secret); + let [initiator, responder] = [b"initiator", b"responder"] + .map(|role| Identity::new(&[&role[..], room.as_bytes()].concat())); + let (pake, message) = match side { + Side::Initiator => Spake2::::start_a(&password, &initiator, &responder), + Side::Responder => Spake2::::start_b(&password, &initiator, &responder), + }; + let ours = Open { + version: VERSION, + room: room.into(), + pake: message, + }; + let message = minicbor::to_vec(&ours).map_err(io::Error::other)?; + Ok(( + Self { + pake, + side, + room: room.to_owned(), + }, + message, + )) } - if theirs.room != room { - return Err(invalid("The peer means another room")); + + pub(super) fn finish(self, bytes: &[u8]) -> io::Result<(Sealer, Sealer)> { + let theirs: Open = minicbor::decode(bytes).map_err(|_| invalid("A malformed opening"))?; + if theirs.version != VERSION { + return Err(io::Error::new( + io::ErrorKind::Unsupported, + Version(theirs.version), + )); + } + if theirs.room != self.room { + return Err(invalid("The peer means another room")); + } + let key = self + .pake + .finish(&theirs.pake) + .map_err(|_| invalid("A malformed key exchange"))?; + let [from_initiator, from_responder] = [ + Sealer::new(&key, b"Snowbound live v1 initiator"), + Sealer::new(&key, b"Snowbound live v1 responder"), + ]; + Ok(match self.side { + Side::Initiator => (from_initiator, from_responder), + Side::Responder => (from_responder, from_initiator), + }) } - let key = pake - .finish(&theirs.pake) - .map_err(|_| invalid("A malformed key exchange"))?; - let [from_initiator, from_responder] = [ - Sealer::new(&key, b"Snowbound live v1 initiator"), - Sealer::new(&key, b"Snowbound live v1 responder"), - ]; - Ok(match side { - Side::Initiator => (from_initiator, from_responder), - Side::Responder => (from_responder, from_initiator), - }) } fn write_block(to: &mut impl Write, bytes: &[u8]) -> io::Result<()> { diff --git a/crates/notebook/src/session.rs b/crates/notebook/src/session.rs index 566499678d3adc3e0465eacb622c67aee93f8d77..5ea758ba37771e6328defdd3dff9967840faac23 100644 --- a/crates/notebook/src/session.rs +++ b/crates/notebook/src/session.rs @@ -1588,9 +1588,9 @@ impl Notebook { let (folder, remote) = (root.clone(), remote.clone()); let bind = move |path: &str| remote(&folder.join(path)); let mut local = discover::Local::open(&root)?; - let list = move |folder: &str| { + let list = move |folder: String| { use discover::Source; - local.entries(folder, LIMITS.entries) + std::future::ready(local.entries(&folder, LIMITS.entries)) }; Ok(((bind, list), watched)) }, diff --git a/crates/notebook/src/sync.rs b/crates/notebook/src/sync.rs index 2d73ab2faf045bcd6d3dd1b672f3c100c2d65f95..49f2260d3d9d15e96fe85c1c77788cd028755e9a 100644 --- a/crates/notebook/src/sync.rs +++ b/crates/notebook/src/sync.rs @@ -14,6 +14,11 @@ use std::{ /// Errors retain publication state; confirmation checks the stamp, flushes, and notifies /// cached readers. pub trait Remote { + /// Completes a non-blocking operation that returned `WouldBlock`; retrying that + /// operation then takes its result. Blocking providers leave this absent. + fn pending(&mut self) -> Option + '_>>> { + None + } fn read(&mut self) -> io::Result>; /// The file's stamp without reading its body or coordinating with writers. While it is /// the last observed image's, synchronization neither reads nor revalidates the file. @@ -54,6 +59,38 @@ pub trait Remote { } } +pub(crate) trait Waiting { + fn waiting(&self) -> bool; +} + +impl Waiting for io::Error { + fn waiting(&self) -> bool { + self.kind() == io::ErrorKind::WouldBlock + } +} + +impl Waiting for CommitError { + fn waiting(&self) -> bool { + self.error.kind() == io::ErrorKind::WouldBlock + } +} + +pub(crate) async fn awaited( + remote: &mut R, + mut operation: impl FnMut(&mut R) -> std::result::Result, +) -> std::result::Result { + loop { + let result = operation(remote); + if result.as_ref().is_err_and(Waiting::waiting) + && let Some(pending) = remote.pending() + { + pending.await; + continue; + } + return result; + } +} + #[derive(Debug, Clone, PartialEq, Eq)] pub enum EditStatus { Pending, @@ -189,29 +226,53 @@ impl Replica { /// Versions the remote keeps beside the file merge into it first, each published as one /// more revision and then retired (`resolve.rs`). pub fn sync_once(&self, remote: &mut impl Remote) -> Result { + crate::task::ready(self.sync_once_async(remote))? + } + + /// The same guarded synchronization step, awaiting non-blocking remote operations. + pub async fn sync_once_async(&self, remote: &mut impl Remote) -> Result { + let result = self.sync_once_inner(remote).await; + crate::fs::durable().await?; + result + } + + /// Ownership uses `try_lock`; the cache mutex is released before every await. + #[allow(clippy::await_holding_lock)] + async fn sync_once_inner(&self, remote: &mut impl Remote) -> Result { let _owner = self.sync_owner()?; - for version in remote.versions().map_err(Error::RemoteIo)? { - let image = remote.version(&version.id).map_err(Error::RemoteIo)?; - let current = remote.read().map_err(Error::RemoteIo)?; + for version in awaited(remote, Remote::versions) + .await + .map_err(Error::RemoteIo)? + { + let image = awaited(remote, |remote| remote.version(&version.id)) + .await + .map_err(Error::RemoteIo)?; + let current = awaited(remote, Remote::read) + .await + .map_err(Error::RemoteIo)?; let device = version.device.as_deref().unwrap_or("Another device"); let keep = match crate::resolve::merge(¤t, &image, device) { Ok(Merged::Held) => false, Ok(Merged::Publish(transaction)) => { - remote.publish(&transaction)?; + awaited(remote, |remote| remote.publish(&transaction)).await?; false } // A version this cannot merge is kept whole rather than lost. Ok(Merged::Foreign) | Err(Error::Document(_) | Error::Rejected(_)) => true, Err(error) => return Err(error), }; - remote.retire(&version.id, keep).map_err(Error::RemoteIo)?; + awaited(remote, |remote| remote.retire(&version.id, keep)) + .await + .map_err(Error::RemoteIo)?; } let state = state(&*self.lock()?)?; let batched = self.section.key.is_none() && state.queued && state.blocked.is_none() && remote.accepts_edits(); - let observed = remote.stamp().map_err(Error::RemoteIo)?; + let observed = awaited(remote, Remote::stamp) + .await + .map_err(Error::RemoteIo)?; if let Some(blocked) = &state.blocked && observed == *state.remote.as_ref().unwrap_or(&state.base) { @@ -229,7 +290,9 @@ impl Replica { let image = match observed { observed if observed == state.base || batched => None, _ => { - let image = remote.read().map_err(Error::RemoteIo)?; + let image = awaited(remote, Remote::read) + .await + .map_err(Error::RemoteIo)?; // Protected elsewhere: a section written anew, which only its key reads. if self.section.key.is_none() && crate::discover::locked(&Store::parse(&image)?) { return Err(Error::RemoteIo(io::Error::new( @@ -297,7 +360,8 @@ impl Replica { .collect(), ) }; - if let Err(error) = remote.confirm(&Stamp::of(&image)?) { + let stamp = Stamp::of(&image)?; + if let Err(error) = awaited(remote, |remote| remote.confirm(&stamp)).await { if error.state == CommitState::Committed { self.acknowledge(*batch, sealed.as_ref(), receipts.as_ref())?; } @@ -351,7 +415,7 @@ impl Replica { let Some(transaction) = transaction else { // Edits that changed nothing are published once the remote's image is durable. let base = base::base_stamp(&*self.lock()?)?; - if let Err(error) = remote.confirm(&base) { + if let Err(error) = awaited(remote, |remote| remote.confirm(&base)).await { if error.state == CommitState::Committed { self.acknowledge(batch, None, None)?; } @@ -364,18 +428,28 @@ impl Replica { changed, }); }; + if let Err(error) = crate::fs::durable().await { + self.lock()? + .execute("UPDATE batches SET attempted=0 WHERE id=?1", [batch])?; + return Err(error.into()); + } let published = if self.section.key.is_none() && remote.accepts_edits() { - let connection = self.lock()?; - let edits = queue::load(&connection, None, Some(batch))?; - let revisions = decode_revisions(&connection.query_row( - "SELECT revisions FROM batches WHERE id=?1", - [batch], - |row| row.get::<_, String>(0), - )?)?; - drop(connection); - remote.publish_edits(&transaction, &edits, &revisions) + let (edits, revisions) = { + let connection = self.lock()?; + let edits = queue::load(&connection, None, Some(batch))?; + let revisions = decode_revisions(&connection.query_row( + "SELECT revisions FROM batches WHERE id=?1", + [batch], + |row| row.get::<_, String>(0), + )?)?; + (edits, revisions) + }; + awaited(remote, |remote| { + remote.publish_edits(&transaction, &edits, &revisions) + }) + .await } else { - remote.publish(&transaction) + awaited(remote, |remote| remote.publish(&transaction)).await }; match published { Ok(()) => {} @@ -393,7 +467,13 @@ impl Replica { self.acknowledge(batch, Some(&transaction), None)?; let revision = self.receipt(id)?; if batched { - changed.extend(self.rebase(Some(remote.read().map_err(Error::RemoteIo)?))?); + changed.extend( + self.rebase(Some( + awaited(remote, Remote::read) + .await + .map_err(Error::RemoteIo)?, + ))?, + ); changed.sort(); changed.dedup(); } @@ -550,15 +630,21 @@ impl Replica { /// Whether nothing is queued and the remote still has the base's stamp, or the queue /// is blocked on a remote that has not changed since; reads neither image and does not /// lock the cache during remote I/O. - pub(crate) fn settled(&self, remote: &mut impl Remote) -> Result { + pub(crate) async fn settled(&self, remote: &mut impl Remote) -> Result { let state = state(&*self.lock()?)?; let expected = match &state.blocked { Some(_) => state.remote.unwrap_or(state.base), None if !state.queued => state.base, None => return Ok(false), }; - Ok(remote.stamp().map_err(Error::RemoteIo)? == expected - && remote.versions().map_err(Error::RemoteIo)?.is_empty()) + Ok(awaited(remote, Remote::stamp) + .await + .map_err(Error::RemoteIo)? + == expected + && awaited(remote, Remote::versions) + .await + .map_err(Error::RemoteIo)? + .is_empty()) } } diff --git a/crates/notebook/src/task.rs b/crates/notebook/src/task.rs index 2ff7bc01631e04d97d0569ad8b5c451524004d97..77f725750996aed7dfd7b775b06e42a54fb25894 100644 --- a/crates/notebook/src/task.rs +++ b/crates/notebook/src/task.rs @@ -1,7 +1,6 @@ //! The sync worker's and background's loops: threads natively, and in the browser, which -//! gives a module one thread, tasks on its event loop. A loop is an `async` block whose only -//! waits are `wait`; natively each wait blocks its thread, so `complete` runs the block to the -//! end in one poll. +//! gives a module one thread, tasks on its event loop. Native waits block their thread, so +//! `complete` runs the block to the end in one poll. use std::{io, time::Duration}; @@ -11,6 +10,22 @@ pub(crate) use std::{ thread::JoinHandle, }; +#[cfg(all(feature = "live", not(target_arch = "wasm32")))] +pub(crate) type Answer = std::sync::mpsc::Sender; + +#[cfg(all(feature = "live", not(target_arch = "wasm32")))] +pub(crate) fn response() -> (Answer, std::sync::mpsc::Receiver) { + std::sync::mpsc::channel() +} + +#[cfg(all(feature = "live", not(target_arch = "wasm32")))] +pub(crate) async fn answered( + receiver: &std::sync::mpsc::Receiver, + timeout: Duration, +) -> Result { + receiver.recv_timeout(timeout) +} + /// A wake that waits, at most one, as `mpsc::sync_channel(1)` keeps one. #[cfg(not(target_arch = "wasm32"))] pub(crate) fn channel() -> (Sender<()>, Receiver<()>) { @@ -50,6 +65,18 @@ pub(crate) fn complete(work: impl Future) -> T { } } +/// Polls a synchronous entry point once; browser I/O must use the async entry point. +pub(crate) fn ready(work: impl Future) -> io::Result { + let mut work = std::pin::pin!(work); + match work + .as_mut() + .poll(&mut std::task::Context::from_waker(std::task::Waker::noop())) + { + std::task::Poll::Ready(output) => Ok(output), + std::task::Poll::Pending => Err(io::ErrorKind::WouldBlock.into()), + } +} + #[cfg(target_arch = "wasm32")] pub(crate) use web::*; @@ -62,6 +89,77 @@ mod web { }; use wasm_bindgen::{JsCast, prelude::*}; + struct Response { + value: Option, + closed: bool, + waiting: Option, + } + + pub(crate) struct Answer(Arc>>); + pub(crate) struct Answered(Arc>>); + + pub(crate) fn response() -> (Answer, Answered) { + let response = Arc::new(Mutex::new(Response { + value: None, + closed: false, + waiting: None, + })); + (Answer(response.clone()), Answered(response)) + } + + impl Answer { + pub(crate) fn send(self, value: T) -> Result<(), ()> { + let wake = { + let mut response = self.0.lock().map_err(|_| ())?; + response.value = Some(value); + response.waiting.take() + }; + if let Some(wake) = wake { + wake.wake(); + } + Ok(()) + } + } + + impl Drop for Answer { + fn drop(&mut self) { + let wake = { + let mut response = self.0.lock().unwrap(); + response.closed = true; + response.waiting.take() + }; + if let Some(wake) = wake { + wake.wake(); + } + } + } + + pub(crate) async fn answered( + receiver: &Answered, + timeout: Duration, + ) -> Result { + let deadline = web_time::Instant::now() + timeout; + let mut armed = false; + std::future::poll_fn(|context| { + let mut response = receiver.0.lock().unwrap(); + if let Some(value) = response.value.take() { + return Poll::Ready(Ok(value)); + } + if response.closed { + return Poll::Ready(Err(std::sync::mpsc::RecvTimeoutError::Disconnected)); + } + if web_time::Instant::now() >= deadline { + return Poll::Ready(Err(std::sync::mpsc::RecvTimeoutError::Timeout)); + } + response.waiting = Some(context.waker().clone()); + if !std::mem::replace(&mut armed, true) { + after(timeout, context.waker().clone()); + } + Poll::Pending + }) + .await + } + #[derive(Default)] struct Bell { rung: bool, diff --git a/crates/notebook/src/worker.rs b/crates/notebook/src/worker.rs index e558ade6c555118ea434bba59d8e13435b0922a2..c47645d84d9c68c6d060b252dc7f93ee0492e961 100644 --- a/crates/notebook/src/worker.rs +++ b/crates/notebook/src/worker.rs @@ -214,13 +214,13 @@ impl Replica { } let result = match remote.as_mut() { Some(remote) => { - if reported && replica.settled(remote).unwrap_or(false) { + if reported && replica.settled(remote).await.unwrap_or(false) { worker_signal.requested.store(false, Ordering::Release); worker_signal.synced.store(crate::now(), Ordering::Release); rest().await; continue; } - let result = replica.sync_once(remote); + let result = replica.sync_once_async(remote).await; reported = result.is_ok(); result } diff --git a/crates/notebook/tests/live_membership.rs b/crates/notebook/tests/live_membership.rs index 835d23afbb6c261384057fc70f97f0cc340bac23..0f3141bc2bf6ccf966396d734c592c8d4884294b 100644 --- a/crates/notebook/tests/live_membership.rs +++ b/crates/notebook/tests/live_membership.rs @@ -119,7 +119,7 @@ fn removal_retires_one_credential_and_other_devices_reconnect_after_a_restart() .secret; host.remove(&bob_key).unwrap(); until("Bob was not removed", || { - bob.stopped() && bob.host().is_none() + bob.ended() == Some(share::Ended::Removed) && bob.host().is_none() }); until("presence did not move to the new room", || { host.guests().len() == 1 @@ -129,7 +129,7 @@ fn removal_retires_one_credential_and_other_devices_reconnect_after_a_restart() assert_eq!(current.members.len(), 1); assert_ne!(original_code, code(&host)); assert!(alice.host().is_some()); - assert!(!alice.stopped()); + assert_eq!(alice.ended(), None); let live = Notebook::open_hosted(Arc::clone(&alice), directory.path().join("alice")).unwrap(); assert_eq!(live.catalog().sections.len(), 2); let restarted: Sharing = diff --git a/crates/notebook/tests/live_share.rs b/crates/notebook/tests/live_share.rs index 6d2967e243eb7435a2df0b602846b224a0e50cec..6488732a92e816e0cb10ab96f4eac96a7d944dfd 100644 --- a/crates/notebook/tests/live_share.rs +++ b/crates/notebook/tests/live_share.rs @@ -398,7 +398,7 @@ fn stopping_lets_every_guest_go_and_retires_the_code() { let (guest, _notebook) = guest("Grace", &code, &url, &directory.path().join("grace")); host.stop(); until("the guest never heard the host stop", || { - guest.stopped() && guest.host().is_none() + guest.ended() == Some(share::Ended::Stopped) && guest.host().is_none() }); assert!(host.code().is_none() && host.guests().is_empty()); assert!(matches!( diff --git a/crates/notebook/tests/sync.rs b/crates/notebook/tests/sync.rs index 0a9ef69555ed098a3727324f6701e9592a4a96cb..1b701132182ecc61bfdbe3734ee4626da8981489 100644 --- a/crates/notebook/tests/sync.rs +++ b/crates/notebook/tests/sync.rs @@ -20,6 +20,106 @@ fn save(cache: &Replica, text: ExGuid, range: Range, replacement: &str) -> .unwrap() } +#[test] +fn asynchronous_io_keeps_the_attempt_and_recovers_a_lost_reply() { + use std::{ + future::Future, + pin::Pin, + task::{Context, Poll, Waker}, + }; + + struct Deferred { + server: Server, + ready: bool, + } + impl Remote for Deferred { + fn pending(&mut self) -> Option + '_>>> { + let mut yielded = false; + Some(Box::pin(std::future::poll_fn(move |context| { + if yielded { + self.ready = true; + Poll::Ready(()) + } else { + yielded = true; + context.waker().wake_by_ref(); + Poll::Pending + } + }))) + } + fn read(&mut self) -> io::Result> { + if !std::mem::take(&mut self.ready) { + return Err(io::ErrorKind::WouldBlock.into()); + } + self.server.read() + } + fn stamp(&mut self) -> io::Result { + if !std::mem::take(&mut self.ready) { + return Err(io::ErrorKind::WouldBlock.into()); + } + self.server.stamp() + } + fn publish(&mut self, transaction: &Transaction) -> Result<(), CommitError> { + if !std::mem::take(&mut self.ready) { + return Err(CommitError { + state: CommitState::NotCommitted, + error: io::ErrorKind::WouldBlock.into(), + }); + } + self.server.publish(transaction) + } + fn confirm(&mut self, base: &onestore::Stamp) -> Result<(), CommitError> { + if !std::mem::take(&mut self.ready) { + return Err(CommitError { + state: CommitState::NotCommitted, + error: io::ErrorKind::WouldBlock.into(), + }); + } + self.server.confirm(base) + } + } + fn run(work: impl Future) -> T { + let mut work = std::pin::pin!(work); + let mut polls = 0; + loop { + polls += 1; + assert!(polls < 100, "An async operation stopped making progress"); + if let Poll::Ready(result) = work.as_mut().poll(&mut Context::from_waker(Waker::noop())) + { + assert!(polls > 1); + return result; + } + } + } + let directory = tempfile::tempdir().unwrap(); + let path = directory.path().join("cache.sqlite"); + let source = onestore::create_section("async.one", "Original", "Fixture").unwrap(); + let (_, object, _) = text(&source); + let cache = Replica::create(&path, &source).unwrap(); + let id = save(&cache, object, 0..0, "Browser ").unwrap(); + let mut remote = Deferred { + server: Server::new(&source), + ready: false, + }; + remote.server.fault = Fault::UnknownAfter; + assert!(matches!( + run(cache.sync_once_async(&mut remote)), + Err(Error::Remote(CommitError { + state: CommitState::Unknown, + .. + })) + )); + assert_eq!(remote.server.publications, 1); + drop(cache); + let cache = Replica::open(&path).unwrap(); + assert!( + matches!(run(cache.sync_once_async(&mut remote)).unwrap().edit, Some((published, EditStatus::Published { .. })) if published == id) + ); + assert_eq!(remote.server.publications, 1); + assert_eq!(text(&remote.server.visible).2, "Browser Original"); + assert_eq!(remote.server.visible, remote.server.durable); + assert!(cache.pending().unwrap().is_empty()); +} + #[test] fn an_unchanged_stamp_publishes_and_settles_without_reading_the_remote() { struct Counted { diff --git a/crates/snowbound/src/commands.rs b/crates/snowbound/src/commands.rs index 46297270be9a13a11a520be307854319f6c86117..fe6547df3af6538e831c088b0bc7735cc0a26b52 100644 --- a/crates/snowbound/src/commands.rs +++ b/crates/snowbound/src/commands.rs @@ -1163,6 +1163,7 @@ impl State { #[cfg(feature = "live")] Id::LiveShare => enabled( !modal + && cfg!(not(target_arch = "wasm32")) && self.notebook().is_some_and(|library| { library.catalog().is_some() && library.joined.is_none() }), diff --git a/crates/snowbound/src/live.rs b/crates/snowbound/src/live.rs index 95adeb6cffe5c5f7afa08f9483c0327c43a0a217..83330811aa40f27925092f8e240638ef816fcd99 100644 --- a/crates/snowbound/src/live.rs +++ b/crates/snowbound/src/live.rs @@ -147,7 +147,7 @@ fn keep(file: &Path, value: &impl serde::Serialize) -> io::Result<()> { let folder = file.parent().unwrap_or(Path::new(".")); notebook::fs::create_dir_all(folder)?; let partial = file.with_extension("partial"); - let mut options = std::fs::OpenOptions::new(); + let mut options = notebook::fs::OpenOptions::new(); options.write(true).create(true).truncate(true); #[cfg(unix)] std::os::unix::fs::OpenOptionsExt::mode(&mut options, 0o600); @@ -217,8 +217,9 @@ impl Joined { ) }) .map_err(|error| refused(&error))?; + #[cfg(not(target_arch = "wasm32"))] if !share.listed { - let since = std::time::Instant::now(); + let since = web_time::Instant::now(); while guest.host().is_none() && since.elapsed() < FIRST_LISTING { std::thread::sleep(std::time::Duration::from_millis(50)); } @@ -257,17 +258,17 @@ impl Joined { } /// Meets whoever shares `code` (and `password`): the share it welcomes this computer to. -pub(crate) fn join( +pub(crate) async fn join( code: &str, password: &str, waiting: impl Fn(bool) -> bool, ) -> Result { let me = hello().map_err(|_| live::share::Refusal::Unreachable(live::Trouble::Other))?; - live::share::join_while(me, code, password, reach(), relay().as_deref(), waiting) + live::share::join_while_async(me, code, password, reach(), relay().as_deref(), waiting).await } /// Keeps `welcome` as the notebook this computer joined: its location. -pub(crate) fn joined(cache: &Path, welcome: live::wire::Welcome) -> io::Result { +pub(crate) async fn joined(cache: &Path, welcome: live::wire::Welcome) -> io::Result { let location = live::share::location(&welcome.share); let file = kept(cache, JOINED); let mut shares: BTreeMap = read_kept(&file); @@ -276,20 +277,32 @@ pub(crate) fn joined(cache: &Path, welcome: live::wire::Welcome) -> io::Result, Option)>, ); @@ -319,9 +332,9 @@ pub(crate) struct Peers { /// Each peer's picture as drawn, cut to a circle, decoded once. pictures: HashMap<[u8; 16], Option>, /// Where each peer's caret was last seen, and when it got there. - moved: HashMap<[u8; 16], (Caret, std::time::Instant)>, + moved: HashMap<[u8; 16], (Caret, web_time::Instant)>, /// Where each peer was last seen, and when they last moved, for ordering the avatars. - active: HashMap<[u8; 16], (Option, std::time::Instant)>, + active: HashMap<[u8; 16], (Option, web_time::Instant)>, /// The peer an avatar's click goes to, since when, and the place already asked to open. going: Option, pub(crate) share: Option, @@ -650,7 +663,7 @@ impl State { if peers.is_empty() { return; } - let now = std::time::Instant::now(); + let now = web_time::Instant::now(); for peer in &peers { let id = peer.hello.peer; let seen = self @@ -921,7 +934,7 @@ impl State { (y * viewport.scale + viewport.origin[1]) / scale, ] }; - let now = std::time::Instant::now(); + let now = web_time::Instant::now(); let (mut placed, mut above, mut below) = (Vec::new(), Vec::new(), Vec::new()); for peer in self.connected() { let Some(caret) = (peer.presence.as_ref()) diff --git a/crates/snowbound/src/main.rs b/crates/snowbound/src/main.rs index 9b704eba4b5dbcd248bcc017c551d8faacfadbdf..edde70e506d0d12c4f35f726697f54fbd094efb7 100644 --- a/crates/snowbound/src/main.rs +++ b/crates/snowbound/src/main.rs @@ -4247,6 +4247,23 @@ impl State { .filter(|(_, paths)| !paths.is_empty()) .collect(); for (library, paths) in &reported { + #[cfg(all(feature = "live", target_arch = "wasm32"))] + if library.joined.is_some() && paths.iter().any(|path| !path.ends_with(".one")) { + let refreshed = Arc::new(library.with(library.reopen()?)); + let shown = self + .session + .as_ref() + .filter(|session| session.library.location == library.location) + .and_then(|session| { + let path = &session.tabs[session.tab].path; + if refreshed.contains(path) { + Some(path.clone()) + } else { + refreshed.first_section() + } + }); + self.adopt(refreshed, shown.as_deref())?; + } self.sections_changed(library, paths.clone()); } let changed: Vec = (reported.iter()) @@ -6184,6 +6201,22 @@ fn spawn(work: impl FnOnce() + Send + 'static) { platform::defer(work); } +/// Runs asynchronous work beside the frame, creating the future on its executor. +#[cfg(feature = "live")] +fn spawn_async + 'static>(work: impl FnOnce() -> F + Send + 'static) { + RUNNING.fetch_add(1, Ordering::Relaxed); + #[cfg(not(target_arch = "wasm32"))] + std::thread::spawn(move || { + pollster::block_on(work()); + RUNNING.fetch_sub(1, Ordering::Release); + }); + #[cfg(target_arch = "wasm32")] + wasm_bindgen_futures::spawn_local(async move { + work().await; + RUNNING.fetch_sub(1, Ordering::Release); + }); +} + /// FILETIME now: when an edit happened, which its modification times record. fn filetime() -> u64 { let unix = web_time::SystemTime::now() diff --git a/crates/snowbound/src/menus.rs b/crates/snowbound/src/menus.rs index 56bc0877cc5728ac471afcbaec5ee1c24e0f0282..b226bf0e5a51cbe0cbdc18bc4164bcdae8b98b1e 100644 --- a/crates/snowbound/src/menus.rs +++ b/crates/snowbound/src/menus.rs @@ -601,7 +601,7 @@ impl State { } actions.extend([ item(Action::CopyLink, "Copy Link to Notebook", false, true), - #[cfg(feature = "live")] + #[cfg(all(feature = "live", not(target_arch = "wasm32")))] item( Action::LiveShare, commands::command(commands::Id::LiveShare).title, diff --git a/crates/snowbound/src/share.rs b/crates/snowbound/src/share.rs index 97933b38a88f644891307e17450d7a9bba3613bf..f212c8b1c120f678d695b4f921ba32c47ab5167a 100644 --- a/crates/snowbound/src/share.rs +++ b/crates/snowbound/src/share.rs @@ -6,7 +6,6 @@ use accesskit::Role; use notebook::live::{ Relayed, Trouble, share::{self, Refusal, Sharing}, - wire::Welcome, }; use std::sync::{ Arc, @@ -59,7 +58,7 @@ pub(crate) struct JoinDialog { code: String, password: String, status: Status, - replies: crate::live::Channel>, + replies: crate::live::Channel>, alive: Arc, approval: Arc, } @@ -607,21 +606,14 @@ impl State { let mut welcomed = None; for reply in dialog.replies.1.try_iter() { match reply { - Ok(welcome) => welcomed = Some(welcome), - Err(refused) => dialog.status = Status::Failed(refusal(&refused)), + Ok(location) => welcomed = Some(location), + Err(message) => dialog.status = Status::Failed(message), } } - if let Some(welcome) = welcomed { - match crate::live::joined(&self.cache, welcome) { - Ok(location) => { - self.ui.close_popup(join_id()); - self.open_notebook(location, None); - return; - } - Err(error) => { - dialog.status = Status::Failed(format!("Couldn’t keep the code: {error}")); - } - } + if let Some(location) = welcomed { + self.ui.close_popup(join_id()); + self.open_notebook(location, None); + return; } let ui = &mut self.ui; let owners = [join_id(), code_field(), join_password()]; @@ -679,15 +671,27 @@ impl State { dialog.status = Status::Waiting; let (replies, redraw) = (dialog.replies.0.clone(), self.redraw.clone()); let password = dialog.password.clone(); + let cache = self.cache.clone(); let (alive, approval) = (Arc::clone(&dialog.alive), Arc::clone(&dialog.approval)); - crate::spawn(move || { + crate::spawn_async(move || async move { let reply = crate::live::join(&code, &password, |waiting| { if approval.swap(waiting, Ordering::AcqRel) != waiting { redraw.wake_by_ref(); } alive.load(Ordering::Acquire) - }); + }) + .await; + approval.store(false, Ordering::Release); + redraw.wake_by_ref(); + let reply = match reply { + Ok(welcome) => { + crate::live::joined(&cache, welcome).await.map_err(|error| { + format!("Couldn’t open the shared notebook: {error}") + }) + } + Err(error) => Err(refusal(&error)), + }; let _ = replies.send(reply); redraw.wake(); }); diff --git a/crates/snowbound/src/sync.rs b/crates/snowbound/src/sync.rs index df25f9c7a7b93f2c17023d36affa0150117aece4..e70f9b9e775e96c65c9ca91a7249395724d84872 100644 --- a/crates/snowbound/src/sync.rs +++ b/crates/snowbound/src/sync.rs @@ -265,8 +265,7 @@ struct Facts<'a> { changes_listed: bool, /// Whether the open section's file can be shown in the file manager. local: bool, - /// The computer sharing this notebook by Live Share stopped sharing it. - stopped: bool, + ended: Option<&'a str>, } /// What the reader picked in the popup. @@ -541,9 +540,9 @@ fn build(ui: &mut Ui, facts: &Facts) -> Picked { let conflict; // A problem's advice stays while working offline, so turning it on moves nothing. let advice = match state { - SyncState::NotConnected if facts.stopped => Some(format!( - "The person sharing this notebook stopped sharing it. Your changes stay on {THIS}." - )), + SyncState::NotConnected if facts.ended.is_some() => facts + .ended + .map(|ended| format!("{ended} Your changes stay on {THIS}.")), SyncState::NotConnected => Some(match sync.queued { 0 => format!("Can’t reach {host}. Sync continues when it’s back."), _ => format!("Can’t reach {host}. {waiting} will sync when it’s back."), @@ -1089,12 +1088,19 @@ impl State { location: &library.location, place: library.place(), #[cfg(feature = "live")] - stopped: library + ended: library .joined .as_ref() - .is_some_and(|joined| joined.guest.stopped()), + .and_then(|joined| match joined.guest.ended()? { + notebook::live::share::Ended::Stopped => { + Some("The person sharing this notebook stopped sharing it.") + } + notebook::live::share::Ended::Removed => Some( + "This device was removed. Ask the person sharing for a new link or code.", + ), + }), #[cfg(not(feature = "live"))] - stopped: false, + ended: None, notice: library.notice.as_deref(), sections: sections(&library, session), conflicts: session @@ -1212,7 +1218,7 @@ mod tests { update: &IDLE, changes_listed: false, local: true, - stopped: false, + ended: None, } } diff --git a/crates/snowbound/src/web.rs b/crates/snowbound/src/web.rs index 66f39afd4f000f8dd6b325462c8346a8a9a27e88..2406947a2f16414592c8462dad99191f1e54c70c 100644 --- a/crates/snowbound/src/web.rs +++ b/crates/snowbound/src/web.rs @@ -1185,6 +1185,18 @@ async fn open( None => state.create_notebook(own.clone())?, } } + #[cfg(feature = "live")] + if let Some(code) = web_sys::window() + .and_then(|window| window.location().search().ok()) + .and_then(|query| { + query + .trim_start_matches('?') + .split('&') + .find_map(|pair| pair.strip_prefix("join=").map(str::to_owned)) + }) + { + state.join_link(code); + } Ok(state) } diff --git a/crates/snowbound/web/glue.js b/crates/snowbound/web/glue.js index 7575cbc566308540b6df2767cfd651ebbc8a798f..ffc726c22da8381c402eb080e9f5301ee21a5ffe 100644 --- a/crates/snowbound/web/glue.js +++ b/crates/snowbound/web/glue.js @@ -23,6 +23,7 @@ function storageWorker() { const handles = new Map(); let root; let queue = Promise.resolve(); + const pending = []; const names = (path) => path.split("/").filter(Boolean); const folder = async (parts) => { @@ -47,7 +48,14 @@ function storageWorker() { }; const write = async (path, length, ranges) => { const open = await handle(path); - for (const [offset, bytes] of ranges) open.write(bytes, { at: offset }); + for (const [offset, bytes] of ranges) { + let written = 0; + while (written < bytes.length) { + const count = open.write(bytes.subarray(written), { at: offset + written }); + if (!count) throw new Error("Browser storage could not finish writing"); + written += count; + } + } open.truncate(length); open.flush(); }; @@ -103,25 +111,34 @@ function storageWorker() { }); }; + const flush = async () => { + while (pending.length) { + for (const [path, length, ranges] of pending[0]) + if (length === undefined) await remove(path); + else if (length === null) await folder(names(path)); + else await write(path, length, ranges); + pending.shift(); + } + }; + onmessage = ({ data }) => { - queue = queue - .then(async () => { + queue = queue.then(async () => { + try { if (data.kind === "load") { root = await navigator.storage.getDirectory(); let files = await list(root, "", []); - if (!files.length) { - await migrate(); - files = await list(root, "", []); - } - postMessage(files, files.flatMap(([, bytes]) => (bytes ? [bytes.buffer] : []))); - } else if (data.kind === "settle") postMessage(null); - else - for (const [path, length, ranges] of data.changes) - if (length === undefined) await remove(path); - else if (length === null) await folder(names(path)); - else await write(path, length, ranges); - }) - .catch((error) => console.error("Keeping files", error)); + if (!files.length) { await migrate(); files = await list(root, "", []); } + postMessage(files, files.flatMap(([, bytes]) => bytes ? [bytes.buffer] : [])); + } else { + if (data.kind === "store") pending.push(data.changes); + await flush(); + if (data.kind === "settle") postMessage({ error: null }); + } + } catch (error) { + if (data.kind === "settle" || data.kind === "load") postMessage({ error: String(error) }); + else console.error("Keeping files", error); + } + }); }; } @@ -132,15 +149,25 @@ export function loadFiles() { storage.postMessage({ kind: "load" }); return new Promise((resolve, reject) => { storage.onmessage = ({ data }) => { - storage.onmessage = () => settling.shift()?.(); - resolve(data); + storage.onmessage = ({ data }) => { + const waiting = settling.shift(); + if (data.error) waiting?.reject(new Error(data.error)); + else waiting?.resolve(); + }; + if (data.error) reject(new Error(data.error)); + else resolve(data); + }; + storage.onerror = (event) => { + storageError = new Error(`The storage worker failed: ${event.message}`); + reject(storageError); + for (const waiting of settling.splice(0)) waiting.reject(storageError); }; - storage.onerror = (error) => reject(new Error(`The storage worker failed: ${error.message}`)); }); } // Callers waiting for the storage worker to finish what it was given. const settling = []; +let storageError; /** Writes `[path]` (removed), `[path, null]` (a folder) and `[path, length, ranges]` entries; * 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) { /** Resolves once every file handed over so far is written. */ function settled() { - const kept = new Promise((resolve) => { - settling.push(resolve); + if (storageError) return Promise.reject(storageError); + const kept = new Promise((resolve, reject) => { + settling.push({ resolve, reject }); storage.postMessage({ kind: "settle" }); }); return Promise.all([kept, writing]); @@ -597,6 +625,7 @@ export function adoptCanvas(fresh) { /** Wires the page's canvas, text area and file input to `module`'s exports. */ export function attach(module) { wasm = module; + globalThis.snowboundFlushStorage = () => { wasm.flush(); return settled(); }; canvas = document.getElementById("page"); input = document.getElementById("input"); picker = document.getElementById("files"); diff --git a/crates/snowbound/web/index.html b/crates/snowbound/web/index.html index daaa0d44dc106d123d335e67e91fe92241666c1b..2c181528ab9aa18c8d1042c35e499c63c4b8214e 100644 --- a/crates/snowbound/web/index.html +++ b/crates/snowbound/web/index.html @@ -205,6 +205,11 @@ ]); status.textContent = "Starting Snowbound"; await snowbound.start(snowbound, fonts, dictionaries); + const opened = new URL(location.href); + if (opened.searchParams.has("join")) { + opened.searchParams.delete("join"); + history.replaceState(null, "", opened); + } sessionStorage.removeItem("snowbound.refetched"); } catch (error) { // A module and JavaScript from different builds can't link: an old cached copy. Once, diff --git a/tools/ci.py b/tools/ci.py index d7d25e41d7174423506a59f3d1b40772ef42e7eb..739dd9d93a8db76c91767af07014eefd07864893 100644 --- a/tools/ci.py +++ b/tools/ci.py @@ -80,7 +80,7 @@ def lanes(): result.append(dict( name='web', minutes=20, packages=['snowbound'], paths=('crates/snowbound/web/', 'tools/release_web.py'), # Clippy builds it; the module itself is linked by release_web.py. - commands=[cargo('clippy', '-p', 'snowbound', '--target', 'wasm32-unknown-unknown', '--no-default-features', '--features', 'wgpu', '--', '-D', 'warnings')], + commands=[cargo('clippy', '-p', 'snowbound', '--target', 'wasm32-unknown-unknown', '--no-default-features', '--features', 'wgpu,live', '--', '-D', 'warnings')], environment={**(wasm or {}), 'CARGO_TARGET_DIR': str(TARGET / 'wasm')}, missing='needs the wasm32-unknown-unknown target' if 'wasm32-unknown-unknown' not in targets.split() else None if wasm else 'needs nix for a clang that builds for wasm32')) diff --git a/tools/release_web.py b/tools/release_web.py index 4768ed6db3013720bae5a548b7342e863c14fc4c..915dcac92a99419d0f1131c886fa1160780bcecf 100644 --- a/tools/release_web.py +++ b/tools/release_web.py @@ -100,7 +100,7 @@ def build(out): """Writes index.html, the module, its JavaScript and the fonts to `out`, the module knowing when it was built.""" built = str(int(time.time())) - run(['cargo', 'build', '--locked', '-p', 'snowbound', '--release', '--target', 'wasm32-unknown-unknown', '--no-default-features', '--features', 'wgpu', + run(['cargo', 'build', '--locked', '-p', 'snowbound', '--release', '--target', 'wasm32-unknown-unknown', '--no-default-features', '--features', 'wgpu,live', *PROFILE], env={**environment(), 'SNOWBOUND_WEB_BUILD': built}) if out.exists(): shutil.rmtree(out)