From a380cb56a0a524f24a9edfa7a153336e4db3201c Mon Sep 17 00:00:00 2001 From: clover caruso Date: Fri, 11 Sep 2026 04:01:57 -0700 Subject: [PATCH] feat: add replica-backed section sessions with offline resume Expose local and embedded SMB sessions over the existing durable replica. Preserve typed transport failures and cache ownership through cancellation. Restrict notebook opens to catalog paths. Validated with notebook tests, live Samba publication, iOS Simulator runtime recovery and a cold OneNote read of a session-published fixture. Device compilation is separate from runtime acceptance. Assisted-by: gpt-6-astra --- crates/notebook/README.md | 38 ++ crates/notebook/src/lib.rs | 1 + crates/notebook/src/session.rs | 297 +++++++++++++++ crates/notebook/tests/session.rs | 603 +++++++++++++++++++++++++++++++ crates/onestore/src/commit.rs | 11 + crates/onestore/src/lib.rs | 4 +- 6 files changed, 953 insertions(+), 1 deletion(-) create mode 100644 crates/notebook/src/session.rs create mode 100644 crates/notebook/tests/session.rs 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; -- 2.54.0