1//! The background worker: waking, backoff, ownership and error propagation.
2
3use notebook::{EditStatus, Error, Remote, Replica};
4use onestore::Stamp;
5use onestore::{
6 CommitError, CommitState, ExGuid, RevisionIndex, Store, Transaction, document::Document,
7};
8use std::io;
9
10#[path = "support/server.rs"]
11mod server;
12use server::*;
13#[path = "support/model_ops.rs"]
14mod model_ops;
15
16use std::{
17 ops::Range,
18 sync::{Arc, Mutex, mpsc},
19 time::{Duration, Instant},
20};
21
22/// How long a wait may take on a loaded machine before it fails: far short of the hourly
23/// poll the waits rule out.
24const PATIENCE: Duration = Duration::from_secs(120);
25
26#[derive(Clone)]
27struct Shared(Arc<Mutex<Server>>);
28
29impl Remote for Shared {
30 fn read(&mut self) -> io::Result<Vec<u8>> {
31 self.0.lock().unwrap().read()
32 }
33 fn stamp(&mut self) -> io::Result<Stamp> {
34 self.0.lock().unwrap().stamp()
35 }
36 fn publish(&mut self, transaction: &Transaction) -> Result<(), CommitError> {
37 self.0.lock().unwrap().publish(transaction)
38 }
39 fn confirm(&mut self, base: &Stamp) -> Result<(), CommitError> {
40 self.0.lock().unwrap().confirm(base)
41 }
42}
43
44/// Saves a replacement of `range` in the page holding `text`.
45fn save(cache: &Replica, text: ExGuid, range: Range<u32>, replacement: &str) -> Option<u64> {
46 model_ops::save(cache, text, |page| {
47 model_ops::replace_text(page, text, range, replacement)
48 })
49 .unwrap()
50}
51
52/// The text of the paragraph holding `object`, wherever the section's pages hold it.
53fn content(bytes: &[u8], object: ExGuid) -> String {
54 let (_, page) = model_ops::locate(bytes, object);
55 model_ops::paragraph_with(&page, object)
56 .unwrap()
57 .text()
58 .unwrap()
59 .text
60 .text()
61 .to_owned()
62}
63
64fn pages(bytes: &[u8]) -> usize {
65 let store = Store::parse(bytes).unwrap();
66 let index = RevisionIndex::parse(&store).unwrap();
67 Document::parse(&index).unwrap().pages().unwrap().len()
68}
69
70#[test]
71fn reconnects_after_connect_read_and_uncertain_publish_without_replaying() {
72 struct Session {
73 shared: Shared,
74 fail_read: bool,
75 }
76 impl Remote for Session {
77 fn read(&mut self) -> io::Result<Vec<u8>> {
78 if self.fail_read {
79 return Err(io::ErrorKind::ConnectionReset.into());
80 }
81 self.shared.read()
82 }
83 fn stamp(&mut self) -> io::Result<Stamp> {
84 if self.fail_read {
85 return Err(io::ErrorKind::ConnectionReset.into());
86 }
87 self.shared.stamp()
88 }
89 fn publish(&mut self, transaction: &Transaction) -> Result<(), CommitError> {
90 self.shared.publish(transaction)
91 }
92 fn confirm(&mut self, base: &Stamp) -> Result<(), CommitError> {
93 self.shared.confirm(base)
94 }
95 }
96 let dir = tempfile::tempdir().unwrap();
97 let path = dir.path().join("cache.sqlite");
98 let source = onestore::create_section("worker.one", "abc", "Fixture").unwrap();
99 let (_, oid, _) = text(&source);
100 let cache = Arc::new(Replica::create(&path, &source).unwrap());
101 let id = save(&cache, oid, 1..2, "🦀").unwrap();
102 let mut server = Server::new(&source);
103 server.fault = Fault::UnknownAfter;
104 let server = Arc::new(Mutex::new(server));
105 let shared = Shared(Arc::clone(&server));
106 let (connected_tx, connected_rx) = mpsc::channel();
107 let (observed_tx, observed_rx) = mpsc::channel();
108 let mut connections = 0;
109 let worker = cache
110 .start_sync(
111 Duration::from_millis(10),
112 move || {
113 connections += 1;
114 connected_tx.send(connections).unwrap();
115 if connections <= 2 {
116 return Err(io::ErrorKind::ConnectionRefused.into());
117 }
118 Ok(Session {
119 shared: shared.clone(),
120 fail_read: connections == 3,
121 })
122 },
123 move |result| {
124 observed_tx
125 .send(
126 result
127 .as_ref()
128 .map(|synced| synced.edit.clone())
129 .map_err(|error| error.to_string()),
130 )
131 .unwrap();
132 },
133 )
134 .unwrap();
135 let mut errors = 0;
136 let published = loop {
137 match observed_rx.recv_timeout(PATIENCE).unwrap() {
138 Err(_) => errors += 1,
139 Ok(Some((actual, status @ EditStatus::Published { .. }))) => {
140 assert_eq!(actual, id);
141 break status;
142 }
143 other => panic!("Unexpected result: {other:?}"),
144 }
145 };
146 worker.stop().unwrap();
147 assert_eq!(errors, 4);
148 assert_eq!(connected_rx.try_iter().collect::<Vec<_>>(), [1, 2, 3, 4, 5]);
149 let server = server.lock().unwrap();
150 assert_eq!(server.publications, 1);
151 assert_eq!(server.confirmations, 1);
152 assert_eq!(text(&server.durable).2, "a🦀c");
153 drop(cache);
154 let cache = Replica::open(&path).unwrap();
155 assert_eq!(cache.status(id).unwrap(), Some(published));
156 let cached = snapshot(&cache);
157 // Confirmation changes reader-notification fields without changing the committed graph.
158 assert_eq!(cached[..212], server.durable[..212]);
159 assert_eq!(cached[252..], server.durable[252..]);
160 let cached_generation = Store::parse(&cached).unwrap().header.generation;
161 let remote_generation = Store::parse(&server.durable).unwrap().header.generation;
162 if cached_generation == remote_generation {
163 assert_eq!(cached, server.durable);
164 } else {
165 assert_eq!(cached_generation.checked_add(1), Some(remote_generation));
166 }
167}
168
169#[test]
170fn local_saves_wake_an_idle_worker_and_publish_every_writer_marker() {
171 let dir = tempfile::tempdir().unwrap();
172 let source = onestore::create_section("worker.one", "abc", "Fixture").unwrap();
173 let (sid, oid, _) = text(&source);
174 let cache = Arc::new(Replica::create(dir.path().join("cache.sqlite"), &source).unwrap());
175 let server = Arc::new(Mutex::new(Server::new(&source)));
176 let shared = Shared(Arc::clone(&server));
177 let (observed_tx, observed_rx) = mpsc::channel();
178 let (resume_tx, resume_rx) = mpsc::channel();
179 let mut first = true;
180 let worker = cache
181 .start_sync(
182 Duration::from_secs(3600),
183 move || Ok(shared.clone()),
184 move |result| {
185 observed_tx
186 .send(result.as_ref().unwrap().edit.clone())
187 .unwrap();
188 if first {
189 first = false;
190 resume_rx.recv_timeout(PATIENCE).unwrap();
191 }
192 },
193 )
194 .unwrap();
195 assert_eq!(observed_rx.recv_timeout(PATIENCE).unwrap(), None);
196 let started = Instant::now();
197 let queued = std::thread::scope(|scope| {
198 (0..12)
199 .map(|writer| {
200 let cache = &cache;
201 // Typed as an op: a page read before another writer's edit would lower stale.
202 scope.spawn(move || {
203 let op = onestore::op::PageOp::Text {
204 text: oid,
205 range: 0..0,
206 with: format!("[{writer}] "),
207 };
208 let ops = vec![onestore::op::Op::Page { space: sid, op }];
209 let edit = onestore::op::Edit {
210 at: model_ops::now(),
211 ops,
212 };
213 cache.apply(model_ops::AUTHOR, edit).unwrap()
214 })
215 })
216 .collect::<Vec<_>>()
217 .into_iter()
218 .map(|join| join.join().unwrap())
219 .collect::<std::collections::BTreeSet<_>>()
220 });
221 resume_tx.send(()).unwrap();
222 while !cache.pending().unwrap().is_empty() {
223 if let Some((_, status)) = observed_rx.recv_timeout(PATIENCE).unwrap() {
224 assert!(matches!(status, EditStatus::Published { .. }), "{status:?}");
225 }
226 }
227 assert!(
228 started.elapsed() < PATIENCE,
229 "Edits waited for the hourly poll"
230 );
231 worker.stop().unwrap();
232 for id in &queued {
233 assert!(matches!(
234 cache.status(*id).unwrap(),
235 Some(EditStatus::Published { .. })
236 ));
237 }
238 let content = text(&snapshot(&cache)).2;
239 for writer in 0..12 {
240 assert_eq!(content.matches(&format!("[{writer}] ")).count(), 1);
241 }
242 assert!(content.ends_with("abc"));
243 let server = server.lock().unwrap();
244 assert_eq!(text(&server.durable).2, content);
245 // Edits queued while the worker was busy publish together.
246 assert_eq!(server.publications, 1);
247 assert_eq!(snapshot(&cache), server.durable);
248}
249
250#[test]
251fn dropping_during_publication_is_nonblocking_and_retains_ownership_until_recovery_is_recorded() {
252 struct Paused {
253 shared: Shared,
254 entered: mpsc::Sender<()>,
255 resume: mpsc::Receiver<()>,
256 }
257 impl Remote for Paused {
258 fn read(&mut self) -> io::Result<Vec<u8>> {
259 self.shared.read()
260 }
261 fn stamp(&mut self) -> io::Result<Stamp> {
262 self.shared.stamp()
263 }
264 fn publish(&mut self, transaction: &Transaction) -> Result<(), CommitError> {
265 self.entered.send(()).unwrap();
266 self.resume.recv_timeout(PATIENCE).unwrap();
267 self.shared.publish(transaction)
268 }
269 fn confirm(&mut self, base: &Stamp) -> Result<(), CommitError> {
270 self.shared.confirm(base)
271 }
272 }
273 let dir = tempfile::tempdir().unwrap();
274 let path = dir.path().join("cache.sqlite");
275 let source = onestore::create_section("worker.one", "abc", "Fixture").unwrap();
276 let (_, oid, _) = text(&source);
277 let cache = Arc::new(Replica::create(&path, &source).unwrap());
278 let id = save(&cache, oid, 0..0, "L ").unwrap();
279 let mut server = Server::new(&source);
280 server.fault = Fault::UnknownAfter;
281 let server = Arc::new(Mutex::new(server));
282 let shared = Shared(Arc::clone(&server));
283 let (entered_tx, entered_rx) = mpsc::channel();
284 let (resume_tx, resume_rx) = mpsc::channel();
285 let mut session = Some(Paused {
286 shared,
287 entered: entered_tx,
288 resume: resume_rx,
289 });
290 let worker = cache
291 .start_sync(
292 Duration::from_secs(3600),
293 move || Ok(session.take().unwrap()),
294 |_| {},
295 )
296 .unwrap();
297 entered_rx.recv_timeout(PATIENCE).unwrap();
298 let stopped = Instant::now();
299 drop(worker);
300 assert!(stopped.elapsed() < PATIENCE);
301 let shared = Shared(Arc::clone(&server));
302 let start = cache.start_sync(
303 Duration::from_secs(3600),
304 move || Ok(shared.clone()),
305 |_| {},
306 );
307 assert!(matches!(start, Err(error) if error.kind() == io::ErrorKind::WouldBlock));
308 drop(cache);
309 assert!(matches!(Replica::open(&path), Err(Error::Database(error))
310 if error.sqlite_error_code() == Some(rusqlite::ErrorCode::DatabaseBusy)));
311 resume_tx.send(()).unwrap();
312 let cache = Arc::new(loop {
313 match Replica::open(&path) {
314 Ok(cache) => break cache,
315 Err(Error::Database(error))
316 if error.sqlite_error_code() == Some(rusqlite::ErrorCode::DatabaseBusy) => {}
317 Err(error) => panic!("{error}"),
318 }
319 assert!(stopped.elapsed() < PATIENCE);
320 std::thread::sleep(Duration::from_millis(1));
321 });
322 assert!(matches!(
323 cache.status(id).unwrap(),
324 Some(EditStatus::AwaitingConfirmation { .. })
325 ));
326 assert_eq!(server.lock().unwrap().publications, 1);
327 let shared = Shared(Arc::clone(&server));
328 let (tx, rx) = mpsc::channel();
329 let worker = cache
330 .start_sync(
331 Duration::from_secs(3600),
332 move || Ok(shared.clone()),
333 move |result| {
334 tx.send(result.as_ref().unwrap().edit.clone()).unwrap();
335 },
336 )
337 .unwrap();
338 assert!(
339 matches!(rx.recv_timeout(PATIENCE).unwrap(), Some((actual, EditStatus::Published { .. })) if actual == id)
340 );
341 worker.stop().unwrap();
342 let server = server.lock().unwrap();
343 assert_eq!(server.publications, 1);
344 assert_eq!(server.confirmations, 1);
345 assert_eq!(text(&server.durable).2, "L abc");
346}
347
348#[test]
349fn cache_failures_stop_retries_and_return_the_error_without_remote_publication() {
350 let dir = tempfile::tempdir().unwrap();
351 let path = dir.path().join("cache.sqlite");
352 let source = onestore::create_section("worker.one", "abc", "Fixture").unwrap();
353 let (_, oid, _) = text(&source);
354 let cache = Replica::create(&path, &source).unwrap();
355 let id = save(&cache, oid, 0..0, "L ").unwrap();
356 drop(cache);
357 let connection = rusqlite::Connection::open(&path).unwrap();
358 connection.execute_batch("CREATE TRIGGER fail_attempt BEFORE UPDATE OF attempted ON batches BEGIN SELECT RAISE(ABORT, 'test cache write failure'); END;").unwrap();
359 drop(connection);
360 let cache = Arc::new(Replica::open(&path).unwrap());
361 let server = Arc::new(Mutex::new(Server::new(&source)));
362 let shared = Shared(Arc::clone(&server));
363 let (tx, rx) = mpsc::channel();
364 let worker = cache
365 .start_sync(
366 Duration::from_millis(1),
367 move || Ok(shared.clone()),
368 move |result| {
369 tx.send(matches!(result, Err(Error::Database(_)))).unwrap();
370 },
371 )
372 .unwrap();
373 assert!(rx.recv_timeout(PATIENCE).unwrap());
374 assert!(matches!(worker.stop(), Err(Error::Database(_))));
375 assert!(rx.try_iter().next().is_none());
376 assert_eq!(server.lock().unwrap().publications, 0);
377 assert_eq!(cache.status(id).unwrap(), Some(EditStatus::Pending));
378 assert_eq!(text(&snapshot(&cache)).2, "L abc");
379}
380
381#[test]
382fn polling_reports_an_uncertain_attempt_once_per_remote_change_without_replay() {
383 let dir = tempfile::tempdir().unwrap();
384 let source = onestore::create_section("worker.one", "abc", "Fixture").unwrap();
385 let (sid, oid, _) = text(&source);
386 let cache = Arc::new(Replica::create(dir.path().join("cache.sqlite"), &source).unwrap());
387 let id = save(&cache, oid, 1..2, "L").unwrap();
388 let local = text(&snapshot(&cache)).2;
389 let mut server = Server::new(&source);
390 server.fault = Fault::UnknownBefore;
391 let server = Arc::new(Mutex::new(server));
392 let shared = Shared(Arc::clone(&server));
393 let (tx, rx) = mpsc::channel();
394 let worker = cache
395 .start_sync(
396 Duration::from_millis(10),
397 move || Ok(shared.clone()),
398 move |result| {
399 tx.send(
400 result
401 .as_ref()
402 .map(|synced| synced.edit.clone())
403 .map_err(|_| ()),
404 )
405 .unwrap();
406 },
407 )
408 .unwrap();
409 assert!(rx.recv_timeout(PATIENCE).unwrap().is_err());
410 for round in 0..3 {
411 let (actual, status) = rx.recv_timeout(PATIENCE).unwrap().unwrap().unwrap();
412 assert_eq!(actual, id);
413 assert!(
414 matches!(status, EditStatus::AwaitingConfirmation { .. }),
415 "{status:?}"
416 );
417 // An unchanged remote is not read or reported again.
418 assert!(rx.recv_timeout(Duration::from_millis(100)).is_err());
419 // Another writer's change elsewhere in the text is read and decided again.
420 let mut server = server.lock().unwrap();
421 let changed = typed(&server.visible, sid, oid, 0..0, &round.to_string());
422 server.visible = changed.clone();
423 server.durable = changed;
424 drop(server);
425 worker.wake();
426 }
427 worker.stop().unwrap();
428 assert_eq!(text(&snapshot(&cache)).2, local);
429 assert_eq!(cache.pending().unwrap().len(), 1);
430 let server = server.lock().unwrap();
431 assert_eq!(server.publications, 1);
432 assert_eq!(server.confirmations, 0);
433}
434
435#[test]
436fn reachability_notification_retries_without_waiting_for_the_poll() {
437 let dir = tempfile::tempdir().unwrap();
438 let source = onestore::create_section("worker.one", "abc", "Fixture").unwrap();
439 let cache = Arc::new(Replica::create(dir.path().join("cache.sqlite"), &source).unwrap());
440 let (tx, rx) = mpsc::channel();
441 let worker = cache
442 .start_sync(
443 Duration::from_secs(3600),
444 || -> io::Result<Shared> { Err(io::ErrorKind::NotConnected.into()) },
445 move |result| {
446 tx.send(matches!(result, Err(Error::RemoteIo(_)))).unwrap();
447 },
448 )
449 .unwrap();
450 assert!(rx.recv_timeout(PATIENCE).unwrap());
451 worker.wake();
452 assert!(rx.recv_timeout(PATIENCE).unwrap());
453 worker.stop().unwrap();
454 assert!(rx.try_iter().next().is_none());
455}
456
457#[test]
458fn cancellation_during_connect_does_not_read_or_report_a_false_refresh() {
459 let dir = tempfile::tempdir().unwrap();
460 let source = onestore::create_section("worker.one", "abc", "Fixture").unwrap();
461 let cache = Arc::new(Replica::create(dir.path().join("cache.sqlite"), &source).unwrap());
462 let (entered_tx, entered_rx) = mpsc::channel();
463 let (resume_tx, resume_rx) = mpsc::channel();
464 let (tx, rx) = mpsc::channel();
465 // An invalid image makes any unexpected read observable as an error callback.
466 let shared = Shared(Arc::new(Mutex::new(Server::new(&[]))));
467 let worker = cache
468 .start_sync(
469 Duration::from_secs(3600),
470 move || {
471 entered_tx.send(()).unwrap();
472 resume_rx.recv_timeout(PATIENCE).unwrap();
473 Ok(shared.clone())
474 },
475 move |_| {
476 tx.send(()).unwrap();
477 },
478 )
479 .unwrap();
480 entered_rx.recv_timeout(PATIENCE).unwrap();
481 drop(worker);
482 let weak = Arc::downgrade(&cache);
483 drop(cache);
484 resume_tx.send(()).unwrap();
485 let started = Instant::now();
486 while weak.upgrade().is_some() {
487 assert!(started.elapsed() < PATIENCE);
488 std::thread::sleep(Duration::from_millis(1));
489 }
490 assert_eq!(
491 rx.recv_timeout(PATIENCE),
492 Err(mpsc::RecvTimeoutError::Disconnected)
493 );
494}
495
496#[test]
497fn invalid_intervals_and_callback_panics_leave_worker_ownership_recoverable() {
498 let dir = tempfile::tempdir().unwrap();
499 let source = onestore::create_section("worker.one", "abc", "Fixture").unwrap();
500 let cache = Arc::new(Replica::create(dir.path().join("cache.sqlite"), &source).unwrap());
501 for interval in [Duration::ZERO, Duration::MAX] {
502 assert!(
503 matches!(cache.start_sync(interval, || -> io::Result<Shared> { panic!("Unexpected connection") }, |_| {}), Err(error) if error.kind() == io::ErrorKind::InvalidInput)
504 );
505 }
506 let shared = Shared(Arc::new(Mutex::new(Server::new(&source))));
507 let first = shared.clone();
508 let (tx, rx) = mpsc::channel();
509 let worker = cache
510 .start_sync(
511 Duration::from_secs(3600),
512 move || Ok(first.clone()),
513 move |_| {
514 tx.send(()).unwrap();
515 panic!("Test observer panic");
516 },
517 )
518 .unwrap();
519 rx.recv_timeout(PATIENCE).unwrap();
520 assert!(matches!(worker.stop(), Err(Error::Io(_))));
521 let (tx, rx) = mpsc::channel();
522 let worker = cache
523 .start_sync(
524 Duration::from_secs(3600),
525 move || Ok(shared.clone()),
526 move |result| {
527 tx.send(result.as_ref().unwrap().edit.clone()).unwrap();
528 },
529 )
530 .unwrap();
531 assert_eq!(rx.recv_timeout(PATIENCE).unwrap(), None);
532 worker.stop().unwrap();
533 assert_eq!(snapshot(&cache), source);
534}
535
536/// A conflict never blocks the worker: one step publishes the remote's version with the
537/// local one as its conflict page.
538#[test]
539fn a_conflict_publishes_its_conflict_page_in_one_worker_step() {
540 let dir = tempfile::tempdir().unwrap();
541 let source = onestore::create_section("resolve.one", "abc", "Fixture").unwrap();
542 let (sid, oid, _) = text(&source);
543 let cache = Arc::new(Replica::create(dir.path().join("cache.sqlite"), &source).unwrap());
544 let id = save(&cache, oid, 1..2, "L").unwrap();
545 let remote = typed(&source, sid, oid, 1..2, "R");
546 let server = Arc::new(Mutex::new(Server::new(&remote)));
547 let shared = Shared(Arc::clone(&server));
548 let (tx, rx) = mpsc::channel();
549 let worker = cache
550 .start_sync(
551 Duration::from_secs(3600),
552 move || Ok(shared.clone()),
553 move |result| {
554 tx.send(result.as_ref().unwrap().edit.clone()).unwrap();
555 },
556 )
557 .unwrap();
558 assert!(
559 matches!(rx.recv_timeout(PATIENCE).unwrap(), Some((actual, EditStatus::Published { .. })) if actual > id)
560 );
561 worker.stop().unwrap();
562 assert!(cache.pending().unwrap().is_empty());
563 let server = server.lock().unwrap();
564 assert_eq!(server.publications, 1);
565 assert_eq!(
566 conflicts(&server.durable),
567 [(
568 sid,
569 vec![(model_ops::AUTHOR.to_owned(), vec!["aLc".to_owned()])]
570 )]
571 );
572 assert_eq!(snapshot(&cache), server.durable);
573}
574
575#[test]
576fn ordinary_read_and_unpublished_write_contention_reuse_the_connection() {
577 struct Busy {
578 server: Server,
579 reads: usize,
580 writes: usize,
581 }
582 impl Remote for Busy {
583 fn read(&mut self) -> io::Result<Vec<u8>> {
584 self.server.read()
585 }
586 fn stamp(&mut self) -> io::Result<Stamp> {
587 self.reads += 1;
588 match self.reads {
589 1 => Err(io::ErrorKind::WouldBlock.into()),
590 2 => Err(io::ErrorKind::ResourceBusy.into()),
591 _ => self.server.stamp(),
592 }
593 }
594 fn publish(&mut self, transaction: &Transaction) -> Result<(), CommitError> {
595 self.writes += 1;
596 match self.writes {
597 1 | 2 => Err(CommitError {
598 state: CommitState::NotCommitted,
599 error: if self.writes == 1 {
600 io::ErrorKind::WouldBlock
601 } else {
602 io::ErrorKind::ResourceBusy
603 }
604 .into(),
605 }),
606 _ => self.server.publish(transaction),
607 }
608 }
609 fn confirm(&mut self, base: &Stamp) -> Result<(), CommitError> {
610 self.server.confirm(base)
611 }
612 }
613 let dir = tempfile::tempdir().unwrap();
614 let source = onestore::create_section("worker.one", "abc", "Fixture").unwrap();
615 let (_, oid, _) = text(&source);
616 let cache = Arc::new(Replica::create(dir.path().join("cache.sqlite"), &source).unwrap());
617 let id = save(&cache, oid, 0..0, "L ").unwrap();
618 let mut remote = Some(Busy {
619 server: Server::new(&source),
620 reads: 0,
621 writes: 0,
622 });
623 let (tx, rx) = mpsc::channel();
624 let worker = cache
625 .start_sync(
626 Duration::from_millis(1),
627 move || Ok(remote.take().expect("Contention caused a reconnect")),
628 move |result| {
629 tx.send(
630 result
631 .as_ref()
632 .map(|synced| synced.edit.clone())
633 .map_err(|error| error.to_string()),
634 )
635 .unwrap();
636 },
637 )
638 .unwrap();
639 for _ in 0..4 {
640 assert!(rx.recv_timeout(PATIENCE).unwrap().is_err());
641 }
642 assert!(
643 matches!(rx.recv_timeout(PATIENCE).unwrap().unwrap(), Some((actual, EditStatus::Published { .. })) if actual == id)
644 );
645 worker.stop().unwrap();
646 assert_eq!(text(&snapshot(&cache)).2, "L abc");
647 assert!(cache.pending().unwrap().is_empty());
648}
649
650#[test]
651fn publication_backoff_drains_local_wakes_without_waiting_for_the_idle_poll() {
652 struct BusyOnce(Server);
653 impl Remote for BusyOnce {
654 fn read(&mut self) -> io::Result<Vec<u8>> {
655 self.0.read()
656 }
657 fn stamp(&mut self) -> io::Result<Stamp> {
658 self.0.stamp()
659 }
660 fn publish(&mut self, transaction: &Transaction) -> Result<(), CommitError> {
661 let server = &mut self.0;
662 if server.publications == 0 {
663 server.publications += 1;
664 return Err(CommitError {
665 state: CommitState::NotCommitted,
666 error: io::ErrorKind::ResourceBusy.into(),
667 });
668 }
669 server.publish(transaction)
670 }
671 fn confirm(&mut self, base: &Stamp) -> Result<(), CommitError> {
672 self.0.confirm(base)
673 }
674 }
675 let dir = tempfile::tempdir().unwrap();
676 let source = onestore::create_section("worker.one", "abc", "Fixture").unwrap();
677 let (_, oid, _) = text(&source);
678 let cache = Arc::new(Replica::create(dir.path().join("cache.sqlite"), &source).unwrap());
679 let first = save(&cache, oid, 0..0, "L ").unwrap();
680 let mut remote = Some(BusyOnce(Server::new(&source)));
681 let observed = Arc::clone(&cache);
682 let (tx, rx) = mpsc::channel();
683 let mut failed = false;
684 let worker = cache
685 .start_sync(
686 Duration::from_secs(3600),
687 move || Ok(remote.take().expect("Contention must retain the session")),
688 move |result| match result {
689 Err(Error::Remote(error)) if error.state == CommitState::NotCommitted => {
690 assert!(!failed);
691 failed = true;
692 let page =
693 onestore::PageCreation::new(None, Some("Queued"), "Fixture").unwrap();
694 let second = section_op(&observed, onestore::op::SectionOp::Create(page));
695 tx.send((second, None)).unwrap();
696 }
697 Ok(notebook::Synced {
698 edit: Some((id, status)),
699 ..
700 }) => tx.send((*id, Some(status.clone()))).unwrap(),
701 Ok(_) => {}
702 other => panic!("Unexpected worker result: {other:?}"),
703 },
704 )
705 .unwrap();
706 let (second, status) = rx.recv_timeout(PATIENCE).unwrap();
707 assert_eq!(status, None);
708 for expected in [first, second] {
709 let (id, status) = rx.recv_timeout(PATIENCE).unwrap();
710 assert_eq!(id, expected);
711 assert!(matches!(status, Some(EditStatus::Published { .. })));
712 }
713 worker.stop().unwrap();
714 assert!(cache.pending().unwrap().is_empty());
715 assert_eq!(content(&snapshot(&cache), oid), "L abc");
716 assert_eq!(pages(&snapshot(&cache)), 2);
717}