1//! The sections of an open notebook that no session holds, kept in sync as OneNote 2010
2//! keeps every section of an open notebook: queued edits publish and other clients' changes
3//! are noticed without the section being open.
4
5use crate::fs;
6use crate::{
7 EditStatus, Error, Remote, Replica, Result,
8 discover::{Entry, Listed},
9 session::{Section, SyncStatus, reached},
10 worker::Signal,
11};
12use onestore::Stamp;
13use std::{
14 collections::{BTreeMap, BTreeSet},
15 io,
16 path::{Path, PathBuf},
17 sync::{
18 Arc, Mutex, Weak,
19 atomic::{AtomicBool, Ordering},
20 },
21 time::Duration,
22};
23use web_time::Instant;
24
25/// How long a section that could not be reached waits to be tried again; OneNote 2010
26/// retries a failing section about every 31 seconds.
27const RETRY: Duration = Duration::from_secs(31);
28/// Between one section's first check and the next's, so opening a notebook reads its files
29/// one after another rather than in a burst.
30const STAGGER: Duration = Duration::from_millis(100);
31/// How long a reported change settles before its section is checked, so the writes of one
32/// commit cost one check.
33const SETTLE: Duration = Duration::from_secs(1);
34
35/// Keeps each watched section in sync: a section is checked when a watch on the notebook's
36/// folder reports it changed, when it could not be reached, and otherwise every interval.
37/// On connecting, and when a watch reports a folder rather than a file, the folders are
38/// listed and only the sections listed otherwise than before are checked, as OneNote 2010
39/// reopens a notebook. A check costs one stamp read; a section whose replica has edits waiting,
40/// or whose file moved past the replica's base, has its replica opened for the synchronization
41/// steps that publish or rebase it, then closed again. A section without a replica gets one,
42/// its offline copy, where the notebook keeps them. A section a session holds is left to that
43/// session's worker, which the watch wakes. Dropping requests cancellation without waiting for
44/// the step in flight.
45pub struct Background(Arc<Shared>, Mutex<Option<crate::task::JoinHandle<()>>>);
46
47/// A section for `Background::watch`, as its notebook knows it.
48pub struct Known {
49 /// The catalog path.
50 pub path: String,
51 pub replica: Option<PathBuf>,
52 /// How the notebook's discovery last found the file: its listing, and its stamp then.
53 pub found: Option<(Listed, Stamp)>,
54 /// The file as that discovery read it, if it did, which the first check takes in place
55 /// of reading the file again while its stamp still matches.
56 pub image: Option<Vec<u8>>,
57}
58
59/// What the background thread and whatever reports changes to it share.
60pub(crate) struct Shared {
61 signal: Arc<Signal>,
62 watched: Mutex<Watched>,
63 /// Asked for by `discard`, once the thread stops.
64 discard: AtomicBool,
65 /// Hears every report of changed paths (`on_touched`).
66 listener: Mutex<Option<Listener>>,
67}
68
69/// Hears the catalog paths a report names.
70pub type Listener = Box<dyn Fn(&[String]) + Send + Sync>;
71
72/// A connection's report of its notebook folder's changes, from a watch armed on connecting.
73#[cfg_attr(not(any(feature = "smb", feature = "live")), allow(dead_code))]
74pub(crate) struct Reports {
75 shared: Weak<Shared>,
76 connection: u64,
77}
78
79#[derive(Default)]
80struct Watched {
81 sections: Vec<Watch>,
82 /// Sections or cached folders changed since `changed` was last asked.
83 changed: Vec<String>,
84 /// Folders a watch reported without naming the file, listed at `relist`.
85 folders: BTreeSet<String>,
86 relist: Option<Instant>,
87 /// Counts connections, so a watch that ended with an earlier one is not this one's.
88 connection: u64,
89 /// The current connection's watch ended, or none was connected: changes since went
90 /// unreported.
91 lost: bool,
92 /// The current connection's watch reports the folder's changes.
93 reported: bool,
94 /// How long a report settles, where not `SETTLE` (`Background::set_settle`).
95 settle: Option<Duration>,
96}
97
98struct Watch {
99 path: String,
100 replica: Option<PathBuf>,
101 /// The file's stamp when last read.
102 stamp: Option<Stamp>,
103 /// How the file was listed when its folder was last listed; the stamp was read since.
104 listed: Option<Listed>,
105 /// The folder's listing since `stamp` shows the file unchanged, so the next check need not
106 /// read it.
107 current: bool,
108 /// A watch reported the file changed since its folder's listing began.
109 reported: bool,
110 status: SyncStatus,
111 /// The file as discovery read it, until the next check takes it.
112 image: Option<Vec<u8>>,
113 /// When the section is next checked.
114 due: Instant,
115 /// The worker of the session that holds the section.
116 held: Option<Weak<Signal>>,
117}
118
119impl Background {
120 /// How often each section is checked while a watch reports its folder's changes, against
121 /// a report that went missing; OneNote 2010, watching, checks nothing on a timer.
122 pub const BACKSTOP: Duration = Duration::from_secs(60 * 60);
123 /// How often each section is checked where nothing reports its folder's changes.
124 pub const UNWATCHED: Duration = Duration::from_secs(15);
125
126 /// Starts the thread. `connect` answers how the connection reaches each section file and
127 /// lists each folder, both by catalog path, arming any watch that reports the folder's
128 /// changes through `Reports`, and whether one does; it runs again after a transport
129 /// failure. `copies` keeps an offline copy of every section.
130 pub(crate) fn start<R, B, L, LF>(
131 copies: bool,
132 mut connect: impl FnMut(Reports) -> io::Result<((B, L), bool)> + Send + 'static,
133 notify: impl Fn() + Send + 'static,
134 ) -> Result<Self>
135 where
136 R: Remote,
137 B: FnMut(&str) -> R,
138 L: FnMut(String) -> LF,
139 LF: Future<Output = io::Result<Vec<Entry>>>,
140 {
141 let (signal, receiver) = Signal::new();
142 let shared = Arc::new(Shared {
143 signal,
144 watched: Mutex::default(),
145 discard: AtomicBool::new(false),
146 listener: Mutex::default(),
147 });
148 let weak = Arc::downgrade(&shared);
149 let owner = Arc::clone(&shared);
150 let thread = crate::task::spawn("onestore-background", move || async move {
151 let signal = &owner.signal;
152 let mut bound: Option<(B, L)> = None;
153 let mut news = false;
154 while !signal.stopped.load(Ordering::Acquire) {
155 if signal.offline.load(Ordering::Acquire)
156 && !signal.requested.load(Ordering::Acquire)
157 {
158 // Working online again lists every folder again.
159 if let Ok(mut watched) = owner.watched.lock() {
160 watched.watching(false);
161 watched.lost = true;
162 }
163 bound = None;
164 crate::task::wait(&receiver, None).await;
165 continue;
166 }
167 let now = Instant::now();
168 let next = {
169 let Ok(mut watched) = owner.watched.lock() else {
170 return;
171 };
172 if std::mem::take(&mut watched.lost) {
173 bound = None;
174 watched.watching(false);
175 for watch in &mut watched.sections {
176 watch.due = now;
177 }
178 }
179 watched.next(now, bound.is_some())
180 };
181 let (path, replica, seen, current, image) = match next {
182 Next::Wait(wait) => {
183 signal.requested.store(false, Ordering::Release);
184 if std::mem::take(&mut news) {
185 notify();
186 }
187 crate::task::wait(&receiver, wait).await;
188 continue;
189 }
190 Next::Connect(connection) => {
191 let reports = Reports {
192 shared: weak.clone(),
193 connection: connection + 1,
194 };
195 let connected = connect(reports);
196 let folders = {
197 let Ok(mut watched) = owner.watched.lock() else {
198 return;
199 };
200 watched.connection += 1;
201 watched.lost = false;
202 match &connected {
203 Ok((_, reported)) => {
204 watched.watching(*reported);
205 Some(watched.rescanning())
206 }
207 Err(error) => {
208 let error = Error::RemoteIo(io::Error::new(
209 error.kind(),
210 error.to_string(),
211 ));
212 let retry = now + RETRY;
213 for watch in &mut watched.sections {
214 news |= watch.fail(&error, None);
215 watch.due = watch.due.max(retry);
216 }
217 watched.relist = watched.relist.map(|due| due.max(retry));
218 None
219 }
220 }
221 };
222 if let (Ok((mut files, _)), Some(folders)) = (connected, folders) {
223 let listings = list(&mut files.1, folders).await;
224 let Ok(mut watched) = owner.watched.lock() else {
225 return;
226 };
227 watched.listed(&listings, Some(copies), Instant::now());
228 bound = Some(files);
229 }
230 continue;
231 }
232 Next::List(folders) => {
233 let Some((_, list_folder)) = &mut bound else {
234 continue;
235 };
236 let listings = list(list_folder, folders).await;
237 let Ok(mut watched) = owner.watched.lock() else {
238 return;
239 };
240 watched.listed(&listings, None, Instant::now());
241 continue;
242 }
243 Next::Check {
244 path,
245 replica,
246 seen,
247 current,
248 image,
249 } => (path, replica, seen, current, image),
250 };
251 let Some((bind, _)) = &mut bound else {
252 continue;
253 };
254 let (queued, outcome) = step_async(
255 &mut bind(&path),
256 replica.as_deref(),
257 seen.as_deref(),
258 current,
259 image,
260 copies,
261 )
262 .await;
263 if outcome.as_ref().is_err_and(disconnected) {
264 bound = None;
265 }
266 let Ok(mut watched) = owner.watched.lock() else {
267 return;
268 };
269 let interval = watched.interval();
270 let Watched {
271 sections, changed, ..
272 } = &mut *watched;
273 let Some(watch) = sections.iter_mut().find(|watch| watch.path == path) else {
274 continue;
275 };
276 let now = Instant::now();
277 watch.current = false;
278 news |= match outcome {
279 Ok((stamp, moved)) => {
280 watch.stamp = Some(stamp);
281 watch.due = now + interval;
282 if moved {
283 changed.push(path);
284 }
285 let before = summary(&watch.status);
286 if let Some(queued) = queued {
287 watch.status = SyncStatus {
288 synced: Some(crate::now()),
289 error: None,
290 queued,
291 };
292 }
293 moved || before != summary(&watch.status)
294 }
295 Err(error) => {
296 watch.due = now + RETRY;
297 watch.fail(&error, queued)
298 }
299 };
300 }
301 if let Ok(mut watched) = owner.watched.lock() {
302 watched.watching(false);
303 if owner.discard.load(Ordering::Acquire) {
304 for replica in watched
305 .sections
306 .iter()
307 .filter_map(|watch| watch.replica.as_ref())
308 {
309 discard(replica);
310 }
311 }
312 }
313 })?;
314 Ok(Self(shared, Mutex::new(Some(thread))))
315 }
316
317 /// Keeps the sections of a notebook on a share in sync while they are not open, with an
318 /// offline copy of each, and one watch on the notebook's folder reporting what changed,
319 /// as OneNote 2010 watches it; `connect` runs again after a transport failure. Paths are
320 /// relative to `root`.
321 #[cfg(feature = "smb")]
322 pub fn smb(
323 root: &str,
324 limit: usize,
325 mut connect: impl FnMut() -> io::Result<crate::smb::Client> + Send + 'static,
326 notify: impl Fn() + Send + 'static,
327 ) -> Result<Self> {
328 let root = root.replace('\\', "/");
329 Self::start(
330 true,
331 move |reports| {
332 let client = Arc::new(connect()?);
333 let reported = match client.watch(&root, move |changed| match changed {
334 Ok(paths) => reports.touched(&paths),
335 Err(_) => reports.lost(),
336 }) {
337 Ok(()) => true,
338 Err(error) if error.kind() == io::ErrorKind::Unsupported => false,
339 Err(error) => return Err(error),
340 };
341 let (bound, folder) = (Arc::clone(&client), root.clone());
342 let bind = move |path: &str| {
343 let file = match folder.as_str() {
344 "" => path.to_owned(),
345 root => format!("{root}/{path}"),
346 };
347 crate::SmbRemote::new(Arc::clone(&bound), file, limit)
348 };
349 let root = root.clone();
350 let list = move |folder: String| {
351 std::future::ready((|| {
352 use crate::discover::Source;
353 crate::discover::Smb::new(&client, &root)?
354 .entries(&folder, crate::session::LIMITS.entries)
355 })())
356 };
357 Ok(((bind, list), reported))
358 },
359 notify,
360 )
361 }
362
363 /// Keeps the sections of a notebook a Live Share host serves in sync while they are not
364 /// open, with an offline copy of each; the host reports what changed. Paths are catalog
365 /// paths.
366 #[cfg(feature = "live")]
367 pub fn hosted(
368 guest: Arc<crate::live::share::Guest>,
369 notify: impl Fn() + Send + 'static,
370 ) -> Result<Self> {
371 let background = Self::start(
372 true,
373 move |reports| {
374 guest.watch(reports)?;
375 let (bound, listing) = (Arc::clone(&guest), Arc::clone(&guest));
376 let bind = move |path: &str| crate::live::share::HostedRemote::new(&bound, path);
377 let list = move |folder: String| {
378 let listing = listing.clone();
379 async move { listing.entries_async(&folder).await }
380 };
381 Ok(((bind, list), true))
382 },
383 notify,
384 )?;
385 background.set_settle(crate::live::share::SETTLE);
386 Ok(background)
387 }
388
389 /// Checks what a report names after `settle` rather than a second, where reports come
390 /// one to a commit, as a Live Share host sends them.
391 pub fn set_settle(&self, settle: Duration) {
392 if let Ok(mut watched) = self.0.watched.lock() {
393 watched.settle = Some(settle);
394 }
395 }
396
397 /// Has `listener` hear every report of changed paths from now on, as a Live Share host
398 /// passes them to its guests; `None` stops.
399 pub fn on_touched(&self, listener: Option<Listener>) {
400 if let Ok(mut held) = self.0.listener.lock() {
401 *held = listener;
402 }
403 }
404
405 /// Watches these sections (`Notebook::replicas`); a section not watched before is first
406 /// checked soon after, the next one a little later.
407 pub fn watch(&self, sections: Vec<Known>) {
408 if let Ok(mut watched) = self.0.watched.lock() {
409 watched.watch(sections, Instant::now());
410 }
411 self.0.signal.wake();
412 }
413
414 /// Leaves the section at catalog `path` to `section`'s worker, which from now on the watch
415 /// on the notebook's folder wakes when the file changes, instead of the worker checking
416 /// it every few seconds. Once the section closes, this takes it over again.
417 pub fn hold(&self, path: &str, section: &Section) {
418 let replica = section.replica();
419 let Some(worker) = replica
420 .section
421 .worker
422 .lock()
423 .ok()
424 .and_then(|worker| worker.upgrade())
425 else {
426 return;
427 };
428 if let Ok(mut watched) = self.0.watched.lock() {
429 let reported = watched.reported;
430 if let Some(watch) = watched.sections.iter_mut().find(|watch| watch.path == path) {
431 worker.watched.store(reported, Ordering::Release);
432 watch.held = Some(Arc::downgrade(&worker));
433 }
434 }
435 let (shared, path) = (Arc::downgrade(&self.0), path.to_owned());
436 if let Ok(mut released) = replica.released.0.lock() {
437 *released = Some(Box::new(move || {
438 if let Some(shared) = shared.upgrade() {
439 shared.released(&path);
440 }
441 }));
442 }
443 }
444
445 /// Checks the sections at or below these catalog paths soon, as a watch on the notebook's
446 /// folder reports them changed: a section's own path checks it, a folder's lists it first
447 /// and checks the sections listed otherwise. `""` names the notebook's folder.
448 pub fn touched(&self, paths: &[String]) {
449 self.0.touched(paths);
450 }
451
452 /// Each watched section's status as its last check left it, in watch order. A section a
453 /// session holds keeps the status it had before.
454 pub fn status(&self) -> Vec<(String, SyncStatus)> {
455 self.0.watched.lock().map_or_else(
456 |_| Vec::new(),
457 |watched| {
458 watched
459 .sections
460 .iter()
461 .map(|watch| {
462 let status = &watch.status;
463 (
464 watch.path.clone(),
465 SyncStatus {
466 synced: status.synced,
467 error: status
468 .error
469 .as_ref()
470 .map(|error| io::Error::new(error.kind(), error.to_string())),
471 queued: status.queued,
472 },
473 )
474 })
475 .collect()
476 },
477 )
478 }
479
480 /// Sections whose file changed since the last call, by catalog path.
481 pub fn changed(&self) -> Vec<String> {
482 self.0
483 .watched
484 .lock()
485 .map(|mut watched| std::mem::take(&mut watched.changed))
486 .unwrap_or_default()
487 }
488
489 /// Reads every section's stamp now, working offline included (Sync Now).
490 pub fn wake(&self) {
491 if let Ok(mut watched) = self.0.watched.lock() {
492 let now = Instant::now();
493 for watch in &mut watched.sections {
494 watch.due = now;
495 watch.current = false;
496 }
497 }
498 self.0.signal.requested.store(true, Ordering::Release);
499 self.0.signal.wake();
500 }
501
502 /// Stops for good, then deletes each section's replica that holds nothing unpublished, as
503 /// OneNote lets go of a notebook it closes; one with edits waiting, or that a session
504 /// holds, stays.
505 pub fn discard(&self) {
506 self.0.discard.store(true, Ordering::Release);
507 self.0.signal.stopped.store(true, Ordering::Release);
508 self.0.signal.wake();
509 }
510
511 /// Stops future checks; native threads finish the current step before returning.
512 pub fn stop(&self) {
513 self.0.signal.stopped.store(true, Ordering::Release);
514 self.0.signal.wake();
515 let thread = self.1.lock().ok().and_then(|mut thread| thread.take());
516 #[cfg(not(target_arch = "wasm32"))]
517 if let Some(thread) = thread {
518 let _ = thread.join();
519 }
520 #[cfg(target_arch = "wasm32")]
521 drop(thread);
522 }
523
524 /// Working offline, nothing is checked until `wake`, or until working online again, which
525 /// lists every folder again.
526 pub fn set_offline(&self, offline: bool) {
527 self.0.signal.offline.store(offline, Ordering::Release);
528 self.0.signal.wake();
529 }
530}
531
532impl Drop for Background {
533 fn drop(&mut self) {
534 self.0.signal.stopped.store(true, Ordering::Release);
535 self.0.signal.wake();
536 }
537}
538
539impl Shared {
540 fn touched(&self, paths: &[String]) {
541 if let Ok(mut watched) = self.watched.lock() {
542 watched.touched(paths, Instant::now());
543 }
544 self.signal.wake();
545 if let Ok(listener) = self.listener.lock()
546 && let Some(listener) = &*listener
547 {
548 listener(paths);
549 }
550 }
551
552 /// The session holding the section at `path` let it go: it is checked now.
553 fn released(&self, path: &str) {
554 if let Ok(mut watched) = self.watched.lock()
555 && let Some(watch) = watched.sections.iter_mut().find(|watch| watch.path == path)
556 {
557 watch.held = None;
558 watch.current = false;
559 watch.due = Instant::now();
560 }
561 self.signal.wake();
562 }
563}
564
565#[cfg_attr(not(any(feature = "smb", feature = "live")), allow(dead_code))]
566impl Reports {
567 #[cfg(all(feature = "live", target_arch = "wasm32"))]
568 pub(crate) fn catalog(&self, folder: &str) {
569 if let Some(shared) = self.shared.upgrade()
570 && let Ok(mut watched) = shared.watched.lock()
571 && !watched.changed.iter().any(|path| path == folder)
572 {
573 watched.changed.push(folder.to_owned());
574 }
575 }
576
577 /// The sections at or below these paths changed.
578 pub(crate) fn touched(&self, paths: &[String]) {
579 if let Some(shared) = self.shared.upgrade() {
580 shared.touched(paths);
581 }
582 }
583
584 /// The watch ended, and with it, as far as anyone can tell, the connection.
585 pub(crate) fn lost(&self) {
586 if let Some(shared) = self.shared.upgrade() {
587 if let Ok(mut watched) = shared.watched.lock()
588 && watched.connection == self.connection
589 {
590 watched.lost = true;
591 }
592 shared.signal.wake();
593 }
594 }
595}
596
597/// What the background thread does next.
598enum Next {
599 Wait(Option<Duration>),
600 /// Connects, as the current connection is `.0`, then lists every folder.
601 Connect(u64),
602 /// Lists these folders, checking the sections listed otherwise.
603 List(BTreeSet<String>),
604 Check {
605 path: String,
606 replica: Option<PathBuf>,
607 seen: Option<Box<Stamp>>,
608 current: bool,
609 image: Option<Vec<u8>>,
610 },
611}
612
613impl Watched {
614 fn watch(&mut self, sections: Vec<Known>, now: Instant) {
615 let mut previous = std::mem::take(&mut self.sections);
616 let mut new = 0;
617 self.sections = sections
618 .into_iter()
619 .map(
620 |known| match previous.iter().position(|watch| watch.path == known.path) {
621 Some(index) => Watch {
622 replica: known.replica,
623 image: known.image,
624 ..previous.swap_remove(index)
625 },
626 None => {
627 new += 1;
628 let (listed, stamp) = known.found.unzip();
629 Watch {
630 path: known.path,
631 replica: known.replica,
632 stamp,
633 listed,
634 current: false,
635 reported: false,
636 status: SyncStatus {
637 synced: None,
638 error: None,
639 queued: 0,
640 },
641 image: known.image,
642 due: now + STAGGER * (new - 1),
643 held: None,
644 }
645 }
646 },
647 )
648 .collect();
649 }
650
651 fn interval(&self) -> Duration {
652 if self.reported {
653 Background::BACKSTOP
654 } else {
655 Background::UNWATCHED
656 }
657 }
658
659 /// Whether a watch reports the folder's changes from now on, as the held sessions' workers
660 /// learn.
661 fn watching(&mut self, reported: bool) {
662 self.reported = reported;
663 for signal in self
664 .sections
665 .iter()
666 .filter_map(|watch| watch.held.as_ref()?.upgrade())
667 {
668 signal.watched.store(reported, Ordering::Release);
669 signal.wake();
670 }
671 }
672
673 fn touched(&mut self, paths: &[String], now: Instant) {
674 let settled = now + self.settle.unwrap_or(SETTLE);
675 for path in paths {
676 let path = path.trim_matches('/');
677 let section = self
678 .sections
679 .iter()
680 .position(|watch| watch.path.eq_ignore_ascii_case(path));
681 match section {
682 Some(index) => {
683 let watch = &mut self.sections[index];
684 watch.due = watch.due.min(settled);
685 watch.current = false;
686 watch.reported = true;
687 }
688 None if self.sections.iter().any(|watch| within(&watch.path, path)) => {
689 self.folders.insert(path.to_owned());
690 self.relist = Some(self.relist.map_or(settled, |due| due.min(settled)));
691 }
692 None => {}
693 }
694 }
695 }
696
697 /// Every folder holding a section, to list them all, as a listing begins.
698 fn rescanning(&mut self) -> BTreeSet<String> {
699 for watch in &mut self.sections {
700 watch.reported = false;
701 }
702 self.parents(&[String::new()])
703 }
704
705 /// The folders holding the sections at or below `paths`.
706 fn parents<'a>(&self, paths: impl IntoIterator<Item = &'a String> + Clone) -> BTreeSet<String> {
707 self.sections
708 .iter()
709 .filter(|watch| {
710 paths
711 .clone()
712 .into_iter()
713 .any(|path| within(&watch.path, path))
714 })
715 .map(|watch| split(&watch.path).0.to_owned())
716 .collect()
717 }
718
719 /// Takes the folders' listings. With `rescan`, as on connecting (whether copies are kept),
720 /// every section of a listed folder is checked, from now: first those whose file lists as
721 /// when its stamp was last read, which need no reading unless for a copy, then the others
722 /// one `STAGGER` after another. Otherwise only the sections listed otherwise are, as are
723 /// those a watch reported meanwhile either way.
724 fn listed(
725 &mut self,
726 listings: &BTreeMap<String, Vec<Entry>>,
727 rescan: Option<bool>,
728 now: Instant,
729 ) {
730 let mut staggered = 0;
731 for watch in &mut self.sections {
732 let (folder, name) = split(&watch.path);
733 let Some(entries) = listings.get(folder) else {
734 continue;
735 };
736 let listed = entries
737 .iter()
738 .find(|entry| entry.name.eq_ignore_ascii_case(name))
739 .map(|entry| entry.listed);
740 let unchanged = watch.stamp.is_some() && listed.is_some() && listed == watch.listed;
741 watch.listed = listed;
742 match rescan {
743 // A report while listing may name a change the listing missed.
744 Some(_) if watch.reported => watch.current = false,
745 Some(copies) => {
746 watch.current = unchanged;
747 let read = !unchanged
748 || copies
749 && watch.image.is_none()
750 && watch
751 .replica
752 .as_ref()
753 .is_some_and(|replica| fs::metadata(replica).is_err());
754 watch.due = if read {
755 staggered += 1;
756 now + STAGGER * (staggered - 1)
757 } else {
758 now
759 };
760 }
761 None if !unchanged => {
762 watch.current = false;
763 watch.due = watch.due.min(now);
764 }
765 None => {}
766 }
767 }
768 }
769
770 /// What to do at `now`, waking the held sections that are due; `bound` while connected.
771 fn next(&mut self, now: Instant, bound: bool) -> Next {
772 let interval = self.interval();
773 loop {
774 let due = (0..self.sections.len()).min_by_key(|&index| self.sections[index].due);
775 let at = due.map(|index| self.sections[index].due);
776 if let Some(relist) = self.relist
777 && relist <= now
778 && at.is_none_or(|at| relist <= at)
779 {
780 if !bound {
781 return Next::Connect(self.connection);
782 }
783 let folders = std::mem::take(&mut self.folders);
784 self.relist = None;
785 return Next::List(self.parents(&folders));
786 }
787 let Some(index) = due.filter(|&index| self.sections[index].due <= now) else {
788 let wait = [at, self.relist].into_iter().flatten().min();
789 return Next::Wait(wait.map(|at| at.saturating_duration_since(now)));
790 };
791 let watch = &mut self.sections[index];
792 if !bound {
793 return Next::Connect(self.connection);
794 }
795 if let Some(worker) = watch.held.as_ref().and_then(Weak::upgrade) {
796 worker.wake();
797 watch.due = now + interval;
798 continue;
799 }
800 return Next::Check {
801 path: watch.path.clone(),
802 replica: watch.replica.clone(),
803 seen: watch.stamp.clone().map(Box::new),
804 current: watch.current,
805 image: watch.image.take(),
806 };
807 }
808 }
809}
810
811/// Lists each of `folders` through `list`; a folder that cannot be listed lists nothing, so
812/// that its sections are checked.
813async fn list<F: Future<Output = io::Result<Vec<Entry>>>>(
814 list: &mut impl FnMut(String) -> F,
815 folders: BTreeSet<String>,
816) -> BTreeMap<String, Vec<Entry>> {
817 let mut listings = BTreeMap::new();
818 for folder in folders {
819 let entries = list(folder.clone()).await.unwrap_or_default();
820 listings.insert(folder, entries);
821 }
822 listings
823}
824
825/// Deletes the replica at `replica` if it holds nothing unpublished and no one holds it.
826fn discard(replica: &Path) {
827 if matches!(
828 crate::closed(replica).and_then(|held| crate::peek(&held)),
829 Ok((_, 0))
830 ) {
831 for suffix in ["-wal", "-shm", ""] {
832 let mut file = replica.as_os_str().to_owned();
833 file.push(suffix);
834 let _ = fs::remove_file(file);
835 }
836 }
837}
838
839/// Whether the catalog path `section` is `path` or lies below it, as a share compares names.
840fn within(section: &str, path: &str) -> bool {
841 let path = path.trim_matches('/');
842 path.is_empty()
843 || section.len() >= path.len()
844 && section.is_char_boundary(path.len())
845 && section[..path.len()].eq_ignore_ascii_case(path)
846 && matches!(section.as_bytes().get(path.len()), None | Some(b'/'))
847}
848
849/// A catalog path's folder and name.
850fn split(path: &str) -> (&str, &str) {
851 path.rsplit_once('/').unwrap_or(("", path))
852}
853
854/// Whether a failed step lost the connection, rather than failing for its one file.
855fn disconnected(error: &Error) -> bool {
856 use io::ErrorKind::*;
857 let kind = match error {
858 Error::RemoteIo(error) => error.kind(),
859 Error::Remote(error) => error.error.kind(),
860 _ => return false,
861 };
862 matches!(
863 kind,
864 NotConnected | TimedOut | ConnectionReset | ConnectionAborted | BrokenPipe
865 )
866}
867
868impl Watch {
869 /// Records a failed step and what it found waiting, answering whether the status shown
870 /// changes.
871 fn fail(&mut self, error: &Error, queued: Option<u64>) -> bool {
872 let before = summary(&self.status);
873 self.status.error = Some(reached(error));
874 self.status.queued = queued.unwrap_or(self.status.queued);
875 before != summary(&self.status)
876 }
877}
878
879/// What the host shows of a status.
880fn summary(status: &SyncStatus) -> (bool, Option<io::ErrorKind>, u64) {
881 (
882 status.synced.is_some(),
883 status.error.as_ref().map(io::Error::kind),
884 status.queued,
885 )
886}
887
888/// One section's step: how many of its edits wait (`None` while unknown, as while a session
889/// holds its replica), then the file's stamp now and whether it changed since `seen`, which
890/// with `current` is the stamp now, unread. With `copies`, a section without a replica gets
891/// one from the file as it is now. `image`, the file as discovery read it, stands in for
892/// reading it while the stamp is still its own.
893async fn step_async<R: Remote>(
894 remote: &mut R,
895 replica: Option<&Path>,
896 seen: Option<&Stamp>,
897 current: bool,
898 image: Option<Vec<u8>>,
899 copies: bool,
900) -> (Option<u64>, Result<(Stamp, bool)>) {
901 let stamp = match seen.filter(|_| current) {
902 Some(seen) => seen.clone(),
903 None => match crate::sync::awaited(remote, Remote::stamp).await {
904 Ok(stamp) => stamp,
905 Err(error) => return (None, Err(Error::RemoteIo(error))),
906 },
907 };
908 let remote = &mut Discovered {
909 image: image.and_then(|image| Some((Stamp::of(&image).ok()?, image))),
910 stamp: stamp.clone(),
911 remote,
912 };
913 let moved = seen.is_some_and(|seen| *seen != stamp);
914 let Some(replica) = replica else {
915 return (Some(0), Ok((stamp, moved)));
916 };
917 if fs::metadata(replica).is_err() {
918 if !copies {
919 return (Some(0), Ok((stamp, moved)));
920 }
921 let copied = (async {
922 let image = crate::sync::awaited(remote, Remote::read)
923 .await
924 .map_err(Error::RemoteIo)?;
925 if let Some(folder) = replica.parent() {
926 fs::create_dir_all(folder)?;
927 }
928 Replica::seed(replica, &image, None)?;
929 fs::durable().await?;
930 Ok(Stamp::of(&image)?)
931 })
932 .await;
933 return match copied {
934 Ok(stamp) => (Some(0), Ok((stamp, moved))),
935 // A session made it first.
936 Err(Error::Io(error)) if error.kind() == io::ErrorKind::AlreadyExists => {
937 (None, Ok((stamp, false)))
938 }
939 Err(error) => (None, Err(error)),
940 };
941 }
942 // A version kept beside the file leaves its stamp as it was.
943 let versions = match crate::sync::awaited(remote, Remote::versions).await {
944 Ok(versions) => versions,
945 Err(error) => return (None, Err(Error::RemoteIo(error))),
946 };
947 match crate::closed(replica).and_then(|held| crate::peek(&held)) {
948 Ok((base, 0)) if base == stamp && versions.is_empty() => {
949 return (Some(0), Ok((stamp, moved)));
950 }
951 Err(error) if error.busy() => return (None, Ok((stamp, false))),
952 _ => {}
953 }
954 let replica = match Replica::open(replica) {
955 Ok(replica) => replica,
956 Err(error) if error.busy() => return (None, Ok((stamp, false))),
957 Err(error) => return (None, Err(error)),
958 };
959 let synced = (async {
960 let mut changed = moved;
961 loop {
962 let synced = replica.sync_once_async(remote).await?;
963 changed |= !synced.changed.is_empty();
964 if !matches!(synced.edit, Some((_, EditStatus::Published { .. }))) {
965 return Ok((
966 crate::sync::awaited(remote, Remote::stamp)
967 .await
968 .map_err(Error::RemoteIo)?,
969 changed,
970 ));
971 }
972 }
973 })
974 .await;
975 let queued = replica
976 .recovery_summary()
977 .ok()
978 .map(|summary| summary.queued_edits);
979 (queued, synced)
980}
981
982/// A remote whose first read, while the file's stamp is still `image`'s, answers `image`.
983struct Discovered<'a, R> {
984 remote: &'a mut R,
985 image: Option<(Stamp, Vec<u8>)>,
986 /// The stamp last read.
987 stamp: Stamp,
988}
989
990impl<R: Remote> Remote for Discovered<'_, R> {
991 fn pending(&mut self) -> Option<std::pin::Pin<Box<dyn Future<Output = ()> + '_>>> {
992 self.remote.pending()
993 }
994
995 fn read(&mut self) -> io::Result<Vec<u8>> {
996 match self.image.take() {
997 Some((stamp, image)) if stamp == self.stamp => Ok(image),
998 _ => self.remote.read(),
999 }
1000 }
1001
1002 fn stamp(&mut self) -> io::Result<Stamp> {
1003 self.stamp = self.remote.stamp()?;
1004 Ok(self.stamp.clone())
1005 }
1006
1007 fn publish(
1008 &mut self,
1009 transaction: &onestore::Transaction,
1010 ) -> std::result::Result<(), onestore::CommitError> {
1011 self.remote.publish(transaction)
1012 }
1013
1014 fn confirm(&mut self, base: &Stamp) -> std::result::Result<(), onestore::CommitError> {
1015 self.remote.confirm(base)
1016 }
1017
1018 fn versions(&mut self) -> io::Result<Vec<crate::Version>> {
1019 self.remote.versions()
1020 }
1021
1022 fn version(&mut self, id: &str) -> io::Result<Vec<u8>> {
1023 self.remote.version(id)
1024 }
1025
1026 fn retire(&mut self, id: &str, keep: bool) -> io::Result<()> {
1027 self.remote.retire(id, keep)
1028 }
1029}
1030
1031#[cfg(test)]
1032mod tests {
1033 use super::*;
1034 use crate::discover::EntryKind;
1035 const BACKSTOP: Duration = Background::BACKSTOP;
1036 use onestore::{CommitError, Transaction};
1037 use std::{
1038 collections::HashMap,
1039 sync::mpsc::{self, Receiver},
1040 };
1041
1042 fn path(n: usize) -> String {
1043 let folder = if n.is_multiple_of(2) { "" } else { "Group/" };
1044 format!("{folder}Section {n:03}.one")
1045 }
1046
1047 fn listed(n: usize) -> Listed {
1048 Listed {
1049 size: 1000 + n as u64,
1050 modified: 7,
1051 }
1052 }
1053
1054 fn stamp(n: usize) -> Stamp {
1055 Stamp {
1056 header: [n as u8; 1024],
1057 length: n as u64,
1058 }
1059 }
1060
1061 /// `count` sections, each found by discovery as `listed` and `stamp` have it, with `warm`.
1062 fn sections(count: usize, warm: bool) -> Vec<Known> {
1063 (0..count)
1064 .map(|n| Known {
1065 path: path(n),
1066 replica: None,
1067 found: warm.then(|| (listed(n), stamp(n))),
1068 image: None,
1069 })
1070 .collect()
1071 }
1072
1073 #[test]
1074 fn a_notebook_with_every_section_held_still_connects_its_watch() {
1075 let now = Instant::now();
1076 let mut watched = Watched::default();
1077 watched.watch(sections(1, true), now);
1078 let (held, woken) = Signal::new();
1079 watched.sections[0].held = Some(Arc::downgrade(&held));
1080 assert!(matches!(watched.next(now, false), Next::Connect(0)));
1081 watched.watching(true);
1082 assert!(held.watched.load(Ordering::Acquire));
1083 assert!(matches!(watched.next(now, true), Next::Wait(_)));
1084 assert!(woken.recv_timeout(Duration::from_millis(50)).is_ok());
1085 }
1086
1087 /// The notebook's folders as a listing shows its first `count` sections.
1088 fn listing(count: usize) -> BTreeMap<String, Vec<Entry>> {
1089 let mut folders: BTreeMap<String, Vec<Entry>> = BTreeMap::new();
1090 for n in 0..count {
1091 let path = path(n);
1092 let (folder, name) = split(&path);
1093 folders.entry(folder.to_owned()).or_default().push(Entry {
1094 name: name.to_owned(),
1095 kind: EntryKind::File,
1096 listed: listed(n),
1097 });
1098 }
1099 folders
1100 }
1101
1102 /// Connects at `now`, a watch reporting, and the folders listing as `listings` has them;
1103 /// `meanwhile` runs while they are listed.
1104 fn connect(
1105 watched: &mut Watched,
1106 listings: &BTreeMap<String, Vec<Entry>>,
1107 now: Instant,
1108 meanwhile: impl FnOnce(&mut Watched),
1109 ) {
1110 watched.watching(true);
1111 let folders = watched.rescanning();
1112 meanwhile(watched);
1113 let listings = folders
1114 .into_iter()
1115 .map(|folder| (folder.clone(), listings[&folder].clone()))
1116 .collect();
1117 watched.listed(&listings, Some(false), now);
1118 }
1119
1120 /// Runs the schedule as the thread would from `now` to `end`, every check succeeding and
1121 /// every folder listing as `listings` has it, answering each check: when, which section,
1122 /// and whether it read the file's stamp.
1123 fn run(
1124 watched: &mut Watched,
1125 mut now: Instant,
1126 end: Instant,
1127 listings: &BTreeMap<String, Vec<Entry>>,
1128 ) -> Vec<(Instant, String, bool)> {
1129 let mut checks = Vec::new();
1130 loop {
1131 match watched.next(now, true) {
1132 Next::Wait(Some(wait)) if now + wait <= end => now += wait,
1133 Next::Wait(_) => return checks,
1134 Next::Connect(_) => unreachable!("connected"),
1135 Next::List(folders) => {
1136 let listings = folders
1137 .into_iter()
1138 .map(|folder| (folder.clone(), listings[&folder].clone()))
1139 .collect();
1140 watched.listed(&listings, None, now);
1141 }
1142 Next::Check { path, current, .. } => {
1143 let interval = watched.interval();
1144 let watch = watched
1145 .sections
1146 .iter_mut()
1147 .find(|watch| watch.path == path)
1148 .unwrap();
1149 watch.current = false;
1150 watch.due = now + interval;
1151 checks.push((now, path, !current));
1152 }
1153 }
1154 }
1155 }
1156
1157 #[test]
1158 fn a_new_notebook_is_read_once_staggered_then_left_alone() {
1159 let start = Instant::now();
1160 let mut watched = Watched::default();
1161 watched.watch(sections(200, false), start);
1162 connect(&mut watched, &listing(200), start, |_| {});
1163 let checks = run(
1164 &mut watched,
1165 start,
1166 start + BACKSTOP - SETTLE,
1167 &listing(200),
1168 );
1169 assert_eq!(checks.len(), 200, "one first check each, then none");
1170 assert!(checks.iter().all(|(.., read)| *read));
1171 for pair in checks.windows(2) {
1172 assert!(pair[1].0 - pair[0].0 >= STAGGER, "never a burst");
1173 }
1174 // Then each section costs one backstop check an interval: over a day, 24 stamp reads
1175 // where a 15-second poll made 5,760.
1176 let day = run(
1177 &mut watched,
1178 start + BACKSTOP - SETTLE,
1179 start + Duration::from_secs(86_400) - SETTLE,
1180 &listing(200),
1181 );
1182 assert_eq!(day.len(), 200 * (86_400 / BACKSTOP.as_secs() as usize - 1));
1183 }
1184
1185 #[test]
1186 fn reopening_reads_only_the_files_listed_otherwise() {
1187 let start = Instant::now();
1188 let mut watched = Watched::default();
1189 watched.watch(sections(200, true), start);
1190 let mut listings = listing(200);
1191 // Another client wrote one section while the notebook was closed.
1192 listings.get_mut("Group").unwrap()[3].listed.modified += 1;
1193 connect(&mut watched, &listings, start, |_| {});
1194 let checks = run(&mut watched, start, start + SETTLE, &listings);
1195 assert_eq!(checks.len(), 200);
1196 let read: Vec<_> = checks.iter().filter(|(.., read)| *read).collect();
1197 assert_eq!(read.len(), 1, "{read:?}");
1198 assert_eq!(read[0].1, "Group/Section 007.one");
1199 assert!(
1200 checks.iter().all(|(at, ..)| *at == start),
1201 "nothing to read, nothing to stagger"
1202 );
1203 }
1204
1205 #[test]
1206 fn a_change_reported_while_listing_is_read_though_the_listing_missed_it() {
1207 let start = Instant::now();
1208 let mut watched = Watched::default();
1209 watched.watch(sections(4, true), start);
1210 connect(&mut watched, &listing(4), start, |watched| {
1211 watched.touched(&[path(2)], start);
1212 });
1213 let checks = run(&mut watched, start, start + SETTLE * 2, &listing(4));
1214 let read: Vec<_> = checks
1215 .into_iter()
1216 .filter(|(.., read)| *read)
1217 .map(|(_, path, _)| path)
1218 .collect();
1219 assert_eq!(read, vec![path(2)]);
1220 }
1221
1222 #[test]
1223 fn a_reported_change_checks_exactly_the_sections_it_names() {
1224 let start = Instant::now();
1225 let mut watched = Watched::default();
1226 watched.watch(sections(200, true), start);
1227 let mut listings = listing(200);
1228 connect(&mut watched, &listings, start, |_| {});
1229 let now = start + Duration::from_secs(60);
1230 run(&mut watched, start, now, &listings);
1231 // A commit's several writes are reported apiece, and settle into one check.
1232 for _ in 0..3 {
1233 watched.touched(&["section 004.ONE".into()], now);
1234 }
1235 watched.touched(
1236 &["Section 004.one".into()],
1237 now + Duration::from_millis(300),
1238 );
1239 let checks = run(&mut watched, now, now + Duration::from_secs(60), &listings);
1240 assert_eq!(
1241 checks,
1242 vec![(now + SETTLE, "Section 004.one".to_owned(), true)]
1243 );
1244 // A folder, as a watch names it when it cannot name the file, is listed, and only the
1245 // section listed otherwise is checked.
1246 let now = now + Duration::from_secs(60);
1247 listings.get_mut("Group").unwrap()[2].listed.size += 1;
1248 watched.touched(&["Group".into()], now);
1249 let group = run(&mut watched, now, now + SETTLE * 2, &listings);
1250 assert_eq!(
1251 group,
1252 vec![(now + SETTLE, "Group/Section 005.one".to_owned(), true)]
1253 );
1254 // `""` lists every folder; nothing listed otherwise, nothing is checked.
1255 watched.touched(&["Grou".into(), String::new()], now + SETTLE * 2);
1256 assert!(run(&mut watched, now + SETTLE * 2, now + SETTLE * 4, &listings).is_empty());
1257 }
1258
1259 /// Each file's image and how many times it was written, by path.
1260 type Images = HashMap<String, (Vec<u8>, u64)>;
1261
1262 /// Section files in memory, each stamp read reported as it happens.
1263 #[derive(Clone)]
1264 struct Files {
1265 images: Arc<Mutex<Images>>,
1266 stamps: mpsc::Sender<String>,
1267 }
1268
1269 impl Files {
1270 fn write(&self, path: &str, image: Vec<u8>) {
1271 let mut images = self.images.lock().unwrap();
1272 let version = images.get(path).map_or(0, |(_, version)| version + 1);
1273 images.insert(path.to_owned(), (image, version));
1274 }
1275
1276 fn list(&self, folder: &str) -> Vec<Entry> {
1277 let images = self.images.lock().unwrap();
1278 images
1279 .iter()
1280 .filter(|(path, _)| split(path).0 == folder)
1281 .map(|(path, (image, version))| Entry {
1282 name: split(path).1.to_owned(),
1283 kind: EntryKind::File,
1284 listed: Listed {
1285 size: image.len() as u64,
1286 modified: *version,
1287 },
1288 })
1289 .collect()
1290 }
1291
1292 fn known(&self, path: &str) -> Known {
1293 let (image, version) = self.images.lock().unwrap()[path].clone();
1294 Known {
1295 path: path.to_owned(),
1296 replica: None,
1297 found: Some((
1298 Listed {
1299 size: image.len() as u64,
1300 modified: version,
1301 },
1302 Stamp::of(&image).unwrap(),
1303 )),
1304 image: None,
1305 }
1306 }
1307 }
1308
1309 struct File(Files, String);
1310
1311 impl Remote for File {
1312 fn read(&mut self) -> io::Result<Vec<u8>> {
1313 let images = self.0.images.lock().unwrap();
1314 Ok(images[&self.1].0.clone())
1315 }
1316
1317 fn stamp(&mut self) -> io::Result<Stamp> {
1318 let _ = self.0.stamps.send(self.1.clone());
1319 Stamp::of(&self.read()?).map_err(io::Error::other)
1320 }
1321
1322 fn publish(&mut self, _: &Transaction) -> std::result::Result<(), CommitError> {
1323 unreachable!("nothing is queued")
1324 }
1325
1326 fn confirm(&mut self, _: &Stamp) -> std::result::Result<(), CommitError> {
1327 unreachable!("nothing is queued")
1328 }
1329 }
1330
1331 fn quiet(stamps: &Receiver<String>, wait: Duration) -> Vec<String> {
1332 let mut read = Vec::new();
1333 while let Ok(path) = stamps.recv_timeout(wait) {
1334 read.push(path);
1335 }
1336 read
1337 }
1338
1339 #[test]
1340 fn a_watch_wakes_the_section_it_reports_and_a_lost_one_lists_again() {
1341 let (sender, stamps) = mpsc::channel();
1342 let files = Files {
1343 images: Arc::default(),
1344 stamps: sender,
1345 };
1346 let mut paths: Vec<String> = (0..6).map(path).collect();
1347 paths.sort();
1348 for path in &paths {
1349 files.write(
1350 path,
1351 onestore::create_section("s.one", path, "Author").unwrap(),
1352 );
1353 }
1354 let (armed, watches) = mpsc::channel();
1355 let remote = files.clone();
1356 let background = Background::start(
1357 false,
1358 move |reports| {
1359 armed.send(reports).unwrap();
1360 let (files, listed) = (remote.clone(), remote.clone());
1361 let bind = move |path: &str| File(files.clone(), path.to_owned());
1362 let list = move |folder: String| std::future::ready(Ok(listed.list(&folder)));
1363 Ok(((bind, list), true))
1364 },
1365 || {},
1366 )
1367 .unwrap();
1368 // A section the catalog read unchanged is not read again; a new one is.
1369 let mut known: Vec<Known> = paths[1..].iter().map(|path| files.known(path)).collect();
1370 known.push(Known {
1371 path: paths[0].clone(),
1372 replica: None,
1373 found: None,
1374 image: None,
1375 });
1376 background.watch(known);
1377 assert_eq!(quiet(&stamps, STAGGER * 4), vec![paths[0].clone()]);
1378 let watch = watches.recv().unwrap();
1379 assert!(quiet(&stamps, SETTLE).is_empty(), "idle reads nothing");
1380
1381 // A folder reported without its file is listed; only the file listed otherwise is read.
1382 let other = &paths[4];
1383 files.write(
1384 other,
1385 onestore::create_section("s.one", "Other", "Author").unwrap(),
1386 );
1387 watch.touched(&[String::new()]);
1388 assert_eq!(quiet(&stamps, SETTLE * 2), vec![other.clone()]);
1389 assert_eq!(background.changed(), vec![other.clone()]);
1390
1391 let changed = &paths[3];
1392 files.write(
1393 changed,
1394 onestore::create_section("s.one", "Changed", "Author").unwrap(),
1395 );
1396 watch.touched(std::slice::from_ref(changed));
1397 assert_eq!(quiet(&stamps, SETTLE * 2), vec![changed.clone()]);
1398 assert_eq!(background.changed(), vec![changed.clone()]);
1399
1400 // A watch that ended leaves changes unreported, so every folder is listed again, and
1401 // the files listed otherwise since the last listing are read.
1402 let third = &paths[5];
1403 files.write(
1404 third,
1405 onestore::create_section("s.one", "Third", "Author").unwrap(),
1406 );
1407 watch.lost();
1408 let rewatch = watches.recv_timeout(SETTLE).unwrap();
1409 let mut read = quiet(&stamps, STAGGER * 4);
1410 read.sort();
1411 assert_eq!(read, vec![changed.clone(), third.clone()]);
1412 // The ended watch's connection is gone; only the current one's loss counts.
1413 watch.lost();
1414 assert!(quiet(&stamps, SETTLE).is_empty());
1415 drop(rewatch);
1416 }
1417
1418 #[test]
1419 fn a_held_section_is_left_to_its_worker_which_the_watch_wakes() {
1420 let (sender, stamps) = mpsc::channel();
1421 let files = Files {
1422 images: Arc::default(),
1423 stamps: sender,
1424 };
1425 let held = path(0);
1426 files.write(
1427 &held,
1428 onestore::create_section("s.one", &held, "Author").unwrap(),
1429 );
1430 let (armed, watches) = mpsc::channel();
1431 let remote = files.clone();
1432 let background = Background::start(
1433 false,
1434 move |reports| {
1435 armed.send(reports).unwrap();
1436 let (files, listed) = (remote.clone(), remote.clone());
1437 let bind = move |path: &str| File(files.clone(), path.to_owned());
1438 let list = move |folder: String| std::future::ready(Ok(listed.list(&folder)));
1439 Ok(((bind, list), true))
1440 },
1441 || {},
1442 )
1443 .unwrap();
1444 background.watch(vec![files.known(&held)]);
1445 let watch = watches.recv().unwrap();
1446 assert!(quiet(&stamps, SETTLE).is_empty());
1447 // As `hold` leaves it to a session's worker.
1448 let (signal, woken) = Signal::new();
1449 if let Ok(mut watched) = background.0.watched.lock() {
1450 watched.sections[0].held = Some(Arc::downgrade(&signal));
1451 signal.watched.store(watched.reported, Ordering::Release);
1452 }
1453 assert!(
1454 signal.watched.load(Ordering::Acquire),
1455 "the watch reports, so the worker need not poll"
1456 );
1457 watch.touched(std::slice::from_ref(&held));
1458 assert!(woken.recv_timeout(SETTLE * 2).is_ok(), "the watch wakes it");
1459 assert!(
1460 quiet(&stamps, SETTLE).is_empty(),
1461 "the worker reads, not this"
1462 );
1463 watch.lost();
1464 assert!(woken.recv_timeout(SETTLE).is_ok());
1465 let rearmed = watches.recv_timeout(SETTLE).unwrap();
1466 while !signal.watched.load(Ordering::Acquire) {
1467 woken.recv_timeout(SETTLE * 2).unwrap();
1468 }
1469 assert!(signal.watched.load(Ordering::Acquire));
1470 rearmed.touched(std::slice::from_ref(&held));
1471 assert!(woken.recv_timeout(SETTLE * 2).is_ok());
1472 assert!(quiet(&stamps, SETTLE).is_empty());
1473 drop(signal);
1474 drop(background);
1475 }
1476}