diff --git a/Cargo.lock b/Cargo.lock index 7aac8ee77b097d29ed8c4008224c767104b79f2c..c5e413c84686efa193d7964110b4e1c634b479ce 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -537,10 +537,12 @@ name = "onestore-offline" version = "0.1.0" dependencies = [ "onestore", + "onestore-notebook", "onestore-smb", "rusqlite", "serde", "serde_json", + "sha2", "tempfile", "thiserror", ] diff --git a/corpus/offline-v4/README.md b/corpus/offline-v4/README.md new file mode 100644 index 0000000000000000000000000000000000000000..efdba39aca39aa6d6e002e1d3289860e16678fec --- /dev/null +++ b/corpus/offline-v4/README.md @@ -0,0 +1,15 @@ +# Original schema-4 offline state + +These databases were produced by `onestore-offline` at revision `2244711b`, +before the external-media table and schema-5 migration existed. Tests copy the +live database before opening it; the archive remains read-only. + +[`producer.rs`](producer.rs) records the original producer. Run it against that +revision's core/offline crates, from a checkout whose `corpus/offline-v4` does +not exist. It uses the public native external-asset fixture, queues a text edit, +simulates a publication with a lost response, queues a dependent text edit, and +exports the complete recovery archive. Generated revision IDs vary between runs. + +[`provenance.json`](provenance.json) records the original source/database hashes +and the retained uncertain and pending edit IDs. No credentials or personal +notebook content are included. diff --git a/corpus/offline-v4/live.sqlite b/corpus/offline-v4/live.sqlite new file mode 100644 index 0000000000000000000000000000000000000000..130c341ba39987df0dfdc217633989160c2f23f6 Binary files /dev/null and b/corpus/offline-v4/live.sqlite differ diff --git a/corpus/offline-v4/producer.rs b/corpus/offline-v4/producer.rs new file mode 100644 index 0000000000000000000000000000000000000000..15f85f587e9811e170c8d63b828967881dc0a029 --- /dev/null +++ b/corpus/offline-v4/producer.rs @@ -0,0 +1,33 @@ +use onestore::{CommitError, CommitState, ExGuid, PreparedEdit, RevisionIndex, Store, document::{Document, Kind}}; +use onestore_offline::{Remote, Replica, Recovery}; +use std::{fs, io, path::Path}; +struct LostReply(Vec); +impl Remote for LostReply { + fn read(&mut self) -> io::Result> { Ok(self.0.clone()) } + fn publish(&mut self, edit: &PreparedEdit<'_>) -> Result<(), CommitError> { + self.0 = edit.as_bytes().to_vec(); + Err(CommitError { state: CommitState::Unknown, error: io::ErrorKind::ConnectionReset.into() }) + } + fn confirm(&mut self, _: &[u8]) -> Result<(), CommitError> { panic!("No confirmation is expected") } +} +fn main() -> Result<(), Box> { + let root = Path::new("corpus/offline-v4"); fs::create_dir(root)?; + let source = fs::read("corpus/native-external-assets/notebook/synthetic.one")?; + let store = Store::parse(&source)?; let index = RevisionIndex::parse(&store)?; let document = Document::parse(&index)?; + let (space, object) = document.spaces.iter().find_map(|(sid, space)| { + let view = &space.revisions[&space.contexts[&ExGuid::default()]]; + view.nodes.iter().find_map(|(oid,node)| matches!(&node.kind,Kind::RichText { text, .. } if text == "Native before 🦀").then_some((*sid,*oid))) + }).unwrap(); + let cache = Replica::create(root.join("live.sqlite"), &source)?; + let first = cache.edit_text(&source, space, object, 0..0, "queued ")?.unwrap(); + assert!(cache.sync_once(&mut LostReply(source.clone())).is_err()); + let working = cache.snapshot()?; + let second = cache.edit_text(&working, space, object, 0..0, "dependent ")?.unwrap(); + cache.export_recovery(root.join("recovery.sqlite"))?; + let recovery = Recovery::open(root.join("recovery.sqlite"))?; + assert_eq!(cache.pending()?, recovery.pending()?); + assert_eq!(recovery.summary()?.queued_edits, 2); + assert_eq!(recovery.summary()?.uncertain_edits, 1); + println!("Original schema-4 producer: IDs {first}, {second}; {:?}", recovery.summary()?); + Ok(()) +} diff --git a/corpus/offline-v4/provenance.json b/corpus/offline-v4/provenance.json new file mode 100644 index 0000000000000000000000000000000000000000..a6c712b2debcba0482044d9da4e8a5ba8a33f585 --- /dev/null +++ b/corpus/offline-v4/provenance.json @@ -0,0 +1,20 @@ +{ + "producer_revision": "2244711b", + "producer": "corpus/offline-v4/producer.rs", + "source_sha256": { + "crates/onestore-offline/src/lib.rs": "f4b6bdc9b8269e43cc8c1d69f5becb3140905f520438da53fb96ce1663352f22", + "crates/onestore-offline/src/schema.rs": "b20d6c127b32dc378c54ef8179583782d37b79496c4056899d053881034b9963", + "crates/onestore-offline/src/sync.rs": "3967d304c42550350c7fe2659317d665b8c4d7cdfd42f3eb013180ca91177442", + "crates/onestore-offline/src/recovery.rs": "1301b649085670e847b4593bc308f9a1e25713a99b5333d4fd217014720d29c3" + }, + "files": { + "recovery.sqlite": "2946f7b2f2636b4dc80b2d3fd930dfe1d5bdf3cafd591d5c33d50be209685a76", + "live.sqlite": "ff0b2ad1141398c10913d3f573ba909aa7cc6e7b3fb809ecd3d9869b689ff970" + }, + "queued_ids": [ + 1, + 2 + ], + "uncertain_id": 1, + "archive_version": 4 +} diff --git a/corpus/offline-v4/recovery.sqlite b/corpus/offline-v4/recovery.sqlite new file mode 100644 index 0000000000000000000000000000000000000000..67f06892d762d34917c2cf6b0520ea303c5c18e1 Binary files /dev/null and b/corpus/offline-v4/recovery.sqlite differ diff --git a/crates/onestore-offline/Cargo.toml b/crates/onestore-offline/Cargo.toml index 2f14c77276dc108b22823d8ccf42efc58fb4ae3a..8fe583780935ca0b1f0ba3c470361f964299b4fa 100644 --- a/crates/onestore-offline/Cargo.toml +++ b/crates/onestore-offline/Cargo.toml @@ -5,16 +5,18 @@ edition = "2024" publish = false [features] -smb = ["dep:onestore-smb"] +smb = ["dep:onestore-smb", "onestore-notebook/smb"] [dependencies] onestore = { path = "../onestore" } +onestore-notebook = { path = "../onestore-notebook" } onestore-smb = { path = "../onestore-smb", optional = true } rusqlite = { version = "=0.40.2", features = ["bundled", "backup"] } thiserror = "2" serde = { version = "1", features = ["derive"] } serde_json = "1" tempfile = "3" +sha2 = "0.11" [[example]] name = "smb_offline_client" diff --git a/crates/onestore-offline/README.md b/crates/onestore-offline/README.md index d25d4d4bbbb379f1c1ab17430153d63789a30e97..e338af2d44764e80cdeb7ad5c72c062cb0648f76 100644 --- a/crates/onestore-offline/README.md +++ b/crates/onestore-offline/README.md @@ -193,7 +193,7 @@ the library continues to preserve conflicts requiring an explicit decision. ## Recovery archives `export_recovery(new_path)` captures both complete notebook images, the typed -queue, uncertain attempts, conflicts, receipts and the edit-ID sequence in one +queue, uncertain attempts, conflicts, receipts, downloaded media and the edit-ID sequence in one SQLite snapshot. It refuses existing destinations and leaves the live queue unchanged. Export to a local directory from a background thread: copying holds the cache mutex while capturing the database. Failure after the final rename can @@ -221,3 +221,25 @@ let receipts = review.receipts()?; counts and byte sizes without notebook text, paths, authors or credentials. Opening an archive validates its schema and images without migration. A recovery archive is evidence for a reviewed recovery decision, not a second active queue. + +## Downloaded media + +`fetch_asset(source, section, filename, limit)` resolves a declared external +payload through `onestore-notebook::Source` and durably caches its exact bytes. +The section path is relative to the source root; the filename comes from a +`FileDataReference::External` in the retained working or remote image. Local and +SMB sources use the same API. Downloads release the cache mutex during network +I/O and recheck the reference before committing; edits can continue meanwhile. + +`cached_asset(filename, limit)` reads previously downloaded bytes without network +access. `None` means never downloaded; `Some(Vec::new())` is a downloaded empty +payload. Each read checks the stored SHA-256 and enforces the byte limit before +loading the payload. Cache contents describe the prior download, not current +server reachability or presence. A failed fetch returns its error without +silently substituting cached bytes. Different bytes for an already cached file +identity return `AssetChanged` and preserve the previous download. + +Downloads do not change pending edits, publication attempts or receipts. Recovery +archives include cached media and expose the same bounded `cached_asset` lookup. +Opening older live caches migrates them transactionally to schema 5; original +schema-4 archives remain readable without migration and contain no media cache. diff --git a/crates/onestore-offline/src/assets.rs b/crates/onestore-offline/src/assets.rs new file mode 100644 index 0000000000000000000000000000000000000000..9f61197f3a188cfb90d3bc60dc65502595b8a7a9 --- /dev/null +++ b/crates/onestore-offline/src/assets.rs @@ -0,0 +1,134 @@ +use super::*; +use onestore::FileDataReference; +use rusqlite::OptionalExtension; +use sha2::{Digest, Sha256}; + +impl Replica { + /// Reads a previously downloaded external payload without network access. + /// Absence is distinct from an empty payload; every returned buffer passes its stored checksum. + pub fn cached_asset(&self, filename: &str, limit: usize) -> Result>> { + let key = key(filename)?; + let connection = self + .connection + .lock() + .map_err(|_| io::Error::other("Cache owner panicked"))?; + cached(&connection, &key, limit) + } + + /// Fetches a declared external payload and durably retains it without changing the edit queue. + /// Different bytes for an already cached identity return `AssetChanged`, preserving the cache. + /// Network I/O does not hold the cache mutex; a stale reference fails before local publication. + pub fn fetch_asset( + &self, + source: &mut impl onestore_notebook::Source, + section: &str, + filename: &str, + limit: usize, + ) -> Result> { + let key = key(filename)?; + { + let connection = self + .connection + .lock() + .map_err(|_| io::Error::other("Cache owner panicked"))?; + if !referenced(&connection, &key)? { + return Err(io::Error::new( + io::ErrorKind::InvalidInput, + "The retained document images do not reference this external payload", + ) + .into()); + } + } + let bytes = onestore_notebook::read_external_asset(source, section, filename, limit)?; + let mut connection = self + .connection + .lock() + .map_err(|_| io::Error::other("Cache owner panicked"))?; + let transaction = connection.transaction_with_behavior(TransactionBehavior::Immediate)?; + if !referenced(&transaction, &key)? { + return Err(io::Error::new( + io::ErrorKind::ResourceBusy, + "The external payload reference changed during download", + ) + .into()); + } + let previous = match cached(&transaction, &key, bytes.len()) { + Err(Error::Io(error)) if error.kind() == io::ErrorKind::FileTooLarge => { + return Err(Error::AssetChanged); + } + other => other?, + }; + if let Some(previous) = previous { + if previous != bytes { + return Err(Error::AssetChanged); + } + } else { + transaction.execute( + "INSERT INTO assets(name,data,sha256) VALUES (?1,?2,?3)", + params![key, &bytes, &Sha256::digest(&bytes)[..]], + )?; + } + transaction.commit()?; + Ok(bytes) + } +} + +pub(crate) fn key(filename: &str) -> Result { + format!("{filename}").parse::()?; + Ok(filename.to_ascii_lowercase()) +} + +fn referenced(connection: &Connection, key: &str) -> Result { + for column in ["working", "base"] { + let image: Vec = connection.query_row( + &format!("SELECT {column} FROM replica WHERE id=1"), + [], + |row| row.get(0), + )?; + let store = Store::parse(&image)?; + let index = RevisionIndex::parse(&store)?; + let document = Document::parse(&index)?; + if document.spaces.values().flat_map(|space| space.revisions.values()) + .flat_map(|revision| revision.nodes.values()).any(|node| { + matches!(&node.kind, Kind::File { reference: FileDataReference::External(name), .. } if name.eq_ignore_ascii_case(key)) + }) { + return Ok(true); + } + } + Ok(false) +} + +pub(crate) fn cached(connection: &Connection, key: &str, limit: usize) -> Result>> { + let version: u32 = connection.pragma_query_value(None, "user_version", |row| row.get(0))?; + if version < 5 { + return Ok(None); + } + let length: Option = connection + .query_row( + "SELECT length(data) FROM assets WHERE name=?1", + [key], + |row| row.get(0), + ) + .optional()?; + let Some(length) = length else { + return Ok(None); + }; + let length = + usize::try_from(length).map_err(|_| io::Error::from(io::ErrorKind::InvalidData))?; + if length > limit { + return Err(io::Error::from(io::ErrorKind::FileTooLarge).into()); + } + let (bytes, expected): (Vec, Vec) = connection.query_row( + "SELECT data,sha256 FROM assets WHERE name=?1", + [key], + |row| Ok((row.get(0)?, row.get(1)?)), + )?; + if bytes.len() != length || Sha256::digest(&bytes)[..] != expected { + return Err(io::Error::new( + io::ErrorKind::InvalidData, + "Cached external payload checksum mismatch", + ) + .into()); + } + Ok(Some(bytes)) +} diff --git a/crates/onestore-offline/src/lib.rs b/crates/onestore-offline/src/lib.rs index 5186c411f1053e8e2646ed2eab27c8549f613157..a21ba1afca74f36ba8932b5c338e82f725192808 100644 --- a/crates/onestore-offline/src/lib.rs +++ b/crates/onestore-offline/src/lib.rs @@ -8,6 +8,7 @@ use onestore::{ use rusqlite::{Connection, OpenFlags, TransactionBehavior, params}; use std::{fs::OpenOptions, io, ops::Range, path::Path, sync::Mutex, time::Duration}; +mod assets; mod formatting; mod rebase; mod recovery; @@ -35,12 +36,16 @@ pub enum Error { Remote(#[from] onestore::CommitError), #[error(transparent)] RemoteIo(io::Error), + #[error(transparent)] + Notebook(#[from] onestore_notebook::Error), + #[error("External payload identity now refers to different bytes")] + AssetChanged, } type Result = std::result::Result; const APPLICATION_ID: u32 = 0x4f4e454f; -const SCHEMA_VERSION: u32 = 4; +const SCHEMA_VERSION: u32 = 5; /// Text and its observed precondition, retained across cache reopen and rebasing. #[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)] diff --git a/crates/onestore-offline/src/recovery.rs b/crates/onestore-offline/src/recovery.rs index 4530d2d1de1aeeebc49aebeb545d2192884af82e..5b87b1c733d53acaba966f24c4fe60ba5df548e7 100644 --- a/crates/onestore-offline/src/recovery.rs +++ b/crates/onestore-offline/src/recovery.rs @@ -13,6 +13,8 @@ pub struct RecoverySummary { pub published_receipts: u64, pub working_bytes: u64, pub remote_bytes: u64, + pub cached_assets: u64, + pub cached_asset_bytes: u64, } /// Read-only recovery evidence; it cannot publish or acknowledge an edit. @@ -29,7 +31,7 @@ impl Recovery { let application: u32 = connection.pragma_query_value(None, "application_id", |row| row.get(0))?; let version: u32 = connection.pragma_query_value(None, "user_version", |row| row.get(0))?; - if application != RECOVERY_ID || version != SCHEMA_VERSION { + if application != RECOVERY_ID || !(4..=SCHEMA_VERSION).contains(&version) { return Err(io::Error::new( io::ErrorKind::InvalidData, "Unrecognized recovery archive or unsupported schema version", @@ -73,6 +75,11 @@ impl Recovery { pub fn receipts(&self) -> Result> { receipts(&self.connection) } + + /// Reads a previously downloaded external payload without accessing its former server. + pub fn cached_asset(&self, filename: &str, limit: usize) -> Result>> { + assets::cached(&self.connection, &assets::key(filename)?, limit) + } } impl Replica { @@ -131,6 +138,16 @@ impl Replica { } fn summary(connection: &Connection) -> Result { + let version: u32 = connection.pragma_query_value(None, "user_version", |row| row.get(0))?; + let (cached_assets, cached_asset_bytes) = if version < 5 { + (0, 0) + } else { + connection.query_row( + "SELECT count(*),coalesce(sum(length(data)),0) FROM assets", + [], + |row| Ok((unsigned(row, 0)?, unsigned(row, 1)?)), + )? + }; Ok(connection.query_row( "SELECT (SELECT count(*) FROM edits), (SELECT count(*) FROM conflicts), (SELECT count(*) FROM attempt), (SELECT count(*) FROM receipts), @@ -144,6 +161,8 @@ fn summary(connection: &Connection) -> Result { published_receipts: unsigned(row, 3)?, working_bytes: unsigned(row, 4)?, remote_bytes: unsigned(row, 5)?, + cached_assets, + cached_asset_bytes, }) }, )?) diff --git a/crates/onestore-offline/src/schema.rs b/crates/onestore-offline/src/schema.rs index a65c19446f13f76bfac63adaed10061da7bfbed8..32b8f454dc8f969a687d9412dd5cd01f74238b56 100644 --- a/crates/onestore-offline/src/schema.rs +++ b/crates/onestore-offline/src/schema.rs @@ -6,6 +6,12 @@ const CONFLICTS: &str = "CREATE TABLE conflicts ( kind INTEGER NOT NULL CHECK(kind BETWEEN 0 AND 3) ) STRICT;"; +const ASSETS: &str = "CREATE TABLE assets ( + name TEXT PRIMARY KEY NOT NULL, + data BLOB NOT NULL, + sha256 BLOB NOT NULL CHECK(length(sha256)=32) +) STRICT;"; + pub(crate) fn create(transaction: &Transaction<'_>) -> Result<()> { transaction.execute_batch( "CREATE TABLE edits ( @@ -24,16 +30,22 @@ pub(crate) fn create(transaction: &Transaction<'_>) -> Result<()> { ) STRICT;", )?; transaction.execute_batch(CONFLICTS)?; + transaction.execute_batch(ASSETS)?; Ok(()) } pub(crate) fn migrate(transaction: &Transaction<'_>, version: u32) -> Result<()> { + if version == 4 { + transaction.execute_batch(ASSETS)?; + return Ok(()); + } if version == 3 { transaction.execute_batch("ALTER TABLE conflicts RENAME TO old_conflicts;")?; transaction.execute_batch(CONFLICTS)?; transaction.execute_batch( "INSERT INTO conflicts SELECT * FROM old_conflicts; DROP TABLE old_conflicts;", )?; + transaction.execute_batch(ASSETS)?; return Ok(()); } diff --git a/crates/onestore-offline/tests/assets.rs b/crates/onestore-offline/tests/assets.rs new file mode 100644 index 0000000000000000000000000000000000000000..3186efe6def0171e62db0640966f925ca9ec0f11 --- /dev/null +++ b/crates/onestore-offline/tests/assets.rs @@ -0,0 +1,449 @@ +use onestore_offline::{EditStatus, Error, Operation, Recovery, Replica}; +use std::{ + fs, io, + path::Path, + sync::{Barrier, mpsc}, + time::Duration, +}; + +const FIXTURE: &str = "../../corpus/native-external-assets/notebook"; +const LEGACY: &str = "../../corpus/offline-v4/live.sqlite"; + +fn copied_cache(root: &Path) -> Replica { + let path = root.join("cache.sqlite"); + assert!(!path.exists()); + fs::copy(LEGACY, &path).unwrap(); + Replica::open(path).unwrap() +} + +fn payload(size: usize) -> (String, Vec) { + fs::read_dir(Path::new(FIXTURE).join("synthetic_onefiles")) + .unwrap() + .map(|entry| entry.unwrap().path()) + .find_map(|path| { + let bytes = fs::read(&path).unwrap(); + (bytes.len() == size) + .then(|| (path.file_name().unwrap().to_str().unwrap().into(), bytes)) + }) + .unwrap() +} + +struct Payload(Option>>); +impl onestore_notebook::Source for Payload { + fn entries(&mut self, _: &str, _: usize) -> io::Result> { + panic!("Unexpected enumeration") + } + fn read(&mut self, _: &str, _: usize) -> io::Result> { + self.0.take().expect("Unexpected repeated payload read") + } +} + +#[test] +fn downloaded_media_survives_migration_reopen_and_recovery_with_the_queue_intact() { + let directory = tempfile::tempdir().unwrap(); + let original = fs::read(LEGACY).unwrap(); + let cache = copied_cache(directory.path()); + let working = cache.snapshot().unwrap(); + let remote = cache.remote_snapshot().unwrap(); + let pending = cache.pending().unwrap(); + assert_eq!( + pending.iter().map(|edit| edit.id).collect::>(), + [1, 2] + ); + let uncertain = cache.status(1).unwrap(); + assert!(matches!( + uncertain, + Some(EditStatus::AwaitingConfirmation { .. }) + )); + let mut source = onestore_notebook::Local::open(FIXTURE).unwrap(); + for size in [0, 1024] { + let (name, bytes) = payload(size); + assert!(cache.cached_asset(&name, size).unwrap().is_none()); + assert_eq!( + cache + .fetch_asset(&mut source, "synthetic.one", &name, size) + .unwrap(), + bytes + ); + assert_eq!( + cache + .cached_asset(&name.to_ascii_lowercase(), size) + .unwrap(), + Some(bytes) + ); + } + let summary = cache.recovery_summary().unwrap(); + assert_eq!( + (summary.cached_assets, summary.cached_asset_bytes), + (2, 1024) + ); + assert_eq!(cache.snapshot().unwrap(), working); + assert_eq!(cache.remote_snapshot().unwrap(), remote); + assert_eq!(cache.pending().unwrap(), pending); + assert_eq!(cache.status(1).unwrap(), uncertain); + cache + .export_recovery(directory.path().join("recovery.sqlite")) + .unwrap(); + drop(cache); + let cache = Replica::open(directory.path().join("cache.sqlite")).unwrap(); + let recovery = Recovery::open(directory.path().join("recovery.sqlite")).unwrap(); + assert_eq!(recovery.summary().unwrap(), summary); + assert_eq!(recovery.pending().unwrap(), pending); + assert_eq!(recovery.status(1).unwrap(), uncertain); + assert!(recovery.receipts().unwrap().is_empty()); + for size in [0, 1024] { + let (name, bytes) = payload(size); + assert_eq!( + cache.cached_asset(&name, size).unwrap(), + Some(bytes.clone()) + ); + assert_eq!(recovery.cached_asset(&name, size).unwrap(), Some(bytes)); + } + assert_eq!(fs::read(LEGACY).unwrap(), original); +} + +#[test] +fn an_original_version_four_archive_remains_read_only_and_has_no_cached_media() { + let path = "../../corpus/offline-v4/recovery.sqlite"; + let before = fs::read(path).unwrap(); + let archive = Recovery::open(path).unwrap(); + assert_eq!(archive.pending().unwrap().len(), 2); + assert!(matches!( + archive.status(1).unwrap(), + Some(EditStatus::AwaitingConfirmation { .. }) + )); + assert_eq!(archive.summary().unwrap().cached_assets, 0); + assert_eq!(archive.summary().unwrap().cached_asset_bytes, 0); + assert!(archive.cached_asset(&payload(0).0, 0).unwrap().is_none()); + drop(archive); + assert_eq!(fs::read(path).unwrap(), before); +} + +#[test] +fn unsuccessful_refreshes_and_changed_identity_data_preserve_the_downloaded_payload() { + let directory = tempfile::tempdir().unwrap(); + let cache = copied_cache(directory.path()); + let (name, bytes) = payload(1024); + let mut source = Payload(Some(Ok(bytes.clone()))); + cache + .fetch_asset(&mut source, "synthetic.one", &name, 1024) + .unwrap(); + let before = fs::read(directory.path().join("cache.sqlite")).unwrap(); + let mut unavailable = Payload(Some(Err(io::ErrorKind::ConnectionReset.into()))); + assert!( + matches!(cache.fetch_asset(&mut unavailable, "synthetic.one", &name, 1024), + Err(Error::Notebook(onestore_notebook::Error::Io { error, .. })) if error.kind() == io::ErrorKind::ConnectionReset) + ); + for size in [0, 1024, 2048] { + let mut changed = Payload(Some(Ok(vec![9; size]))); + assert!(matches!( + cache.fetch_asset(&mut changed, "synthetic.one", &name, 2048), + Err(Error::AssetChanged) + )); + } + assert!( + matches!(cache.cached_asset(&name, 1023), Err(Error::Io(error)) if error.kind() == io::ErrorKind::FileTooLarge) + ); + assert_eq!(cache.cached_asset(&name, 1024).unwrap(), Some(bytes)); + assert_eq!( + fs::read(directory.path().join("cache.sqlite")).unwrap(), + before + ); +} + +#[test] +fn unreferenced_payloads_are_rejected_before_io_or_local_changes() { + let directory = tempfile::tempdir().unwrap(); + let cache = copied_cache(directory.path()); + let mut unused = Payload(None); + assert!( + matches!(cache.fetch_asset(&mut unused, "synthetic.one", "00000000-0000-0000-0000-000000000001.onebin", 100), + Err(Error::Io(error)) if error.kind() == io::ErrorKind::InvalidInput) + ); + assert!( + cache + .fetch_asset(&mut unused, "synthetic.one", "../payload.onebin", 100) + .is_err() + ); + assert_eq!(cache.recovery_summary().unwrap().cached_assets, 0); +} + +#[test] +fn download_network_wait_does_not_block_local_edits() { + struct Waiting { + entered: mpsc::Sender<()>, + released: mpsc::Receiver<()>, + bytes: Vec, + } + impl onestore_notebook::Source for Waiting { + fn entries(&mut self, _: &str, _: usize) -> io::Result> { + panic!("Unexpected enumeration") + } + fn read(&mut self, _: &str, _: usize) -> io::Result> { + self.entered.send(()).unwrap(); + self.released.recv_timeout(Duration::from_secs(5)).unwrap(); + Ok(self.bytes.clone()) + } + } + let directory = tempfile::tempdir().unwrap(); + let cache = copied_cache(directory.path()); + let (name, bytes) = payload(1024); + let (entered, waiting) = mpsc::channel(); + let (release, released) = mpsc::channel(); + let mut source = Waiting { + entered, + released, + bytes: bytes.clone(), + }; + std::thread::scope(|scope| { + let download = scope.spawn(|| { + cache + .fetch_asset(&mut source, "synthetic.one", &name, 1024) + .unwrap() + }); + waiting.recv_timeout(Duration::from_secs(5)).unwrap(); + let pending = cache.pending().unwrap(); + let Operation::Text(edit) = &pending[0].operation else { + panic!() + }; + let id = cache + .edit_text( + &cache.snapshot().unwrap(), + pending[0].space, + edit.object, + 0..0, + "during download ", + ) + .unwrap() + .unwrap(); + release.send(()).unwrap(); + assert_eq!(download.join().unwrap(), bytes); + assert_eq!(cache.status(id).unwrap(), Some(EditStatus::Pending)); + }); +} + +#[test] +fn concurrent_downloads_publish_one_immutable_cache_entry() { + let directory = tempfile::tempdir().unwrap(); + let cache = copied_cache(directory.path()); + let (name, bytes) = payload(1024); + let ready = Barrier::new(8); + std::thread::scope(|scope| { + for _ in 0..8 { + scope.spawn(|| { + let mut source = onestore_notebook::Local::open(FIXTURE).unwrap(); + ready.wait(); + assert_eq!( + cache + .fetch_asset(&mut source, "synthetic.one", &name, 1024) + .unwrap(), + bytes + ); + }); + } + }); + assert_eq!(cache.recovery_summary().unwrap().cached_assets, 1); +} + +#[test] +fn a_native_refresh_removing_the_reference_rejects_an_inflight_download() { + struct Native; + impl onestore_offline::Remote for Native { + fn read(&mut self) -> io::Result> { + fs::read("../../corpus/native-external-assets/native/synthetic.one") + } + fn publish(&mut self, _: &onestore::PreparedEdit<'_>) -> Result<(), onestore::CommitError> { + panic!("Unexpected publication") + } + fn confirm(&mut self, _: &[u8]) -> Result<(), onestore::CommitError> { + panic!("Unexpected confirmation") + } + } + struct Refresh<'a>(&'a Replica); + impl onestore_notebook::Source for Refresh<'_> { + fn entries(&mut self, _: &str, _: usize) -> io::Result> { + panic!("Unexpected enumeration") + } + fn read(&mut self, _: &str, _: usize) -> io::Result> { + assert_eq!(self.0.sync_once(&mut Native).unwrap(), None); + Ok(payload(1024).1) + } + } + let root = tempfile::tempdir().unwrap(); + let source = fs::read(Path::new(FIXTURE).join("synthetic.one")).unwrap(); + let cache = Replica::create(root.path().join("cache.sqlite"), &source).unwrap(); + let (name, _) = payload(1024); + assert!( + matches!(cache.fetch_asset(&mut Refresh(&cache), "synthetic.one", &name, 1024), + Err(Error::Io(error)) if error.kind() == io::ErrorKind::ResourceBusy) + ); + assert!(cache.cached_asset(&name, 1024).unwrap().is_none()); + assert_eq!( + cache.snapshot().unwrap(), + fs::read("../../corpus/native-external-assets/native/synthetic.one").unwrap() + ); +} + +#[test] +fn local_failure_and_bad_cached_bytes_never_become_successful_downloads() { + let directory = tempfile::tempdir().unwrap(); + drop(copied_cache(directory.path())); + let path = directory.path().join("cache.sqlite"); + let connection = rusqlite::Connection::open(&path).unwrap(); + connection.execute_batch("CREATE TRIGGER fail_asset BEFORE INSERT ON assets BEGIN SELECT RAISE(ABORT,'Test asset failure'); END").unwrap(); + drop(connection); + let (name, bytes) = payload(1024); + let cache = Replica::open(&path).unwrap(); + let mut source = onestore_notebook::Local::open(FIXTURE).unwrap(); + assert!(matches!( + cache.fetch_asset(&mut source, "synthetic.one", &name, 1024), + Err(Error::Database(_)) + )); + assert!(cache.cached_asset(&name, 1024).unwrap().is_none()); + assert_eq!(cache.pending().unwrap().len(), 2); + drop(cache); + let connection = rusqlite::Connection::open(&path).unwrap(); + connection.execute_batch("DROP TRIGGER fail_asset").unwrap(); + drop(connection); + let cache = Replica::open(&path).unwrap(); + cache + .fetch_asset(&mut source, "synthetic.one", &name, 1024) + .unwrap(); + drop(cache); + let connection = rusqlite::Connection::open(&path).unwrap(); + let mut damaged = bytes; + damaged[0] ^= 1; + connection + .execute("UPDATE assets SET data=?1", [damaged]) + .unwrap(); + drop(connection); + let cache = Replica::open(&path).unwrap(); + assert!( + matches!(cache.cached_asset(&name, 1024), Err(Error::Io(error)) if error.kind() == io::ErrorKind::InvalidData) + ); + assert!( + matches!(cache.fetch_asset(&mut source, "synthetic.one", &name, 1024), Err(Error::Io(error)) if error.kind() == io::ErrorKind::InvalidData) + ); + cache + .export_recovery(directory.path().join("damaged.sqlite")) + .unwrap(); + let archive = Recovery::open(directory.path().join("damaged.sqlite")).unwrap(); + assert!(archive.cached_asset(&name, 1024).is_err()); + assert_eq!(archive.pending().unwrap().len(), 2); +} + +#[test] +fn abrupt_process_exit_retains_only_completed_downloads_and_archives() { + const CHILD: &str = "ONESTORE_ASSET_EXIT_CASE"; + if let Ok(phase) = std::env::var(CHILD) { + struct ExitDuringRead; + impl onestore_notebook::Source for ExitDuringRead { + fn entries(&mut self, _: &str, _: usize) -> io::Result> { + panic!("Unexpected enumeration") + } + fn read(&mut self, _: &str, _: usize) -> io::Result> { + std::process::exit(83) + } + } + let root = std::env::var("ONESTORE_ASSET_EXIT_ROOT").unwrap(); + let cache = copied_cache(Path::new(&root)); + let (name, bytes) = payload(1024); + if phase == "download" { + cache + .fetch_asset(&mut ExitDuringRead, "synthetic.one", &name, bytes.len()) + .unwrap(); + panic!("Read returned after process exit"); + } + cache + .fetch_asset( + &mut onestore_notebook::Local::open(FIXTURE).unwrap(), + "synthetic.one", + &name, + bytes.len(), + ) + .unwrap(); + if phase == "archive" { + cache + .export_recovery(Path::new(&root).join("recovery.sqlite")) + .unwrap(); + } + std::process::exit(83); + } + for phase in ["download", "cached", "archive"] { + let root = tempfile::tempdir().unwrap(); + let output = std::process::Command::new(std::env::current_exe().unwrap()) + .args([ + "--exact", + "abrupt_process_exit_retains_only_completed_downloads_and_archives", + ]) + .env(CHILD, phase) + .env("ONESTORE_ASSET_EXIT_ROOT", root.path()) + .output() + .unwrap(); + assert_eq!( + output.status.code(), + Some(83), + "{}", + String::from_utf8_lossy(&output.stderr) + ); + let cache = Replica::open(root.path().join("cache.sqlite")).unwrap(); + let (name, bytes) = payload(1024); + let expected = (phase != "download").then_some(bytes); + assert_eq!(cache.cached_asset(&name, 1024).unwrap(), expected); + assert_eq!(cache.pending().unwrap().len(), 2); + assert!(matches!( + cache.status(1).unwrap(), + Some(EditStatus::AwaitingConfirmation { .. }) + )); + if phase == "archive" { + let recovery = Recovery::open(root.path().join("recovery.sqlite")).unwrap(); + assert_eq!(recovery.cached_asset(&name, 1024).unwrap(), expected); + assert_eq!(recovery.pending().unwrap(), cache.pending().unwrap()); + assert_eq!(recovery.status(1).unwrap(), cache.status(1).unwrap()); + } + } +} + +#[test] +#[cfg(feature = "smb")] +#[ignore = "requires a disposable Samba mirror at ONESTORE_SMB_NOTEBOOK"] +fn live_smb_downloads_survive_disconnect_and_cache_reopen() { + let root = tempfile::tempdir().unwrap(); + let cache = copied_cache(root.path()); + let client = onestore_smb::Client::connect( + &std::env::var("ONESTORE_SMB_LAB").unwrap(), + "agent", + onestore_smb::Credentials::default(), + Duration::from_secs(5), + ) + .unwrap(); + let notebook = std::env::var("ONESTORE_SMB_NOTEBOOK").unwrap(); + let mut source = onestore_notebook::Smb::new(&client, ¬ebook).unwrap(); + for size in [0, 771, 1024] { + let (name, bytes) = payload(size); + assert_eq!( + cache + .fetch_asset(&mut source, "synthetic.one", &name, size) + .unwrap(), + bytes + ); + } + drop(client); + cache + .export_recovery(root.path().join("recovery.sqlite")) + .unwrap(); + drop(cache); + let cache = Replica::open(root.path().join("cache.sqlite")).unwrap(); + let recovery = Recovery::open(root.path().join("recovery.sqlite")).unwrap(); + for size in [0, 771, 1024] { + let (name, bytes) = payload(size); + assert_eq!( + cache.cached_asset(&name, size).unwrap(), + Some(bytes.clone()) + ); + assert_eq!(recovery.cached_asset(&name, size).unwrap(), Some(bytes)); + } + assert_eq!(cache.pending().unwrap(), recovery.pending().unwrap()); + assert_eq!(cache.status(1).unwrap(), recovery.status(1).unwrap()); + assert_eq!(cache.recovery_summary().unwrap().cached_assets, 3); +} diff --git a/crates/onestore-offline/tests/cache.rs b/crates/onestore-offline/tests/cache.rs index 7e762d412e19449e22aa147f432f5ef0adde3858..31f727a97413e28cc255bcae111b8ef8bc8f580a 100644 --- a/crates/onestore-offline/tests/cache.rs +++ b/crates/onestore-offline/tests/cache.rs @@ -734,7 +734,7 @@ fn unrecognized_persisted_operations_are_rejected_without_dropping_fields() { db.execute("UPDATE edits SET operation=?1", [value.to_string()]) .unwrap(); if operation != "Format" { - db.execute_batch("DROP TABLE conflicts; CREATE TABLE conflicts (edit_id INTEGER PRIMARY KEY REFERENCES edits(id) ON DELETE CASCADE, kind INTEGER NOT NULL CHECK(kind BETWEEN 0 AND 2)) STRICT; PRAGMA user_version=3;").unwrap(); + db.execute_batch("DROP TABLE assets; DROP TABLE conflicts; CREATE TABLE conflicts (edit_id INTEGER PRIMARY KEY REFERENCES edits(id) ON DELETE CASCADE, kind INTEGER NOT NULL CHECK(kind BETWEEN 0 AND 2)) STRICT; PRAGMA user_version=3;").unwrap(); } drop(db); let before = fs::read(&path).unwrap(); diff --git a/crates/onestore-offline/tests/sync.rs b/crates/onestore-offline/tests/sync.rs index c6f52ae896486d9f64cf93d52ba334f75e2180b5..5128b17829f84f71916c3c548d437ee866feca4b 100644 --- a/crates/onestore-offline/tests/sync.rs +++ b/crates/onestore-offline/tests/sync.rs @@ -197,6 +197,8 @@ fn recovery_archive_preserves_typed_queue_uncertainty_and_receipts_without_becom published_receipts: 1, working_bytes: working.len() as u64, remote_bytes: remote.len() as u64, + cached_assets: 0, + cached_asset_bytes: 0, }; assert_eq!(cache.recovery_summary().unwrap(), summary); cache.export_recovery(&archive_path).unwrap(); @@ -712,7 +714,7 @@ fn version_one_cache_migration_preserves_images_intents_and_local_ids() { drop(cache); let db = rusqlite::Connection::open(&path).unwrap(); db.execute_batch( - "DROP TABLE attempt; DROP TABLE conflicts; DROP TABLE receipts; DROP TABLE edits; + "DROP TABLE attempt; DROP TABLE conflicts; DROP TABLE receipts; DROP TABLE edits; DROP TABLE assets; CREATE TABLE edits ( id INTEGER PRIMARY KEY AUTOINCREMENT CHECK(id>0), space TEXT NOT NULL, object TEXT NOT NULL, before_text TEXT NOT NULL, @@ -737,7 +739,7 @@ fn version_one_cache_migration_preserves_images_intents_and_local_ids() { assert_eq!( db.pragma_query_value(None, "user_version", |row| row.get::<_, u32>(0)) .unwrap(), - 4 + 5 ); }