1use notebook::{EditStatus, Error, Recovery, Replica};
2use onestore::{
3 ExGuid, RevisionIndex, Store,
4 document::Document,
5 page::{Page, PageObject},
6};
7use std::{
8 fs, io,
9 path::Path,
10 sync::{Barrier, mpsc},
11 time::Duration,
12};
13
14#[path = "support/model_ops.rs"]
15mod model_ops;
16#[path = "support/server.rs"]
17mod server;
18use server::{remote_snapshot, snapshot};
19
20const FIXTURE: &str = "../../corpus/native-external-assets/notebook";
21
22/// The first body text object of the section referencing the fixture's media.
23fn first_text(bytes: &[u8]) -> ExGuid {
24 let store = Store::parse(bytes).unwrap();
25 let index = RevisionIndex::parse(&store).unwrap();
26 let document = Document::parse(&index).unwrap();
27 document
28 .pages()
29 .unwrap()
30 .into_iter()
31 .find_map(|(space, _)| {
32 Page::from_space(&document, space)
33 .ok()?
34 .objects
35 .iter()
36 .find_map(|object| match object {
37 PageObject::Outline(outline) => outline
38 .paragraphs
39 .iter()
40 .find_map(|p| p.text().map(|t| t.id)),
41 _ => None,
42 })
43 })
44 .unwrap()
45}
46
47/// A cache of the media fixture with a queued text save and a queued page creation.
48fn queued_cache(root: &Path) -> Replica {
49 let path = root.join("cache.sqlite");
50 assert!(!path.exists());
51 let source = fs::read(Path::new(FIXTURE).join("synthetic.one")).unwrap();
52 let cache = Replica::create(path, &source).unwrap();
53 let text = first_text(&source);
54 model_ops::save(&cache, text, |page| {
55 model_ops::replace_text(page, text, 0..0, "Queued ")
56 })
57 .unwrap()
58 .unwrap();
59 let page = onestore::PageCreation::new(None, Some("Queued page"), "Author").unwrap();
60 server::section_op(&cache, onestore::op::SectionOp::Create(page));
61 cache
62}
63
64fn payload(size: usize) -> (String, Vec<u8>) {
65 fs::read_dir(Path::new(FIXTURE).join("synthetic_onefiles"))
66 .unwrap()
67 .map(|entry| entry.unwrap().path())
68 .find_map(|path| {
69 let bytes = fs::read(&path).unwrap();
70 (bytes.len() == size)
71 .then(|| (path.file_name().unwrap().to_str().unwrap().into(), bytes))
72 })
73 .unwrap()
74}
75
76struct Payload(Option<io::Result<Vec<u8>>>);
77impl notebook::discover::Source for Payload {
78 fn entries(&mut self, _: &str, _: usize) -> io::Result<Vec<notebook::discover::Entry>> {
79 panic!("Unexpected enumeration")
80 }
81 fn read(&mut self, _: &str, _: usize) -> io::Result<Vec<u8>> {
82 self.0.take().expect("Unexpected repeated payload read")
83 }
84}
85
86#[test]
87fn downloaded_media_survives_reopen_and_recovery_with_the_queue_intact() {
88 let directory = tempfile::tempdir().unwrap();
89 let cache = queued_cache(directory.path());
90 let working = snapshot(&cache);
91 let remote = remote_snapshot(&cache);
92 let pending = cache.pending().unwrap();
93 assert_eq!(
94 pending.iter().map(|edit| edit.id).collect::<Vec<_>>(),
95 [1, 2]
96 );
97 let uncertain = cache.status(1).unwrap();
98 assert_eq!(uncertain, Some(EditStatus::Pending));
99 let mut source = notebook::discover::Local::open(FIXTURE).unwrap();
100 for size in [0, 1024] {
101 let (name, bytes) = payload(size);
102 assert!(cache.cached_asset(&name, size).unwrap().is_none());
103 assert_eq!(
104 cache
105 .fetch_asset(&mut source, "synthetic.one", &name, size)
106 .unwrap(),
107 bytes
108 );
109 assert_eq!(
110 cache
111 .cached_asset(&name.to_ascii_lowercase(), size)
112 .unwrap(),
113 Some(bytes)
114 );
115 }
116 let summary = cache.recovery_summary().unwrap();
117 assert_eq!(
118 (summary.cached_assets, summary.cached_asset_bytes),
119 (2, 1024)
120 );
121 assert_eq!(server::pages(&snapshot(&cache)), server::pages(&working));
122 assert_eq!(remote_snapshot(&cache), remote);
123 assert_eq!(cache.pending().unwrap(), pending);
124 assert_eq!(cache.status(1).unwrap(), uncertain);
125 cache
126 .export_recovery(directory.path().join("recovery.sqlite"))
127 .unwrap();
128 drop(cache);
129 let cache = Replica::open(directory.path().join("cache.sqlite")).unwrap();
130 let recovery = Recovery::open(directory.path().join("recovery.sqlite")).unwrap();
131 assert_eq!(recovery.summary().unwrap(), summary);
132 assert_eq!(recovery.pending().unwrap(), pending);
133 assert_eq!(recovery.status(1).unwrap(), uncertain);
134 assert!(recovery.receipts().unwrap().is_empty());
135 for size in [0, 1024] {
136 let (name, bytes) = payload(size);
137 assert_eq!(
138 cache.cached_asset(&name, size).unwrap(),
139 Some(bytes.clone())
140 );
141 assert_eq!(recovery.cached_asset(&name, size).unwrap(), Some(bytes));
142 }
143}
144
145#[test]
146fn unsuccessful_refreshes_and_changed_identity_data_preserve_the_downloaded_payload() {
147 let directory = tempfile::tempdir().unwrap();
148 let cache = queued_cache(directory.path());
149 let (name, bytes) = payload(1024);
150 let mut source = Payload(Some(Ok(bytes.clone())));
151 cache
152 .fetch_asset(&mut source, "synthetic.one", &name, 1024)
153 .unwrap();
154 let before = fs::read(directory.path().join("cache.sqlite")).unwrap();
155 let mut unavailable = Payload(Some(Err(io::ErrorKind::ConnectionReset.into())));
156 assert!(
157 matches!(cache.fetch_asset(&mut unavailable, "synthetic.one", &name, 1024),
158 Err(Error::Discovery(notebook::discover::Error::Io { error, .. })) if error.kind() == io::ErrorKind::ConnectionReset)
159 );
160 for size in [0, 1024, 2048] {
161 let mut changed = Payload(Some(Ok(vec![9; size])));
162 assert!(matches!(
163 cache.fetch_asset(&mut changed, "synthetic.one", &name, 2048),
164 Err(Error::AssetChanged)
165 ));
166 }
167 assert!(
168 matches!(cache.cached_asset(&name, 1023), Err(Error::Io(error)) if error.kind() == io::ErrorKind::FileTooLarge)
169 );
170 assert_eq!(cache.cached_asset(&name, 1024).unwrap(), Some(bytes));
171 assert_eq!(
172 fs::read(directory.path().join("cache.sqlite")).unwrap(),
173 before
174 );
175}
176
177#[test]
178fn unreferenced_payloads_are_rejected_before_io_or_local_changes() {
179 let directory = tempfile::tempdir().unwrap();
180 let cache = queued_cache(directory.path());
181 let mut unused = Payload(None);
182 assert!(
183 matches!(cache.fetch_asset(&mut unused, "synthetic.one", "00000000-0000-0000-0000-000000000001.onebin", 100),
184 Err(Error::Io(error)) if error.kind() == io::ErrorKind::InvalidInput)
185 );
186 assert!(
187 cache
188 .fetch_asset(&mut unused, "synthetic.one", "../payload.onebin", 100)
189 .is_err()
190 );
191 assert_eq!(cache.recovery_summary().unwrap().cached_assets, 0);
192}
193
194#[test]
195fn download_network_wait_does_not_block_local_edits() {
196 struct Waiting {
197 entered: mpsc::Sender<()>,
198 released: mpsc::Receiver<()>,
199 bytes: Vec<u8>,
200 }
201 impl notebook::discover::Source for Waiting {
202 fn entries(&mut self, _: &str, _: usize) -> io::Result<Vec<notebook::discover::Entry>> {
203 panic!("Unexpected enumeration")
204 }
205 fn read(&mut self, _: &str, _: usize) -> io::Result<Vec<u8>> {
206 self.entered.send(()).unwrap();
207 self.released
208 .recv_timeout(Duration::from_secs(120))
209 .unwrap();
210 Ok(self.bytes.clone())
211 }
212 }
213 let directory = tempfile::tempdir().unwrap();
214 let cache = queued_cache(directory.path());
215 let (name, bytes) = payload(1024);
216 let (entered, waiting) = mpsc::channel();
217 let (release, released) = mpsc::channel();
218 let mut source = Waiting {
219 entered,
220 released,
221 bytes: bytes.clone(),
222 };
223 std::thread::scope(|scope| {
224 let download = scope.spawn(|| {
225 cache
226 .fetch_asset(&mut source, "synthetic.one", &name, 1024)
227 .unwrap()
228 });
229 waiting.recv_timeout(Duration::from_secs(120)).unwrap();
230 let text = first_text(&snapshot(&cache));
231 let id = model_ops::save(&cache, text, |page| {
232 model_ops::replace_text(page, text, 0..0, "during download ")
233 })
234 .unwrap()
235 .unwrap();
236 release.send(()).unwrap();
237 assert_eq!(download.join().unwrap(), bytes);
238 assert_eq!(cache.status(id).unwrap(), Some(EditStatus::Pending));
239 assert!(matches!(
240 &cache.pending().unwrap()[2].edit.ops[..],
241 [onestore::op::Op::Page { op: onestore::op::PageOp::Text { with, .. }, .. }]
242 if with == "during download "
243 ));
244 });
245}
246
247#[test]
248fn concurrent_downloads_publish_one_immutable_cache_entry() {
249 let directory = tempfile::tempdir().unwrap();
250 let cache = queued_cache(directory.path());
251 let (name, bytes) = payload(1024);
252 let ready = Barrier::new(8);
253 std::thread::scope(|scope| {
254 for _ in 0..8 {
255 scope.spawn(|| {
256 let mut source = notebook::discover::Local::open(FIXTURE).unwrap();
257 ready.wait();
258 assert_eq!(
259 cache
260 .fetch_asset(&mut source, "synthetic.one", &name, 1024)
261 .unwrap(),
262 bytes
263 );
264 });
265 }
266 });
267 assert_eq!(cache.recovery_summary().unwrap().cached_assets, 1);
268}
269
270#[test]
271fn a_native_refresh_removing_the_reference_rejects_an_inflight_download() {
272 struct Native;
273 impl notebook::Remote for Native {
274 fn read(&mut self) -> io::Result<Vec<u8>> {
275 let mut image = fs::read("../../corpus/native-external-assets/native/synthetic.one")?;
276 // The fixture keeps the old header; a native commit writes a new file version.
277 image[212] ^= 1;
278 Ok(image)
279 }
280 fn publish(&mut self, _: &onestore::Transaction) -> Result<(), onestore::CommitError> {
281 panic!("Unexpected publication")
282 }
283 fn stamp(&mut self) -> io::Result<onestore::Stamp> {
284 onestore::Stamp::of(&self.read()?).map_err(io::Error::other)
285 }
286 fn confirm(&mut self, _: &onestore::Stamp) -> Result<(), onestore::CommitError> {
287 panic!("Unexpected confirmation")
288 }
289 }
290 struct Refresh<'a>(&'a Replica);
291 impl notebook::discover::Source for Refresh<'_> {
292 fn entries(&mut self, _: &str, _: usize) -> io::Result<Vec<notebook::discover::Entry>> {
293 panic!("Unexpected enumeration")
294 }
295 fn read(&mut self, _: &str, _: usize) -> io::Result<Vec<u8>> {
296 assert_eq!(self.0.sync_once(&mut Native).unwrap().edit, None);
297 Ok(payload(1024).1)
298 }
299 }
300 let root = tempfile::tempdir().unwrap();
301 let source = fs::read(Path::new(FIXTURE).join("synthetic.one")).unwrap();
302 let cache = Replica::create(root.path().join("cache.sqlite"), &source).unwrap();
303 let (name, _) = payload(1024);
304 assert!(
305 matches!(cache.fetch_asset(&mut Refresh(&cache), "synthetic.one", &name, 1024),
306 Err(Error::Io(error)) if error.kind() == io::ErrorKind::ResourceBusy)
307 );
308 assert!(cache.cached_asset(&name, 1024).unwrap().is_none());
309 assert_eq!(
310 snapshot(&cache),
311 notebook::Remote::read(&mut Native).unwrap()
312 );
313}
314
315#[test]
316fn local_failure_and_bad_cached_bytes_never_become_successful_downloads() {
317 let directory = tempfile::tempdir().unwrap();
318 drop(queued_cache(directory.path()));
319 let path = directory.path().join("cache.sqlite");
320 let connection = rusqlite::Connection::open(&path).unwrap();
321 connection.execute_batch("CREATE TRIGGER fail_asset BEFORE INSERT ON assets BEGIN SELECT RAISE(ABORT,'Test asset failure'); END").unwrap();
322 drop(connection);
323 let (name, bytes) = payload(1024);
324 let cache = Replica::open(&path).unwrap();
325 let mut source = notebook::discover::Local::open(FIXTURE).unwrap();
326 assert!(matches!(
327 cache.fetch_asset(&mut source, "synthetic.one", &name, 1024),
328 Err(Error::Database(_))
329 ));
330 assert!(cache.cached_asset(&name, 1024).unwrap().is_none());
331 assert_eq!(cache.pending().unwrap().len(), 2);
332 drop(cache);
333 let connection = rusqlite::Connection::open(&path).unwrap();
334 connection.execute_batch("DROP TRIGGER fail_asset").unwrap();
335 drop(connection);
336 let cache = Replica::open(&path).unwrap();
337 cache
338 .fetch_asset(&mut source, "synthetic.one", &name, 1024)
339 .unwrap();
340 drop(cache);
341 let connection = rusqlite::Connection::open(&path).unwrap();
342 let mut damaged = bytes;
343 damaged[0] ^= 1;
344 connection
345 .execute("UPDATE assets SET data=?1", [damaged])
346 .unwrap();
347 drop(connection);
348 let cache = Replica::open(&path).unwrap();
349 assert!(
350 matches!(cache.cached_asset(&name, 1024), Err(Error::Io(error)) if error.kind() == io::ErrorKind::InvalidData)
351 );
352 assert!(
353 matches!(cache.fetch_asset(&mut source, "synthetic.one", &name, 1024), Err(Error::Io(error)) if error.kind() == io::ErrorKind::InvalidData)
354 );
355 cache
356 .export_recovery(directory.path().join("damaged.sqlite"))
357 .unwrap();
358 let archive = Recovery::open(directory.path().join("damaged.sqlite")).unwrap();
359 assert!(archive.cached_asset(&name, 1024).is_err());
360 assert_eq!(archive.pending().unwrap().len(), 2);
361}
362
363#[test]
364fn abrupt_process_exit_retains_only_completed_downloads_and_archives() {
365 const CHILD: &str = "ONESTORE_ASSET_EXIT_CASE";
366 if let Ok(phase) = std::env::var(CHILD) {
367 struct ExitDuringRead;
368 impl notebook::discover::Source for ExitDuringRead {
369 fn entries(&mut self, _: &str, _: usize) -> io::Result<Vec<notebook::discover::Entry>> {
370 panic!("Unexpected enumeration")
371 }
372 fn read(&mut self, _: &str, _: usize) -> io::Result<Vec<u8>> {
373 std::process::exit(83)
374 }
375 }
376 let root = std::env::var("ONESTORE_ASSET_EXIT_ROOT").unwrap();
377 let cache = queued_cache(Path::new(&root));
378 let (name, bytes) = payload(1024);
379 if phase == "download" {
380 cache
381 .fetch_asset(&mut ExitDuringRead, "synthetic.one", &name, bytes.len())
382 .unwrap();
383 panic!("Read returned after process exit");
384 }
385 cache
386 .fetch_asset(
387 &mut notebook::discover::Local::open(FIXTURE).unwrap(),
388 "synthetic.one",
389 &name,
390 bytes.len(),
391 )
392 .unwrap();
393 if phase == "archive" {
394 cache
395 .export_recovery(Path::new(&root).join("recovery.sqlite"))
396 .unwrap();
397 }
398 std::process::exit(83);
399 }
400 for phase in ["download", "cached", "archive"] {
401 let root = tempfile::tempdir().unwrap();
402 let output = std::process::Command::new(std::env::current_exe().unwrap())
403 .args([
404 "--exact",
405 "abrupt_process_exit_retains_only_completed_downloads_and_archives",
406 ])
407 .env(CHILD, phase)
408 .env("ONESTORE_ASSET_EXIT_ROOT", root.path())
409 .output()
410 .unwrap();
411 assert_eq!(
412 output.status.code(),
413 Some(83),
414 "{}",
415 String::from_utf8_lossy(&output.stderr)
416 );
417 let cache = Replica::open(root.path().join("cache.sqlite")).unwrap();
418 let (name, bytes) = payload(1024);
419 let expected = (phase != "download").then_some(bytes);
420 assert_eq!(cache.cached_asset(&name, 1024).unwrap(), expected);
421 assert_eq!(cache.pending().unwrap().len(), 2);
422 assert_eq!(cache.status(1).unwrap(), Some(EditStatus::Pending));
423 if phase == "archive" {
424 let recovery = Recovery::open(root.path().join("recovery.sqlite")).unwrap();
425 assert_eq!(recovery.cached_asset(&name, 1024).unwrap(), expected);
426 assert_eq!(recovery.pending().unwrap(), cache.pending().unwrap());
427 assert_eq!(recovery.status(1).unwrap(), cache.status(1).unwrap());
428 }
429 }
430}
431
432#[test]
433#[cfg(feature = "smb")]
434#[ignore = "requires a disposable Samba mirror at ONESTORE_SMB_NOTEBOOK"]
435fn live_smb_downloads_survive_disconnect_and_cache_reopen() {
436 let root = tempfile::tempdir().unwrap();
437 let cache = queued_cache(root.path());
438 let client = notebook::smb::Client::connect(
439 &std::env::var("ONESTORE_SMB_LAB").unwrap(),
440 "agent",
441 notebook::smb::Credentials::default(),
442 Duration::from_secs(5),
443 )
444 .unwrap();
445 let notebook = std::env::var("ONESTORE_SMB_NOTEBOOK").unwrap();
446 let mut source = notebook::discover::Smb::new(&client, &notebook).unwrap();
447 for size in [0, 771, 1024] {
448 let (name, bytes) = payload(size);
449 assert_eq!(
450 cache
451 .fetch_asset(&mut source, "synthetic.one", &name, size)
452 .unwrap(),
453 bytes
454 );
455 }
456 drop(client);
457 cache
458 .export_recovery(root.path().join("recovery.sqlite"))
459 .unwrap();
460 drop(cache);
461 let cache = Replica::open(root.path().join("cache.sqlite")).unwrap();
462 let recovery = Recovery::open(root.path().join("recovery.sqlite")).unwrap();
463 for size in [0, 771, 1024] {
464 let (name, bytes) = payload(size);
465 assert_eq!(
466 cache.cached_asset(&name, size).unwrap(),
467 Some(bytes.clone())
468 );
469 assert_eq!(recovery.cached_asset(&name, size).unwrap(), Some(bytes));
470 }
471 assert_eq!(cache.pending().unwrap(), recovery.pending().unwrap());
472 assert_eq!(cache.status(1).unwrap(), recovery.status(1).unwrap());
473 assert_eq!(cache.recovery_summary().unwrap().cached_assets, 3);
474}