authorgravatar for git@paperclover.netclover caruso <git@paperclover.net> 2026-09-29 09:45:05-07:00
committergravatar for git@paperclover.netclover caruso <git@paperclover.net> 2026-09-29 10:00:31-07:00
logf4e1afda77c93656bcf5c1e0e634aa9ca735d9ec
treee32ef5880f4ca7924fae722a64029adb731c7655
parent24912d53294b686bf8bfebdd7ab2e17ca14aba1f
signature Signed by SSH key SHA256:52mNGHRsVFBDED9IAX5pe+LRWUefqTbxEReunq21QvU

feat(sync): sync closed sections in the background, as OneNote 2010 does

Each open notebook gets a background thread that polls every closed section's stamp every 15 s (immediately on Sync Now or when a section closes). A section with queued edits or a moved file runs the existing sync step, publishing or merging as one appended revision. Work Offline stops it. The sync popup lists every section, and the top-right status shows the notebook's worst state. Assisted-by: claude-opus-5.5

12 files changed, 842 insertions(+), 142 deletions(-)

arc/sync.md+10
......@@ -167,6 +167,16 @@ Working offline on purpose, as OneNote's Work Offline does, is the same state
167167chosen: the sync thread stops stepping until Sync Now or until the user works
168168online again.
169169
170## Sections that aren't open
171
172OneNote keeps every section of an open notebook in sync, not just the one on screen: it
173reads a section another client changed within seconds and publishes a closed section's
174offline edits the moment the share is back. Snowbound's `session::Background` does the
175same with one thread per notebook. Each round reads every section file's stamp. Only a
176section whose replica has edits waiting, or whose file moved past the replica's base,
177has its replica opened for the usual sync steps, and it is closed again afterwards. The
178open section is left to its own session.
179
170180## Notebook structure
171181
172182Sections, groups and the notebook's own colour live in the `.onetoc2` files,
crates/notebook/README.md+11
......@@ -64,6 +64,17 @@ edit names is the host's: the app passes the account's full name.
6464and how many edits wait for it. `set_offline(true)` works offline as OneNote does: the
6565worker stops connecting and edits queue until `wake()` (Sync Now) or `set_offline(false)`.
6666
67`session::Background` keeps the sections no session holds in sync, as OneNote 2010 keeps
68every section of an open notebook: `Notebook::background(interval, notify)` for a mounted
69notebook, `Background::smb(root, limit, interval, connect, notify)` on a share, then
70`watch(notebook.replicas())`. Each round reads every watched file's stamp; a section whose
71replica has edits waiting, or whose file moved past the replica's base, has the replica
72opened for the synchronization steps that publish or rebase it and closed again. A replica
73a session holds is skipped (`Error::busy`), so opening a section may wait out one step.
74`status()` gives each section's `SyncStatus`, `changed()` the sections another client
75changed, and `wake` and `set_offline` follow Sync Now and Work Offline.
76`Notebook::replica_path(path)` names a section's replica for either kind of notebook.
77
6778`Section::resume(file, replica, notify)` starts from an owned `Replica` without
6879consulting the remote file; with the `smb` feature, `Section::resume_smb(path,
6980replica, limit, connect, notify)` binds a share-relative path, `connect` running on
crates/notebook/examples/scratch_bgsync.rs deleted-42
......@@ -1,42 +0,0 @@
1//! Scratch lab writer (delete before reporting): appends TEXT to the first body paragraph of
2//! the first page of a share-relative section. `scratch_bgsync ADDRESS PATH TEXT`
3use notebook::{
4 Replica, SmbRemote,
5 smb::{Client, Credentials},
6};
7use onestore::op::{Edit, Op, PageOp};
8use std::time::Duration;
9
10fn main() -> Result<(), Box<dyn std::error::Error>> {
11 let args: Vec<_> = std::env::args().skip(1).collect();
12 let [address, path, with] = &args[..] else {
13 return Err("scratch_bgsync ADDRESS PATH TEXT".into());
14 };
15 let connect = || Client::connect(address, "agent", Credentials::default(), Duration::from_secs(5));
16 let source = connect()?.read_storage(path, 1 << 26)?;
17 let directory = tempfile::tempdir()?;
18 let replica = Replica::create(directory.path().join("cache.sqlite"), &source)?;
19 let space = replica.pages()?[0].0;
20 let page = replica.page(space)?;
21 let (text, len) = page
22 .objects
23 .iter()
24 .find_map(|object| match object {
25 onestore::page::PageObject::Outline(outline) => outline
26 .paragraphs
27 .iter()
28 .find_map(|p| p.text().map(|t| (t.id, t.text.text().encode_utf16().count() as u32))),
29 _ => None,
30 })
31 .ok_or("no text")?;
32 replica.apply(
33 "Lab",
34 Edit {
35 at: 134_000_000_000_000_000,
36 ops: vec![Op::Page { space, op: PageOp::Text { text, range: len..len, with: with.clone() } }],
37 },
38 )?;
39 let mut remote = SmbRemote::new(connect()?, path.clone(), 1 << 26);
40 println!("{:?}", replica.sync_once(&mut remote)?.edit);
41 Ok(())
42}
crates/notebook/src/background.rs created+326
......@@ -0,0 +1,326 @@
1//! The sections of an open notebook that no session holds, kept in sync as OneNote 2010
2//! keeps every section of an open notebook: queued edits publish and other clients' changes
3//! are noticed without the section being open.
4
5use crate::{
6 EditStatus, Error, Remote, Replica, Result,
7 session::{SyncStatus, reached},
8 worker::Signal,
9};
10use onestore::Stamp;
11use std::{
12 io,
13 path::{Path, PathBuf},
14 sync::{Arc, Mutex, atomic::Ordering},
15 thread,
16 time::Duration,
17};
18
19/// Polls each watched section file's stamp every interval. A section whose replica has
20/// edits waiting, or whose file moved past the replica's base, has its replica opened for
21/// the synchronization steps that publish or rebase it, then closed again; any other costs
22/// one stamp read. A replica a session holds is left to that session's worker. Dropping
23/// requests cancellation without waiting for the step in flight.
24pub struct Background {
25 signal: Arc<Signal>,
26 watched: Arc<Mutex<Watched>>,
27}
28
29#[derive(Default)]
30struct Watched {
31 sections: Vec<Watch>,
32 /// Sections whose file changed since `changed` was last asked.
33 changed: Vec<String>,
34}
35
36struct Watch {
37 path: String,
38 replica: Option<PathBuf>,
39 /// The file's stamp when last reached.
40 stamp: Option<Stamp>,
41 status: SyncStatus,
42}
43
44impl Background {
45 pub(crate) fn start<R, B>(
46 interval: Duration,
47 mut connect: impl FnMut() -> io::Result<B> + Send + 'static,
48 notify: impl Fn() + Send + 'static,
49 ) -> Result<Self>
50 where
51 R: Remote,
52 B: FnMut(&str) -> R,
53 {
54 let (signal, receiver) = Signal::new();
55 let watched = Arc::new(Mutex::new(Watched::default()));
56 let (shared, sections) = (Arc::clone(&signal), Arc::clone(&watched));
57 thread::Builder::new()
58 .name("onestore-background".into())
59 .spawn(move || {
60 let mut bound: Option<B> = None;
61 while !shared.stopped.load(Ordering::Acquire) {
62 if shared.offline.load(Ordering::Acquire)
63 && !shared.requested.load(Ordering::Acquire)
64 {
65 bound = None;
66 let _ = receiver.recv();
67 continue;
68 }
69 shared.requested.store(false, Ordering::Release);
70 let round: Vec<_> = match sections.lock() {
71 Ok(watched) => watched
72 .sections
73 .iter()
74 .map(|watch| {
75 (
76 watch.path.clone(),
77 watch.replica.clone(),
78 watch.stamp.clone(),
79 )
80 })
81 .collect(),
82 Err(_) => return,
83 };
84 let mut news = false;
85 for (path, replica, seen) in round {
86 if shared.stopped.load(Ordering::Acquire) {
87 return;
88 }
89 let bind = match &mut bound {
90 Some(bind) => bind,
91 None => match connect() {
92 Ok(bind) => bound.insert(bind),
93 Err(error) => {
94 // Nothing on the share is reachable this round.
95 let Ok(mut watched) = sections.lock() else {
96 return;
97 };
98 for watch in &mut watched.sections {
99 news |= watch.fail(
100 &Error::RemoteIo(io::Error::new(
101 error.kind(),
102 error.to_string(),
103 )),
104 None,
105 );
106 }
107 break;
108 }
109 },
110 };
111 let (queued, outcome) =
112 step(&mut bind(&path), replica.as_deref(), seen.as_ref());
113 if matches!(outcome, Err(Error::RemoteIo(_) | Error::Remote(_))) {
114 bound = None;
115 }
116 let Ok(mut watched) = sections.lock() else {
117 return;
118 };
119 let Watched { sections, changed } = &mut *watched;
120 let Some(watch) = sections.iter_mut().find(|watch| watch.path == path)
121 else {
122 continue;
123 };
124 news |= match outcome {
125 Ok((stamp, moved)) => {
126 watch.stamp = Some(stamp);
127 if moved {
128 changed.push(path);
129 }
130 let before = summary(&watch.status);
131 if let Some(queued) = queued {
132 watch.status = SyncStatus {
133 synced: Some(crate::now()),
134 error: None,
135 queued,
136 };
137 }
138 moved || before != summary(&watch.status)
139 }
140 Err(error) => watch.fail(&error, queued),
141 };
142 }
143 if news {
144 notify();
145 }
146 let _ = receiver.recv_timeout(interval);
147 }
148 })?;
149 Ok(Self { signal, watched })
150 }
151
152 /// Keeps the sections of a notebook on a share in sync while they are not open;
153 /// `connect` runs again after a transport failure. Paths are relative to `root`.
154 #[cfg(feature = "smb")]
155 pub fn smb(
156 root: &str,
157 limit: usize,
158 interval: Duration,
159 mut connect: impl FnMut() -> io::Result<crate::smb::Client> + Send + 'static,
160 notify: impl Fn() + Send + 'static,
161 ) -> Result<Self> {
162 let root = root.replace('\\', "/");
163 Self::start(
164 interval,
165 move || {
166 let client = Arc::new(connect()?);
167 let root = root.clone();
168 Ok(move |path: &str| {
169 let file = match root.as_str() {
170 "" => path.to_owned(),
171 root => format!("{root}/{path}"),
172 };
173 crate::SmbRemote::new(Arc::clone(&client), file, limit)
174 })
175 },
176 notify,
177 )
178 }
179
180 /// Watches these sections, by catalog path and replica (`Notebook::replicas`), from
181 /// the next round, which starts now.
182 pub fn watch(&self, sections: Vec<(String, Option<PathBuf>)>) {
183 if let Ok(mut watched) = self.watched.lock() {
184 let mut previous = std::mem::take(&mut watched.sections);
185 watched.sections = sections
186 .into_iter()
187 .map(
188 |(path, replica)| match previous.iter().position(|watch| watch.path == path) {
189 Some(index) => Watch {
190 replica,
191 ..previous.swap_remove(index)
192 },
193 None => Watch {
194 path,
195 replica,
196 stamp: None,
197 status: SyncStatus {
198 synced: None,
199 error: None,
200 queued: 0,
201 },
202 },
203 },
204 )
205 .collect();
206 }
207 self.signal.wake();
208 }
209
210 /// Each watched section's status as its last round left it, in watch order. A section a
211 /// session holds keeps the status it had before.
212 pub fn status(&self) -> Vec<(String, SyncStatus)> {
213 self.watched.lock().map_or_else(
214 |_| Vec::new(),
215 |watched| {
216 watched
217 .sections
218 .iter()
219 .map(|watch| {
220 let status = &watch.status;
221 (
222 watch.path.clone(),
223 SyncStatus {
224 synced: status.synced,
225 error: status
226 .error
227 .as_ref()
228 .map(|error| io::Error::new(error.kind(), error.to_string())),
229 queued: status.queued,
230 },
231 )
232 })
233 .collect()
234 },
235 )
236 }
237
238 /// Sections whose file changed since the last call, by catalog path.
239 pub fn changed(&self) -> Vec<String> {
240 self.watched
241 .lock()
242 .map(|mut watched| std::mem::take(&mut watched.changed))
243 .unwrap_or_default()
244 }
245
246 /// Runs a round now, working offline included (Sync Now).
247 pub fn wake(&self) {
248 self.signal.requested.store(true, Ordering::Release);
249 self.signal.wake();
250 }
251
252 /// Working offline, no round runs until `wake`, or until working online again.
253 pub fn set_offline(&self, offline: bool) {
254 self.signal.offline.store(offline, Ordering::Release);
255 self.signal.wake();
256 }
257}
258
259impl Drop for Background {
260 fn drop(&mut self) {
261 self.signal.stopped.store(true, Ordering::Release);
262 self.signal.wake();
263 }
264}
265
266impl Watch {
267 /// Records a failed step and what it found waiting, answering whether the status shown
268 /// changes.
269 fn fail(&mut self, error: &Error, queued: Option<u64>) -> bool {
270 let before = summary(&self.status);
271 self.status.error = Some(reached(error));
272 self.status.queued = queued.unwrap_or(self.status.queued);
273 before != summary(&self.status)
274 }
275}
276
277/// What the host shows of a status.
278fn summary(status: &SyncStatus) -> (bool, Option<io::ErrorKind>, u64) {
279 (
280 status.synced.is_some(),
281 status.error.as_ref().map(io::Error::kind),
282 status.queued,
283 )
284}
285
286/// One section's step: how many of its edits wait (`None` while unknown, as while a session
287/// holds its replica), then the file's stamp now and whether it changed since `seen`.
288fn step<R: Remote>(
289 remote: &mut R,
290 replica: Option<&Path>,
291 seen: Option<&Stamp>,
292) -> (Option<u64>, Result<(Stamp, bool)>) {
293 let stamp = match remote.stamp() {
294 Ok(stamp) => stamp,
295 Err(error) => return (None, Err(Error::RemoteIo(error))),
296 };
297 let moved = seen.is_some_and(|seen| *seen != stamp);
298 let Some(replica) = replica.filter(|replica| replica.exists()) else {
299 return (Some(0), Ok((stamp, moved)));
300 };
301 match crate::peek(replica) {
302 Ok((base, 0)) if base == stamp => return (Some(0), Ok((stamp, moved))),
303 Err(error) if error.busy() => return (None, Ok((stamp, false))),
304 _ => {}
305 }
306 let replica = match Replica::open(replica) {
307 Ok(replica) => replica,
308 Err(error) if error.busy() => return (None, Ok((stamp, false))),
309 Err(error) => return (None, Err(error)),
310 };
311 let synced = (|| {
312 let mut changed = moved;
313 loop {
314 let synced = replica.sync_once(remote)?;
315 changed |= !synced.changed.is_empty();
316 if !matches!(synced.edit, Some((_, EditStatus::Published { .. }))) {
317 return Ok((remote.stamp().map_err(Error::RemoteIo)?, changed));
318 }
319 }
320 })();
321 let queued = replica
322 .recovery_summary()
323 .ok()
324 .map(|summary| summary.queued_edits);
325 (queued, synced)
326}
crates/notebook/src/lib.rs+24
......@@ -20,6 +20,7 @@ use std::{
2020};
2121
2222mod assets;
23mod background;
2324mod base;
2425mod merge;
2526mod migrate;
......@@ -60,6 +61,14 @@ pub enum Error {
6061 Protected(#[from] onestore::protected::Error),
6162}
6263
64impl Error {
65 /// Whether the replica is open elsewhere, as in another section session of this process.
66 pub fn busy(&self) -> bool {
67 matches!(self, Self::Database(rusqlite::Error::SqliteFailure(error, _))
68 if error.code == rusqlite::ErrorCode::DatabaseBusy)
69 }
70}
71
6372type Result<T> = std::result::Result<T, Error>;
6473
6574const APPLICATION_ID: u32 = 0x4f4e454f;
......@@ -307,6 +316,21 @@ fn validate(source: &[u8]) -> Result<ExGuid> {
307316 Ok(onestore::Section::open(&arena, source.to_vec())?.root())
308317}
309318
319/// A closed cache's base stamp and how many edits wait, without opening the replica.
320fn peek(path: &Path) -> Result<(onestore::Stamp, u64)> {
321 let connection = cache_connection(path)?;
322 let application: u32 =
323 connection.pragma_query_value(None, "application_id", |row| row.get(0))?;
324 let version: u32 = connection.pragma_query_value(None, "user_version", |row| row.get(0))?;
325 if application != APPLICATION_ID || version != schema::VERSION {
326 return Err(io::Error::from(io::ErrorKind::InvalidData).into());
327 }
328 let base = base::stamp(&connection, base::Image::Base)?
329 .ok_or_else(|| io::Error::new(io::ErrorKind::InvalidData, "The cache has no base image"))?;
330 let queued: i64 = connection.query_row("SELECT count(*) FROM edits", [], |row| row.get(0))?;
331 Ok((base, unsigned(queued)?))
332}
333
310334fn pending(connection: &Connection) -> Result<Vec<PendingEdit>> {
311335 queue::load(connection, None)
312336}
crates/notebook/src/session.rs+73-11
......@@ -1,6 +1,7 @@
11//! The application's view of a notebook: sections opened through a local replica that
22//! publishes their edits to the section file in the background.
33
4pub use crate::background::Background;
45use crate::{
56 EditStatus, Error, PendingEdit, Remote, Replica, Resolution, Result, SyncWorker, discover,
67};
......@@ -881,6 +882,61 @@ impl Notebook {
881882 Ok(())
882883 }
883884
885 /// Where the replica of the section at catalog `path` lives: named by the section's
886 /// document identity in a mounted notebook (`Section::open`), by its file identity under
887 /// `smb` on a share. It exists once the section has been opened.
888 pub fn replica_path(&self, path: &str) -> Result<PathBuf> {
889 let section = self.section_path(path)?;
890 Ok(match &self.root {
891 Some(_) => {
892 let image = self.storage.read(&section.path)?;
893 let root = RevisionIndex::parse(&Store::parse(&image)?)?.root;
894 replica_file(&self.cache, &root.guid)
895 }
896 None => replica_file(&self.cache.join("smb"), &section.file_id),
897 })
898 }
899
900 /// Every readable section's catalog path and replica, for `Background::watch`.
901 pub fn replicas(&self) -> Vec<(String, Option<PathBuf>)> {
902 let mut sections = Vec::new();
903 let mut folders = vec![&self.catalog];
904 while let Some(folder) = folders.pop() {
905 for section in &folder.sections {
906 if matches!(section.state, discover::SectionState::Readable { .. }) {
907 let replica = self.replica_path(&section.path).ok();
908 sections.push((section.path.clone(), replica));
909 }
910 }
911 folders.extend(folder.groups.iter().rev());
912 }
913 sections
914 }
915
916 /// Keeps the sections of a mounted notebook in sync while they are not open
917 /// (`Background`), polling each file every `interval` once watched.
918 pub fn background(
919 &self,
920 interval: Duration,
921 notify: impl Fn() + Send + 'static,
922 ) -> Result<Background> {
923 let Some(root) = self.root.clone() else {
924 return Err(io::Error::new(
925 io::ErrorKind::Unsupported,
926 "Sections on a share sync through Background::smb",
927 )
928 .into());
929 };
930 Background::start(
931 interval,
932 move || {
933 let root = root.clone();
934 Ok(move |path: &str| FileRemote(root.join(path)))
935 },
936 notify,
937 )
938 }
939
884940 /// Opens a section of a mounted notebook by its catalog path.
885941 pub fn section(&self, path: &str, notify: impl Fn() + Send + 'static) -> Result<Section> {
886942 self.section_with(path, |file| Ok(FileRemote(file.to_owned())), notify)
......@@ -1010,6 +1066,21 @@ fn catalog_path(folder: &str, name: &str) -> String {
10101066 }
10111067}
10121068
1069/// The replica in `cache` of the section `identity` names.
1070fn replica_file(cache: &Path, identity: &[u8; 16]) -> PathBuf {
1071 let name: String = identity.iter().map(|byte| format!("{byte:02x}")).collect();
1072 cache.join(format!("{name}.sqlite"))
1073}
1074
1075/// Why a synchronization step did not reach the section file, as the host shows it.
1076pub(crate) fn reached(error: &Error) -> io::Error {
1077 match error {
1078 Error::RemoteIo(error) => io::Error::new(error.kind(), error.to_string()),
1079 Error::Remote(error) => io::Error::new(error.error.kind(), error.to_string()),
1080 error => io::Error::other(error.to_string()),
1081 }
1082}
1083
10131084/// What happened to the section since the last poll.
10141085#[derive(Debug)]
10151086pub enum Event {
......@@ -1079,13 +1150,8 @@ impl Section {
10791150 let source = connect(&file)?.read()?;
10801151 let store = Store::parse(&source)?;
10811152 let identity = RevisionIndex::parse(&store)?.root;
1082 let name: String = identity
1083 .guid
1084 .iter()
1085 .map(|byte| format!("{byte:02x}"))
1086 .collect();
10871153 std::fs::create_dir_all(&cache)?;
1088 let cache = cache.as_ref().join(format!("{name}.sqlite"));
1154 let cache = replica_file(cache.as_ref(), &identity.guid);
10891155 let replica = if cache.exists() {
10901156 Replica::open(&cache)?
10911157 } else {
......@@ -1149,11 +1215,7 @@ impl Section {
11491215 let (sender, notify) = (sender.clone(), Arc::clone(&notify));
11501216 let observed = Arc::clone(&observed);
11511217 replica.start_sync(Duration::from_secs(2), connect, move |result| {
1152 let error = result.as_ref().err().map(|error| match error {
1153 Error::RemoteIo(error) => io::Error::new(error.kind(), error.to_string()),
1154 Error::Remote(error) => io::Error::new(error.error.kind(), error.to_string()),
1155 error => io::Error::other(error.to_string()),
1156 });
1218 let error = result.as_ref().err().map(reached);
11571219 // Reaching the file as the last attempt did is no news to the host.
11581220 let mut changed = false;
11591221 if let Ok(mut observed) = observed.lock() {
crates/notebook/src/smb/remote.rs+5-4
......@@ -1,20 +1,21 @@
11use super::Client;
22use crate::Remote;
33use onestore::{CommitError, Stamp, Transaction};
4use std::io;
4use std::{io, sync::Arc};
55
66/// Binds every reconciliation operation to one share-relative file and read limit.
77/// Connection loss retires the client; reconnect before subsequent sync attempts.
88pub struct SmbRemote {
9 client: Client,
9 client: Arc<Client>,
1010 path: String,
1111 limit: usize,
1212}
1313
1414impl SmbRemote {
15 pub fn new(client: Client, path: impl Into<String>, limit: usize) -> Self {
15 /// A client shared between remotes serves each of their files over one connection.
16 pub fn new(client: impl Into<Arc<Client>>, path: impl Into<String>, limit: usize) -> Self {
1617 Self {
17 client,
18 client: client.into(),
1819 path: path.into(),
1920 limit,
2021 }
crates/notebook/src/worker.rs+17-12
......@@ -12,16 +12,28 @@ use std::{
1212};
1313
1414pub(super) struct Signal {
15 stopped: AtomicBool,
16 offline: AtomicBool,
15 pub(crate) stopped: AtomicBool,
16 pub(crate) offline: AtomicBool,
1717 /// FILETIME of the last step or poll that reached the remote; 0 before one has.
1818 synced: AtomicU64,
19 /// Asked for by `SyncWorker::wake`: steps run while offline until one leaves nothing to do.
20 requested: AtomicBool,
19 /// Asked for by `wake`: steps run while offline until one leaves nothing to do.
20 pub(crate) requested: AtomicBool,
2121 sender: SyncSender<()>,
2222}
2323
2424impl Signal {
25 pub(crate) fn new() -> (Arc<Self>, mpsc::Receiver<()>) {
26 let (sender, receiver) = mpsc::sync_channel(1);
27 let signal = Self {
28 stopped: AtomicBool::new(false),
29 offline: AtomicBool::new(false),
30 synced: AtomicU64::new(0),
31 requested: AtomicBool::new(false),
32 sender,
33 };
34 (Arc::new(signal), receiver)
35 }
36
2537 pub(super) fn wake(&self) {
2638 // One retained notification covers edits that arrive during network I/O.
2739 let _ = self.sender.try_send(());
......@@ -110,14 +122,7 @@ impl Replica {
110122 if owner.upgrade().is_some() {
111123 return Err(io::ErrorKind::WouldBlock.into());
112124 }
113 let (sender, receiver) = mpsc::sync_channel(1);
114 let signal = Arc::new(Signal {
115 stopped: AtomicBool::new(false),
116 offline: AtomicBool::new(false),
117 synced: AtomicU64::new(0),
118 requested: AtomicBool::new(false),
119 sender,
120 });
125 let (signal, receiver) = Signal::new();
121126 let replica = Arc::clone(self);
122127 let worker_signal = Arc::clone(&signal);
123128 let thread = thread::Builder::new()
crates/notebook/tests/session.rs+126
......@@ -999,3 +999,129 @@ fn a_missing_section_file_reports_its_error_until_it_returns() {
999999 assert!(notified.load(Ordering::SeqCst) > before);
10001000 section.close().unwrap();
10011001}
1002
1003/// A notebook folder holding `First.one` and `Second.one`, and a cache beside it.
1004fn two_sections(directory: &Path) -> Notebook {
1005 let root = directory.join("Shared");
1006 std::fs::create_dir(&root).unwrap();
1007 for name in ["First", "Second"] {
1008 let file = format!("{name}.one");
1009 std::fs::write(
1010 root.join(&file),
1011 onestore::create_section(&file, name, "Author").unwrap(),
1012 )
1013 .unwrap();
1014 }
1015 Notebook::open(&root, directory.join("cache")).unwrap()
1016}
1017
1018fn until(what: &str, mut accept: impl FnMut() -> bool) {
1019 let deadline = Instant::now() + Duration::from_secs(20);
1020 while !accept() {
1021 assert!(Instant::now() < deadline, "{what}");
1022 std::thread::sleep(Duration::from_millis(20));
1023 }
1024}
1025
1026#[test]
1027fn a_closed_sections_queued_edits_publish_in_the_background_when_online() {
1028 let directory = tempfile::tempdir().unwrap();
1029 let notebook = two_sections(directory.path());
1030 let file = directory.path().join("Shared/Second.one");
1031 let background = notebook
1032 .background(Duration::from_millis(20), || {})
1033 .unwrap();
1034 background.set_offline(true);
1035 background.watch(notebook.replicas());
1036 // The round started with the background ends before the edit exists.
1037 std::thread::sleep(Duration::from_millis(200));
1038 let section = notebook.section("Second.one", || {}).unwrap();
1039 section.set_offline(true);
1040 let space = section.pages().unwrap()[0].0;
1041 let before = section.page(space).unwrap();
1042 typed(&section, space, &before, 0..0, "Closed ");
1043 section.close().unwrap();
1044 std::thread::sleep(Duration::from_millis(300));
1045 assert_eq!(
1046 stored_page(&file, space),
1047 before,
1048 "working offline publishes nothing"
1049 );
1050 background.set_offline(false);
1051 let mut after = edited(&before, "Closed ");
1052 until("the closed section published", || {
1053 let stored = stored_page(&file, space);
1054 after.title = stored.title.clone();
1055 stored == after
1056 });
1057 until("the status shows nothing waiting", || {
1058 background.status().iter().any(|(path, status)| {
1059 path == "Second.one" && status.synced.is_some() && status.queued == 0
1060 })
1061 });
1062 drop(background);
1063 let mut reopened = None;
1064 until("the background released the replica", || {
1065 reopened = notebook.section("Second.one", || {}).ok();
1066 reopened.is_some()
1067 });
1068 let reopened = reopened.unwrap();
1069 assert!(reopened.pending().unwrap().is_empty());
1070 reopened.close().unwrap();
1071}
1072
1073#[test]
1074fn a_remote_change_to_a_closed_section_is_noticed_and_rebases_its_replica() {
1075 let directory = tempfile::tempdir().unwrap();
1076 let notebook = two_sections(directory.path());
1077 // Second has a replica from being opened once; First has never been opened.
1078 notebook
1079 .section("Second.one", || {})
1080 .unwrap()
1081 .close()
1082 .unwrap();
1083 let notified = Arc::new(AtomicUsize::new(0));
1084 let counter = Arc::clone(&notified);
1085 let background = notebook
1086 .background(Duration::from_millis(20), move || {
1087 counter.fetch_add(1, Ordering::SeqCst);
1088 })
1089 .unwrap();
1090 background.watch(notebook.replicas());
1091 until("both sections were reached", || {
1092 background
1093 .status()
1094 .iter()
1095 .all(|(_, status)| status.synced.is_some() && status.error.is_none())
1096 });
1097 assert!(background.changed().is_empty());
1098 let mut changed = Vec::new();
1099 for name in ["First.one", "Second.one"] {
1100 let file = directory.path().join("Shared").join(name);
1101 let bytes = onestore::read_file(&file).unwrap();
1102 let space = notebook::session::stored_pages(&bytes).unwrap()[0].space;
1103 let native = edited(&model_ops::page_of(&bytes, space), "Native ");
1104 ops::save(&bytes, space, &native)
1105 .unwrap()
1106 .commit_file(&file)
1107 .unwrap();
1108 changed.push((name, space, native));
1109 }
1110 let mut noticed = Vec::new();
1111 until("both changes were noticed", || {
1112 noticed.extend(background.changed());
1113 changed
1114 .iter()
1115 .all(|(name, ..)| noticed.iter().any(|path| path == name))
1116 });
1117 assert!(notified.load(Ordering::SeqCst) > 0);
1118 drop(background);
1119 let (_, space, native) = &changed[1];
1120 let replica = notebook.replica_path("Second.one").unwrap();
1121 let mut opened = None;
1122 until("the background released the replica", || {
1123 opened = notebook::Replica::open(&replica).ok();
1124 opened.is_some()
1125 });
1126 assert_same(opened.unwrap().page(*space).unwrap(), native);
1127}
crates/snowbound/src/library.rs+99-48
......@@ -1,17 +1,17 @@
11//! Open notebooks: where each lives and the sections its tabs offer.
22
33use notebook::discover::{Folder, SectionState};
4use notebook::session::{Notebook, Section};
4use notebook::session::{Background, Notebook, Section};
55use notebook::smb::{Client, Credentials};
66use std::{
77 error::Error,
88 io,
99 path::{Path, PathBuf},
1010 sync::{
11 Arc,
11 Arc, Mutex, OnceLock,
1212 atomic::{AtomicBool, Ordering},
1313 },
14 time::Duration,
14 time::{Duration, Instant},
1515};
1616
1717/// OneNote's Work Offline, which like OneNote's holds for every notebook.
......@@ -26,6 +26,23 @@ pub fn set_offline(offline: bool) {
2626 OFFLINE.store(offline, Ordering::Relaxed);
2727}
2828
29/// Wakes the app when a notebook's closed sections report, once set.
30static NOTIFY: OnceLock<Box<dyn Fn() + Send + Sync>> = OnceLock::new();
31
32pub fn on_background(notify: impl Fn() + Send + Sync + 'static) {
33 let _ = NOTIFY.set(Box::new(notify));
34}
35
36/// How often the sections no tab shows are polled; OneNote 2010 read a closed section
37/// fifteen seconds after another client changed it.
38const BACKGROUND: Duration = Duration::from_secs(15);
39
40fn notify_background() {
41 if let Some(notify) = NOTIFY.get() {
42 notify();
43 }
44}
45
2946/// How long an SMB request may take before the share counts as unreachable.
3047const TIMEOUT: Duration = Duration::from_secs(10);
3148/// The largest section file read whole over SMB.
......@@ -151,6 +168,9 @@ pub struct Library {
151168 /// Why a notebook on a mounted SMB share opened through the mount instead.
152169 pub notice: Option<String>,
153170 cache: PathBuf,
171 /// Syncs the notebook's sections while no tab shows them, as OneNote syncs every
172 /// section of an open notebook.
173 pub background: Option<Arc<Background>>,
154174}
155175
156176impl Library {
......@@ -171,12 +191,12 @@ impl Library {
171191 }
172192 }
173193 }
194 let notebook = Notebook::open(location, cache);
174195 Self {
175196 location: location.to_owned(),
176197 name: file_name(Path::new(location)),
177 notebook: Notebook::open(location, cache)
178 .map(Some)
179 .map_err(|error| error.to_string()),
198 background: notebook.as_ref().ok().and_then(local_background),
199 notebook: notebook.map(Some).map_err(|error| error.to_string()),
180200 server: None,
181201 notice,
182202 cache: cache.to_owned(),
......@@ -195,6 +215,17 @@ impl Library {
195215 let client = server.connect().map_err(|error| error.to_string())?;
196216 let notebook = Notebook::open_smb(Arc::new(client), &server.mount.root, cache)
197217 .map_err(|error| error.to_string())?;
218 let connect = Arc::clone(&server);
219 let background = Background::smb(
220 &server.mount.root,
221 LIMIT,
222 BACKGROUND,
223 move || connect.connect(),
224 notify_background,
225 )
226 .map_err(|error| error.to_string())?;
227 background.set_offline(offline());
228 background.watch(notebook.replicas());
198229 Ok(Self {
199230 location: location.to_owned(),
200231 name: file_name(Path::new(location)),
......@@ -202,6 +233,7 @@ impl Library {
202233 server: Some(server),
203234 notice: None,
204235 cache: cache.to_owned(),
236 background: Some(Arc::new(background)),
205237 })
206238 }
207239
......@@ -217,6 +249,9 @@ impl Library {
217249
218250 /// This notebook as `notebook`, read again after a change.
219251 pub fn with(&self, notebook: Notebook) -> Self {
252 if let Some(background) = &self.background {
253 background.watch(notebook.replicas());
254 }
220255 Self {
221256 location: self.location.clone(),
222257 name: self.name.clone(),
......@@ -224,6 +259,7 @@ impl Library {
224259 server: self.server.clone(),
225260 notice: self.notice.clone(),
226261 cache: self.cache.clone(),
262 background: self.background.clone(),
227263 }
228264 }
229265
......@@ -232,6 +268,7 @@ impl Library {
232268 Self {
233269 location: location.to_owned(),
234270 name: file_name(Path::new(location)),
271 background: local_background(&notebook),
235272 notebook: Ok(Some(notebook)),
236273 server: None,
237274 notice: None,
......@@ -248,6 +285,7 @@ impl Library {
248285 server: None,
249286 notice: None,
250287 cache: cache.to_owned(),
288 background: None,
251289 }
252290 }
253291
......@@ -258,31 +296,50 @@ impl Library {
258296 path: &str,
259297 notify: impl Fn() + Send + 'static,
260298 ) -> Result<Section, Box<dyn Error>> {
261 let section = match (&self.notebook, &self.server) {
262 (Ok(Some(notebook)), Some(server)) => {
263 let identity = self
264 .catalog_section(notebook.catalog(), path)
265 .ok_or("The notebook doesn’t list this section")?;
266 let file = match server.mount.root.as_str() {
267 "" => path.to_owned(),
268 root => format!("{root}/{path}"),
269 };
270 let cache = self.cache.join("smb").join(format!("{identity}.sqlite"));
271 std::fs::create_dir_all(self.cache.join("smb"))?;
272 let replica = if cache.exists() {
273 notebook::Replica::open(&cache)?
274 } else {
275 notebook::Replica::create(
276 &cache,
277 &server.connect()?.read_storage(&file, LIMIT)?,
278 )?
279 };
280 let server = Arc::clone(server);
281 Section::resume_smb(file, replica, LIMIT, move || server.connect(), notify)?
299 // The background may hold the replica for a step, which takes a network round trip.
300 let deadline = Instant::now() + TIMEOUT * 3;
301 let notify = Arc::new(Mutex::new(notify));
302 let notifier = || {
303 let notify = Arc::clone(&notify);
304 move || {
305 if let Ok(notify) = notify.lock() {
306 notify();
307 }
308 }
309 };
310 let section = loop {
311 let opened = match (&self.notebook, &self.server) {
312 (Ok(Some(notebook)), Some(server)) => {
313 let file = match server.mount.root.as_str() {
314 "" => path.to_owned(),
315 root => format!("{root}/{path}"),
316 };
317 let cache = notebook.replica_path(path)?;
318 std::fs::create_dir_all(self.cache.join("smb"))?;
319 let replica = if cache.exists() {
320 notebook::Replica::open(&cache)
321 } else {
322 notebook::Replica::create(
323 &cache,
324 &server.connect()?.read_storage(&file, LIMIT)?,
325 )
326 };
327 replica.and_then(|replica| {
328 let server = Arc::clone(server);
329 let connect = move || server.connect();
330 Section::resume_smb(file, replica, LIMIT, connect, notifier())
331 })
332 }
333 (Ok(Some(notebook)), None) => notebook.section(path, notifier()),
334 (Ok(None), _) => Section::open(path, &self.cache, notifier()),
335 (Err(error), _) => return Err(error.clone().into()),
336 };
337 match opened {
338 Err(error) if error.busy() && Instant::now() < deadline => {
339 std::thread::sleep(Duration::from_millis(50));
340 }
341 opened => break opened?,
282342 }
283 (Ok(Some(notebook)), None) => notebook.section(path, notify)?,
284 (Ok(None), _) => Section::open(path, &self.cache, notify)?,
285 (Err(error), _) => return Err(error.clone().into()),
286343 };
287344 section.set_offline(offline());
288345 Ok(section)
......@@ -311,24 +368,6 @@ impl Library {
311368 }
312369 }
313370
314 /// The file identity, in hex, of the section at catalog `path`.
315 fn catalog_section(&self, catalog: &Folder, path: &str) -> Option<String> {
316 let mut folders = vec![catalog];
317 while let Some(folder) = folders.pop() {
318 if let Some(section) = folder.sections.iter().find(|section| section.path == path) {
319 return Some(
320 section
321 .file_id
322 .iter()
323 .map(|byte| format!("{byte:02x}"))
324 .collect(),
325 );
326 }
327 folders.extend(&folder.groups);
328 }
329 None
330 }
331
332371 /// Names the section at catalog `path` across every open notebook, for remembering
333372 /// its pages.
334373 pub fn key(&self, path: &str) -> String {
......@@ -407,6 +446,17 @@ pub fn recycle_bin(path: &str) -> bool {
407446 path.rsplit('/').next() == Some("OneNote_RecycleBin")
408447}
409448
449/// A mounted notebook's background sync, following Work Offline.
450fn local_background(notebook: &Notebook) -> Option<Arc<Background>> {
451 let background = notebook
452 .background(BACKGROUND, notify_background)
453 .inspect_err(|error| eprintln!("Background sync did not start: {error}"))
454 .ok()?;
455 background.set_offline(offline());
456 background.watch(notebook.replicas());
457 Some(Arc::new(background))
458}
459
410460pub fn file_name(path: &Path) -> String {
411461 path.canonicalize()
412462 .ok()
......@@ -617,6 +667,7 @@ mod tests {
617667 })),
618668 notice: None,
619669 cache: PathBuf::new(),
670 background: None,
620671 };
621672 let shown = |root, file| library(root).local(Path::new(file));
622673 let under = Some(PathBuf::from("/Volumes/agent/lab/Group/New Section 1.one"));
crates/snowbound/src/main.rs+31-6
......@@ -662,6 +662,10 @@ impl State {
662662 }
663663 let layouts = Arc::new(Mutex::new(engine.clone()));
664664 let temporary = matches!(input, Input::Notes { .. } | Input::Page(_));
665 let background = proxy.clone();
666 library::on_background(move || {
667 let _ = background.send_event(UserEvent::Sync);
668 });
665669 let mut notebooks = Vec::new();
666670 let mut session = None;
667671 let mut sectionless = None;
......@@ -2429,7 +2433,19 @@ impl State {
24292433 Some(library) => *library = Arc::clone(&session.library),
24302434 None => self.notebooks.push(Arc::clone(&session.library)),
24312435 }
2432 self.session = Some(*session);
2436 // The notebook's background sync takes the section over once it is closed.
2437 if let Some(previous) = self.session.replace(*session) {
2438 let background = previous.library.background.clone();
2439 let section = previous.section;
2440 std::thread::spawn(move || {
2441 if let Err(error) = section.close() {
2442 eprintln!("Synchronization stopped: {error}");
2443 }
2444 if let Some(background) = background {
2445 background.wake();
2446 }
2447 });
2448 }
24332449 self.sectionless = None;
24342450 if other_notebook {
24352451 self.save_settings();
......@@ -2769,12 +2785,24 @@ impl State {
27692785 Ok(())
27702786 }
27712787
2772 /// Applies what the section reported since the last poll.
2788 /// Applies what the section, and each notebook's closed sections, reported since the
2789 /// last poll.
27732790 fn synced(&mut self) -> Result<(), Box<dyn Error>> {
2791 // Every notebook's changes are taken, so none is reported again.
2792 let noticed = self
2793 .notebooks
2794 .iter()
2795 .filter_map(|library| library.background.as_ref())
2796 .filter(|background| !background.changed().is_empty())
2797 .count();
2798 if noticed > 0 {
2799 self.sync_index(true);
2800 }
2801 // Reports are rare: a status changed, or a section did.
2802 self.window.request_redraw();
27742803 let Some(session) = &mut self.session else {
27752804 return Ok(());
27762805 };
2777 let shown = sync::label(&session.sync);
27782806 let mut listed = false;
27792807 let mut changed = false;
27802808 let mut rejected = None;
......@@ -2802,9 +2830,6 @@ impl State {
28022830 }
28032831 }
28042832 session.sync = session.section.sync_status()?;
2805 if listed || changed || shown != sync::label(&session.sync) {
2806 self.window.request_redraw();
2807 }
28082833 if listed {
28092834 let back = session.shown;
28102835 session.pages = session.section.pages()?;
crates/snowbound/src/sync.rs+120-19
......@@ -1,5 +1,5 @@
11//! The sync status at the top right, and the popup it opens: OneNote's Shared Notebook
2//! Synchronization for the open section.
2//! Synchronization for the open section's notebook and each of its sections.
33
44use crate::{Session, State, art, filetime, library, platform};
55use notebook::session::SyncStatus;
......@@ -56,8 +56,71 @@ fn describe(sync: &SyncStatus) -> (&'static str, &'static [&'static str], Option
5656 }
5757}
5858
59pub(crate) fn label(sync: &SyncStatus) -> &'static str {
60 describe(sync).0
59/// Labels from the best state to the worst, which a notebook's status shows.
60const ORDER: [&str; 5] = [
61 "Up to date",
62 "Syncing…",
63 "Section in use",
64 "Not connected",
65 "Unable to sync",
66];
67
68fn rank(sync: &SyncStatus) -> usize {
69 let label = describe(sync).0;
70 ORDER.iter().position(|shown| *shown == label).unwrap_or(0)
71}
72
73/// Every section of the open section's notebook with its status, in catalog order: the
74/// open one's from its session, the others' from the notebook's background sync.
75fn sections(session: &Session) -> Vec<(String, SyncStatus)> {
76 let open = &session.tabs[session.tab].path;
77 let copy = |sync: &SyncStatus| SyncStatus {
78 synced: sync.synced,
79 error: sync
80 .error
81 .as_ref()
82 .map(|error| std::io::Error::new(error.kind(), error.to_string())),
83 queued: sync.queued,
84 };
85 let mut sections = session
86 .library
87 .background
88 .as_ref()
89 .map(|background| background.status())
90 .unwrap_or_default();
91 match sections.iter_mut().find(|(path, _)| path == open) {
92 Some((_, sync)) => *sync = copy(&session.sync),
93 None => sections.insert(0, (open.clone(), copy(&session.sync))),
94 }
95 sections
96}
97
98/// The notebook as a whole: the worst section's error, every waiting change, and the
99/// oldest last sync.
100fn overall(sections: &[(String, SyncStatus)]) -> SyncStatus {
101 let worst = sections
102 .iter()
103 .map(|(_, sync)| sync)
104 .filter(|sync| sync.error.is_some())
105 .max_by_key(|sync| rank(sync));
106 SyncStatus {
107 synced: sections
108 .iter()
109 .map(|(_, sync)| sync.synced)
110 .collect::<Option<Vec<_>>>()
111 .and_then(|times| times.into_iter().min()),
112 error: worst
113 .and_then(|sync| sync.error.as_ref())
114 .map(|error| std::io::Error::new(error.kind(), error.to_string())),
115 queued: sections.iter().map(|(_, sync)| sync.queued).sum(),
116 }
117}
118
119fn changes(count: u64) -> String {
120 match count {
121 1 => "1 change".to_owned(),
122 count => format!("{count} changes"),
123 }
61124}
62125
63126/// A FILETIME as OneNote's "Last sync": the time, and the date too before today.
......@@ -73,8 +136,9 @@ fn when(time: u64) -> String {
73136
74137/// The status's icon in the toolbar, named in its tooltip, which opens the popup.
75138pub(crate) fn control(ui: &mut Ui, session: &Session, theme: &Theme) {
76 let strong = ui.popup_open(id()) || session.sync.error.is_some() && !library::offline();
77 let (label, icon, _) = describe(&session.sync);
139 let sync = overall(&sections(session));
140 let strong = ui.popup_open(id()) || sync.error.is_some() && !library::offline();
141 let (label, icon, _) = describe(&sync);
78142 ui.open_as(
79143 button(),
80144 Spec {
......@@ -102,7 +166,7 @@ pub(crate) fn control(ui: &mut Ui, session: &Session, theme: &Theme) {
102166impl State {
103167 /// Builds the sync status popup while it is open.
104168 pub(crate) fn sync_popup(&mut self) -> Result<(), Box<dyn Error>> {
105 let ui = &mut self.ui;
169 let (ui, notebooks) = (&mut self.ui, &self.notebooks);
106170 let Some(session) = &mut self.session else {
107171 ui.close_popup(id());
108172 return Ok(());
......@@ -114,7 +178,9 @@ impl State {
114178 session.sync = session.section.sync_status()?;
115179 let theme = ui.theme.clone();
116180 let offline = library::offline();
117 let (progress, _, advice) = describe(&session.sync);
181 let sections = sections(session);
182 let sync = overall(&sections);
183 let (progress, _, advice) = describe(&sync);
118184 ui.open_as(
119185 id(),
120186 Spec {
......@@ -177,23 +243,54 @@ impl State {
177243 row(ui, "Notebook", &session.library.name);
178244 row(ui, "Location", &session.library.location);
179245 row(ui, "Connection", &session.library.transport());
180 row(ui, "Section", &session.tabs[session.tab].name);
181246 row(ui, "Progress", progress);
182 if let Some(synced) = session.sync.synced {
247 if let Some(synced) = sync.synced {
183248 row(ui, "Last sync", &when(synced));
184249 }
185 if session.sync.queued > 0 {
186 let queued = match session.sync.queued {
187 1 => "1 change".to_owned(),
188 count => format!("{count} changes"),
189 };
190 row(ui, "Not yet synced", &queued);
250 if sync.queued > 0 {
251 row(ui, "Not yet synced", &changes(sync.queued));
191252 }
192253 if let Some(advice) = advice {
193254 text(ui, "advice", advice, theme.text, false);
194255 }
195 if let Some(error) = &session.sync.error {
196 text(ui, "error", &error.to_string(), theme.text_dim, false);
256 text(ui, "sections", "Sections", theme.text, true);
257 for (index, (path, sync)) in sections.iter().enumerate() {
258 ui.open(
259 format!("section-{index}"),
260 Spec {
261 size: [fill(), children()],
262 gap: 8.0,
263 ..Spec::default()
264 },
265 );
266 let name = path.strip_suffix(".one").unwrap_or(path);
267 ui.leaf(
268 "name",
269 Spec {
270 size: [fill(), px(theme.font_size * 1.6)],
271 text: Some(name),
272 overflow: Overflow::Ellipsis,
273 ..Spec::default()
274 },
275 );
276 let status = match sync.queued {
277 0 => describe(sync).0.to_owned(),
278 queued => format!("{}, {}", describe(sync).0, changes(queued)),
279 };
280 ui.leaf(
281 "status",
282 Spec {
283 size: [fit(), px(theme.font_size * 1.6)],
284 text: Some(&status),
285 color: Some(theme.text_dim),
286 ..Spec::default()
287 },
288 );
289 ui.close();
290 if let Some(error) = &sync.error {
291 let part = format!("section-{index}-error");
292 text(ui, &part, &error.to_string(), theme.text_dim, false);
293 }
197294 }
198295 if let Some(notice) = &session.library.notice {
199296 let notice = format!("Snowbound’s SMB client couldn’t sign in: {notice}");
......@@ -239,12 +336,16 @@ impl State {
239336 let now = ui::button(ui, "sync-now", "Sync Now").clicked;
240337 ui.close();
241338 ui.close();
339 let backgrounds = notebooks
340 .iter()
341 .filter_map(|library| library.background.as_ref());
242342 if toggled {
243343 library::set_offline(!offline);
244344 session.section.set_offline(!offline);
245 }
246 if now {
345 backgrounds.for_each(|background| background.set_offline(!offline));
346 } else if now {
247347 session.section.wake();
348 backgrounds.for_each(|background| background.wake());
248349 }
249350 if let Some(file) = file.filter(|_| show) {
250351 platform::show_file(&file);