| 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 | |
| 5 | use crate::fs; |
| 6 | use crate::{ |
| 7 | EditStatus, Error, Remote, Replica, Result, |
| 8 | discover::{Entry, Listed}, |
| 9 | session::{Section, SyncStatus, reached}, |
| 10 | worker::Signal, |
| 11 | }; |
| 12 | use onestore::Stamp; |
| 13 | use 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 | }; |
| 23 | use 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. |
| 27 | const 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. |
| 30 | const 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. |
| 33 | const 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. |
| 45 | pub struct Background(Arc<Shared>, Mutex<Option<crate::task::JoinHandle<()>>>); |
| 46 | |
| 47 | /// A section for `Background::watch`, as its notebook knows it. |
| 48 | pub 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. |
| 60 | pub(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. |
| 70 | pub 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))] |
| 74 | pub(crate) struct Reports { |
| 75 | shared: Weak<Shared>, |
| 76 | connection: u64, |
| 77 | } |
| 78 | |
| 79 | #[derive(Default)] |
| 80 | struct 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 | |
| 98 | struct 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 | |
| 119 | impl 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 | |
| 532 | impl 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 | |
| 539 | impl 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))] |
| 566 | impl 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. |
| 598 | enum 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 | |
| 613 | impl 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. |
| 813 | async 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. |
| 826 | fn 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. |
| 840 | fn 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. |
| 850 | fn 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. |
| 855 | fn 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 | |
| 868 | impl 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. |
| 880 | fn 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. |
| 893 | async 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`. |
| 983 | struct Discovered<'a, R> { |
| 984 | remote: &'a mut R, |
| 985 | image: Option<(Stamp, Vec<u8>)>, |
| 986 | /// The stamp last read. |
| 987 | stamp: Stamp, |
| 988 | } |
| 989 | |
| 990 | impl<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)] |
| 1032 | mod 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 | } |