diff --git a/arc/sync.md b/arc/sync.md index 76ae7fc71e2e8f9876049ae5c2e56e3c881d4e67..0d53c438c69479bb488f8060cbce3eba52042830 100644 --- a/arc/sync.md +++ b/arc/sync.md @@ -167,6 +167,16 @@ Working offline on purpose, as OneNote's Work Offline does, is the same state chosen: the sync thread stops stepping until Sync Now or until the user works online again. +## Sections that aren't open + +OneNote keeps every section of an open notebook in sync, not just the one on screen: it +reads a section another client changed within seconds and publishes a closed section's +offline edits the moment the share is back. Snowbound's `session::Background` does the +same with one thread per notebook. Each round reads every section file's stamp. Only a +section whose replica has edits waiting, or whose file moved past the replica's base, +has its replica opened for the usual sync steps, and it is closed again afterwards. The +open section is left to its own session. + ## Notebook structure Sections, groups and the notebook's own colour live in the `.onetoc2` files, diff --git a/crates/notebook/README.md b/crates/notebook/README.md index 588e4860384735a2ff38b03e5616b67f99b7abab..61ba64a1660588da42bd6f3846b4771a472355b0 100644 --- a/crates/notebook/README.md +++ b/crates/notebook/README.md @@ -64,6 +64,17 @@ edit names is the host's: the app passes the account's full name. and how many edits wait for it. `set_offline(true)` works offline as OneNote does: the worker stops connecting and edits queue until `wake()` (Sync Now) or `set_offline(false)`. +`session::Background` keeps the sections no session holds in sync, as OneNote 2010 keeps +every section of an open notebook: `Notebook::background(interval, notify)` for a mounted +notebook, `Background::smb(root, limit, interval, connect, notify)` on a share, then +`watch(notebook.replicas())`. Each round reads every watched file's stamp; a section whose +replica has edits waiting, or whose file moved past the replica's base, has the replica +opened for the synchronization steps that publish or rebase it and closed again. A replica +a session holds is skipped (`Error::busy`), so opening a section may wait out one step. +`status()` gives each section's `SyncStatus`, `changed()` the sections another client +changed, and `wake` and `set_offline` follow Sync Now and Work Offline. +`Notebook::replica_path(path)` names a section's replica for either kind of notebook. + `Section::resume(file, replica, notify)` starts from an owned `Replica` without consulting the remote file; with the `smb` feature, `Section::resume_smb(path, replica, limit, connect, notify)` binds a share-relative path, `connect` running on diff --git a/crates/notebook/examples/scratch_bgsync.rs b/crates/notebook/examples/scratch_bgsync.rs deleted file mode 100644 index 5608de2f4f315edb01272dbf7c3fc75e6fc08283..0000000000000000000000000000000000000000 --- a/crates/notebook/examples/scratch_bgsync.rs +++ /dev/null @@ -1,42 +0,0 @@ -//! Scratch lab writer (delete before reporting): appends TEXT to the first body paragraph of -//! the first page of a share-relative section. `scratch_bgsync ADDRESS PATH TEXT` -use notebook::{ - Replica, SmbRemote, - smb::{Client, Credentials}, -}; -use onestore::op::{Edit, Op, PageOp}; -use std::time::Duration; - -fn main() -> Result<(), Box> { - let args: Vec<_> = std::env::args().skip(1).collect(); - let [address, path, with] = &args[..] else { - return Err("scratch_bgsync ADDRESS PATH TEXT".into()); - }; - let connect = || Client::connect(address, "agent", Credentials::default(), Duration::from_secs(5)); - let source = connect()?.read_storage(path, 1 << 26)?; - let directory = tempfile::tempdir()?; - let replica = Replica::create(directory.path().join("cache.sqlite"), &source)?; - let space = replica.pages()?[0].0; - let page = replica.page(space)?; - let (text, len) = page - .objects - .iter() - .find_map(|object| match object { - onestore::page::PageObject::Outline(outline) => outline - .paragraphs - .iter() - .find_map(|p| p.text().map(|t| (t.id, t.text.text().encode_utf16().count() as u32))), - _ => None, - }) - .ok_or("no text")?; - replica.apply( - "Lab", - Edit { - at: 134_000_000_000_000_000, - ops: vec![Op::Page { space, op: PageOp::Text { text, range: len..len, with: with.clone() } }], - }, - )?; - let mut remote = SmbRemote::new(connect()?, path.clone(), 1 << 26); - println!("{:?}", replica.sync_once(&mut remote)?.edit); - Ok(()) -} diff --git a/crates/notebook/src/background.rs b/crates/notebook/src/background.rs new file mode 100644 index 0000000000000000000000000000000000000000..1c327727876c1099208f13396682bca462be5fa8 --- /dev/null +++ b/crates/notebook/src/background.rs @@ -0,0 +1,326 @@ +//! The sections of an open notebook that no session holds, kept in sync as OneNote 2010 +//! keeps every section of an open notebook: queued edits publish and other clients' changes +//! are noticed without the section being open. + +use crate::{ + EditStatus, Error, Remote, Replica, Result, + session::{SyncStatus, reached}, + worker::Signal, +}; +use onestore::Stamp; +use std::{ + io, + path::{Path, PathBuf}, + sync::{Arc, Mutex, atomic::Ordering}, + thread, + time::Duration, +}; + +/// Polls each watched section file's stamp every interval. A section whose replica has +/// edits waiting, or whose file moved past the replica's base, has its replica opened for +/// the synchronization steps that publish or rebase it, then closed again; any other costs +/// one stamp read. A replica a session holds is left to that session's worker. Dropping +/// requests cancellation without waiting for the step in flight. +pub struct Background { + signal: Arc, + watched: Arc>, +} + +#[derive(Default)] +struct Watched { + sections: Vec, + /// Sections whose file changed since `changed` was last asked. + changed: Vec, +} + +struct Watch { + path: String, + replica: Option, + /// The file's stamp when last reached. + stamp: Option, + status: SyncStatus, +} + +impl Background { + pub(crate) fn start( + interval: Duration, + mut connect: impl FnMut() -> io::Result + Send + 'static, + notify: impl Fn() + Send + 'static, + ) -> Result + where + R: Remote, + B: FnMut(&str) -> R, + { + let (signal, receiver) = Signal::new(); + let watched = Arc::new(Mutex::new(Watched::default())); + let (shared, sections) = (Arc::clone(&signal), Arc::clone(&watched)); + thread::Builder::new() + .name("onestore-background".into()) + .spawn(move || { + let mut bound: Option = None; + while !shared.stopped.load(Ordering::Acquire) { + if shared.offline.load(Ordering::Acquire) + && !shared.requested.load(Ordering::Acquire) + { + bound = None; + let _ = receiver.recv(); + continue; + } + shared.requested.store(false, Ordering::Release); + let round: Vec<_> = match sections.lock() { + Ok(watched) => watched + .sections + .iter() + .map(|watch| { + ( + watch.path.clone(), + watch.replica.clone(), + watch.stamp.clone(), + ) + }) + .collect(), + Err(_) => return, + }; + let mut news = false; + for (path, replica, seen) in round { + if shared.stopped.load(Ordering::Acquire) { + return; + } + let bind = match &mut bound { + Some(bind) => bind, + None => match connect() { + Ok(bind) => bound.insert(bind), + Err(error) => { + // Nothing on the share is reachable this round. + let Ok(mut watched) = sections.lock() else { + return; + }; + for watch in &mut watched.sections { + news |= watch.fail( + &Error::RemoteIo(io::Error::new( + error.kind(), + error.to_string(), + )), + None, + ); + } + break; + } + }, + }; + let (queued, outcome) = + step(&mut bind(&path), replica.as_deref(), seen.as_ref()); + if matches!(outcome, Err(Error::RemoteIo(_) | Error::Remote(_))) { + bound = None; + } + let Ok(mut watched) = sections.lock() else { + return; + }; + let Watched { sections, changed } = &mut *watched; + let Some(watch) = sections.iter_mut().find(|watch| watch.path == path) + else { + continue; + }; + news |= match outcome { + Ok((stamp, moved)) => { + watch.stamp = Some(stamp); + if moved { + changed.push(path); + } + let before = summary(&watch.status); + if let Some(queued) = queued { + watch.status = SyncStatus { + synced: Some(crate::now()), + error: None, + queued, + }; + } + moved || before != summary(&watch.status) + } + Err(error) => watch.fail(&error, queued), + }; + } + if news { + notify(); + } + let _ = receiver.recv_timeout(interval); + } + })?; + Ok(Self { signal, watched }) + } + + /// Keeps the sections of a notebook on a share in sync while they are not open; + /// `connect` runs again after a transport failure. Paths are relative to `root`. + #[cfg(feature = "smb")] + pub fn smb( + root: &str, + limit: usize, + interval: Duration, + mut connect: impl FnMut() -> io::Result + Send + 'static, + notify: impl Fn() + Send + 'static, + ) -> Result { + let root = root.replace('\\', "/"); + Self::start( + interval, + move || { + let client = Arc::new(connect()?); + let root = root.clone(); + Ok(move |path: &str| { + let file = match root.as_str() { + "" => path.to_owned(), + root => format!("{root}/{path}"), + }; + crate::SmbRemote::new(Arc::clone(&client), file, limit) + }) + }, + notify, + ) + } + + /// Watches these sections, by catalog path and replica (`Notebook::replicas`), from + /// the next round, which starts now. + pub fn watch(&self, sections: Vec<(String, Option)>) { + if let Ok(mut watched) = self.watched.lock() { + let mut previous = std::mem::take(&mut watched.sections); + watched.sections = sections + .into_iter() + .map( + |(path, replica)| match previous.iter().position(|watch| watch.path == path) { + Some(index) => Watch { + replica, + ..previous.swap_remove(index) + }, + None => Watch { + path, + replica, + stamp: None, + status: SyncStatus { + synced: None, + error: None, + queued: 0, + }, + }, + }, + ) + .collect(); + } + self.signal.wake(); + } + + /// Each watched section's status as its last round left it, in watch order. A section a + /// session holds keeps the status it had before. + pub fn status(&self) -> Vec<(String, SyncStatus)> { + self.watched.lock().map_or_else( + |_| Vec::new(), + |watched| { + watched + .sections + .iter() + .map(|watch| { + let status = &watch.status; + ( + watch.path.clone(), + SyncStatus { + synced: status.synced, + error: status + .error + .as_ref() + .map(|error| io::Error::new(error.kind(), error.to_string())), + queued: status.queued, + }, + ) + }) + .collect() + }, + ) + } + + /// Sections whose file changed since the last call, by catalog path. + pub fn changed(&self) -> Vec { + self.watched + .lock() + .map(|mut watched| std::mem::take(&mut watched.changed)) + .unwrap_or_default() + } + + /// Runs a round now, working offline included (Sync Now). + pub fn wake(&self) { + self.signal.requested.store(true, Ordering::Release); + self.signal.wake(); + } + + /// Working offline, no round runs until `wake`, or until working online again. + pub fn set_offline(&self, offline: bool) { + self.signal.offline.store(offline, Ordering::Release); + self.signal.wake(); + } +} + +impl Drop for Background { + fn drop(&mut self) { + self.signal.stopped.store(true, Ordering::Release); + self.signal.wake(); + } +} + +impl Watch { + /// Records a failed step and what it found waiting, answering whether the status shown + /// changes. + fn fail(&mut self, error: &Error, queued: Option) -> bool { + let before = summary(&self.status); + self.status.error = Some(reached(error)); + self.status.queued = queued.unwrap_or(self.status.queued); + before != summary(&self.status) + } +} + +/// What the host shows of a status. +fn summary(status: &SyncStatus) -> (bool, Option, u64) { + ( + status.synced.is_some(), + status.error.as_ref().map(io::Error::kind), + status.queued, + ) +} + +/// One section's step: how many of its edits wait (`None` while unknown, as while a session +/// holds its replica), then the file's stamp now and whether it changed since `seen`. +fn step( + remote: &mut R, + replica: Option<&Path>, + seen: Option<&Stamp>, +) -> (Option, Result<(Stamp, bool)>) { + let stamp = match remote.stamp() { + Ok(stamp) => stamp, + Err(error) => return (None, Err(Error::RemoteIo(error))), + }; + let moved = seen.is_some_and(|seen| *seen != stamp); + let Some(replica) = replica.filter(|replica| replica.exists()) else { + return (Some(0), Ok((stamp, moved))); + }; + match crate::peek(replica) { + Ok((base, 0)) if base == stamp => return (Some(0), Ok((stamp, moved))), + Err(error) if error.busy() => return (None, Ok((stamp, false))), + _ => {} + } + let replica = match Replica::open(replica) { + Ok(replica) => replica, + Err(error) if error.busy() => return (None, Ok((stamp, false))), + Err(error) => return (None, Err(error)), + }; + let synced = (|| { + let mut changed = moved; + loop { + let synced = replica.sync_once(remote)?; + changed |= !synced.changed.is_empty(); + if !matches!(synced.edit, Some((_, EditStatus::Published { .. }))) { + return Ok((remote.stamp().map_err(Error::RemoteIo)?, changed)); + } + } + })(); + let queued = replica + .recovery_summary() + .ok() + .map(|summary| summary.queued_edits); + (queued, synced) +} diff --git a/crates/notebook/src/lib.rs b/crates/notebook/src/lib.rs index fe6db93334d5e794ee2dea17df25dd97b65d941e..7c7a39fa46e31bc7f50dde9ac6e05557e25ae7a3 100644 --- a/crates/notebook/src/lib.rs +++ b/crates/notebook/src/lib.rs @@ -20,6 +20,7 @@ use std::{ }; mod assets; +mod background; mod base; mod merge; mod migrate; @@ -60,6 +61,14 @@ pub enum Error { Protected(#[from] onestore::protected::Error), } +impl Error { + /// Whether the replica is open elsewhere, as in another section session of this process. + pub fn busy(&self) -> bool { + matches!(self, Self::Database(rusqlite::Error::SqliteFailure(error, _)) + if error.code == rusqlite::ErrorCode::DatabaseBusy) + } +} + type Result = std::result::Result; const APPLICATION_ID: u32 = 0x4f4e454f; @@ -307,6 +316,21 @@ fn validate(source: &[u8]) -> Result { Ok(onestore::Section::open(&arena, source.to_vec())?.root()) } +/// A closed cache's base stamp and how many edits wait, without opening the replica. +fn peek(path: &Path) -> Result<(onestore::Stamp, u64)> { + let connection = cache_connection(path)?; + let application: u32 = + connection.pragma_query_value(None, "application_id", |row| row.get(0))?; + let version: u32 = connection.pragma_query_value(None, "user_version", |row| row.get(0))?; + if application != APPLICATION_ID || version != schema::VERSION { + return Err(io::Error::from(io::ErrorKind::InvalidData).into()); + } + let base = base::stamp(&connection, base::Image::Base)? + .ok_or_else(|| io::Error::new(io::ErrorKind::InvalidData, "The cache has no base image"))?; + let queued: i64 = connection.query_row("SELECT count(*) FROM edits", [], |row| row.get(0))?; + Ok((base, unsigned(queued)?)) +} + fn pending(connection: &Connection) -> Result> { queue::load(connection, None) } diff --git a/crates/notebook/src/session.rs b/crates/notebook/src/session.rs index 977238d3233776afdaf4e444e6d646197328ee0a..f1973f933a74a0703eb65d394e5c3ac75ef665eb 100644 --- a/crates/notebook/src/session.rs +++ b/crates/notebook/src/session.rs @@ -1,6 +1,7 @@ //! The application's view of a notebook: sections opened through a local replica that //! publishes their edits to the section file in the background. +pub use crate::background::Background; use crate::{ EditStatus, Error, PendingEdit, Remote, Replica, Resolution, Result, SyncWorker, discover, }; @@ -881,6 +882,61 @@ impl Notebook { Ok(()) } + /// Where the replica of the section at catalog `path` lives: named by the section's + /// document identity in a mounted notebook (`Section::open`), by its file identity under + /// `smb` on a share. It exists once the section has been opened. + pub fn replica_path(&self, path: &str) -> Result { + let section = self.section_path(path)?; + Ok(match &self.root { + Some(_) => { + let image = self.storage.read(§ion.path)?; + let root = RevisionIndex::parse(&Store::parse(&image)?)?.root; + replica_file(&self.cache, &root.guid) + } + None => replica_file(&self.cache.join("smb"), §ion.file_id), + }) + } + + /// Every readable section's catalog path and replica, for `Background::watch`. + pub fn replicas(&self) -> Vec<(String, Option)> { + let mut sections = Vec::new(); + let mut folders = vec![&self.catalog]; + while let Some(folder) = folders.pop() { + for section in &folder.sections { + if matches!(section.state, discover::SectionState::Readable { .. }) { + let replica = self.replica_path(§ion.path).ok(); + sections.push((section.path.clone(), replica)); + } + } + folders.extend(folder.groups.iter().rev()); + } + sections + } + + /// Keeps the sections of a mounted notebook in sync while they are not open + /// (`Background`), polling each file every `interval` once watched. + pub fn background( + &self, + interval: Duration, + notify: impl Fn() + Send + 'static, + ) -> Result { + let Some(root) = self.root.clone() else { + return Err(io::Error::new( + io::ErrorKind::Unsupported, + "Sections on a share sync through Background::smb", + ) + .into()); + }; + Background::start( + interval, + move || { + let root = root.clone(); + Ok(move |path: &str| FileRemote(root.join(path))) + }, + notify, + ) + } + /// Opens a section of a mounted notebook by its catalog path. pub fn section(&self, path: &str, notify: impl Fn() + Send + 'static) -> Result
{ self.section_with(path, |file| Ok(FileRemote(file.to_owned())), notify) @@ -1010,6 +1066,21 @@ fn catalog_path(folder: &str, name: &str) -> String { } } +/// The replica in `cache` of the section `identity` names. +fn replica_file(cache: &Path, identity: &[u8; 16]) -> PathBuf { + let name: String = identity.iter().map(|byte| format!("{byte:02x}")).collect(); + cache.join(format!("{name}.sqlite")) +} + +/// Why a synchronization step did not reach the section file, as the host shows it. +pub(crate) fn reached(error: &Error) -> io::Error { + match error { + Error::RemoteIo(error) => io::Error::new(error.kind(), error.to_string()), + Error::Remote(error) => io::Error::new(error.error.kind(), error.to_string()), + error => io::Error::other(error.to_string()), + } +} + /// What happened to the section since the last poll. #[derive(Debug)] pub enum Event { @@ -1079,13 +1150,8 @@ impl Section { let source = connect(&file)?.read()?; let store = Store::parse(&source)?; let identity = RevisionIndex::parse(&store)?.root; - let name: String = identity - .guid - .iter() - .map(|byte| format!("{byte:02x}")) - .collect(); std::fs::create_dir_all(&cache)?; - let cache = cache.as_ref().join(format!("{name}.sqlite")); + let cache = replica_file(cache.as_ref(), &identity.guid); let replica = if cache.exists() { Replica::open(&cache)? } else { @@ -1149,11 +1215,7 @@ impl Section { let (sender, notify) = (sender.clone(), Arc::clone(¬ify)); let observed = Arc::clone(&observed); replica.start_sync(Duration::from_secs(2), connect, move |result| { - let error = result.as_ref().err().map(|error| match error { - Error::RemoteIo(error) => io::Error::new(error.kind(), error.to_string()), - Error::Remote(error) => io::Error::new(error.error.kind(), error.to_string()), - error => io::Error::other(error.to_string()), - }); + let error = result.as_ref().err().map(reached); // Reaching the file as the last attempt did is no news to the host. let mut changed = false; if let Ok(mut observed) = observed.lock() { diff --git a/crates/notebook/src/smb/remote.rs b/crates/notebook/src/smb/remote.rs index b1178eb2115de1f5478754686127724f74713299..6cecf8c55ee45cf4fa020676e17dc307e6e36947 100644 --- a/crates/notebook/src/smb/remote.rs +++ b/crates/notebook/src/smb/remote.rs @@ -1,20 +1,21 @@ use super::Client; use crate::Remote; use onestore::{CommitError, Stamp, Transaction}; -use std::io; +use std::{io, sync::Arc}; /// Binds every reconciliation operation to one share-relative file and read limit. /// Connection loss retires the client; reconnect before subsequent sync attempts. pub struct SmbRemote { - client: Client, + client: Arc, path: String, limit: usize, } impl SmbRemote { - pub fn new(client: Client, path: impl Into, limit: usize) -> Self { + /// A client shared between remotes serves each of their files over one connection. + pub fn new(client: impl Into>, path: impl Into, limit: usize) -> Self { Self { - client, + client: client.into(), path: path.into(), limit, } diff --git a/crates/notebook/src/worker.rs b/crates/notebook/src/worker.rs index a683045309d6e0de2e50bf3ae7f195474ae69422..ed9c9e707f58f5ee8c747cd6df454093f283a5f2 100644 --- a/crates/notebook/src/worker.rs +++ b/crates/notebook/src/worker.rs @@ -12,16 +12,28 @@ use std::{ }; pub(super) struct Signal { - stopped: AtomicBool, - offline: AtomicBool, + pub(crate) stopped: AtomicBool, + pub(crate) offline: AtomicBool, /// FILETIME of the last step or poll that reached the remote; 0 before one has. synced: AtomicU64, - /// Asked for by `SyncWorker::wake`: steps run while offline until one leaves nothing to do. - requested: AtomicBool, + /// Asked for by `wake`: steps run while offline until one leaves nothing to do. + pub(crate) requested: AtomicBool, sender: SyncSender<()>, } impl Signal { + pub(crate) fn new() -> (Arc, mpsc::Receiver<()>) { + let (sender, receiver) = mpsc::sync_channel(1); + let signal = Self { + stopped: AtomicBool::new(false), + offline: AtomicBool::new(false), + synced: AtomicU64::new(0), + requested: AtomicBool::new(false), + sender, + }; + (Arc::new(signal), receiver) + } + pub(super) fn wake(&self) { // One retained notification covers edits that arrive during network I/O. let _ = self.sender.try_send(()); @@ -110,14 +122,7 @@ impl Replica { if owner.upgrade().is_some() { return Err(io::ErrorKind::WouldBlock.into()); } - let (sender, receiver) = mpsc::sync_channel(1); - let signal = Arc::new(Signal { - stopped: AtomicBool::new(false), - offline: AtomicBool::new(false), - synced: AtomicU64::new(0), - requested: AtomicBool::new(false), - sender, - }); + let (signal, receiver) = Signal::new(); let replica = Arc::clone(self); let worker_signal = Arc::clone(&signal); let thread = thread::Builder::new() diff --git a/crates/notebook/tests/session.rs b/crates/notebook/tests/session.rs index a63954015f49167dc63246a4c4ed89b743220d4f..0620ce0a6d3aa30a5a657104612649b32bcbc4eb 100644 --- a/crates/notebook/tests/session.rs +++ b/crates/notebook/tests/session.rs @@ -999,3 +999,129 @@ fn a_missing_section_file_reports_its_error_until_it_returns() { assert!(notified.load(Ordering::SeqCst) > before); section.close().unwrap(); } + +/// A notebook folder holding `First.one` and `Second.one`, and a cache beside it. +fn two_sections(directory: &Path) -> Notebook { + let root = directory.join("Shared"); + std::fs::create_dir(&root).unwrap(); + for name in ["First", "Second"] { + let file = format!("{name}.one"); + std::fs::write( + root.join(&file), + onestore::create_section(&file, name, "Author").unwrap(), + ) + .unwrap(); + } + Notebook::open(&root, directory.join("cache")).unwrap() +} + +fn until(what: &str, mut accept: impl FnMut() -> bool) { + let deadline = Instant::now() + Duration::from_secs(20); + while !accept() { + assert!(Instant::now() < deadline, "{what}"); + std::thread::sleep(Duration::from_millis(20)); + } +} + +#[test] +fn a_closed_sections_queued_edits_publish_in_the_background_when_online() { + let directory = tempfile::tempdir().unwrap(); + let notebook = two_sections(directory.path()); + let file = directory.path().join("Shared/Second.one"); + let background = notebook + .background(Duration::from_millis(20), || {}) + .unwrap(); + background.set_offline(true); + background.watch(notebook.replicas()); + // The round started with the background ends before the edit exists. + std::thread::sleep(Duration::from_millis(200)); + let section = notebook.section("Second.one", || {}).unwrap(); + section.set_offline(true); + let space = section.pages().unwrap()[0].0; + let before = section.page(space).unwrap(); + typed(§ion, space, &before, 0..0, "Closed "); + section.close().unwrap(); + std::thread::sleep(Duration::from_millis(300)); + assert_eq!( + stored_page(&file, space), + before, + "working offline publishes nothing" + ); + background.set_offline(false); + let mut after = edited(&before, "Closed "); + until("the closed section published", || { + let stored = stored_page(&file, space); + after.title = stored.title.clone(); + stored == after + }); + until("the status shows nothing waiting", || { + background.status().iter().any(|(path, status)| { + path == "Second.one" && status.synced.is_some() && status.queued == 0 + }) + }); + drop(background); + let mut reopened = None; + until("the background released the replica", || { + reopened = notebook.section("Second.one", || {}).ok(); + reopened.is_some() + }); + let reopened = reopened.unwrap(); + assert!(reopened.pending().unwrap().is_empty()); + reopened.close().unwrap(); +} + +#[test] +fn a_remote_change_to_a_closed_section_is_noticed_and_rebases_its_replica() { + let directory = tempfile::tempdir().unwrap(); + let notebook = two_sections(directory.path()); + // Second has a replica from being opened once; First has never been opened. + notebook + .section("Second.one", || {}) + .unwrap() + .close() + .unwrap(); + let notified = Arc::new(AtomicUsize::new(0)); + let counter = Arc::clone(¬ified); + let background = notebook + .background(Duration::from_millis(20), move || { + counter.fetch_add(1, Ordering::SeqCst); + }) + .unwrap(); + background.watch(notebook.replicas()); + until("both sections were reached", || { + background + .status() + .iter() + .all(|(_, status)| status.synced.is_some() && status.error.is_none()) + }); + assert!(background.changed().is_empty()); + let mut changed = Vec::new(); + for name in ["First.one", "Second.one"] { + let file = directory.path().join("Shared").join(name); + let bytes = onestore::read_file(&file).unwrap(); + let space = notebook::session::stored_pages(&bytes).unwrap()[0].space; + let native = edited(&model_ops::page_of(&bytes, space), "Native "); + ops::save(&bytes, space, &native) + .unwrap() + .commit_file(&file) + .unwrap(); + changed.push((name, space, native)); + } + let mut noticed = Vec::new(); + until("both changes were noticed", || { + noticed.extend(background.changed()); + changed + .iter() + .all(|(name, ..)| noticed.iter().any(|path| path == name)) + }); + assert!(notified.load(Ordering::SeqCst) > 0); + drop(background); + let (_, space, native) = &changed[1]; + let replica = notebook.replica_path("Second.one").unwrap(); + let mut opened = None; + until("the background released the replica", || { + opened = notebook::Replica::open(&replica).ok(); + opened.is_some() + }); + assert_same(opened.unwrap().page(*space).unwrap(), native); +} diff --git a/crates/snowbound/src/library.rs b/crates/snowbound/src/library.rs index c01fb3cb2456537f486114433aca50ccb4ca1eb4..d5ffe0f22f05b120971698491cc4ef16139f6f1d 100644 --- a/crates/snowbound/src/library.rs +++ b/crates/snowbound/src/library.rs @@ -1,17 +1,17 @@ //! Open notebooks: where each lives and the sections its tabs offer. use notebook::discover::{Folder, SectionState}; -use notebook::session::{Notebook, Section}; +use notebook::session::{Background, Notebook, Section}; use notebook::smb::{Client, Credentials}; use std::{ error::Error, io, path::{Path, PathBuf}, sync::{ - Arc, + Arc, Mutex, OnceLock, atomic::{AtomicBool, Ordering}, }, - time::Duration, + time::{Duration, Instant}, }; /// OneNote's Work Offline, which like OneNote's holds for every notebook. @@ -26,6 +26,23 @@ pub fn set_offline(offline: bool) { OFFLINE.store(offline, Ordering::Relaxed); } +/// Wakes the app when a notebook's closed sections report, once set. +static NOTIFY: OnceLock> = OnceLock::new(); + +pub fn on_background(notify: impl Fn() + Send + Sync + 'static) { + let _ = NOTIFY.set(Box::new(notify)); +} + +/// How often the sections no tab shows are polled; OneNote 2010 read a closed section +/// fifteen seconds after another client changed it. +const BACKGROUND: Duration = Duration::from_secs(15); + +fn notify_background() { + if let Some(notify) = NOTIFY.get() { + notify(); + } +} + /// How long an SMB request may take before the share counts as unreachable. const TIMEOUT: Duration = Duration::from_secs(10); /// The largest section file read whole over SMB. @@ -151,6 +168,9 @@ pub struct Library { /// Why a notebook on a mounted SMB share opened through the mount instead. pub notice: Option, cache: PathBuf, + /// Syncs the notebook's sections while no tab shows them, as OneNote syncs every + /// section of an open notebook. + pub background: Option>, } impl Library { @@ -171,12 +191,12 @@ impl Library { } } } + let notebook = Notebook::open(location, cache); Self { location: location.to_owned(), name: file_name(Path::new(location)), - notebook: Notebook::open(location, cache) - .map(Some) - .map_err(|error| error.to_string()), + background: notebook.as_ref().ok().and_then(local_background), + notebook: notebook.map(Some).map_err(|error| error.to_string()), server: None, notice, cache: cache.to_owned(), @@ -195,6 +215,17 @@ impl Library { let client = server.connect().map_err(|error| error.to_string())?; let notebook = Notebook::open_smb(Arc::new(client), &server.mount.root, cache) .map_err(|error| error.to_string())?; + let connect = Arc::clone(&server); + let background = Background::smb( + &server.mount.root, + LIMIT, + BACKGROUND, + move || connect.connect(), + notify_background, + ) + .map_err(|error| error.to_string())?; + background.set_offline(offline()); + background.watch(notebook.replicas()); Ok(Self { location: location.to_owned(), name: file_name(Path::new(location)), @@ -202,6 +233,7 @@ impl Library { server: Some(server), notice: None, cache: cache.to_owned(), + background: Some(Arc::new(background)), }) } @@ -217,6 +249,9 @@ impl Library { /// This notebook as `notebook`, read again after a change. pub fn with(&self, notebook: Notebook) -> Self { + if let Some(background) = &self.background { + background.watch(notebook.replicas()); + } Self { location: self.location.clone(), name: self.name.clone(), @@ -224,6 +259,7 @@ impl Library { server: self.server.clone(), notice: self.notice.clone(), cache: self.cache.clone(), + background: self.background.clone(), } } @@ -232,6 +268,7 @@ impl Library { Self { location: location.to_owned(), name: file_name(Path::new(location)), + background: local_background(¬ebook), notebook: Ok(Some(notebook)), server: None, notice: None, @@ -248,6 +285,7 @@ impl Library { server: None, notice: None, cache: cache.to_owned(), + background: None, } } @@ -258,31 +296,50 @@ impl Library { path: &str, notify: impl Fn() + Send + 'static, ) -> Result> { - let section = match (&self.notebook, &self.server) { - (Ok(Some(notebook)), Some(server)) => { - let identity = self - .catalog_section(notebook.catalog(), path) - .ok_or("The notebook doesn’t list this section")?; - let file = match server.mount.root.as_str() { - "" => path.to_owned(), - root => format!("{root}/{path}"), - }; - let cache = self.cache.join("smb").join(format!("{identity}.sqlite")); - std::fs::create_dir_all(self.cache.join("smb"))?; - let replica = if cache.exists() { - notebook::Replica::open(&cache)? - } else { - notebook::Replica::create( - &cache, - &server.connect()?.read_storage(&file, LIMIT)?, - )? - }; - let server = Arc::clone(server); - Section::resume_smb(file, replica, LIMIT, move || server.connect(), notify)? + // The background may hold the replica for a step, which takes a network round trip. + let deadline = Instant::now() + TIMEOUT * 3; + let notify = Arc::new(Mutex::new(notify)); + let notifier = || { + let notify = Arc::clone(¬ify); + move || { + if let Ok(notify) = notify.lock() { + notify(); + } + } + }; + let section = loop { + let opened = match (&self.notebook, &self.server) { + (Ok(Some(notebook)), Some(server)) => { + let file = match server.mount.root.as_str() { + "" => path.to_owned(), + root => format!("{root}/{path}"), + }; + let cache = notebook.replica_path(path)?; + std::fs::create_dir_all(self.cache.join("smb"))?; + let replica = if cache.exists() { + notebook::Replica::open(&cache) + } else { + notebook::Replica::create( + &cache, + &server.connect()?.read_storage(&file, LIMIT)?, + ) + }; + replica.and_then(|replica| { + let server = Arc::clone(server); + let connect = move || server.connect(); + Section::resume_smb(file, replica, LIMIT, connect, notifier()) + }) + } + (Ok(Some(notebook)), None) => notebook.section(path, notifier()), + (Ok(None), _) => Section::open(path, &self.cache, notifier()), + (Err(error), _) => return Err(error.clone().into()), + }; + match opened { + Err(error) if error.busy() && Instant::now() < deadline => { + std::thread::sleep(Duration::from_millis(50)); + } + opened => break opened?, } - (Ok(Some(notebook)), None) => notebook.section(path, notify)?, - (Ok(None), _) => Section::open(path, &self.cache, notify)?, - (Err(error), _) => return Err(error.clone().into()), }; section.set_offline(offline()); Ok(section) @@ -311,24 +368,6 @@ impl Library { } } - /// The file identity, in hex, of the section at catalog `path`. - fn catalog_section(&self, catalog: &Folder, path: &str) -> Option { - let mut folders = vec![catalog]; - while let Some(folder) = folders.pop() { - if let Some(section) = folder.sections.iter().find(|section| section.path == path) { - return Some( - section - .file_id - .iter() - .map(|byte| format!("{byte:02x}")) - .collect(), - ); - } - folders.extend(&folder.groups); - } - None - } - /// Names the section at catalog `path` across every open notebook, for remembering /// its pages. pub fn key(&self, path: &str) -> String { @@ -407,6 +446,17 @@ pub fn recycle_bin(path: &str) -> bool { path.rsplit('/').next() == Some("OneNote_RecycleBin") } +/// A mounted notebook's background sync, following Work Offline. +fn local_background(notebook: &Notebook) -> Option> { + let background = notebook + .background(BACKGROUND, notify_background) + .inspect_err(|error| eprintln!("Background sync did not start: {error}")) + .ok()?; + background.set_offline(offline()); + background.watch(notebook.replicas()); + Some(Arc::new(background)) +} + pub fn file_name(path: &Path) -> String { path.canonicalize() .ok() @@ -617,6 +667,7 @@ mod tests { })), notice: None, cache: PathBuf::new(), + background: None, }; let shown = |root, file| library(root).local(Path::new(file)); let under = Some(PathBuf::from("/Volumes/agent/lab/Group/New Section 1.one")); diff --git a/crates/snowbound/src/main.rs b/crates/snowbound/src/main.rs index 29f523d3b770f19bfde553e3d26fd671b6de8f73..025a25adf18cc9376699b36a6bf180593ad167c6 100644 --- a/crates/snowbound/src/main.rs +++ b/crates/snowbound/src/main.rs @@ -662,6 +662,10 @@ impl State { } let layouts = Arc::new(Mutex::new(engine.clone())); let temporary = matches!(input, Input::Notes { .. } | Input::Page(_)); + let background = proxy.clone(); + library::on_background(move || { + let _ = background.send_event(UserEvent::Sync); + }); let mut notebooks = Vec::new(); let mut session = None; let mut sectionless = None; @@ -2429,7 +2433,19 @@ impl State { Some(library) => *library = Arc::clone(&session.library), None => self.notebooks.push(Arc::clone(&session.library)), } - self.session = Some(*session); + // The notebook's background sync takes the section over once it is closed. + if let Some(previous) = self.session.replace(*session) { + let background = previous.library.background.clone(); + let section = previous.section; + std::thread::spawn(move || { + if let Err(error) = section.close() { + eprintln!("Synchronization stopped: {error}"); + } + if let Some(background) = background { + background.wake(); + } + }); + } self.sectionless = None; if other_notebook { self.save_settings(); @@ -2769,12 +2785,24 @@ impl State { Ok(()) } - /// Applies what the section reported since the last poll. + /// Applies what the section, and each notebook's closed sections, reported since the + /// last poll. fn synced(&mut self) -> Result<(), Box> { + // Every notebook's changes are taken, so none is reported again. + let noticed = self + .notebooks + .iter() + .filter_map(|library| library.background.as_ref()) + .filter(|background| !background.changed().is_empty()) + .count(); + if noticed > 0 { + self.sync_index(true); + } + // Reports are rare: a status changed, or a section did. + self.window.request_redraw(); let Some(session) = &mut self.session else { return Ok(()); }; - let shown = sync::label(&session.sync); let mut listed = false; let mut changed = false; let mut rejected = None; @@ -2802,9 +2830,6 @@ impl State { } } session.sync = session.section.sync_status()?; - if listed || changed || shown != sync::label(&session.sync) { - self.window.request_redraw(); - } if listed { let back = session.shown; session.pages = session.section.pages()?; diff --git a/crates/snowbound/src/sync.rs b/crates/snowbound/src/sync.rs index 4e498fb752bf33dbe4622997f63acf36aeae7a38..9c2a8d132fde1ac1f5f4e20c5632d95d198d1709 100644 --- a/crates/snowbound/src/sync.rs +++ b/crates/snowbound/src/sync.rs @@ -1,5 +1,5 @@ //! The sync status at the top right, and the popup it opens: OneNote's Shared Notebook -//! Synchronization for the open section. +//! Synchronization for the open section's notebook and each of its sections. use crate::{Session, State, art, filetime, library, platform}; use notebook::session::SyncStatus; @@ -56,8 +56,71 @@ fn describe(sync: &SyncStatus) -> (&'static str, &'static [&'static str], Option } } -pub(crate) fn label(sync: &SyncStatus) -> &'static str { - describe(sync).0 +/// Labels from the best state to the worst, which a notebook's status shows. +const ORDER: [&str; 5] = [ + "Up to date", + "Syncing…", + "Section in use", + "Not connected", + "Unable to sync", +]; + +fn rank(sync: &SyncStatus) -> usize { + let label = describe(sync).0; + ORDER.iter().position(|shown| *shown == label).unwrap_or(0) +} + +/// Every section of the open section's notebook with its status, in catalog order: the +/// open one's from its session, the others' from the notebook's background sync. +fn sections(session: &Session) -> Vec<(String, SyncStatus)> { + let open = &session.tabs[session.tab].path; + let copy = |sync: &SyncStatus| SyncStatus { + synced: sync.synced, + error: sync + .error + .as_ref() + .map(|error| std::io::Error::new(error.kind(), error.to_string())), + queued: sync.queued, + }; + let mut sections = session + .library + .background + .as_ref() + .map(|background| background.status()) + .unwrap_or_default(); + match sections.iter_mut().find(|(path, _)| path == open) { + Some((_, sync)) => *sync = copy(&session.sync), + None => sections.insert(0, (open.clone(), copy(&session.sync))), + } + sections +} + +/// The notebook as a whole: the worst section's error, every waiting change, and the +/// oldest last sync. +fn overall(sections: &[(String, SyncStatus)]) -> SyncStatus { + let worst = sections + .iter() + .map(|(_, sync)| sync) + .filter(|sync| sync.error.is_some()) + .max_by_key(|sync| rank(sync)); + SyncStatus { + synced: sections + .iter() + .map(|(_, sync)| sync.synced) + .collect::>>() + .and_then(|times| times.into_iter().min()), + error: worst + .and_then(|sync| sync.error.as_ref()) + .map(|error| std::io::Error::new(error.kind(), error.to_string())), + queued: sections.iter().map(|(_, sync)| sync.queued).sum(), + } +} + +fn changes(count: u64) -> String { + match count { + 1 => "1 change".to_owned(), + count => format!("{count} changes"), + } } /// A FILETIME as OneNote's "Last sync": the time, and the date too before today. @@ -73,8 +136,9 @@ fn when(time: u64) -> String { /// The status's icon in the toolbar, named in its tooltip, which opens the popup. pub(crate) fn control(ui: &mut Ui, session: &Session, theme: &Theme) { - let strong = ui.popup_open(id()) || session.sync.error.is_some() && !library::offline(); - let (label, icon, _) = describe(&session.sync); + let sync = overall(§ions(session)); + let strong = ui.popup_open(id()) || sync.error.is_some() && !library::offline(); + let (label, icon, _) = describe(&sync); ui.open_as( button(), Spec { @@ -102,7 +166,7 @@ pub(crate) fn control(ui: &mut Ui, session: &Session, theme: &Theme) { impl State { /// Builds the sync status popup while it is open. pub(crate) fn sync_popup(&mut self) -> Result<(), Box> { - let ui = &mut self.ui; + let (ui, notebooks) = (&mut self.ui, &self.notebooks); let Some(session) = &mut self.session else { ui.close_popup(id()); return Ok(()); @@ -114,7 +178,9 @@ impl State { session.sync = session.section.sync_status()?; let theme = ui.theme.clone(); let offline = library::offline(); - let (progress, _, advice) = describe(&session.sync); + let sections = sections(session); + let sync = overall(§ions); + let (progress, _, advice) = describe(&sync); ui.open_as( id(), Spec { @@ -177,23 +243,54 @@ impl State { row(ui, "Notebook", &session.library.name); row(ui, "Location", &session.library.location); row(ui, "Connection", &session.library.transport()); - row(ui, "Section", &session.tabs[session.tab].name); row(ui, "Progress", progress); - if let Some(synced) = session.sync.synced { + if let Some(synced) = sync.synced { row(ui, "Last sync", &when(synced)); } - if session.sync.queued > 0 { - let queued = match session.sync.queued { - 1 => "1 change".to_owned(), - count => format!("{count} changes"), - }; - row(ui, "Not yet synced", &queued); + if sync.queued > 0 { + row(ui, "Not yet synced", &changes(sync.queued)); } if let Some(advice) = advice { text(ui, "advice", advice, theme.text, false); } - if let Some(error) = &session.sync.error { - text(ui, "error", &error.to_string(), theme.text_dim, false); + text(ui, "sections", "Sections", theme.text, true); + for (index, (path, sync)) in sections.iter().enumerate() { + ui.open( + format!("section-{index}"), + Spec { + size: [fill(), children()], + gap: 8.0, + ..Spec::default() + }, + ); + let name = path.strip_suffix(".one").unwrap_or(path); + ui.leaf( + "name", + Spec { + size: [fill(), px(theme.font_size * 1.6)], + text: Some(name), + overflow: Overflow::Ellipsis, + ..Spec::default() + }, + ); + let status = match sync.queued { + 0 => describe(sync).0.to_owned(), + queued => format!("{}, {}", describe(sync).0, changes(queued)), + }; + ui.leaf( + "status", + Spec { + size: [fit(), px(theme.font_size * 1.6)], + text: Some(&status), + color: Some(theme.text_dim), + ..Spec::default() + }, + ); + ui.close(); + if let Some(error) = &sync.error { + let part = format!("section-{index}-error"); + text(ui, &part, &error.to_string(), theme.text_dim, false); + } } if let Some(notice) = &session.library.notice { let notice = format!("Snowbound’s SMB client couldn’t sign in: {notice}"); @@ -239,12 +336,16 @@ impl State { let now = ui::button(ui, "sync-now", "Sync Now").clicked; ui.close(); ui.close(); + let backgrounds = notebooks + .iter() + .filter_map(|library| library.background.as_ref()); if toggled { library::set_offline(!offline); session.section.set_offline(!offline); - } - if now { + backgrounds.for_each(|background| background.set_offline(!offline)); + } else if now { session.section.wake(); + backgrounds.for_each(|background| background.wake()); } if let Some(file) = file.filter(|_| show) { platform::show_file(&file);