1// One Win32 call, where std has no safe form: `session`'s hidden attribute.
2#![deny(unsafe_code)]
3#![doc = include_str!("../README.md")]
4
5pub mod discover;
6pub mod fs;
7#[cfg(feature = "live")]
8#[cfg_attr(target_arch = "wasm32", path = "live/web.rs")]
9pub mod live;
10#[cfg(feature = "smb")]
11pub mod smb;
12
13use fs::OpenOptions;
14use onestore::{
15 ExGuid,
16 op::{Edit, OpError},
17 page::Page,
18 protected::Key,
19};
20use rusqlite::{Connection, OpenFlags, TransactionBehavior, params};
21use std::{
22 io,
23 path::{Path, PathBuf},
24 sync::{Arc, Mutex, MutexGuard, mpsc},
25 time::Duration,
26};
27
28mod assets;
29mod background;
30mod base;
31pub mod location;
32mod merge;
33mod migrate;
34pub mod package;
35mod queue;
36mod recovery;
37mod resolve;
38pub use resolve::Version;
39mod schema;
40pub mod session;
41pub use recovery::{Recovery, RecoverySummary};
42pub mod sidecar;
43mod sync;
44pub use sync::{EditStatus, Remote, Synced};
45mod task;
46mod worker;
47mod working;
48#[cfg(feature = "smb")]
49pub use smb::SmbRemote;
50pub use worker::SyncWorker;
51
52/// The largest file a notebook session or Live Share reads or writes whole.
53pub const MAX_FILE_BYTES: usize = 256 << 20;
54
55#[derive(Debug, thiserror::Error)]
56pub enum Error {
57 #[error(transparent)]
58 Database(#[from] rusqlite::Error),
59 #[error(transparent)]
60 Io(#[from] io::Error),
61 #[error(transparent)]
62 Document(#[from] onestore::Error),
63 #[error(transparent)]
64 Remote(#[from] onestore::CommitError),
65 #[error(transparent)]
66 RemoteIo(io::Error),
67 #[error(transparent)]
68 Discovery(#[from] discover::Error),
69 /// The section refused an edit; the pages it names are as they were.
70 #[error(transparent)]
71 Rejected(#[from] OpError),
72 #[error("External payload identity now refers to different bytes")]
73 AssetChanged,
74 /// A password-protected section's password did not match, or its format is unknown.
75 #[error(transparent)]
76 Protected(#[from] onestore::protected::Error),
77}
78
79impl Error {
80 /// Whether the replica is open elsewhere, as in another section session of this process.
81 pub fn busy(&self) -> bool {
82 matches!(self, Self::Database(rusqlite::Error::SqliteFailure(error, _))
83 if error.code == rusqlite::ErrorCode::DatabaseBusy)
84 }
85}
86
87type Result<T> = std::result::Result<T, Error>;
88
89const APPLICATION_ID: u32 = 0x4f4e454f;
90
91/// A locally durable edit; its ID remains stable across cache reopen.
92#[derive(Debug, Clone, PartialEq)]
93pub struct PendingEdit {
94 pub id: u64,
95 pub author: String,
96 pub edit: Edit,
97}
98
99/// How an uncertain attempt ends after review: `Mine` publishes it again, `Theirs` drops
100/// the whole unpublished branch.
101#[derive(Debug, Clone, Copy, PartialEq, Eq)]
102pub enum Resolution {
103 Mine,
104 Theirs,
105}
106
107/// Owns one local cache. Share this handle between threads; a second open fails busy.
108/// Edits apply on the cache's section thread, which keeps the section parsed; SQLite's
109/// exclusive connection retains ownership between local transactions.
110pub struct Replica {
111 section: Arc<working::Thread>,
112 synchronization: Mutex<()>,
113 /// The section's root object space, which names the document.
114 root: ExGuid,
115 /// Last, so that it runs once the fields above have closed the cache.
116 released: Released,
117}
118
119/// Runs once the replica is released and its cache closed, as `Background::hold` asks.
120#[derive(Default)]
121struct Released(Mutex<Option<Box<dyn FnOnce() + Send>>>);
122
123impl Drop for Released {
124 fn drop(&mut self) {
125 if let Some(released) = self.0.get_mut().ok().and_then(Option::take) {
126 released();
127 }
128 }
129}
130
131impl Replica {
132 /// Seeds a new cache from a validated section image, refusing any existing path.
133 /// An initialization error preserves the created file for inspection.
134 pub fn create(path: impl AsRef<Path>, source: &[u8]) -> Result<Self> {
135 Self::seed(path.as_ref(), source, None)?;
136 Self::start(cache_connection(path.as_ref())?, None)
137 }
138
139 /// Opens the cache at `path`, a protected section's under `key`, first creating it from
140 /// the image `source` reads where there is none. A cache another thread creates meanwhile,
141 /// as the background makes an offline copy, is opened instead.
142 pub fn open_or_create(
143 path: impl AsRef<Path>,
144 key: Option<&Key>,
145 source: impl FnOnce() -> Result<Vec<u8>>,
146 ) -> Result<Self> {
147 let path = path.as_ref();
148 if fs::metadata(path).is_err() {
149 match Self::seed(path, &source()?, key) {
150 Err(Error::Io(error)) if error.kind() == io::ErrorKind::AlreadyExists => {}
151 seeded => seeded?,
152 }
153 }
154 Self::open_with(path, key.cloned())
155 }
156
157 /// `create` without opening the cache it made, as for an offline copy.
158 pub(crate) fn seed(path: &Path, source: &[u8], key: Option<&Key>) -> Result<()> {
159 validate(source, key)?;
160 // Built under a name of its own and linked into place whole, so that a cache that
161 // exists is complete, and two threads making the same one never share a build.
162 static BUILDS: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
163 let build = BUILDS.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
164 let mut building = path.as_os_str().to_owned();
165 building.push(format!(".creating-{}-{build}", fs::process_id()));
166 let building = PathBuf::from(building);
167 let built = Self::build(&building, source);
168 let linked = built.and_then(|()| Ok(fs::hard_link(&building, path)?));
169 for suffix in ["", "-wal", "-shm"] {
170 let mut file = building.as_os_str().to_owned();
171 file.push(suffix);
172 let _ = fs::remove_file(file);
173 }
174 linked
175 }
176
177 fn build(path: &Path, source: &[u8]) -> Result<()> {
178 let mut options = OpenOptions::new();
179 options.read(true).write(true).create_new(true);
180 #[cfg(unix)]
181 {
182 use std::os::unix::fs::OpenOptionsExt;
183 options.mode(0o600);
184 }
185 drop(options.open(path)?);
186 let mut connection = cache_connection(path)?;
187 write_ahead(&connection)?;
188 let transaction = connection.transaction_with_behavior(TransactionBehavior::Exclusive)?;
189 let tables: i64 =
190 transaction.query_row("SELECT count(*) FROM sqlite_schema", [], |row| row.get(0))?;
191 let application: u32 =
192 transaction.pragma_query_value(None, "application_id", |row| row.get(0))?;
193 if application != 0 || tables != 0 {
194 return Err(io::Error::new(
195 io::ErrorKind::AlreadyExists,
196 "Cache initialization found an existing database",
197 )
198 .into());
199 }
200 transaction.pragma_update(None, "application_id", APPLICATION_ID)?;
201 transaction.pragma_update(None, "user_version", schema::VERSION)?;
202 schema::create(&transaction)?;
203 base::write(&transaction, base::Image::Base, source)?;
204 transaction.commit()?;
205 Ok(connection.close().map_err(|(_, error)| error)?)
206 }
207
208 /// Reopens an existing cache and its durable pending edits without network access.
209 /// A schema-14 cache is converted once, after exporting it to `<path>.v14-recovery`;
210 /// a conversion that cannot reproduce every queued page leaves it untouched.
211 /// Unrecognized databases and unsupported journal modes are rejected.
212 pub fn open(path: impl AsRef<Path>) -> Result<Self> {
213 Self::open_with(path.as_ref(), None)
214 }
215
216 fn open_with(path: &Path, key: Option<Key>) -> Result<Self> {
217 let mut connection = cache_connection(path)?;
218 let application: u32 =
219 connection.pragma_query_value(None, "application_id", |row| row.get(0))?;
220 if application != APPLICATION_ID {
221 return Err(io::Error::new(io::ErrorKind::InvalidData, "Not a notebook cache").into());
222 }
223 let integrity: String = connection.query_row("PRAGMA quick_check", [], |row| row.get(0))?;
224 if integrity != "ok" {
225 return Err(
226 io::Error::new(io::ErrorKind::InvalidData, "Cache integrity check failed").into(),
227 );
228 }
229 let version: u32 = connection.pragma_query_value(None, "user_version", |row| row.get(0))?;
230 // A schema-14 cache converts in its rollback journal, so a failed conversion leaves
231 // the file as it was.
232 match version {
233 migrate::VERSION => migrate::migrate(&mut connection, path)?,
234 schema::PREVIOUS => schema::upgrade(&mut connection)?,
235 schema::VERSION => {}
236 _ => {
237 return Err(io::Error::new(
238 io::ErrorKind::InvalidData,
239 format!(
240 "Cache schema version {version} is not the supported version {}",
241 schema::VERSION
242 ),
243 )
244 .into());
245 }
246 }
247 write_ahead(&connection)?;
248 Self::start(connection, key)
249 }
250
251 fn start(connection: Connection, key: Option<Key>) -> Result<Self> {
252 let (section, root) = working::spawn(connection, key)?;
253 Ok(Self {
254 section,
255 synchronization: Mutex::new(()),
256 root,
257 released: Released::default(),
258 })
259 }
260
261 fn lock(&self) -> Result<MutexGuard<'_, Connection>> {
262 lock(&self.section.connection)
263 }
264
265 /// Hands `request` to the section thread.
266 fn send(&self, request: working::Request) -> Result<()> {
267 self.section.send(request)
268 }
269
270 /// Asks the section thread and waits for its answer.
271 fn ask<T: Send + 'static>(
272 &self,
273 request: impl FnOnce(working::Reply<T>) -> working::Request,
274 ) -> Result<T> {
275 let (sender, receiver) = mpsc::sync_channel(1);
276 self.send(request(Box::new(move |result| {
277 let _ = sender.send(result);
278 })))?;
279 receiver
280 .recv()
281 .map_err(|_| io::Error::other("The section thread stopped"))?
282 }
283
284 /// Applies an edit and queues it durably for publication, returning its id once
285 /// written. A refused edit returns `Rejected` and leaves every page as it was.
286 pub fn apply(&self, author: &str, edit: Edit) -> Result<u64> {
287 self.ask(|reply| working::Request::Apply {
288 author: author.to_owned(),
289 edit,
290 reply,
291 })
292 }
293
294 /// `apply` without waiting: `reply` runs on the section thread once the edit is
295 /// durable or refused.
296 pub(crate) fn submit(
297 &self,
298 author: &str,
299 edit: Edit,
300 reply: working::Reply<u64>,
301 ) -> Result<()> {
302 self.send(working::Request::Apply {
303 author: author.to_owned(),
304 edit,
305 reply,
306 })
307 }
308
309 /// The page in `space` as the queued edits leave it; O(page), for opening and reloading.
310 pub fn page(&self, space: ExGuid) -> Result<Page> {
311 self.ask(|reply| working::Request::Page { space, reply })
312 }
313
314 /// Page spaces, titles and outline levels (1 at the top) in section order.
315 pub fn pages(&self) -> Result<Vec<(ExGuid, String, u32)>> {
316 self.ask(|reply| working::Request::Pages { reply })
317 }
318
319 /// The conflict pages of each page that has them (`onestore::Section::conflicts`).
320 pub fn conflicts(&self) -> Result<Vec<(ExGuid, Vec<onestore::ConflictPage>)>> {
321 self.ask(|reply| working::Request::Conflicts { reply })
322 }
323
324 /// The versions of each page that has them (`onestore::Section::versions`).
325 pub fn versions(&self) -> Result<Vec<(ExGuid, Vec<onestore::PageVersion>)>> {
326 self.ask(|reply| working::Request::Versions { reply })
327 }
328
329 /// A page as one of its versions holds it; O(section).
330 pub fn version(&self, space: ExGuid, version: ExGuid) -> Result<Page> {
331 self.ask(|reply| working::Request::Version {
332 space,
333 version,
334 reply,
335 })
336 }
337
338 /// The section image the queued edits leave, the unsealed ones sealed as one more
339 /// revision whose identities differ per call: O(section).
340 pub fn snapshot(&self) -> Result<Vec<u8>> {
341 self.written()?;
342 working::image(&*self.lock()?, self.section.key.as_ref())
343 }
344
345 /// Waits until the edits applied before it are written to the queue, as reads answer
346 /// before an open burst of edits is.
347 pub fn written(&self) -> Result<()> {
348 self.ask(|reply| working::Request::Flush { reply })
349 }
350
351 /// The section file's identity, which internal links name as `section-id`.
352 pub fn identity(&self) -> Result<[u8; 16]> {
353 let stamp = base::base_stamp(&*self.lock()?)?;
354 Ok(onestore::Header::parse(&stamp.header)?.file_id)
355 }
356
357 /// Queued edits, oldest first.
358 pub fn pending(&self) -> Result<Vec<PendingEdit>> {
359 pending(&*self.lock()?, self.section.key.as_ref())
360 }
361
362 /// The edits putting each page the queue changed, as the queue leaves it, into another
363 /// section as a copy under fresh identities, by the author of the first edit changing it;
364 /// content outside the page model stays behind.
365 pub(crate) fn copies(&self) -> Result<Vec<(String, Edit)>> {
366 let arena = onestore::Arena::default();
367 let base = base::base(&*self.lock()?)?;
368 let base = working::open(&arena, base, self.section.key.as_ref())?;
369 let mut seen = std::collections::BTreeSet::from([self.root]);
370 let mut copies = Vec::new();
371 for queued in self.pending()? {
372 for space in queue::spaces(&queued.edit, self.root) {
373 if !seen.insert(space) {
374 continue;
375 }
376 let Ok(mut page) = self.page(space) else {
377 continue;
378 };
379 if base.page(space).is_ok_and(|base| base == page) {
380 continue;
381 }
382 page.objects
383 .retain(|object| !matches!(object, onestore::page::PageObject::Unsupported(_)));
384 let creation =
385 onestore::PageCreation::new(None, Some(&page.title), &queued.author)?;
386 let ops = vec![onestore::op::Op::Section(onestore::op::SectionOp::Import {
387 creation,
388 page: page.copy()?,
389 })];
390 copies.push((queued.author.clone(), Edit { at: now(), ops }));
391 }
392 }
393 Ok(copies)
394 }
395
396 fn wake_sync(&self) {
397 wake(&self.section.worker);
398 }
399}
400
401impl Drop for Replica {
402 fn drop(&mut self) {
403 self.section.stop();
404 }
405}
406
407fn lock(connection: &Mutex<Connection>) -> Result<MutexGuard<'_, Connection>> {
408 connection
409 .lock()
410 .map_err(|_| io::Error::other("Cache owner panicked").into())
411}
412
413/// Tells the worker a burst of local edits is durable.
414fn edited(worker: &Mutex<std::sync::Weak<worker::Signal>>) {
415 if let Ok(worker) = worker.lock()
416 && let Some(worker) = worker.upgrade()
417 {
418 worker.edited();
419 }
420}
421
422fn wake(worker: &Mutex<std::sync::Weak<worker::Signal>>) {
423 if let Ok(worker) = worker.lock()
424 && let Some(worker) = worker.upgrade()
425 {
426 worker.wake();
427 }
428}
429
430/// Fully validates a section image, a protected one under `key`, returning its root object
431/// space.
432fn validate(source: &[u8], key: Option<&Key>) -> Result<ExGuid> {
433 let arena = onestore::Arena::default();
434 Ok(match key {
435 Some(key) => onestore::Section::unlock(&arena, source.to_vec(), key)?.root(),
436 None => onestore::Section::open(&arena, source.to_vec())?.root(),
437 })
438}
439
440/// A closed cache, without opening the replica.
441fn closed(path: &Path) -> Result<Connection> {
442 let connection = cache_connection(path)?;
443 let application: u32 =
444 connection.pragma_query_value(None, "application_id", |row| row.get(0))?;
445 let version: u32 = connection.pragma_query_value(None, "user_version", |row| row.get(0))?;
446 if application != APPLICATION_ID || version != schema::VERSION {
447 return Err(io::Error::from(io::ErrorKind::InvalidData).into());
448 }
449 Ok(connection)
450}
451
452/// A closed cache's base stamp and how many edits wait.
453fn peek(connection: &Connection) -> Result<(onestore::Stamp, u64)> {
454 let base = base::base_stamp(connection)?;
455 let queued: i64 = connection.query_row("SELECT count(*) FROM edits", [], |row| row.get(0))?;
456 Ok((base, unsigned(queued)?))
457}
458
459fn pending(connection: &Connection, key: Option<&Key>) -> Result<Vec<PendingEdit>> {
460 queue::load(connection, key, None)
461}
462
463fn unsigned(value: i64) -> Result<u64> {
464 Ok(u64::try_from(value).map_err(|_| io::Error::from(io::ErrorKind::InvalidData))?)
465}
466
467fn signed(value: u64) -> Result<i64> {
468 Ok(i64::try_from(value).map_err(|_| io::Error::from(io::ErrorKind::InvalidInput))?)
469}
470
471/// Opens a cache without writing to it: exclusive locking, a full sync of every commit
472/// (`F_FULLFSYNC` on macOS, the WAL's syncs included) and foreign keys, each queried back.
473fn cache_connection(path: &Path) -> Result<Connection> {
474 #[cfg(target_arch = "wasm32")]
475 fs::install_sqlite();
476 let connection = Connection::open_with_flags(path, OpenFlags::SQLITE_OPEN_READ_WRITE)?;
477 connection.busy_timeout(Duration::ZERO)?;
478 connection.execute_batch(
479 "PRAGMA locking_mode=EXCLUSIVE; PRAGMA synchronous=FULL; PRAGMA fullfsync=ON; PRAGMA foreign_keys=ON;",
480 )?;
481 let locking: String = connection.pragma_query_value(None, "locking_mode", |row| row.get(0))?;
482 let journal: String = connection.pragma_query_value(None, "journal_mode", |row| row.get(0))?;
483 if locking != "exclusive" || !["wal", "delete"].contains(&journal.as_str()) {
484 return Err(io::Error::new(
485 io::ErrorKind::InvalidData,
486 "Unsupported cache locking or journal mode",
487 )
488 .into());
489 }
490 for (name, expected) in [("synchronous", 2), ("fullfsync", 1), ("foreign_keys", 1)] {
491 let actual: i64 = connection.pragma_query_value(None, name, |row| row.get(0))?;
492 if actual != expected {
493 return Err(io::Error::new(
494 io::ErrorKind::Unsupported,
495 "Required cache synchronization is unavailable",
496 )
497 .into());
498 }
499 }
500 Ok(connection)
501}
502
503/// Switches a cache to write-ahead logging: a commit appends its pages to `<cache>-wal`
504/// instead of copying the originals to a rollback journal. Under exclusive locking the WAL
505/// index lives in memory, so the WAL is the only file beside the cache; it is checkpointed
506/// and removed on close, and replayed on the next open after a crash.
507fn write_ahead(connection: &Connection) -> Result<()> {
508 let mode: String =
509 connection.pragma_update_and_check(None, "journal_mode", "WAL", |row| row.get(0))?;
510 if mode != "wal" {
511 return Err(io::Error::new(
512 io::ErrorKind::Unsupported,
513 "The cache cannot use write-ahead logging",
514 )
515 .into());
516 }
517 Ok(())
518}
519
520/// FILETIME now, as an edit's `at`: later than any this process took before, as Windows 7's
521/// clock ticks only every 15.6 ms and a later edit must still order after an earlier one.
522pub(crate) fn now() -> u64 {
523 use std::sync::atomic::{AtomicU64, Ordering};
524 static LAST: AtomicU64 = AtomicU64::new(0);
525 let unix = web_time::SystemTime::now()
526 .duration_since(web_time::UNIX_EPOCH)
527 .unwrap_or_default();
528 let clock =
529 (unix.as_secs() + 11_644_473_600) * 10_000_000 + u64::from(unix.subsec_nanos() / 100);
530 let previous = LAST
531 .fetch_update(Ordering::Relaxed, Ordering::Relaxed, |last| {
532 Some(clock.max(last + 1))
533 })
534 .unwrap_or_default();
535 clock.max(previous + 1)
536}