1use std::collections::BTreeMap;
17use std::path::PathBuf;
18use std::time::{Duration, Instant};
19
20use crate::engine_contract::{Commit, EffectiveChange, EntryKind, Error, Result};
21use crate::index::IndexHandle;
22use crate::query::{Basis, Delivery, Query, Report, Request, Selection, WatchDelivery, report};
23use crate::scan::ScanConfig;
24use crate::watch::{WatchConfig, Watcher};
25
26#[derive(Clone, Debug, PartialEq, Eq)]
28pub struct Change {
29 pub path: PathBuf,
31 pub kind: ChangeKind,
33 pub entry_kind: Option<EntryKind>,
35 pub bytes: Option<u64>,
37 pub allocated: Option<u64>,
39 pub mtime_ns: Option<i64>,
41 pub ignored: Option<bool>,
50 pub clock: u64,
52}
53
54#[derive(Clone, Copy, PartialEq, Eq, Debug)]
56pub enum ChangeKind {
57 Upsert,
59 Remove,
61 Invalidate,
66}
67
68#[derive(Clone, Debug, Default)]
70pub struct Batch {
71 pub changes: Vec<Change>,
73 pub dirty: bool,
79}
80
81struct BatchFacts {
83 ignored: Option<BTreeMap<PathBuf, bool>>,
89 reclassified: BTreeMap<PathBuf, EntryFacts>,
92}
93
94impl BatchFacts {
95 fn is_ignored(&self, path: &std::path::Path) -> Option<bool> {
98 self.ignored.as_ref()?.get(path).copied()
99 }
100}
101
102#[derive(Clone, Copy)]
104struct EntryFacts {
105 kind: EntryKind,
106 bytes: u64,
107 allocated: u64,
108 mtime_ns: i64,
109}
110
111#[derive(Debug)]
113pub enum SaveOutcome {
114 Written,
116 Skipped,
118 Failed(Error),
120}
121
122struct Persistence {
123 pending: bool,
124 last_attempt: Instant,
125}
126
127impl Persistence {
128 fn persist_due(
129 &mut self,
130 now: Instant,
131 interval: Duration,
132 save: impl FnOnce() -> Result<bool>,
133 ) -> SaveOutcome {
134 if !save_is_due(self.pending, now.saturating_duration_since(self.last_attempt), interval) {
135 return SaveOutcome::Skipped;
136 }
137 let outcome = match save() {
138 Ok(true) => SaveOutcome::Written,
139 Ok(false) => SaveOutcome::Skipped,
140 Err(error) => SaveOutcome::Failed(error),
141 };
142 self.pending = pending_after(&outcome);
143 self.last_attempt = now;
145 outcome
146 }
147}
148
149fn save_is_due(pending: bool, since_last_save: Duration, interval: Duration) -> bool {
150 pending && since_last_save >= interval
151}
152
153fn pending_after(outcome: &SaveOutcome) -> bool {
154 !matches!(outcome, SaveOutcome::Written)
155}
156
157pub struct Session {
159 index: IndexHandle,
160 watcher: Watcher,
161 scan: ScanConfig,
162 request: Request,
163 plan: crate::Plan,
164 persistence: Persistence,
165 startup_save_error: Option<Error>,
166}
167
168impl Session {
169 pub fn start(request: Request, delivery: Delivery) -> Result<Self> {
175 Self::start_observed(request, delivery, None)
176 }
177
178 pub fn start_with_progress(
190 request: Request,
191 delivery: Delivery,
192 progress: &crate::Progress,
193 ) -> Result<Self> {
194 Self::start_observed(request, delivery, Some(progress))
195 }
196
197 fn start_observed(
198 request: Request,
199 mut delivery: Delivery,
200 progress: Option<&crate::Progress>,
201 ) -> Result<Self> {
202 delivery.watch.get_or_insert_with(WatchDelivery::default);
203 let plan =
204 crate::plan(&request, &delivery, crate::Route::Watch).map_err(Error::InvalidRequest)?;
205 let (index, report, pending, _diagnostics) =
206 crate::execute(&plan, &request.basis, false, progress)?;
207 let startup_save_error = pending.join().err();
208 let index = std::sync::Arc::into_inner(index)
209 .expect("the joined writer released the only other reference");
210 let mut session = Self::new_observed(
211 IndexHandle::new(index),
212 request,
213 &delivery,
214 WatchConfig::default(),
215 progress,
216 )?;
217 session.persistence.pending |= startup_save_error.is_some() || !report.is_complete();
218 session.startup_save_error = startup_save_error;
219 Ok(session)
220 }
221
222 pub fn persist_due(&mut self, now: Instant) -> SaveOutcome {
228 if let Some(error) = self.startup_save_error.take() {
229 self.persistence.last_attempt = now;
230 return SaveOutcome::Failed(error);
231 }
232 let interval = self.plan.delivery().watch.expect("watch plan").interval;
233 self.persistence.persist_due(now, interval, || {
234 if !self.plan.persists() || self.plan.delivery().cache_path.is_none() {
235 return Ok(false);
236 }
237 let index = self.index.snapshot()?;
238 crate::persist_index(&index, &self.plan)
239 })
240 }
241
242 pub fn new(
269 index: IndexHandle,
270 request: Request,
271 delivery: &Delivery,
272 watch: WatchConfig,
273 ) -> Result<Self> {
274 Self::new_observed(index, request, delivery, watch, None)
275 }
276
277 fn new_observed(
278 index: IndexHandle,
279 request: Request,
280 delivery: &Delivery,
281 watch: WatchConfig,
282 progress: Option<&crate::Progress>,
283 ) -> Result<Self> {
284 let root = index.root_path()?;
285 let scan = request.basis.scope.scan_config(delivery);
286 let delivery =
290 Delivery { watch: Some(delivery.watch.unwrap_or_default()), ..delivery.clone() };
291 crate::plan(&request, &delivery, crate::Route::Watch).map_err(Error::InvalidRequest)?;
292 crate::validate_basis_root(&root, &request.basis)?;
293 scan.validate_for_scope(index.scope()?)?;
296 let held = Basis {
299 root: root.clone(),
300 scope: scan.clone().into(),
301 content: index.read_with(crate::Index::content_set)?,
302 };
303 request.validate_read(&held).map_err(Error::InvalidRequest)?;
304 let watcher = Watcher::new(&root, watch)?;
308 Self::finish_initial_handoff(index, request, &delivery, watcher, scan, progress)
309 }
310
311 fn finish_initial_handoff(
321 index: IndexHandle,
322 request: Request,
323 delivery: &Delivery,
324 watcher: Watcher,
325 scan: ScanConfig,
326 progress: Option<&crate::Progress>,
327 ) -> Result<Self> {
328 if let Some(progress) = progress {
329 progress.begin_pass(crate::ProgressPhase::Revalidating);
330 }
331 let observed = ScanConfig { progress: progress.cloned(), ..scan.clone() };
332 let mut dirty = false;
333 let reconciliation = crate::scan::reconcile_handle(&index, &observed, &mut |commit| {
334 dirty |= !commit.changes.is_empty();
335 })?;
336 if !reconciliation.scan.is_complete() && !delivery.accept_partial {
337 return Err(Error::ObservationHandoffIncomplete);
338 }
339 dirty |= drain_initial_capture(&watcher, &index, &observed)?;
340 if !delivery.accept_partial
341 && !index.read_with(|index| crate::query::TreeStatus::of(index, &request).complete)?
342 {
343 return Err(Error::ObservationHandoffIncomplete);
344 }
345 let plan =
346 crate::plan(&request, delivery, crate::Route::Watch).map_err(Error::InvalidRequest)?;
347 Ok(Self {
348 index,
349 watcher,
350 scan,
351 request,
352 plan,
353 persistence: Persistence { pending: dirty, last_attempt: Instant::now() },
354 startup_save_error: None,
355 })
356 }
357
358 pub fn request(&self) -> &Request {
360 &self.request
361 }
362
363 pub fn query(&self) -> &Query {
365 &self.request.query
366 }
367
368 pub fn report(&self, generated_at: std::time::SystemTime) -> Result<Report> {
373 let index = self.index.snapshot()?;
374 report(&index, &self.request, generated_at)
375 }
376
377 pub fn index_snapshot(&self) -> Result<crate::Index> {
381 self.index.snapshot()
382 }
383
384 pub fn next_batch(&mut self, timeout: Duration) -> Result<Option<Batch>> {
392 let mut commits: Vec<Commit> = Vec::new();
393 let outcome =
394 self.watcher.apply_next(&self.index, &self.scan, timeout, &mut |commit: &Commit| {
395 commits.push(commit.clone());
396 });
397
398 self.persistence.pending |= commits.iter().any(|commit| !commit.changes.is_empty());
399 let Some(_report) = outcome? else {
400 return Ok(None);
401 };
402
403 let mut batch = Batch {
404 changes: Vec::new(),
405 dirty: commits.iter().any(|commit| !commit.changes.is_empty()),
406 };
407 let facts = self.batch_facts(&commits)?;
408 for commit in &commits {
409 for effective in &commit.changes {
410 if let Some(change) = self.change_for(effective, commit.clock.0, &facts) {
411 batch.changes.push(change);
412 }
413 }
414 }
415 Ok(Some(batch))
416 }
417
418 fn batch_facts(&self, commits: &[Commit]) -> Result<BatchFacts> {
426 let mut touched: Vec<&PathBuf> = Vec::new();
427 let mut removed: Vec<(&PathBuf, crate::EntryKind)> = Vec::new();
428 let mut reclassified: Vec<&PathBuf> = Vec::new();
429 for effective in commits.iter().flat_map(|commit| &commit.changes) {
430 match effective {
431 EffectiveChange::Inserted { path, .. } | EffectiveChange::Updated { path, .. } => {
432 touched.push(path);
433 }
434 EffectiveChange::Reclassified { path, .. } => {
435 reclassified.push(path);
436 }
437 EffectiveChange::Removed { path, kind, .. } => removed.push((path, *kind)),
438 EffectiveChange::Invalidated { .. }
439 | EffectiveChange::ControlUpdated { .. }
440 | EffectiveChange::ControlRefusalUpdated { .. } => {}
441 }
442 }
443
444 self.index.read_with(|index| {
445 let observed = index.observes_controls();
446 let entries = reclassified
447 .into_iter()
448 .filter_map(|path| {
449 let id = index.lookup(path)?;
450 let attrs = index.attrs_of(id)?;
451 Some((
452 path.clone(),
453 EntryFacts {
454 kind: index.kind_of(id)?,
455 bytes: attrs.size,
456 allocated: attrs.allocated,
457 mtime_ns: attrs.mtime_ns,
458 },
459 ))
460 })
461 .collect();
462 BatchFacts {
463 ignored: observed.then(|| {
464 let mut ignored = touched
465 .into_iter()
466 .filter_map(|path| match index.is_ignored(path) {
467 Ok(Some(ignored)) => Some((path.clone(), ignored)),
468 Ok(None) | Err(_) => None,
471 })
472 .collect::<BTreeMap<_, _>>();
473 for (path, kind) in removed {
474 if index.control_classification_known(path) {
475 ignored.insert(
476 path.clone(),
477 index.control_table().is_ignored(path, kind.is_dir()),
478 );
479 }
480 }
481 ignored
482 }),
483 reclassified: entries,
484 }
485 })
486 }
487
488 fn change_for(
490 &self,
491 effective: &EffectiveChange,
492 clock: u64,
493 facts: &BatchFacts,
494 ) -> Option<Change> {
495 match effective {
496 EffectiveChange::Inserted { path, kind, attrs } => {
497 let name = path.file_name()?.to_string_lossy().into_owned();
498 let candidate = crate::query::Candidate {
499 relative: path,
500 name: &name,
501 kind: *kind,
502 bytes: attrs.size,
503 allocated: attrs.allocated,
504 mtime_ns: attrs.mtime_ns,
505 ignored: facts.is_ignored(path).unwrap_or(false),
506 };
507 self.selection().admits(&candidate).then(|| Change {
508 path: path.clone(),
509 kind: ChangeKind::Upsert,
510 entry_kind: Some(*kind),
511 bytes: Some(attrs.size),
512 allocated: Some(attrs.allocated),
513 mtime_ns: Some(attrs.mtime_ns),
514 ignored: facts.is_ignored(path),
515 clock,
516 })
517 }
518 EffectiveChange::Updated { path, kind, previous: _, current } => {
519 let name = path.file_name()?.to_string_lossy().into_owned();
520 let ignored = facts.is_ignored(path).unwrap_or(false);
521 let candidate = crate::query::Candidate {
522 relative: path,
523 name: &name,
524 kind: *kind,
525 bytes: current.size,
526 allocated: current.allocated,
527 mtime_ns: current.mtime_ns,
528 ignored,
529 };
530 if self.selection().admits(&candidate) {
531 Some(Change {
532 path: path.clone(),
533 kind: ChangeKind::Upsert,
534 entry_kind: Some(*kind),
535 bytes: Some(current.size),
536 allocated: Some(current.allocated),
537 mtime_ns: Some(current.mtime_ns),
538 ignored: facts.is_ignored(path),
539 clock,
540 })
541 } else if self.admits_by_path(path, &name) {
542 Some(Change {
543 path: path.clone(),
544 kind: ChangeKind::Remove,
545 entry_kind: None,
546 bytes: None,
547 allocated: None,
548 mtime_ns: None,
549 ignored: None,
550 clock,
551 })
552 } else {
553 None
554 }
555 }
556 EffectiveChange::Removed { path, .. } => {
562 let name = path.file_name()?.to_string_lossy().into_owned();
563 self.admits_by_path(path, &name).then(|| Change {
564 path: path.clone(),
565 kind: ChangeKind::Remove,
566 entry_kind: None,
567 bytes: None,
568 allocated: None,
569 mtime_ns: None,
570 ignored: facts.is_ignored(path),
571 clock,
572 })
573 }
574 EffectiveChange::Invalidated { path, .. } => Some(Change {
577 path: path.clone(),
578 kind: ChangeKind::Invalidate,
579 entry_kind: None,
580 bytes: None,
581 allocated: None,
582 mtime_ns: None,
583 ignored: None,
584 clock,
585 }),
586 EffectiveChange::Reclassified { path, previous_ignored, current_ignored } => {
591 let name = path.file_name()?.to_string_lossy().into_owned();
592 let entry = facts.reclassified.get(path)?;
597 let admits = |ignored: bool| {
598 self.selection().admits(&crate::query::Candidate {
599 relative: path,
600 name: &name,
601 kind: entry.kind,
602 bytes: entry.bytes,
603 allocated: entry.allocated,
604 mtime_ns: entry.mtime_ns,
605 ignored,
606 })
607 };
608 match (admits(*previous_ignored), admits(*current_ignored)) {
609 (true, false) => Some(Change {
610 path: path.clone(),
611 kind: ChangeKind::Remove,
612 entry_kind: None,
613 bytes: None,
614 allocated: None,
615 mtime_ns: None,
616 ignored: Some(*current_ignored),
617 clock,
618 }),
619 (_, true) => Some(Change {
620 path: path.clone(),
621 kind: ChangeKind::Upsert,
622 entry_kind: Some(entry.kind),
623 bytes: Some(entry.bytes),
624 allocated: Some(entry.allocated),
625 mtime_ns: Some(entry.mtime_ns),
626 ignored: Some(*current_ignored),
627 clock,
628 }),
629 _ => None,
630 }
631 }
632 EffectiveChange::ControlUpdated { .. }
636 | EffectiveChange::ControlRefusalUpdated { .. } => None,
637 }
638 }
639
640 fn admits_by_path(&self, path: &std::path::Path, name: &str) -> bool {
642 let selection = self.selection();
643 if selection.exclude.iter().any(|pattern| pattern.matches(path, name)) {
644 return false;
645 }
646 selection.include.is_empty()
647 || selection.include.iter().any(|pattern| pattern.matches(path, name))
648 }
649
650 fn selection(&self) -> &Selection {
651 &self.request.query.selection
652 }
653}
654
655fn drain_initial_capture(
656 watcher: &Watcher,
657 index: &IndexHandle,
658 scan: &ScanConfig,
659) -> Result<bool> {
660 let mut dirty = false;
661 for _ in 0..2 {
662 watcher.flush_capture()?;
663 let mut drained = false;
664 for _ in 0..=watcher.capture_backlog_bound() {
665 if watcher
666 .apply_next(index, scan, Duration::ZERO, &mut |commit| {
667 dirty |= !commit.changes.is_empty();
668 })?
669 .is_none()
670 {
671 drained = true;
672 break;
673 }
674 }
675 if !drained {
676 return Err(Error::ObservationHandoffIncomplete);
677 }
678 }
679 Ok(dirty)
680}
681
682#[cfg(test)]
683mod tests {
684 use super::*;
685
686 #[test]
693 fn a_retained_session_rejects_a_request_for_another_root_before_binding() {
694 let a = tempfile::tempdir().expect("root a");
695 let b = tempfile::tempdir().expect("root b");
696 let basis = Basis {
697 root: a.path().into(),
698 scope: crate::query::Scope::default(),
699 content: crate::content::AnalysisSet::NONE,
700 };
701 let delivery = Delivery::new(crate::CachePolicy::Off, None);
702 let (index, _) = crate::open(&basis, &delivery).expect("open a");
703 let handle = IndexHandle::new(index);
704 let before = handle.clock().expect("clock");
705 let request = Request::new(
706 Basis { root: b.path().into(), ..basis },
707 Query::default(),
708 std::time::SystemTime::now(),
709 );
710 assert!(matches!(
711 Session::new(handle.clone(), request, &delivery, WatchConfig::default()),
712 Err(Error::InvalidRequest(crate::query::RequestError::RootMismatch { .. }))
713 ));
714 assert_eq!(handle.clock().expect("clock"), before);
715 }
716
717 #[test]
718 fn a_save_is_due_only_when_a_change_is_pending_and_the_throttle_has_elapsed() {
719 let interval = Duration::from_secs(1);
720 let cases = [
721 (true, Duration::from_secs(2), true, "pending and past the interval"),
723 (true, interval, true, "pending, exactly at the interval: inclusive"),
724 (true, Duration::from_millis(1), false, "pending but throttled"),
727 (false, Duration::from_secs(60), false, "nothing pending, however long it has been"),
728 (false, Duration::ZERO, false, "nothing pending and just saved"),
729 ];
730
731 for (pending, since, want, case) in cases {
732 assert_eq!(save_is_due(pending, since, interval), want, "{case}");
733 }
734 }
735
736 #[test]
738 fn only_a_completed_write_clears_the_pending_change() {
739 assert!(!pending_after(&SaveOutcome::Written), "a completed write persists the change");
742 assert!(
743 pending_after(&SaveOutcome::Skipped),
744 "a skipped save wrote nothing, so the change is still owed to disk",
745 );
746 assert!(
747 pending_after(&SaveOutcome::Failed(Error::Snapshot("failed".into()))),
748 "a failed save must be retried, not forgotten"
749 );
750 }
751
752 #[test]
754 fn a_burst_then_a_quiet_tree_still_persists() {
755 let interval = Duration::from_secs(1);
756
757 let mut pending = true;
759 assert!(!save_is_due(pending, Duration::from_millis(50), interval));
760 assert!(pending, "the throttle must not consume the change");
761
762 assert!(save_is_due(pending, Duration::from_secs(3), interval));
765
766 pending = pending_after(&SaveOutcome::Skipped);
769 assert!(pending);
770 pending = pending_after(&SaveOutcome::Written);
771 assert!(!pending, "once written, the loop stops rewriting an unchanged index");
772 }
773
774 #[test]
775 fn skips_and_failures_retry_only_after_another_interval() {
776 let start = Instant::now();
777 let interval = Duration::from_secs(2);
778 let mut persistence = Persistence { pending: true, last_attempt: start };
779 assert!(matches!(
780 persistence.persist_due(start + interval / 2, interval, || panic!("throttled")),
781 SaveOutcome::Skipped
782 ));
783 assert!(matches!(
784 persistence.persist_due(start + interval, interval, || Ok(false)),
785 SaveOutcome::Skipped
786 ));
787 assert!(persistence.pending);
788 assert!(matches!(
789 persistence.persist_due(start + interval, interval, || panic!("skip was throttled")),
790 SaveOutcome::Skipped
791 ));
792 assert!(matches!(
793 persistence.persist_due(start + interval * 2, interval, || {
794 Err(Error::Snapshot("disk unavailable".into()))
795 }),
796 SaveOutcome::Failed(_)
797 ));
798 assert!(persistence.pending);
799 assert!(matches!(
800 persistence
801 .persist_due(start + interval * 2, interval, || panic!("failure was throttled")),
802 SaveOutcome::Skipped
803 ));
804 assert!(matches!(
805 persistence.persist_due(start + interval * 3, interval, || Ok(true)),
806 SaveOutcome::Written
807 ));
808 assert!(!persistence.pending);
809 assert!(matches!(
810 persistence.persist_due(start + interval * 4, interval, || panic!("already persisted")),
811 SaveOutcome::Skipped
812 ));
813 }
814
815 #[test]
816 fn handoff_changes_are_persisted_after_the_tree_goes_quiet() {
817 let root = tempfile::tempdir().expect("root");
818 let cache = tempfile::tempdir().expect("cache");
819 let cache_path = cache.path().join("snapshot");
820 let scan = ScanConfig::default();
821 std::fs::write(root.path().join("before.txt"), b"before").expect("before");
822 let (index, _) = crate::scan::scan_into_index(root.path(), &scan).expect("scan");
823 let request = Request::new(
824 Basis {
825 root: root.path().to_path_buf(),
826 scope: scan.clone().into(),
827 content: crate::content::AnalysisSet::NONE,
828 },
829 Query::default(),
830 std::time::SystemTime::now(),
831 );
832 let interval = Duration::from_secs(2);
833 let delivery = Delivery {
834 stale_ok: false,
835 cache: crate::CachePolicy::Auto,
836 cache_path: Some(cache_path.clone()),
837 accept_partial: false,
838 watch: Some(WatchDelivery { interval }),
839 workers: crate::query::Workers::default(),
840 batch_size: ScanConfig::default().batch_size,
841 order: crate::scan::ScanOrder::default(),
842 };
843 let script = tempfile::NamedTempFile::new().expect("script");
844 let (watcher, _sender) =
845 Watcher::scripted(root.path(), WatchConfig::default(), script.path()).expect("watcher");
846 std::fs::write(root.path().join("during-handoff.txt"), b"handoff").expect("handoff change");
847 let mut session = Session::finish_initial_handoff(
848 IndexHandle::new(index),
849 request,
850 &delivery,
851 watcher,
852 scan,
853 None,
854 )
855 .expect("handoff");
856 let started = session.persistence.last_attempt;
857 assert!(session.persistence.pending, "handoff changes need persistence too");
858 assert!(matches!(session.persist_due(started), SaveOutcome::Skipped));
859 assert!(!cache_path.exists(), "throttle delays the write");
860 assert!(matches!(session.persist_due(started + interval), SaveOutcome::Written));
861 let restored = crate::snapshot::load(&cache_path).expect("load").expect("saved");
862 assert!(matches!(
863 restored.path_state(std::path::Path::new("during-handoff.txt")),
864 crate::PathState::Present { .. }
865 ));
866 assert!(matches!(session.persist_due(started + interval * 2), SaveOutcome::Skipped));
867 }
868
869 #[test]
870 fn startup_save_failure_keeps_the_session_live_and_retries() {
871 let root = tempfile::tempdir().expect("root");
872 let cache = tempfile::tempdir().expect("cache");
873 let parent = cache.path().join("blocked");
879 #[cfg(unix)]
880 std::os::unix::fs::symlink(cache.path().join("missing"), &parent)
881 .expect("dangling parent blocks the write");
882 #[cfg(not(unix))]
883 std::fs::write(&parent, b"").expect("file parent blocks the write");
884 let restore = || {
885 #[cfg(unix)]
886 std::fs::create_dir(cache.path().join("missing")).expect("restore the parent");
887 #[cfg(not(unix))]
888 std::fs::remove_file(&parent).expect("restore the parent");
889 };
890 let cache_path = parent.join("snapshot");
891 std::fs::write(root.path().join("file.txt"), b"content").expect("file");
892 let interval = Duration::from_secs(2);
893 let request = Request::new(
894 Basis {
895 root: root.path().to_path_buf(),
896 scope: ScanConfig::default().into(),
897 content: crate::content::AnalysisSet::NONE,
898 },
899 Query::default(),
900 std::time::SystemTime::now(),
901 );
902 let delivery = Delivery {
903 stale_ok: false,
904 cache: crate::CachePolicy::On,
905 cache_path: Some(cache_path.clone()),
906 accept_partial: false,
907 watch: Some(WatchDelivery { interval }),
908 workers: crate::query::Workers::default(),
909 batch_size: ScanConfig::default().batch_size,
910 order: crate::scan::ScanOrder::default(),
911 };
912 let mut session = Session::start(request, delivery).expect("save failure is nonfatal");
913 assert!(session.report(std::time::SystemTime::now()).expect("live report").status.complete);
914 let now = Instant::now();
915 assert!(matches!(session.persist_due(now), SaveOutcome::Failed(_)));
916 restore();
917 assert!(matches!(session.persist_due(now), SaveOutcome::Skipped));
918 assert!(matches!(session.persist_due(now + interval), SaveOutcome::Written));
919 assert!(crate::snapshot::load(&cache_path).expect("read snapshot").is_some());
920 }
921
922 #[test]
923 fn an_update_that_leaves_attribute_selection_emits_remove() {
924 let root = tempfile::tempdir().expect("tempdir");
925 std::fs::write(root.path().join("file.txt"), b"12345678").expect("fixture");
926 let scan = ScanConfig::default();
927 let (index, report) = crate::scan::scan_into_index(root.path(), &scan).expect("scan");
928 assert!(report.is_complete());
929 let request = Request::new(
930 Basis {
931 root: root.path().to_path_buf(),
932 scope: scan.clone().into(),
933 content: crate::content::AnalysisSet::NONE,
934 },
935 Query {
936 selection: Selection { min_size: Some(4), ..Selection::default() },
937 ..Query::default()
938 },
939 std::time::SystemTime::now(),
940 );
941 let delivery = Delivery {
942 stale_ok: false,
943 cache: crate::CachePolicy::Off,
944 cache_path: None,
945 accept_partial: false,
946 watch: Some(WatchDelivery { interval: Duration::from_millis(50) }),
947 workers: crate::query::Workers::default(),
948 batch_size: ScanConfig::default().batch_size,
949 order: crate::scan::ScanOrder::default(),
950 };
951 let session =
952 Session::new(IndexHandle::new(index), request, &delivery, WatchConfig::default())
953 .expect("session");
954 let path = PathBuf::from("file.txt");
955 let change = session
956 .change_for(
957 &EffectiveChange::Updated {
958 path: path.clone(),
959 kind: EntryKind::File,
960 previous: crate::Attrs { size: 8, allocated: 8, ..crate::Attrs::default() },
961 current: crate::Attrs { size: 1, allocated: 1, ..crate::Attrs::default() },
962 },
963 1,
964 &BatchFacts {
965 ignored: Some(BTreeMap::from([(path, false)])),
966 reclassified: BTreeMap::new(),
967 },
968 )
969 .expect("membership transition");
970 assert_eq!(change.kind, ChangeKind::Remove);
971 }
972
973 #[test]
974 fn initial_handoff_drains_a_sticky_overflow_after_a_full_intent_queue() {
975 let root = tempfile::tempdir().expect("root");
976 let script = tempfile::NamedTempFile::new().expect("script");
977 let scan = ScanConfig::default();
978 let (index, report) = crate::scan::scan_into_index(root.path(), &scan).expect("scan");
979 assert!(report.is_complete());
980 let handle = IndexHandle::new(index);
981 let config = WatchConfig {
982 settle: Duration::from_millis(1),
983 max_hold: Duration::from_millis(2),
984 event_capacity: 8,
985 batch_path_capacity: 1,
986 intent_capacity: 1,
987 ..WatchConfig::default()
988 };
989 let (watcher, sender) =
990 Watcher::scripted(root.path(), config, script.path()).expect("scripted watcher");
991 std::fs::write(root.path().join("a.txt"), b"a").expect("a");
992 sender.send("create\ta.txt\n").expect("first event");
993 watcher.flush_capture().expect("first barrier fills the intent queue");
994 std::fs::write(root.path().join("b.txt"), b"b").expect("b");
995 sender.send("create\tb.txt\n").expect("second event");
996 watcher.flush_capture().expect("second barrier retains sticky overflow");
997
998 drain_initial_capture(&watcher, &handle, &scan).expect("bounded handoff");
999
1000 assert!(
1001 handle.snapshot().expect("snapshot").lookup(std::path::Path::new("a.txt")).is_some()
1002 );
1003 assert!(
1004 handle.snapshot().expect("snapshot").lookup(std::path::Path::new("b.txt")).is_some()
1005 );
1006 assert!(
1007 watcher
1008 .apply_next(&handle, &scan, Duration::ZERO, &mut |_| {})
1009 .expect("proof poll")
1010 .is_none(),
1011 "no queued or sticky pre-handoff work remains"
1012 );
1013 }
1014
1015 #[test]
1016 fn initial_handoff_enforces_partial_acceptance_after_its_reconciliation() {
1017 let root = tempfile::tempdir().expect("root");
1018 std::fs::write(root.path().join("kept.txt"), b"kept").expect("fixture");
1019 let scan = ScanConfig::default();
1020 let (index, report) = crate::scan::scan_into_index(root.path(), &scan).expect("scan");
1021 assert!(report.is_complete());
1022 let request = Request::new(
1023 Basis {
1024 root: root.path().to_path_buf(),
1025 scope: scan.clone().into(),
1026 content: crate::content::AnalysisSet::NONE,
1027 },
1028 Query::default(),
1029 std::time::SystemTime::now(),
1030 );
1031 let _fault = crate::scan::install_walk_hook(root.path(), |_| {
1032 Some(std::io::Error::new(
1033 std::io::ErrorKind::PermissionDenied,
1034 "deterministic handoff refusal",
1035 ))
1036 });
1037 let delivery = Delivery {
1038 stale_ok: false,
1039 cache: crate::CachePolicy::Off,
1040 cache_path: None,
1041 accept_partial: false,
1042 watch: Some(WatchDelivery { interval: Duration::from_millis(50) }),
1043 workers: crate::query::Workers::default(),
1044 batch_size: ScanConfig::default().batch_size,
1045 order: crate::scan::ScanOrder::default(),
1046 };
1047
1048 let Err(error) = Session::new(
1049 IndexHandle::new(index.clone()),
1050 request.clone(),
1051 &delivery,
1052 WatchConfig::default(),
1053 ) else {
1054 panic!("a partial handoff is refused");
1055 };
1056 assert!(matches!(error, Error::ObservationHandoffIncomplete));
1057
1058 let accepted = Session::new(
1059 IndexHandle::new(index),
1060 request,
1061 &Delivery { accept_partial: true, ..delivery },
1062 WatchConfig::default(),
1063 )
1064 .expect("the caller explicitly accepts a partial handoff");
1065 assert!(
1066 !accepted.report(std::time::SystemTime::now()).expect("partial report").status.complete
1067 );
1068 }
1069
1070 #[test]
1071 fn initial_handoff_rechecks_partial_acceptance_after_draining_capture() {
1072 use std::sync::Arc;
1073 use std::sync::atomic::{AtomicUsize, Ordering};
1074
1075 let root = tempfile::tempdir().expect("root");
1076 std::fs::write(root.path().join("kept.txt"), b"kept").expect("fixture");
1077 let scan = ScanConfig::default();
1078 let (index, report) = crate::scan::scan_into_index(root.path(), &scan).expect("scan");
1079 assert!(report.is_complete());
1080 let request = Request::new(
1081 Basis {
1082 root: root.path().to_path_buf(),
1083 scope: scan.clone().into(),
1084 content: crate::content::AnalysisSet::NONE,
1085 },
1086 Query::default(),
1087 std::time::SystemTime::now(),
1088 );
1089 let delivery = Delivery {
1090 stale_ok: false,
1091 cache: crate::CachePolicy::Off,
1092 cache_path: None,
1093 accept_partial: false,
1094 watch: Some(WatchDelivery { interval: Duration::from_millis(50) }),
1095 workers: crate::query::Workers::default(),
1096 batch_size: ScanConfig::default().batch_size,
1097 order: crate::scan::ScanOrder::default(),
1098 };
1099 let script = tempfile::NamedTempFile::new().expect("script");
1100 let (watcher, sender) =
1101 Watcher::scripted(root.path(), WatchConfig::default(), script.path()).expect("watcher");
1102 sender.send("rescan\t.\n").expect("queue initial gap");
1103
1104 let attempts = Arc::new(AtomicUsize::new(0));
1107 let hook_attempts = Arc::clone(&attempts);
1108 let _fault = crate::scan::install_walk_hook(root.path(), move |_| {
1109 (hook_attempts.fetch_add(1, Ordering::SeqCst) > 0).then(|| {
1110 std::io::Error::new(
1111 std::io::ErrorKind::PermissionDenied,
1112 "deterministic drain-only refusal",
1113 )
1114 })
1115 });
1116
1117 let Err(error) = Session::finish_initial_handoff(
1118 IndexHandle::new(index),
1119 request,
1120 &delivery,
1121 watcher,
1122 scan,
1123 None,
1124 ) else {
1125 panic!("a partial state created while draining is refused");
1126 };
1127 assert!(matches!(error, Error::ObservationHandoffIncomplete));
1128 assert!(attempts.load(Ordering::SeqCst) > 1, "the drain ran after startup reconciliation");
1129 }
1130
1131 #[test]
1136 fn the_handoff_pass_restarts_the_counts() {
1137 let root = tempfile::tempdir().expect("root");
1138 let mut bytes = 0;
1139 for directory in 0..2 {
1140 let dir = root.path().join(format!("d{directory}"));
1141 std::fs::create_dir(&dir).expect("directory");
1142 for file in 0..3 {
1143 let size = directory * 3 + file + 1;
1144 std::fs::write(dir.join(format!("f{file}.txt")), vec![b'.'; size]).expect("file");
1145 bytes += size as u64;
1146 }
1147 }
1148 let scan = ScanConfig::default();
1149 let (index, report) = crate::scan::scan_into_index(root.path(), &scan).expect("scan");
1150 assert!(report.is_complete());
1151 let request = Request::new(
1152 Basis {
1153 root: root.path().to_path_buf(),
1154 scope: scan.clone().into(),
1155 content: crate::content::AnalysisSet::NONE,
1156 },
1157 Query { views: vec![crate::query::ViewSpec::Summary], ..Query::default() },
1158 std::time::SystemTime::now(),
1159 );
1160 let delivery = Delivery {
1161 stale_ok: false,
1162 cache: crate::CachePolicy::Off,
1163 cache_path: None,
1164 accept_partial: false,
1165 watch: Some(WatchDelivery { interval: Duration::from_millis(50) }),
1166 workers: crate::query::Workers::default(),
1167 batch_size: ScanConfig::default().batch_size,
1168 order: crate::scan::ScanOrder::default(),
1169 };
1170 let script = tempfile::NamedTempFile::new().expect("script");
1171 let (watcher, _sender) =
1172 Watcher::scripted(root.path(), WatchConfig::default(), script.path()).expect("watcher");
1173 let progress = crate::Progress::new();
1174 progress.add_walked(100, 100, 100, 100);
1175 progress.enter(crate::ProgressPhase::Saving);
1176
1177 let session = Session::finish_initial_handoff(
1178 IndexHandle::new(index.clone()),
1179 request.clone(),
1180 &delivery,
1181 watcher,
1182 scan.clone(),
1183 Some(&progress),
1184 )
1185 .expect("handoff");
1186
1187 let snapshot = progress.snapshot();
1188 assert_eq!(snapshot.phase, crate::ProgressPhase::Revalidating);
1189 assert_eq!(
1190 (snapshot.directories, snapshot.files, snapshot.bytes),
1191 (3, 6, bytes),
1192 "one walk of the root and its two directories, the first pass not added in"
1193 );
1194 assert_eq!(snapshot.allocated, report.allocated_walked, "allocated restarts with them");
1195
1196 let (plain_watcher, _plain_sender) =
1197 Watcher::scripted(root.path(), WatchConfig::default(), script.path()).expect("watcher");
1198 let plain = Session::finish_initial_handoff(
1199 IndexHandle::new(index),
1200 request,
1201 &delivery,
1202 plain_watcher,
1203 scan,
1204 None,
1205 )
1206 .expect("plain handoff");
1207 let generated_at = std::time::SystemTime::now();
1208 let observed_report = session.report(generated_at).expect("observed report");
1209 let mut plain_report = plain.report(generated_at).expect("plain report");
1210 plain_report.provenance = observed_report.provenance.clone();
1211 let json = |report: &Report| {
1212 crate::report_format::render(report, crate::report_format::Format::Json, false)
1213 .expect("render")
1214 };
1215 assert_eq!(json(&plain_report), json(&observed_report));
1216 assert_eq!(progress.snapshot(), snapshot, "the second handoff and reads were not observed");
1217 }
1218
1219 #[test]
1227 fn a_record_claims_a_classification_only_for_an_entry_the_index_still_holds() {
1228 let unobserved = BatchFacts { ignored: None, reclassified: BTreeMap::new() };
1229 assert_eq!(unobserved.is_ignored(std::path::Path::new("any.txt")), None);
1230
1231 let observed = BatchFacts {
1232 ignored: Some(BTreeMap::from([
1233 (PathBuf::from("build/out.bin"), true),
1234 (PathBuf::from("src/main.rs"), false),
1235 ])),
1236 reclassified: BTreeMap::new(),
1237 };
1238 assert_eq!(observed.is_ignored(std::path::Path::new("build/out.bin")), Some(true));
1239 assert_eq!(observed.is_ignored(std::path::Path::new("src/main.rs")), Some(false));
1240 assert_eq!(
1241 observed.is_ignored(std::path::Path::new("gone.tmp")),
1242 None,
1243 "an entry the batch removed is in neither partition, not in the unignored one"
1244 );
1245 }
1246
1247 #[test]
1255 fn a_started_session_reports_its_second_pass_and_then_nothing() {
1256 let root = tempfile::tempdir().expect("root");
1257 let cache = tempfile::tempdir().expect("cache");
1258 let mut bytes = 0;
1259 for directory in 0..4 {
1260 let dir = root.path().join(format!("d{directory}"));
1261 std::fs::create_dir(&dir).expect("directory");
1262 for file in 0..3 {
1263 let size = directory * 3 + file + 1;
1264 std::fs::write(dir.join(format!("f{file}.txt")), vec![b'.'; size]).expect("file");
1265 bytes += size as u64;
1266 }
1267 }
1268 let request = || {
1269 Request::new(
1270 Basis {
1271 root: root.path().to_path_buf(),
1272 scope: ScanConfig::default().into(),
1273 content: crate::content::AnalysisSet::NONE,
1274 },
1275 Query::default(),
1276 std::time::UNIX_EPOCH,
1277 )
1278 };
1279 let delivery = Delivery {
1280 stale_ok: false,
1281 cache: crate::CachePolicy::Auto,
1282 cache_path: Some(cache.path().join("snapshot")),
1283 accept_partial: false,
1284 watch: Some(WatchDelivery { interval: Duration::from_secs(2) }),
1285 workers: crate::query::Workers::default(),
1286 batch_size: 4,
1287 order: crate::scan::ScanOrder::default(),
1288 };
1289
1290 let progress = crate::Progress::new();
1291 let session = Session::start_with_progress(request(), delivery.clone(), &progress)
1292 .expect("observed start");
1293 let after_start = progress.snapshot();
1294 assert_eq!(
1295 after_start.phase,
1296 crate::ProgressPhase::Revalidating,
1297 "the handoff revalidation follows the joined save"
1298 );
1299 assert!(after_start.directories >= 5, "the root and four children: {after_start:?}");
1303 assert!(after_start.files >= 12, "{after_start:?}");
1304 assert!(after_start.bytes >= bytes, "{after_start:?}");
1305 assert_eq!(after_start.analysis, None);
1306
1307 let plain = Session::start(request(), delivery).expect("plain start");
1308 let generated_at = std::time::SystemTime::now();
1309 let observed_report = session.report(generated_at).expect("observed report");
1310 let plain_report = plain.report(generated_at).expect("plain report");
1311 assert!(observed_report.status.complete);
1312 assert!(plain_report.status.complete);
1313 assert_eq!(progress.snapshot(), after_start, "the second start was not observed");
1314 }
1315}