diff --git a/crates/notebook/README.md b/crates/notebook/README.md index 80df0d3718d76a36c2b1418e779d4b2c281e6fbe..8641c24faa3afdd6c54e3d058f60aa6dd6a13adb 100644 --- a/crates/notebook/README.md +++ b/crates/notebook/README.md @@ -51,6 +51,44 @@ filesystem: the connection holds exclusive ownership between transactions, and a second open fails busy. No network wait occurs in a local save. After a database error, reopen and inspect the durable state before retrying. +## Sessions + +`session::Notebook::open(root, cache_dir)` discovers a notebook directory and +`session::Section::open(file, cache_dir, notify)` opens one section file through +a replica stored under the cache directory, named by the section's document +identity so the same file reopens the same queue after a relaunch. A section +publishes in the background to the file itself under OneNote-compatible +exclusion. `pages()` lists page spaces and titles from the local image, +`page(space)` returns the model to edit, and `save(space, before, after, +author)` queues the edited model: `Save::Queued(id)` is durable locally, +`Save::Unchanged` means the model equals the stored page, and `Save::Stale` +means the stored page no longer matches `before` because the section changed +underneath the editor, so the page must be reloaded before saving again. +`events()` drains what the synchronization thread reported since the last +poll: refreshes, attempt outcomes and unreachable files; `notify` runs on that +thread whenever an event is available so the application can wake its event +loop. `close()` stops publication. + +`Event::Unreachable` retains an `io::Error`: callers can distinguish permission +denial, missing targets, timeouts, and connection failures through `kind()` without +parsing display text. Document/cache failures are reported as `Event::Failed` and +stop the worker. Durable publication outcomes remain available through `status`. + +`Section::resume(file, replica, notify)` starts a session from an owned +`Replica::open(cache_file)` without consulting the remote file. The caller retains +the cache location and publication path for offline relaunch. Local pages and +saves remain available while the target is absent; synchronization verifies the +document identity before adopting or publishing remote content. A different +document stops the worker and retains the pending local edits. `Section::open` +remains the online convenience constructor that discovers the identity from the +file and creates or reopens its cache. + +With the `smb` feature, `Section::resume_smb(path, replica, limit, connect, +notify)` binds the same session operations to a share-relative path. `connect` +returns a new SMB client on the synchronization worker and is called again after +transport failure. The caller supplies credentials there, outside the replica +schema; construction and local saving do not wait for a network connection. + ## Pages of a section `create_page` accepts the core `PageCreation` intent and queues its page space diff --git a/crates/notebook/src/lib.rs b/crates/notebook/src/lib.rs index c4acdec369f3e9d36c16bc7ddab3c327473b6272..9793056e4ec3e3f25e196febc341e672e31f3477 100644 --- a/crates/notebook/src/lib.rs +++ b/crates/notebook/src/lib.rs @@ -25,6 +25,7 @@ mod pages; mod rebase; mod recovery; mod schema; +pub mod session; pub use pages::PageEdits; pub use recovery::{Recovery, RecoverySummary}; mod sync; diff --git a/crates/notebook/src/session.rs b/crates/notebook/src/session.rs new file mode 100644 index 0000000000000000000000000000000000000000..4278c063709873d1c02f6c17a8cdb59afa445aa4 --- /dev/null +++ b/crates/notebook/src/session.rs @@ -0,0 +1,297 @@ +//! The application's view of a notebook: sections opened through a local replica that +//! publishes page saves to the section file in the background. + +use crate::{EditStatus, Error, PendingEdit, Remote, Replica, Result, SyncWorker, discover}; +use onestore::{ + CommitError, ExGuid, PreparedEdit, RevisionIndex, Store, document::Document, page::Page, +}; +use std::{ + io, + path::{Path, PathBuf}, + sync::{ + Arc, + mpsc::{self, Receiver}, + }, + time::Duration, +}; + +/// A notebook directory and the cache directory holding its section replicas. +pub struct Notebook { + root: PathBuf, + cache: PathBuf, + catalog: discover::Folder, +} + +impl Notebook { + pub fn open(root: impl AsRef, cache: impl AsRef) -> Result { + let root = root.as_ref().canonicalize()?; + let cache = cache.as_ref().to_path_buf(); + std::fs::create_dir_all(&cache)?; + let catalog = discover::discover( + &mut discover::Local::open(&root)?, + discover::Limits { + entries: 100_000, + bytes_per_file: 256 * 1024 * 1024, + depth: 64, + }, + )?; + Ok(Self { + root, + cache, + catalog, + }) + } + + pub fn catalog(&self) -> &discover::Folder { + &self.catalog + } + + /// Opens a section by its catalog path. + pub fn section(&self, path: &str, notify: impl Fn() + Send + 'static) -> Result
{ + let mut folders = vec![&self.catalog]; + while let Some(folder) = folders.pop() { + if folder.sections.iter().any(|section| section.path == path) { + let file = self.root.join(path).canonicalize()?; + if !file.starts_with(&self.root) { + return Err(io::Error::from(io::ErrorKind::PermissionDenied).into()); + } + return Section::open(file, &self.cache, notify); + } + folders.extend(&folder.groups); + } + Err(io::Error::from(io::ErrorKind::NotFound).into()) + } +} + +/// What happened on the synchronization thread since the last poll. +#[derive(Debug)] +pub enum Event { + /// The working image was refreshed from the section file with no local edit pending. + Refreshed, + /// A publication attempt finished with this durable state. + Attempt { id: u64, status: EditStatus }, + /// The section file could not be reached; the replica keeps its state. + Unreachable(io::Error), + /// The replica itself failed; the worker has stopped. + Failed(String), +} + +/// The outcome of saving an edited page. +#[derive(Debug, PartialEq, Eq)] +pub enum Save { + /// The model equals the stored page. + Unchanged, + /// The edit is durable locally under this id. + Queued(u64), + /// The stored page no longer matches `before`; reload it before saving again. + Stale, +} + +/// A section file with its replica and background publication. +pub struct Section { + file: PathBuf, + replica: Arc, + worker: Option, + events: Receiver, +} + +impl Section { + /// Opens the section file through a replica in `cache`, creating the replica from the + /// file on first use. `notify` runs on the synchronization thread whenever an event is + /// available. + pub fn open( + file: impl AsRef, + cache: impl AsRef, + notify: impl Fn() + Send + 'static, + ) -> Result { + let file = file.as_ref().canonicalize()?; + let source = onestore::read_file(&file)?; + 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 replica = if cache.exists() { + Replica::open(&cache)? + } else { + Replica::create(&cache, &source)? + }; + Self::resume(file, replica, notify) + } + + /// Resumes an owned local replica without reading the publication target. + /// A target with a different document identity is rejected by synchronization. + pub fn resume( + file: impl AsRef, + replica: Replica, + notify: impl Fn() + Send + 'static, + ) -> Result { + let file = std::path::absolute(file)?; + let remote = file.clone(); + Self::start( + file, + replica, + move || Ok(FileRemote(remote.clone())), + notify, + ) + } + + /// Resumes a replica against a share-relative section path. Credentials remain in + /// `connect`, which is called on the worker again after transport failures. + #[cfg(feature = "smb")] + pub fn resume_smb( + path: String, + replica: Replica, + limit: usize, + mut connect: impl FnMut() -> io::Result + Send + 'static, + notify: impl Fn() + Send + 'static, + ) -> Result { + Self::start( + PathBuf::from(&path), + replica, + move || Ok(crate::SmbRemote::new(connect()?, path.clone(), limit)), + notify, + ) + } + + fn start( + file: PathBuf, + replica: Replica, + connect: impl FnMut() -> io::Result + Send + 'static, + notify: impl Fn() + Send + 'static, + ) -> Result { + let replica = Arc::new(replica); + let (sender, events) = mpsc::channel(); + let worker = replica.start_sync(Duration::from_secs(2), connect, move |result| { + let event = match result { + Ok(None) => Event::Refreshed, + Ok(Some((id, status))) => Event::Attempt { + id: *id, + status: *status, + }, + Err(Error::RemoteIo(error)) => Event::Unreachable(match error.raw_os_error() { + Some(code) => io::Error::from_raw_os_error(code), + None => io::Error::new(error.kind(), error.to_string()), + }), + Err(Error::Remote(error)) => { + Event::Unreachable(io::Error::new(error.error.kind(), error.to_string())) + } + Err(error) => Event::Failed(error.to_string()), + }; + if sender.send(event).is_ok() { + notify(); + } + })?; + Ok(Self { + file, + replica, + worker: Some(worker), + events, + }) + } + + /// The absolute local path, or share-relative path for an SMB session. + pub fn file(&self) -> &Path { + &self.file + } + + /// Page spaces and titles in section order, from the local working image. + pub fn pages(&self) -> Result> { + let snapshot = self.replica.snapshot()?; + let store = Store::parse(&snapshot)?; + let index = RevisionIndex::parse(&store)?; + let document = Document::parse(&index)?; + document + .pages()? + .into_iter() + .map(|(space, _)| Ok((space, Page::from_space(&document, space)?.title))) + .collect() + } + + pub fn page(&self, space: ExGuid) -> Result { + let snapshot = self.replica.snapshot()?; + let store = Store::parse(&snapshot)?; + let index = RevisionIndex::parse(&store)?; + Ok(Page::from_space(&Document::parse(&index)?, space)?) + } + + /// Saves an edited page. `before` is the model the edit started from; a stored page + /// that differs from it means the section changed underneath the editor. + pub fn save(&self, space: ExGuid, before: &Page, after: &Page, author: &str) -> Result { + loop { + let snapshot = self.replica.snapshot()?; + let store = Store::parse(&snapshot)?; + let index = RevisionIndex::parse(&store)?; + if Page::from_space(&Document::parse(&index)?, space)? != *before { + return Ok(Save::Stale); + } + match self.replica.save(&snapshot, space, after, author) { + Ok(Some(id)) => return Ok(Save::Queued(id)), + Ok(None) => return Ok(Save::Unchanged), + Err(Error::Io(error)) if error.kind() == io::ErrorKind::ResourceBusy => {} + Err(error) => return Err(error), + } + } + } + + pub fn status(&self, id: u64) -> Result> { + self.replica.status(id) + } + + pub fn pending(&self) -> Result> { + self.replica.pending() + } + + /// Events since the last poll, oldest first. + pub fn events(&self) -> Vec { + self.events.try_iter().collect() + } + + /// Requests a synchronization attempt now. + pub fn wake(&self) { + if let Some(worker) = &self.worker { + worker.wake(); + } + } + + pub fn replica(&self) -> &Replica { + &self.replica + } + + /// Waits for the in-flight operation and callback before releasing the replica. + /// Dropping instead requests cancellation without waiting; the worker retains + /// cache ownership until that operation finishes. Remote calls must be bounded. + pub fn close(mut self) -> Result<()> { + match self.worker.take() { + Some(worker) => worker.stop(), + None => Ok(()), + } + } +} + +/// The section file itself as the publication target, under OneNote-compatible exclusion. +struct FileRemote(PathBuf); + +impl Remote for FileRemote { + fn read(&mut self) -> io::Result> { + onestore::read_file(&self.0) + } + + fn publish(&mut self, edit: &PreparedEdit<'_>) -> std::result::Result<(), CommitError> { + edit.commit_file(&self.0) + } + + fn confirm(&mut self, snapshot: &[u8]) -> std::result::Result<(), CommitError> { + onestore::confirm_file_snapshot(&self.0, snapshot) + } +} + +impl Drop for Section { + fn drop(&mut self) { + drop(self.worker.take()); + } +} diff --git a/crates/notebook/tests/session.rs b/crates/notebook/tests/session.rs new file mode 100644 index 0000000000000000000000000000000000000000..609d289f56719ec769c5a19542122a47db618002 --- /dev/null +++ b/crates/notebook/tests/session.rs @@ -0,0 +1,603 @@ +use notebook::{ + EditStatus, + session::{Event, Notebook, Save, Section}, +}; +use onestore::{ExGuid, PreparedEdit, page::Page}; +use std::{ + path::Path, + sync::{ + Arc, + atomic::{AtomicUsize, Ordering}, + }, + time::{Duration, Instant}, +}; + +#[path = "support/model_ops.rs"] +mod model_ops; + +fn first_text(page: &Page) -> ExGuid { + 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)), + _ => None, + }) + .unwrap() +} + +fn edited(page: &Page, text: &str) -> Page { + let mut after = page.clone(); + let id = first_text(page); + model_ops::replace_text(&mut after, id, 0..0, text); + after +} + +fn open(file: &Path, cache: &Path) -> (Section, Arc) { + let notified = Arc::new(AtomicUsize::new(0)); + let counter = Arc::clone(¬ified); + let section = Section::open(file, cache, move || { + counter.fetch_add(1, Ordering::SeqCst); + }) + .unwrap(); + (section, notified) +} + +fn wait(section: &Section, mut accept: impl FnMut(&Event) -> bool) { + let deadline = Instant::now() + Duration::from_secs(20); + loop { + for event in section.events() { + if accept(&event) { + return; + } + } + assert!( + Instant::now() < deadline, + "the expected event did not arrive" + ); + std::thread::sleep(Duration::from_millis(20)); + } +} + +fn published(section: &Section, id: u64) { + wait( + section, + |event| matches!(event, Event::Attempt { id: n, status: EditStatus::Published { .. } } if *n == id), + ); +} + +fn stored_page(file: &Path, space: ExGuid) -> Page { + model_ops::page_of(&onestore::read_file(file).unwrap(), space) +} + +/// The page title is derived from its content on read; an edited model keeps the old one. +fn assert_same(actual: Page, expected: &Page) { + let mut expected = expected.clone(); + expected.title = actual.title.clone(); + assert_eq!(actual, expected); +} + +#[test] +fn a_section_opens_through_its_replica_and_a_save_reaches_the_file_and_survives_relaunch() { + let directory = tempfile::tempdir().unwrap(); + let file = directory.path().join("notes.one"); + let cache = directory.path().join("cache"); + std::fs::write( + &file, + onestore::create_section("notes.one", "Original", "Author").unwrap(), + ) + .unwrap(); + let (section, notified) = open(&file, &cache); + let pages = section.pages().unwrap(); + assert_eq!(pages.len(), 1); + let space = pages[0].0; + let before = section.page(space).unwrap(); + assert_eq!( + section.save(space, &before, &before, "Editor").unwrap(), + Save::Unchanged + ); + let after = edited(&before, "Saved "); + let Save::Queued(id) = section.save(space, &before, &after, "Editor").unwrap() else { + panic!() + }; + assert_same(section.page(space).unwrap(), &after); + published(§ion, id); + assert!(notified.load(Ordering::SeqCst) > 0); + assert_same(stored_page(&file, space), &after); + assert!(matches!( + section.status(id).unwrap(), + Some(EditStatus::Published { .. }) + )); + section.close().unwrap(); + let (section, _) = open(&file, &cache); + assert_same(section.page(space).unwrap(), &after); + assert!(section.pending().unwrap().is_empty()); + assert!(matches!( + section.status(id).unwrap(), + Some(EditStatus::Published { .. }) + )); + section.close().unwrap(); + if let Some(destination) = std::env::var_os("ONESTORE_SESSION_NATIVE_EXPORT") { + let destination = std::path::PathBuf::from(destination); + std::fs::create_dir(&destination).unwrap(); + let bytes = onestore::read_file(&file).unwrap(); + let identity = onestore::Store::parse(&bytes).unwrap().header.file_id; + std::fs::write(destination.join("notes.one"), bytes).unwrap(); + std::fs::write( + destination.join("Open Notebook.onetoc2"), + onestore::create_table_of_contents("Open Notebook.onetoc2", &[("notes.one", identity)]) + .unwrap(), + ) + .unwrap(); + } +} + +#[test] +fn saves_wait_for_an_unreachable_file_and_publish_after_relaunch() { + let directory = tempfile::tempdir().unwrap(); + let file = directory.path().join("notes.one"); + let cache = directory.path().join("cache"); + std::fs::write( + &file, + onestore::create_section("notes.one", "Original", "Author").unwrap(), + ) + .unwrap(); + let (section, _) = open(&file, &cache); + let space = section.pages().unwrap()[0].0; + let before = section.page(space).unwrap(); + use std::os::unix::fs::PermissionsExt; + let permissions = std::fs::metadata(&file).unwrap().permissions(); + std::fs::set_permissions(&file, std::fs::Permissions::from_mode(0o444)).unwrap(); + let after = edited(&before, "Offline "); + let Save::Queued(id) = section.save(space, &before, &after, "Editor").unwrap() else { + panic!() + }; + wait(§ion, |event| matches!(event, Event::Unreachable(_))); + section.close().unwrap(); + assert_eq!( + notebook::Replica::open( + std::fs::read_dir(&cache) + .unwrap() + .next() + .unwrap() + .unwrap() + .path() + ) + .unwrap() + .status(id) + .unwrap(), + Some(EditStatus::Pending) + ); + let (section, _) = open(&file, &cache); + assert_same(section.page(space).unwrap(), &after); + assert_eq!(section.pending().unwrap().len(), 1); + assert_eq!(stored_page(&file, space), before); + std::fs::set_permissions(&file, permissions).unwrap(); + section.wake(); + published(§ion, id); + assert_same(stored_page(&file, space), &after); + section.close().unwrap(); +} + +#[test] +fn an_external_change_refreshes_the_page_and_stales_a_save_from_the_old_model() { + let directory = tempfile::tempdir().unwrap(); + let file = directory.path().join("notes.one"); + let cache = directory.path().join("cache"); + std::fs::write( + &file, + onestore::create_section("notes.one", "Original", "Author").unwrap(), + ) + .unwrap(); + let (section, _) = open(&file, &cache); + let space = section.pages().unwrap()[0].0; + let before = section.page(space).unwrap(); + let native = edited(&before, "Native "); + let deadline = Instant::now() + Duration::from_secs(20); + loop { + let bytes = onestore::read_file(&file).unwrap(); + match PreparedEdit::page(&bytes, space, &native, "Native") + .unwrap() + .commit_file(&file) + { + Ok(()) => break, + Err(error) if error.error.kind() == std::io::ErrorKind::WouldBlock => { + assert!(Instant::now() < deadline); + std::thread::sleep(Duration::from_millis(20)); + } + Err(error) => panic!("{error:?}"), + } + } + section.wake(); + let deadline = Instant::now() + Duration::from_secs(20); + while section.page(space).unwrap().title == before.title { + assert!(Instant::now() < deadline, "the refresh did not arrive"); + wait(§ion, |event| matches!(event, Event::Refreshed)); + } + assert_same(section.page(space).unwrap(), &native); + let native = section.page(space).unwrap(); + let stale = edited(&before, "Local "); + assert_eq!( + section.save(space, &before, &stale, "Editor").unwrap(), + Save::Stale + ); + let after = edited(&native, "Local "); + let Save::Queued(id) = section.save(space, &native, &after, "Editor").unwrap() else { + panic!() + }; + published(§ion, id); + assert_same(stored_page(&file, space), &after); + section.close().unwrap(); +} + +#[test] +fn a_notebook_directory_lists_its_sections_and_opens_them() { + let directory = tempfile::tempdir().unwrap(); + let root = directory.path().join("Personal"); + std::fs::create_dir_all(root.join("Group")).unwrap(); + for (name, text) in [ + ("Personal/First.one", "One"), + ("Personal/Group/Second.one", "Two"), + ] { + std::fs::write( + directory.path().join(name), + onestore::create_section( + Path::new(name).file_name().unwrap().to_str().unwrap(), + text, + "Author", + ) + .unwrap(), + ) + .unwrap(); + } + let notebook = Notebook::open(&root, directory.path().join("cache")).unwrap(); + let catalog = notebook.catalog(); + assert_eq!(catalog.sections.len(), 1); + assert_eq!(catalog.groups.len(), 1); + assert_eq!(catalog.groups[0].sections.len(), 1); + let path = catalog.groups[0].sections[0].path.clone(); + let section = notebook.section(&path, || {}).unwrap(); + let (space, _) = section.pages().unwrap()[0]; + let page = section.page(space).unwrap(); + let text = first_text(&page); + assert_eq!( + model_ops::paragraph_with(&page, text) + .unwrap() + .text() + .unwrap() + .text + .text(), + "Two" + ); + section.close().unwrap(); +} + +#[test] +fn a_notebook_only_opens_discovered_section_paths() { + let directory = tempfile::tempdir().unwrap(); + let root = directory.path().join("Notebook"); + let cache = directory.path().join("cache"); + std::fs::create_dir(&root).unwrap(); + let source = onestore::create_section("notes.one", "Original", "Author").unwrap(); + std::fs::write(root.join("notes.one"), &source).unwrap(); + let outside = directory.path().join("outside.one"); + std::fs::write(&outside, &source).unwrap(); + let notebook = Notebook::open(&root, &cache).unwrap(); + std::fs::write(root.join("added.one"), &source).unwrap(); + for path in ["../outside.one", outside.to_str().unwrap(), "added.one"] { + assert!(matches!( + notebook.section(path, || {}), + Err(notebook::Error::Io(error)) if error.kind() == std::io::ErrorKind::NotFound + )); + } + assert_eq!(std::fs::read_dir(&cache).unwrap().count(), 0); + assert_eq!(std::fs::read(&outside).unwrap(), source); +} + +#[cfg(unix)] +#[test] +fn a_catalog_section_replaced_by_an_outside_symlink_is_rejected() { + let directory = tempfile::tempdir().unwrap(); + let root = directory.path().join("Notebook"); + let cache = directory.path().join("cache"); + std::fs::create_dir(&root).unwrap(); + let file = root.join("notes.one"); + let outside = directory.path().join("outside.one"); + let source = onestore::create_section("notes.one", "Original", "Author").unwrap(); + std::fs::write(&file, &source).unwrap(); + std::fs::write(&outside, &source).unwrap(); + let notebook = Notebook::open(&root, &cache).unwrap(); + std::fs::remove_file(&file).unwrap(); + std::os::unix::fs::symlink(&outside, &file).unwrap(); + assert!(matches!( + notebook.section("notes.one", || {}), + Err(notebook::Error::Io(error)) if error.kind() == std::io::ErrorKind::PermissionDenied + )); + assert_eq!(std::fs::read_dir(&cache).unwrap().count(), 0); + assert_eq!(std::fs::read(&outside).unwrap(), source); +} + +#[test] +fn an_absent_remote_does_not_prevent_local_relaunch_or_further_saves() { + let directory = tempfile::tempdir().unwrap(); + let file = directory.path().join("remote.one"); + let cache = directory.path().join("replica.sqlite"); + let source = onestore::create_section("remote.one", "Original", "Author").unwrap(); + let replica = notebook::Replica::create(&cache, &source).unwrap(); + let section = Section::resume(&file, replica, || {}).unwrap(); + let space = section.pages().unwrap()[0].0; + let before = section.page(space).unwrap(); + let after = edited(&before, "First "); + let Save::Queued(id) = section.save(space, &before, &after, "Editor").unwrap() else { + panic!() + }; + wait( + §ion, + |event| matches!(event, Event::Unreachable(error) if error.kind() == std::io::ErrorKind::NotFound), + ); + assert_eq!(section.status(id).unwrap(), Some(EditStatus::Pending)); + section.close().unwrap(); + assert!(!file.exists()); + + let section = Section::resume(&file, notebook::Replica::open(&cache).unwrap(), || {}).unwrap(); + assert_same(section.page(space).unwrap(), &after); + assert_eq!(section.status(id).unwrap(), Some(EditStatus::Pending)); + let before = section.page(space).unwrap(); + let after = edited(&before, "Second "); + let Save::Queued(next) = section.save(space, &before, &after, "Editor").unwrap() else { + panic!() + }; + assert!(!file.exists()); + let restored = directory.path().join("restored.one"); + std::fs::write(&restored, &source).unwrap(); + std::fs::rename(restored, &file).unwrap(); + section.wake(); + published(§ion, next); + assert_same(stored_page(&file, space), &after); + section.close().unwrap(); + let section = Section::resume(&file, notebook::Replica::open(&cache).unwrap(), || {}).unwrap(); + assert_same(section.page(space).unwrap(), &after); + assert!(section.pending().unwrap().is_empty()); + section.close().unwrap(); +} + +#[test] +fn resuming_against_another_document_preserves_both_remote_and_pending_edit() { + let directory = tempfile::tempdir().unwrap(); + let file = directory.path().join("other.one"); + let cache = directory.path().join("replica.sqlite"); + let source = onestore::create_section("notes.one", "Original", "Author").unwrap(); + let replica = notebook::Replica::create(&cache, &source).unwrap(); + let section = Section::resume(&file, replica, || {}).unwrap(); + let space = section.pages().unwrap()[0].0; + let before = section.page(space).unwrap(); + let after = edited(&before, "Local "); + let Save::Queued(id) = section.save(space, &before, &after, "Editor").unwrap() else { + panic!() + }; + wait(§ion, |event| matches!(event, Event::Unreachable(_))); + section.close().unwrap(); + + let other = onestore::create_section("other.one", "Unrelated", "Other").unwrap(); + std::fs::write(&file, &other).unwrap(); + let section = Section::resume(&file, notebook::Replica::open(&cache).unwrap(), || {}).unwrap(); + wait(§ion, |event| matches!(event, Event::Failed(_))); + assert_eq!(section.status(id).unwrap(), Some(EditStatus::Pending)); + assert_same(section.page(space).unwrap(), &after); + assert_eq!(std::fs::read(&file).unwrap(), other); + assert!(section.close().is_err()); + let replica = notebook::Replica::open(&cache).unwrap(); + assert_eq!(replica.status(id).unwrap(), Some(EditStatus::Pending)); + assert_eq!(std::fs::read(&file).unwrap(), other); +} + +#[test] +fn offline_save_process() { + let Some(root) = std::env::var_os("ONENOTE_SESSION_CHILD") else { + return; + }; + let root = std::path::PathBuf::from(root); + let replica = notebook::Replica::open(root.join("replica.sqlite")).unwrap(); + let section = Section::resume(root.join("absent.one"), replica, || {}).unwrap(); + let space = section.pages().unwrap()[0].0; + let before = section.page(space).unwrap(); + let after = edited(&before, "Durable "); + let Save::Queued(id) = section.save(space, &before, &after, "Editor").unwrap() else { + panic!() + }; + std::fs::write(root.join("acknowledged"), id.to_string()).unwrap(); + // Terminate without running Session or SQLite destructors after acknowledgement. + std::process::exit(0); +} + +#[test] +fn an_acknowledged_offline_save_survives_process_exit_without_cleanup() { + let directory = tempfile::tempdir().unwrap(); + let cache = directory.path().join("replica.sqlite"); + let source = onestore::create_section("absent.one", "Original", "Author").unwrap(); + drop(notebook::Replica::create(&cache, &source).unwrap()); + let output = std::process::Command::new(std::env::current_exe().unwrap()) + .args(["--exact", "offline_save_process", "--nocapture"]) + .env("ONENOTE_SESSION_CHILD", directory.path()) + .output() + .unwrap(); + assert!(output.status.success(), "{output:?}"); + let id: u64 = std::fs::read_to_string(directory.path().join("acknowledged")) + .unwrap() + .parse() + .unwrap(); + let file = directory.path().join("absent.one"); + assert!(!file.exists()); + let section = Section::resume(&file, notebook::Replica::open(&cache).unwrap(), || {}).unwrap(); + let space = section.pages().unwrap()[0].0; + let original = model_ops::page_of(&source, space); + let after = edited(&original, "Durable "); + assert_same(section.page(space).unwrap(), &after); + assert_eq!(section.status(id).unwrap(), Some(EditStatus::Pending)); + section.close().unwrap(); +} + +#[cfg(feature = "smb")] +#[test] +fn smb_connection_failure_keeps_the_session_locally_editable_and_retries() { + let directory = tempfile::tempdir().unwrap(); + let cache = directory.path().join("replica.sqlite"); + let source = onestore::create_section("notes.one", "Original", "Author").unwrap(); + let replica = notebook::Replica::create(&cache, &source).unwrap(); + let attempts = Arc::new(AtomicUsize::new(0)); + let counter = Arc::clone(&attempts); + let section = Section::resume_smb( + "Folder/notes.one".into(), + replica, + 1024 * 1024, + move || { + counter.fetch_add(1, Ordering::SeqCst); + Err(std::io::ErrorKind::PermissionDenied.into()) + }, + || {}, + ) + .unwrap(); + wait( + §ion, + |event| matches!(event, Event::Unreachable(error) if error.kind() == std::io::ErrorKind::PermissionDenied), + ); + let space = section.pages().unwrap()[0].0; + let before = section.page(space).unwrap(); + let after = edited(&before, "Offline SMB "); + let Save::Queued(id) = section.save(space, &before, &after, "Editor").unwrap() else { + panic!() + }; + section.wake(); + wait( + §ion, + |event| matches!(event, Event::Unreachable(error) if error.kind() == std::io::ErrorKind::PermissionDenied), + ); + assert!(attempts.load(Ordering::SeqCst) >= 2); + assert_eq!(section.status(id).unwrap(), Some(EditStatus::Pending)); + assert_same(section.page(space).unwrap(), &after); + assert_eq!(section.file(), Path::new("Folder/notes.one")); + section.close().unwrap(); + assert_eq!( + notebook::Replica::open(&cache).unwrap().status(id).unwrap(), + Some(EditStatus::Pending) + ); +} + +#[cfg(feature = "smb")] +#[test] +#[ignore = "requires ONESTORE_SESSION_SMB address and a disposable session.one on its agent share"] +fn live_smb_session_save_publishes_and_reopens() { + use notebook::smb::{Client, Credentials}; + let address = std::env::var("ONESTORE_SESSION_SMB").unwrap(); + let connect = move || { + Client::connect( + &address, + "agent", + Credentials::default(), + Duration::from_secs(5), + ) + }; + let client = connect().unwrap(); + let source = client.read("session.one", 1024 * 1024).unwrap(); + let directory = tempfile::tempdir().unwrap(); + let cache = directory.path().join("replica.sqlite"); + let replica = notebook::Replica::create(&cache, &source).unwrap(); + let section = Section::resume_smb( + "session.one".into(), + replica, + 1024 * 1024, + connect.clone(), + || {}, + ) + .unwrap(); + let space = section.pages().unwrap()[0].0; + let before = section.page(space).unwrap(); + let after = edited(&before, "Session SMB "); + let Save::Queued(id) = section.save(space, &before, &after, "Editor").unwrap() else { + panic!() + }; + published(§ion, id); + let remote = client.read("session.one", 1024 * 1024).unwrap(); + assert_same(model_ops::page_of(&remote, space), &after); + section.close().unwrap(); + let section = Section::resume_smb( + "session.one".into(), + notebook::Replica::open(&cache).unwrap(), + 1024 * 1024, + connect, + || {}, + ) + .unwrap(); + assert_same(section.page(space).unwrap(), &after); + assert!(section.pending().unwrap().is_empty()); + assert!(matches!( + section.status(id).unwrap(), + Some(EditStatus::Published { .. }) + )); + section.close().unwrap(); +} + +#[cfg(feature = "smb")] +#[test] +fn dropping_during_connection_keeps_cache_owned_until_the_worker_finishes() { + use std::sync::mpsc; + let directory = tempfile::tempdir().unwrap(); + let cache = directory.path().join("replica.sqlite"); + let source = onestore::create_section("notes.one", "Original", "Author").unwrap(); + let replica = notebook::Replica::create(&cache, &source).unwrap(); + let (entered, connecting) = mpsc::channel(); + let (release, stalled) = mpsc::channel(); + let section = Section::resume_smb( + "notes.one".into(), + replica, + 1024 * 1024, + move || { + entered.send(()).unwrap(); + stalled.recv_timeout(Duration::from_secs(10)).unwrap(); + Err(std::io::ErrorKind::TimedOut.into()) + }, + || {}, + ) + .unwrap(); + connecting.recv_timeout(Duration::from_secs(5)).unwrap(); + let space = section.pages().unwrap()[0].0; + let before = section.page(space).unwrap(); + let after = edited(&before, "Saved during connection "); + let Save::Queued(id) = section.save(space, &before, &after, "Editor").unwrap() else { + panic!() + }; + let (dropped, finished) = mpsc::channel(); + let owner = std::thread::spawn(move || { + drop(section); + dropped.send(()).unwrap(); + }); + finished.recv_timeout(Duration::from_secs(2)).unwrap(); + assert!(matches!( + notebook::Replica::open(&cache), + Err(notebook::Error::Database(error)) + if error.sqlite_error_code() == Some(rusqlite::ErrorCode::DatabaseBusy) + )); + release.send(()).unwrap(); + owner.join().unwrap(); + let deadline = Instant::now() + Duration::from_secs(5); + let replica = loop { + match notebook::Replica::open(&cache) { + Ok(replica) => break replica, + Err(notebook::Error::Database(error)) + if error.sqlite_error_code() == Some(rusqlite::ErrorCode::DatabaseBusy) => + { + assert!(Instant::now() < deadline, "worker retained cache ownership"); + std::thread::sleep(Duration::from_millis(10)); + } + Err(error) => panic!("{error}"), + } + }; + assert_eq!(replica.status(id).unwrap(), Some(EditStatus::Pending)); + assert_same( + model_ops::page_of(&replica.snapshot().unwrap(), space), + &after, + ); + assert!(connecting.try_recv().is_err()); +} diff --git a/crates/onestore/src/commit.rs b/crates/onestore/src/commit.rs index d0022820e849a4cd374b269fb254855bf9db6fc1..ad2abf55a7f16bff33ae5a85f753a617095c0f12 100644 --- a/crates/onestore/src/commit.rs +++ b/crates/onestore/src/commit.rs @@ -345,6 +345,17 @@ impl<'a> PreparedEdit<'a> { } } +/// `confirm_snapshot` under the same whole-file exclusion `commit_file` uses. +#[cfg(any(unix, windows))] +pub fn confirm_file_snapshot(path: impl AsRef, source: &[u8]) -> Result<(), CommitError> { + let mut io = FileIo::open(path, true).map_err(|error| CommitError { + state: CommitState::NotCommitted, + error, + })?; + let result = confirm_snapshot(&mut io, source); + io.finish(result) +} + /// Compares and flushes a snapshot, then refreshes its header version metadata. /// No revision is added; reread before using the snapshot for another physical commit. /// The caller must hold OneNote-compatible exclusion and independently establish which diff --git a/crates/onestore/src/lib.rs b/crates/onestore/src/lib.rs index 9a4deb50081cd1c1bf9f4fe581e78727ca5e1ac8..06a63f2e28bb773430ac9fd77976a81aaf077298 100644 --- a/crates/onestore/src/lib.rs +++ b/crates/onestore/src/lib.rs @@ -29,7 +29,9 @@ pub use commit::{ confirm_snapshot, }; #[cfg(any(unix, windows))] -pub use commit::{commit_file_property, commit_file_text, read_file, read_file_limited}; +pub use commit::{ + commit_file_property, commit_file_text, confirm_file_snapshot, read_file, read_file_limited, +}; pub use create::{create_section, create_table_of_contents}; pub use edit::replace_text; pub use files::FileDataReference;