authorgravatar for git@paperclover.netclover caruso <git@paperclover.net> 2026-09-07 19:32:30-07:00
committergravatar for git@paperclover.netclover caruso <git@paperclover.net> 2026-09-07 19:32:30-07:00
logaabbadfaf1e956a3daaa6d59d0a2213d4cf78ffe
tree5c836d3d86aa38e84e649d57aef97818fe12898f
parent30d02b72f6a5b8dd8b72cea593bdc9c883afea1e
signature Signed by SSH key SHA256:52mNGHRsVFBDED9IAX5pe+LRWUefqTbxEReunq21QvU

feat: retain verified external media in offline recovery caches

Download declared external payloads without blocking local edits, then cache exact bytes and SHA-256 transactionally. Reject changed identities, stale references, and damaged cached data while preserving the edit queue. Include media in read-only recovery exports. Migrate existing live caches to schema 5 and retain original schema-4 archive compatibility. Verify actual legacy cache migration, concurrent download/edit behavior, abrupt process exit, live Samba downloads followed by offline reopen, full workspace regressions, and iOS device/simulator linking. Assisted-by: gpt-6-astra

15 files changed, 722 insertions(+), 7 deletions(-)

Cargo.lock+2
......@@ -537,10 +537,12 @@ name = "onestore-offline"
537537version = "0.1.0"
538538dependencies = [
539539 "onestore",
540 "onestore-notebook",
540541 "onestore-smb",
541542 "rusqlite",
542543 "serde",
543544 "serde_json",
545 "sha2",
544546 "tempfile",
545547 "thiserror",
546548]
corpus/offline-v4/README.md created+15
......@@ -0,0 +1,15 @@
1# Original schema-4 offline state
2
3These databases were produced by `onestore-offline` at revision `2244711b`,
4before the external-media table and schema-5 migration existed. Tests copy the
5live database before opening it; the archive remains read-only.
6
7[`producer.rs`](producer.rs) records the original producer. Run it against that
8revision's core/offline crates, from a checkout whose `corpus/offline-v4` does
9not exist. It uses the public native external-asset fixture, queues a text edit,
10simulates a publication with a lost response, queues a dependent text edit, and
11exports the complete recovery archive. Generated revision IDs vary between runs.
12
13[`provenance.json`](provenance.json) records the original source/database hashes
14and the retained uncertain and pending edit IDs. No credentials or personal
15notebook content are included.
corpus/offline-v4/live.sqlite created
Binary files /dev/null and b/corpus/offline-v4/live.sqlite differ
corpus/offline-v4/producer.rs created+33
......@@ -0,0 +1,33 @@
1use onestore::{CommitError, CommitState, ExGuid, PreparedEdit, RevisionIndex, Store, document::{Document, Kind}};
2use onestore_offline::{Remote, Replica, Recovery};
3use std::{fs, io, path::Path};
4struct LostReply(Vec<u8>);
5impl Remote for LostReply {
6 fn read(&mut self) -> io::Result<Vec<u8>> { Ok(self.0.clone()) }
7 fn publish(&mut self, edit: &PreparedEdit<'_>) -> Result<(), CommitError> {
8 self.0 = edit.as_bytes().to_vec();
9 Err(CommitError { state: CommitState::Unknown, error: io::ErrorKind::ConnectionReset.into() })
10 }
11 fn confirm(&mut self, _: &[u8]) -> Result<(), CommitError> { panic!("No confirmation is expected") }
12}
13fn main() -> Result<(), Box<dyn std::error::Error>> {
14 let root = Path::new("corpus/offline-v4"); fs::create_dir(root)?;
15 let source = fs::read("corpus/native-external-assets/notebook/synthetic.one")?;
16 let store = Store::parse(&source)?; let index = RevisionIndex::parse(&store)?; let document = Document::parse(&index)?;
17 let (space, object) = document.spaces.iter().find_map(|(sid, space)| {
18 let view = &space.revisions[&space.contexts[&ExGuid::default()]];
19 view.nodes.iter().find_map(|(oid,node)| matches!(&node.kind,Kind::RichText { text, .. } if text == "Native before 🦀").then_some((*sid,*oid)))
20 }).unwrap();
21 let cache = Replica::create(root.join("live.sqlite"), &source)?;
22 let first = cache.edit_text(&source, space, object, 0..0, "queued ")?.unwrap();
23 assert!(cache.sync_once(&mut LostReply(source.clone())).is_err());
24 let working = cache.snapshot()?;
25 let second = cache.edit_text(&working, space, object, 0..0, "dependent ")?.unwrap();
26 cache.export_recovery(root.join("recovery.sqlite"))?;
27 let recovery = Recovery::open(root.join("recovery.sqlite"))?;
28 assert_eq!(cache.pending()?, recovery.pending()?);
29 assert_eq!(recovery.summary()?.queued_edits, 2);
30 assert_eq!(recovery.summary()?.uncertain_edits, 1);
31 println!("Original schema-4 producer: IDs {first}, {second}; {:?}", recovery.summary()?);
32 Ok(())
33}
corpus/offline-v4/provenance.json created+20
......@@ -0,0 +1,20 @@
1{
2 "producer_revision": "2244711b",
3 "producer": "corpus/offline-v4/producer.rs",
4 "source_sha256": {
5 "crates/onestore-offline/src/lib.rs": "f4b6bdc9b8269e43cc8c1d69f5becb3140905f520438da53fb96ce1663352f22",
6 "crates/onestore-offline/src/schema.rs": "b20d6c127b32dc378c54ef8179583782d37b79496c4056899d053881034b9963",
7 "crates/onestore-offline/src/sync.rs": "3967d304c42550350c7fe2659317d665b8c4d7cdfd42f3eb013180ca91177442",
8 "crates/onestore-offline/src/recovery.rs": "1301b649085670e847b4593bc308f9a1e25713a99b5333d4fd217014720d29c3"
9 },
10 "files": {
11 "recovery.sqlite": "2946f7b2f2636b4dc80b2d3fd930dfe1d5bdf3cafd591d5c33d50be209685a76",
12 "live.sqlite": "ff0b2ad1141398c10913d3f573ba909aa7cc6e7b3fb809ecd3d9869b689ff970"
13 },
14 "queued_ids": [
15 1,
16 2
17 ],
18 "uncertain_id": 1,
19 "archive_version": 4
20}
corpus/offline-v4/recovery.sqlite created
Binary files /dev/null and b/corpus/offline-v4/recovery.sqlite differ
crates/onestore-offline/Cargo.toml+3-1
......@@ -5,16 +5,18 @@ edition = "2024"
55publish = false
66
77[features]
8smb = ["dep:onestore-smb"]
8smb = ["dep:onestore-smb", "onestore-notebook/smb"]
99
1010[dependencies]
1111onestore = { path = "../onestore" }
12onestore-notebook = { path = "../onestore-notebook" }
1213onestore-smb = { path = "../onestore-smb", optional = true }
1314rusqlite = { version = "=0.40.2", features = ["bundled", "backup"] }
1415thiserror = "2"
1516serde = { version = "1", features = ["derive"] }
1617serde_json = "1"
1718tempfile = "3"
19sha2 = "0.11"
1820
1921[[example]]
2022name = "smb_offline_client"
crates/onestore-offline/README.md+23-1
......@@ -193,7 +193,7 @@ the library continues to preserve conflicts requiring an explicit decision.
193193## Recovery archives
194194
195195`export_recovery(new_path)` captures both complete notebook images, the typed
196queue, uncertain attempts, conflicts, receipts and the edit-ID sequence in one
196queue, uncertain attempts, conflicts, receipts, downloaded media and the edit-ID sequence in one
197197SQLite snapshot. It refuses existing destinations and leaves the live queue
198198unchanged. Export to a local directory from a background thread: copying holds
199199the cache mutex while capturing the database. Failure after the final rename can
......@@ -221,3 +221,25 @@ let receipts = review.receipts()?;
221221counts and byte sizes without notebook text, paths, authors or credentials.
222222Opening an archive validates its schema and images without migration. A recovery
223223archive is evidence for a reviewed recovery decision, not a second active queue.
224
225## Downloaded media
226
227`fetch_asset(source, section, filename, limit)` resolves a declared external
228payload through `onestore-notebook::Source` and durably caches its exact bytes.
229The section path is relative to the source root; the filename comes from a
230`FileDataReference::External` in the retained working or remote image. Local and
231SMB sources use the same API. Downloads release the cache mutex during network
232I/O and recheck the reference before committing; edits can continue meanwhile.
233
234`cached_asset(filename, limit)` reads previously downloaded bytes without network
235access. `None` means never downloaded; `Some(Vec::new())` is a downloaded empty
236payload. Each read checks the stored SHA-256 and enforces the byte limit before
237loading the payload. Cache contents describe the prior download, not current
238server reachability or presence. A failed fetch returns its error without
239silently substituting cached bytes. Different bytes for an already cached file
240identity return `AssetChanged` and preserve the previous download.
241
242Downloads do not change pending edits, publication attempts or receipts. Recovery
243archives include cached media and expose the same bounded `cached_asset` lookup.
244Opening older live caches migrates them transactionally to schema 5; original
245schema-4 archives remain readable without migration and contain no media cache.
crates/onestore-offline/src/assets.rs created+134
......@@ -0,0 +1,134 @@
1use super::*;
2use onestore::FileDataReference;
3use rusqlite::OptionalExtension;
4use sha2::{Digest, Sha256};
5
6impl Replica {
7 /// Reads a previously downloaded external payload without network access.
8 /// Absence is distinct from an empty payload; every returned buffer passes its stored checksum.
9 pub fn cached_asset(&self, filename: &str, limit: usize) -> Result<Option<Vec<u8>>> {
10 let key = key(filename)?;
11 let connection = self
12 .connection
13 .lock()
14 .map_err(|_| io::Error::other("Cache owner panicked"))?;
15 cached(&connection, &key, limit)
16 }
17
18 /// Fetches a declared external payload and durably retains it without changing the edit queue.
19 /// Different bytes for an already cached identity return `AssetChanged`, preserving the cache.
20 /// Network I/O does not hold the cache mutex; a stale reference fails before local publication.
21 pub fn fetch_asset(
22 &self,
23 source: &mut impl onestore_notebook::Source,
24 section: &str,
25 filename: &str,
26 limit: usize,
27 ) -> Result<Vec<u8>> {
28 let key = key(filename)?;
29 {
30 let connection = self
31 .connection
32 .lock()
33 .map_err(|_| io::Error::other("Cache owner panicked"))?;
34 if !referenced(&connection, &key)? {
35 return Err(io::Error::new(
36 io::ErrorKind::InvalidInput,
37 "The retained document images do not reference this external payload",
38 )
39 .into());
40 }
41 }
42 let bytes = onestore_notebook::read_external_asset(source, section, filename, limit)?;
43 let mut connection = self
44 .connection
45 .lock()
46 .map_err(|_| io::Error::other("Cache owner panicked"))?;
47 let transaction = connection.transaction_with_behavior(TransactionBehavior::Immediate)?;
48 if !referenced(&transaction, &key)? {
49 return Err(io::Error::new(
50 io::ErrorKind::ResourceBusy,
51 "The external payload reference changed during download",
52 )
53 .into());
54 }
55 let previous = match cached(&transaction, &key, bytes.len()) {
56 Err(Error::Io(error)) if error.kind() == io::ErrorKind::FileTooLarge => {
57 return Err(Error::AssetChanged);
58 }
59 other => other?,
60 };
61 if let Some(previous) = previous {
62 if previous != bytes {
63 return Err(Error::AssetChanged);
64 }
65 } else {
66 transaction.execute(
67 "INSERT INTO assets(name,data,sha256) VALUES (?1,?2,?3)",
68 params![key, &bytes, &Sha256::digest(&bytes)[..]],
69 )?;
70 }
71 transaction.commit()?;
72 Ok(bytes)
73 }
74}
75
76pub(crate) fn key(filename: &str) -> Result<String> {
77 format!("<file>{filename}").parse::<FileDataReference>()?;
78 Ok(filename.to_ascii_lowercase())
79}
80
81fn referenced(connection: &Connection, key: &str) -> Result<bool> {
82 for column in ["working", "base"] {
83 let image: Vec<u8> = connection.query_row(
84 &format!("SELECT {column} FROM replica WHERE id=1"),
85 [],
86 |row| row.get(0),
87 )?;
88 let store = Store::parse(&image)?;
89 let index = RevisionIndex::parse(&store)?;
90 let document = Document::parse(&index)?;
91 if document.spaces.values().flat_map(|space| space.revisions.values())
92 .flat_map(|revision| revision.nodes.values()).any(|node| {
93 matches!(&node.kind, Kind::File { reference: FileDataReference::External(name), .. } if name.eq_ignore_ascii_case(key))
94 }) {
95 return Ok(true);
96 }
97 }
98 Ok(false)
99}
100
101pub(crate) fn cached(connection: &Connection, key: &str, limit: usize) -> Result<Option<Vec<u8>>> {
102 let version: u32 = connection.pragma_query_value(None, "user_version", |row| row.get(0))?;
103 if version < 5 {
104 return Ok(None);
105 }
106 let length: Option<i64> = connection
107 .query_row(
108 "SELECT length(data) FROM assets WHERE name=?1",
109 [key],
110 |row| row.get(0),
111 )
112 .optional()?;
113 let Some(length) = length else {
114 return Ok(None);
115 };
116 let length =
117 usize::try_from(length).map_err(|_| io::Error::from(io::ErrorKind::InvalidData))?;
118 if length > limit {
119 return Err(io::Error::from(io::ErrorKind::FileTooLarge).into());
120 }
121 let (bytes, expected): (Vec<u8>, Vec<u8>) = connection.query_row(
122 "SELECT data,sha256 FROM assets WHERE name=?1",
123 [key],
124 |row| Ok((row.get(0)?, row.get(1)?)),
125 )?;
126 if bytes.len() != length || Sha256::digest(&bytes)[..] != expected {
127 return Err(io::Error::new(
128 io::ErrorKind::InvalidData,
129 "Cached external payload checksum mismatch",
130 )
131 .into());
132 }
133 Ok(Some(bytes))
134}
crates/onestore-offline/src/lib.rs+6-1
......@@ -8,6 +8,7 @@ use onestore::{
88use rusqlite::{Connection, OpenFlags, TransactionBehavior, params};
99use std::{fs::OpenOptions, io, ops::Range, path::Path, sync::Mutex, time::Duration};
1010
11mod assets;
1112mod formatting;
1213mod rebase;
1314mod recovery;
......@@ -35,12 +36,16 @@ pub enum Error {
3536 Remote(#[from] onestore::CommitError),
3637 #[error(transparent)]
3738 RemoteIo(io::Error),
39 #[error(transparent)]
40 Notebook(#[from] onestore_notebook::Error),
41 #[error("External payload identity now refers to different bytes")]
42 AssetChanged,
3843}
3944
4045type Result<T> = std::result::Result<T, Error>;
4146
4247const APPLICATION_ID: u32 = 0x4f4e454f;
43const SCHEMA_VERSION: u32 = 4;
48const SCHEMA_VERSION: u32 = 5;
4449
4550/// Text and its observed precondition, retained across cache reopen and rebasing.
4651#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
crates/onestore-offline/src/recovery.rs+20-1
......@@ -13,6 +13,8 @@ pub struct RecoverySummary {
1313 pub published_receipts: u64,
1414 pub working_bytes: u64,
1515 pub remote_bytes: u64,
16 pub cached_assets: u64,
17 pub cached_asset_bytes: u64,
1618}
1719
1820/// Read-only recovery evidence; it cannot publish or acknowledge an edit.
......@@ -29,7 +31,7 @@ impl Recovery {
2931 let application: u32 =
3032 connection.pragma_query_value(None, "application_id", |row| row.get(0))?;
3133 let version: u32 = connection.pragma_query_value(None, "user_version", |row| row.get(0))?;
32 if application != RECOVERY_ID || version != SCHEMA_VERSION {
34 if application != RECOVERY_ID || !(4..=SCHEMA_VERSION).contains(&version) {
3335 return Err(io::Error::new(
3436 io::ErrorKind::InvalidData,
3537 "Unrecognized recovery archive or unsupported schema version",
......@@ -73,6 +75,11 @@ impl Recovery {
7375 pub fn receipts(&self) -> Result<BTreeMap<u64, ExGuid>> {
7476 receipts(&self.connection)
7577 }
78
79 /// Reads a previously downloaded external payload without accessing its former server.
80 pub fn cached_asset(&self, filename: &str, limit: usize) -> Result<Option<Vec<u8>>> {
81 assets::cached(&self.connection, &assets::key(filename)?, limit)
82 }
7683}
7784
7885impl Replica {
......@@ -131,6 +138,16 @@ impl Replica {
131138}
132139
133140fn summary(connection: &Connection) -> Result<RecoverySummary> {
141 let version: u32 = connection.pragma_query_value(None, "user_version", |row| row.get(0))?;
142 let (cached_assets, cached_asset_bytes) = if version < 5 {
143 (0, 0)
144 } else {
145 connection.query_row(
146 "SELECT count(*),coalesce(sum(length(data)),0) FROM assets",
147 [],
148 |row| Ok((unsigned(row, 0)?, unsigned(row, 1)?)),
149 )?
150 };
134151 Ok(connection.query_row(
135152 "SELECT (SELECT count(*) FROM edits), (SELECT count(*) FROM conflicts),
136153 (SELECT count(*) FROM attempt), (SELECT count(*) FROM receipts),
......@@ -144,6 +161,8 @@ fn summary(connection: &Connection) -> Result<RecoverySummary> {
144161 published_receipts: unsigned(row, 3)?,
145162 working_bytes: unsigned(row, 4)?,
146163 remote_bytes: unsigned(row, 5)?,
164 cached_assets,
165 cached_asset_bytes,
147166 })
148167 },
149168 )?)
crates/onestore-offline/src/schema.rs+12
......@@ -6,6 +6,12 @@ const CONFLICTS: &str = "CREATE TABLE conflicts (
66 kind INTEGER NOT NULL CHECK(kind BETWEEN 0 AND 3)
77) STRICT;";
88
9const ASSETS: &str = "CREATE TABLE assets (
10 name TEXT PRIMARY KEY NOT NULL,
11 data BLOB NOT NULL,
12 sha256 BLOB NOT NULL CHECK(length(sha256)=32)
13) STRICT;";
14
915pub(crate) fn create(transaction: &Transaction<'_>) -> Result<()> {
1016 transaction.execute_batch(
1117 "CREATE TABLE edits (
......@@ -24,16 +30,22 @@ pub(crate) fn create(transaction: &Transaction<'_>) -> Result<()> {
2430 ) STRICT;",
2531 )?;
2632 transaction.execute_batch(CONFLICTS)?;
33 transaction.execute_batch(ASSETS)?;
2734 Ok(())
2835}
2936
3037pub(crate) fn migrate(transaction: &Transaction<'_>, version: u32) -> Result<()> {
38 if version == 4 {
39 transaction.execute_batch(ASSETS)?;
40 return Ok(());
41 }
3142 if version == 3 {
3243 transaction.execute_batch("ALTER TABLE conflicts RENAME TO old_conflicts;")?;
3344 transaction.execute_batch(CONFLICTS)?;
3445 transaction.execute_batch(
3546 "INSERT INTO conflicts SELECT * FROM old_conflicts; DROP TABLE old_conflicts;",
3647 )?;
48 transaction.execute_batch(ASSETS)?;
3749 return Ok(());
3850 }
3951
crates/onestore-offline/tests/assets.rs created+449
......@@ -0,0 +1,449 @@
1use onestore_offline::{EditStatus, Error, Operation, Recovery, Replica};
2use std::{
3 fs, io,
4 path::Path,
5 sync::{Barrier, mpsc},
6 time::Duration,
7};
8
9const FIXTURE: &str = "../../corpus/native-external-assets/notebook";
10const LEGACY: &str = "../../corpus/offline-v4/live.sqlite";
11
12fn copied_cache(root: &Path) -> Replica {
13 let path = root.join("cache.sqlite");
14 assert!(!path.exists());
15 fs::copy(LEGACY, &path).unwrap();
16 Replica::open(path).unwrap()
17}
18
19fn payload(size: usize) -> (String, Vec<u8>) {
20 fs::read_dir(Path::new(FIXTURE).join("synthetic_onefiles"))
21 .unwrap()
22 .map(|entry| entry.unwrap().path())
23 .find_map(|path| {
24 let bytes = fs::read(&path).unwrap();
25 (bytes.len() == size)
26 .then(|| (path.file_name().unwrap().to_str().unwrap().into(), bytes))
27 })
28 .unwrap()
29}
30
31struct Payload(Option<io::Result<Vec<u8>>>);
32impl onestore_notebook::Source for Payload {
33 fn entries(&mut self, _: &str, _: usize) -> io::Result<Vec<onestore_notebook::Entry>> {
34 panic!("Unexpected enumeration")
35 }
36 fn read(&mut self, _: &str, _: usize) -> io::Result<Vec<u8>> {
37 self.0.take().expect("Unexpected repeated payload read")
38 }
39}
40
41#[test]
42fn downloaded_media_survives_migration_reopen_and_recovery_with_the_queue_intact() {
43 let directory = tempfile::tempdir().unwrap();
44 let original = fs::read(LEGACY).unwrap();
45 let cache = copied_cache(directory.path());
46 let working = cache.snapshot().unwrap();
47 let remote = cache.remote_snapshot().unwrap();
48 let pending = cache.pending().unwrap();
49 assert_eq!(
50 pending.iter().map(|edit| edit.id).collect::<Vec<_>>(),
51 [1, 2]
52 );
53 let uncertain = cache.status(1).unwrap();
54 assert!(matches!(
55 uncertain,
56 Some(EditStatus::AwaitingConfirmation { .. })
57 ));
58 let mut source = onestore_notebook::Local::open(FIXTURE).unwrap();
59 for size in [0, 1024] {
60 let (name, bytes) = payload(size);
61 assert!(cache.cached_asset(&name, size).unwrap().is_none());
62 assert_eq!(
63 cache
64 .fetch_asset(&mut source, "synthetic.one", &name, size)
65 .unwrap(),
66 bytes
67 );
68 assert_eq!(
69 cache
70 .cached_asset(&name.to_ascii_lowercase(), size)
71 .unwrap(),
72 Some(bytes)
73 );
74 }
75 let summary = cache.recovery_summary().unwrap();
76 assert_eq!(
77 (summary.cached_assets, summary.cached_asset_bytes),
78 (2, 1024)
79 );
80 assert_eq!(cache.snapshot().unwrap(), working);
81 assert_eq!(cache.remote_snapshot().unwrap(), remote);
82 assert_eq!(cache.pending().unwrap(), pending);
83 assert_eq!(cache.status(1).unwrap(), uncertain);
84 cache
85 .export_recovery(directory.path().join("recovery.sqlite"))
86 .unwrap();
87 drop(cache);
88 let cache = Replica::open(directory.path().join("cache.sqlite")).unwrap();
89 let recovery = Recovery::open(directory.path().join("recovery.sqlite")).unwrap();
90 assert_eq!(recovery.summary().unwrap(), summary);
91 assert_eq!(recovery.pending().unwrap(), pending);
92 assert_eq!(recovery.status(1).unwrap(), uncertain);
93 assert!(recovery.receipts().unwrap().is_empty());
94 for size in [0, 1024] {
95 let (name, bytes) = payload(size);
96 assert_eq!(
97 cache.cached_asset(&name, size).unwrap(),
98 Some(bytes.clone())
99 );
100 assert_eq!(recovery.cached_asset(&name, size).unwrap(), Some(bytes));
101 }
102 assert_eq!(fs::read(LEGACY).unwrap(), original);
103}
104
105#[test]
106fn an_original_version_four_archive_remains_read_only_and_has_no_cached_media() {
107 let path = "../../corpus/offline-v4/recovery.sqlite";
108 let before = fs::read(path).unwrap();
109 let archive = Recovery::open(path).unwrap();
110 assert_eq!(archive.pending().unwrap().len(), 2);
111 assert!(matches!(
112 archive.status(1).unwrap(),
113 Some(EditStatus::AwaitingConfirmation { .. })
114 ));
115 assert_eq!(archive.summary().unwrap().cached_assets, 0);
116 assert_eq!(archive.summary().unwrap().cached_asset_bytes, 0);
117 assert!(archive.cached_asset(&payload(0).0, 0).unwrap().is_none());
118 drop(archive);
119 assert_eq!(fs::read(path).unwrap(), before);
120}
121
122#[test]
123fn unsuccessful_refreshes_and_changed_identity_data_preserve_the_downloaded_payload() {
124 let directory = tempfile::tempdir().unwrap();
125 let cache = copied_cache(directory.path());
126 let (name, bytes) = payload(1024);
127 let mut source = Payload(Some(Ok(bytes.clone())));
128 cache
129 .fetch_asset(&mut source, "synthetic.one", &name, 1024)
130 .unwrap();
131 let before = fs::read(directory.path().join("cache.sqlite")).unwrap();
132 let mut unavailable = Payload(Some(Err(io::ErrorKind::ConnectionReset.into())));
133 assert!(
134 matches!(cache.fetch_asset(&mut unavailable, "synthetic.one", &name, 1024),
135 Err(Error::Notebook(onestore_notebook::Error::Io { error, .. })) if error.kind() == io::ErrorKind::ConnectionReset)
136 );
137 for size in [0, 1024, 2048] {
138 let mut changed = Payload(Some(Ok(vec![9; size])));
139 assert!(matches!(
140 cache.fetch_asset(&mut changed, "synthetic.one", &name, 2048),
141 Err(Error::AssetChanged)
142 ));
143 }
144 assert!(
145 matches!(cache.cached_asset(&name, 1023), Err(Error::Io(error)) if error.kind() == io::ErrorKind::FileTooLarge)
146 );
147 assert_eq!(cache.cached_asset(&name, 1024).unwrap(), Some(bytes));
148 assert_eq!(
149 fs::read(directory.path().join("cache.sqlite")).unwrap(),
150 before
151 );
152}
153
154#[test]
155fn unreferenced_payloads_are_rejected_before_io_or_local_changes() {
156 let directory = tempfile::tempdir().unwrap();
157 let cache = copied_cache(directory.path());
158 let mut unused = Payload(None);
159 assert!(
160 matches!(cache.fetch_asset(&mut unused, "synthetic.one", "00000000-0000-0000-0000-000000000001.onebin", 100),
161 Err(Error::Io(error)) if error.kind() == io::ErrorKind::InvalidInput)
162 );
163 assert!(
164 cache
165 .fetch_asset(&mut unused, "synthetic.one", "../payload.onebin", 100)
166 .is_err()
167 );
168 assert_eq!(cache.recovery_summary().unwrap().cached_assets, 0);
169}
170
171#[test]
172fn download_network_wait_does_not_block_local_edits() {
173 struct Waiting {
174 entered: mpsc::Sender<()>,
175 released: mpsc::Receiver<()>,
176 bytes: Vec<u8>,
177 }
178 impl onestore_notebook::Source for Waiting {
179 fn entries(&mut self, _: &str, _: usize) -> io::Result<Vec<onestore_notebook::Entry>> {
180 panic!("Unexpected enumeration")
181 }
182 fn read(&mut self, _: &str, _: usize) -> io::Result<Vec<u8>> {
183 self.entered.send(()).unwrap();
184 self.released.recv_timeout(Duration::from_secs(5)).unwrap();
185 Ok(self.bytes.clone())
186 }
187 }
188 let directory = tempfile::tempdir().unwrap();
189 let cache = copied_cache(directory.path());
190 let (name, bytes) = payload(1024);
191 let (entered, waiting) = mpsc::channel();
192 let (release, released) = mpsc::channel();
193 let mut source = Waiting {
194 entered,
195 released,
196 bytes: bytes.clone(),
197 };
198 std::thread::scope(|scope| {
199 let download = scope.spawn(|| {
200 cache
201 .fetch_asset(&mut source, "synthetic.one", &name, 1024)
202 .unwrap()
203 });
204 waiting.recv_timeout(Duration::from_secs(5)).unwrap();
205 let pending = cache.pending().unwrap();
206 let Operation::Text(edit) = &pending[0].operation else {
207 panic!()
208 };
209 let id = cache
210 .edit_text(
211 &cache.snapshot().unwrap(),
212 pending[0].space,
213 edit.object,
214 0..0,
215 "during download ",
216 )
217 .unwrap()
218 .unwrap();
219 release.send(()).unwrap();
220 assert_eq!(download.join().unwrap(), bytes);
221 assert_eq!(cache.status(id).unwrap(), Some(EditStatus::Pending));
222 });
223}
224
225#[test]
226fn concurrent_downloads_publish_one_immutable_cache_entry() {
227 let directory = tempfile::tempdir().unwrap();
228 let cache = copied_cache(directory.path());
229 let (name, bytes) = payload(1024);
230 let ready = Barrier::new(8);
231 std::thread::scope(|scope| {
232 for _ in 0..8 {
233 scope.spawn(|| {
234 let mut source = onestore_notebook::Local::open(FIXTURE).unwrap();
235 ready.wait();
236 assert_eq!(
237 cache
238 .fetch_asset(&mut source, "synthetic.one", &name, 1024)
239 .unwrap(),
240 bytes
241 );
242 });
243 }
244 });
245 assert_eq!(cache.recovery_summary().unwrap().cached_assets, 1);
246}
247
248#[test]
249fn a_native_refresh_removing_the_reference_rejects_an_inflight_download() {
250 struct Native;
251 impl onestore_offline::Remote for Native {
252 fn read(&mut self) -> io::Result<Vec<u8>> {
253 fs::read("../../corpus/native-external-assets/native/synthetic.one")
254 }
255 fn publish(&mut self, _: &onestore::PreparedEdit<'_>) -> Result<(), onestore::CommitError> {
256 panic!("Unexpected publication")
257 }
258 fn confirm(&mut self, _: &[u8]) -> Result<(), onestore::CommitError> {
259 panic!("Unexpected confirmation")
260 }
261 }
262 struct Refresh<'a>(&'a Replica);
263 impl onestore_notebook::Source for Refresh<'_> {
264 fn entries(&mut self, _: &str, _: usize) -> io::Result<Vec<onestore_notebook::Entry>> {
265 panic!("Unexpected enumeration")
266 }
267 fn read(&mut self, _: &str, _: usize) -> io::Result<Vec<u8>> {
268 assert_eq!(self.0.sync_once(&mut Native).unwrap(), None);
269 Ok(payload(1024).1)
270 }
271 }
272 let root = tempfile::tempdir().unwrap();
273 let source = fs::read(Path::new(FIXTURE).join("synthetic.one")).unwrap();
274 let cache = Replica::create(root.path().join("cache.sqlite"), &source).unwrap();
275 let (name, _) = payload(1024);
276 assert!(
277 matches!(cache.fetch_asset(&mut Refresh(&cache), "synthetic.one", &name, 1024),
278 Err(Error::Io(error)) if error.kind() == io::ErrorKind::ResourceBusy)
279 );
280 assert!(cache.cached_asset(&name, 1024).unwrap().is_none());
281 assert_eq!(
282 cache.snapshot().unwrap(),
283 fs::read("../../corpus/native-external-assets/native/synthetic.one").unwrap()
284 );
285}
286
287#[test]
288fn local_failure_and_bad_cached_bytes_never_become_successful_downloads() {
289 let directory = tempfile::tempdir().unwrap();
290 drop(copied_cache(directory.path()));
291 let path = directory.path().join("cache.sqlite");
292 let connection = rusqlite::Connection::open(&path).unwrap();
293 connection.execute_batch("CREATE TRIGGER fail_asset BEFORE INSERT ON assets BEGIN SELECT RAISE(ABORT,'Test asset failure'); END").unwrap();
294 drop(connection);
295 let (name, bytes) = payload(1024);
296 let cache = Replica::open(&path).unwrap();
297 let mut source = onestore_notebook::Local::open(FIXTURE).unwrap();
298 assert!(matches!(
299 cache.fetch_asset(&mut source, "synthetic.one", &name, 1024),
300 Err(Error::Database(_))
301 ));
302 assert!(cache.cached_asset(&name, 1024).unwrap().is_none());
303 assert_eq!(cache.pending().unwrap().len(), 2);
304 drop(cache);
305 let connection = rusqlite::Connection::open(&path).unwrap();
306 connection.execute_batch("DROP TRIGGER fail_asset").unwrap();
307 drop(connection);
308 let cache = Replica::open(&path).unwrap();
309 cache
310 .fetch_asset(&mut source, "synthetic.one", &name, 1024)
311 .unwrap();
312 drop(cache);
313 let connection = rusqlite::Connection::open(&path).unwrap();
314 let mut damaged = bytes;
315 damaged[0] ^= 1;
316 connection
317 .execute("UPDATE assets SET data=?1", [damaged])
318 .unwrap();
319 drop(connection);
320 let cache = Replica::open(&path).unwrap();
321 assert!(
322 matches!(cache.cached_asset(&name, 1024), Err(Error::Io(error)) if error.kind() == io::ErrorKind::InvalidData)
323 );
324 assert!(
325 matches!(cache.fetch_asset(&mut source, "synthetic.one", &name, 1024), Err(Error::Io(error)) if error.kind() == io::ErrorKind::InvalidData)
326 );
327 cache
328 .export_recovery(directory.path().join("damaged.sqlite"))
329 .unwrap();
330 let archive = Recovery::open(directory.path().join("damaged.sqlite")).unwrap();
331 assert!(archive.cached_asset(&name, 1024).is_err());
332 assert_eq!(archive.pending().unwrap().len(), 2);
333}
334
335#[test]
336fn abrupt_process_exit_retains_only_completed_downloads_and_archives() {
337 const CHILD: &str = "ONESTORE_ASSET_EXIT_CASE";
338 if let Ok(phase) = std::env::var(CHILD) {
339 struct ExitDuringRead;
340 impl onestore_notebook::Source for ExitDuringRead {
341 fn entries(&mut self, _: &str, _: usize) -> io::Result<Vec<onestore_notebook::Entry>> {
342 panic!("Unexpected enumeration")
343 }
344 fn read(&mut self, _: &str, _: usize) -> io::Result<Vec<u8>> {
345 std::process::exit(83)
346 }
347 }
348 let root = std::env::var("ONESTORE_ASSET_EXIT_ROOT").unwrap();
349 let cache = copied_cache(Path::new(&root));
350 let (name, bytes) = payload(1024);
351 if phase == "download" {
352 cache
353 .fetch_asset(&mut ExitDuringRead, "synthetic.one", &name, bytes.len())
354 .unwrap();
355 panic!("Read returned after process exit");
356 }
357 cache
358 .fetch_asset(
359 &mut onestore_notebook::Local::open(FIXTURE).unwrap(),
360 "synthetic.one",
361 &name,
362 bytes.len(),
363 )
364 .unwrap();
365 if phase == "archive" {
366 cache
367 .export_recovery(Path::new(&root).join("recovery.sqlite"))
368 .unwrap();
369 }
370 std::process::exit(83);
371 }
372 for phase in ["download", "cached", "archive"] {
373 let root = tempfile::tempdir().unwrap();
374 let output = std::process::Command::new(std::env::current_exe().unwrap())
375 .args([
376 "--exact",
377 "abrupt_process_exit_retains_only_completed_downloads_and_archives",
378 ])
379 .env(CHILD, phase)
380 .env("ONESTORE_ASSET_EXIT_ROOT", root.path())
381 .output()
382 .unwrap();
383 assert_eq!(
384 output.status.code(),
385 Some(83),
386 "{}",
387 String::from_utf8_lossy(&output.stderr)
388 );
389 let cache = Replica::open(root.path().join("cache.sqlite")).unwrap();
390 let (name, bytes) = payload(1024);
391 let expected = (phase != "download").then_some(bytes);
392 assert_eq!(cache.cached_asset(&name, 1024).unwrap(), expected);
393 assert_eq!(cache.pending().unwrap().len(), 2);
394 assert!(matches!(
395 cache.status(1).unwrap(),
396 Some(EditStatus::AwaitingConfirmation { .. })
397 ));
398 if phase == "archive" {
399 let recovery = Recovery::open(root.path().join("recovery.sqlite")).unwrap();
400 assert_eq!(recovery.cached_asset(&name, 1024).unwrap(), expected);
401 assert_eq!(recovery.pending().unwrap(), cache.pending().unwrap());
402 assert_eq!(recovery.status(1).unwrap(), cache.status(1).unwrap());
403 }
404 }
405}
406
407#[test]
408#[cfg(feature = "smb")]
409#[ignore = "requires a disposable Samba mirror at ONESTORE_SMB_NOTEBOOK"]
410fn live_smb_downloads_survive_disconnect_and_cache_reopen() {
411 let root = tempfile::tempdir().unwrap();
412 let cache = copied_cache(root.path());
413 let client = onestore_smb::Client::connect(
414 &std::env::var("ONESTORE_SMB_LAB").unwrap(),
415 "agent",
416 onestore_smb::Credentials::default(),
417 Duration::from_secs(5),
418 )
419 .unwrap();
420 let notebook = std::env::var("ONESTORE_SMB_NOTEBOOK").unwrap();
421 let mut source = onestore_notebook::Smb::new(&client, &notebook).unwrap();
422 for size in [0, 771, 1024] {
423 let (name, bytes) = payload(size);
424 assert_eq!(
425 cache
426 .fetch_asset(&mut source, "synthetic.one", &name, size)
427 .unwrap(),
428 bytes
429 );
430 }
431 drop(client);
432 cache
433 .export_recovery(root.path().join("recovery.sqlite"))
434 .unwrap();
435 drop(cache);
436 let cache = Replica::open(root.path().join("cache.sqlite")).unwrap();
437 let recovery = Recovery::open(root.path().join("recovery.sqlite")).unwrap();
438 for size in [0, 771, 1024] {
439 let (name, bytes) = payload(size);
440 assert_eq!(
441 cache.cached_asset(&name, size).unwrap(),
442 Some(bytes.clone())
443 );
444 assert_eq!(recovery.cached_asset(&name, size).unwrap(), Some(bytes));
445 }
446 assert_eq!(cache.pending().unwrap(), recovery.pending().unwrap());
447 assert_eq!(cache.status(1).unwrap(), recovery.status(1).unwrap());
448 assert_eq!(cache.recovery_summary().unwrap().cached_assets, 3);
449}
crates/onestore-offline/tests/cache.rs+1-1
......@@ -734,7 +734,7 @@ fn unrecognized_persisted_operations_are_rejected_without_dropping_fields() {
734734 db.execute("UPDATE edits SET operation=?1", [value.to_string()])
735735 .unwrap();
736736 if operation != "Format" {
737 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();
737 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();
738738 }
739739 drop(db);
740740 let before = fs::read(&path).unwrap();
crates/onestore-offline/tests/sync.rs+4-2
......@@ -197,6 +197,8 @@ fn recovery_archive_preserves_typed_queue_uncertainty_and_receipts_without_becom
197197 published_receipts: 1,
198198 working_bytes: working.len() as u64,
199199 remote_bytes: remote.len() as u64,
200 cached_assets: 0,
201 cached_asset_bytes: 0,
200202 };
201203 assert_eq!(cache.recovery_summary().unwrap(), summary);
202204 cache.export_recovery(&archive_path).unwrap();
......@@ -712,7 +714,7 @@ fn version_one_cache_migration_preserves_images_intents_and_local_ids() {
712714 drop(cache);
713715 let db = rusqlite::Connection::open(&path).unwrap();
714716 db.execute_batch(
715 "DROP TABLE attempt; DROP TABLE conflicts; DROP TABLE receipts; DROP TABLE edits;
717 "DROP TABLE attempt; DROP TABLE conflicts; DROP TABLE receipts; DROP TABLE edits; DROP TABLE assets;
716718 CREATE TABLE edits (
717719 id INTEGER PRIMARY KEY AUTOINCREMENT CHECK(id>0), space TEXT NOT NULL,
718720 object TEXT NOT NULL, before_text TEXT NOT NULL,
......@@ -737,7 +739,7 @@ fn version_one_cache_migration_preserves_images_intents_and_local_ids() {
737739 assert_eq!(
738740 db.pragma_query_value(None, "user_version", |row| row.get::<_, u32>(0))
739741 .unwrap(),
740 4
742 5
741743 );
742744}
743745