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.delivery().cache.writes() || 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 reclassified: Vec<&PathBuf> = Vec::new();
428 for effective in commits.iter().flat_map(|commit| &commit.changes) {
429 match effective {
430 EffectiveChange::Inserted { path, .. } | EffectiveChange::Updated { path, .. } => {
431 touched.push(path);
432 }
433 EffectiveChange::Reclassified { path, .. } => {
434 reclassified.push(path);
435 }
436 EffectiveChange::Removed { .. }
437 | EffectiveChange::Invalidated { .. }
438 | EffectiveChange::ControlUpdated { .. }
439 | EffectiveChange::ControlRefusalUpdated { .. } => {}
440 }
441 }
442
443 self.index.read_with(|index| {
444 let observed = index.observes_controls();
445 let entries = reclassified
446 .into_iter()
447 .filter_map(|path| {
448 let id = index.lookup(path)?;
449 let attrs = index.attrs_of(id)?;
450 Some((
451 path.clone(),
452 EntryFacts {
453 kind: index.kind_of(id)?,
454 bytes: attrs.size,
455 allocated: attrs.allocated,
456 mtime_ns: attrs.mtime_ns,
457 },
458 ))
459 })
460 .collect();
461 BatchFacts {
462 ignored: observed.then(|| {
463 touched
464 .into_iter()
465 .filter_map(|path| match index.is_ignored(path) {
466 Ok(Some(ignored)) => Some((path.clone(), ignored)),
467 Ok(None) | Err(_) => None,
470 })
471 .collect()
472 }),
473 reclassified: entries,
474 }
475 })
476 }
477
478 fn change_for(
480 &self,
481 effective: &EffectiveChange,
482 clock: u64,
483 facts: &BatchFacts,
484 ) -> Option<Change> {
485 match effective {
486 EffectiveChange::Inserted { path, kind, attrs } => {
487 let name = path.file_name()?.to_string_lossy().into_owned();
488 let candidate = crate::query::Candidate {
489 relative: path,
490 name: &name,
491 kind: *kind,
492 bytes: attrs.size,
493 allocated: attrs.allocated,
494 mtime_ns: attrs.mtime_ns,
495 ignored: facts.is_ignored(path).unwrap_or(false),
496 };
497 self.selection().admits(&candidate).then(|| Change {
498 path: path.clone(),
499 kind: ChangeKind::Upsert,
500 entry_kind: Some(*kind),
501 bytes: Some(attrs.size),
502 allocated: Some(attrs.allocated),
503 mtime_ns: Some(attrs.mtime_ns),
504 ignored: facts.is_ignored(path),
505 clock,
506 })
507 }
508 EffectiveChange::Updated { path, kind, previous: _, current } => {
509 let name = path.file_name()?.to_string_lossy().into_owned();
510 let ignored = facts.is_ignored(path).unwrap_or(false);
511 let candidate = crate::query::Candidate {
512 relative: path,
513 name: &name,
514 kind: *kind,
515 bytes: current.size,
516 allocated: current.allocated,
517 mtime_ns: current.mtime_ns,
518 ignored,
519 };
520 if self.selection().admits(&candidate) {
521 Some(Change {
522 path: path.clone(),
523 kind: ChangeKind::Upsert,
524 entry_kind: Some(*kind),
525 bytes: Some(current.size),
526 allocated: Some(current.allocated),
527 mtime_ns: Some(current.mtime_ns),
528 ignored: facts.is_ignored(path),
529 clock,
530 })
531 } else if self.admits_by_path(path, &name) {
532 Some(Change {
533 path: path.clone(),
534 kind: ChangeKind::Remove,
535 entry_kind: None,
536 bytes: None,
537 allocated: None,
538 mtime_ns: None,
539 ignored: None,
540 clock,
541 })
542 } else {
543 None
544 }
545 }
546 EffectiveChange::Removed { path, .. } => {
551 let name = path.file_name()?.to_string_lossy().into_owned();
552 self.admits_by_path(path, &name).then(|| Change {
553 path: path.clone(),
554 kind: ChangeKind::Remove,
555 entry_kind: None,
556 bytes: None,
557 allocated: None,
558 mtime_ns: None,
559 ignored: None,
560 clock,
561 })
562 }
563 EffectiveChange::Invalidated { path, .. } => Some(Change {
566 path: path.clone(),
567 kind: ChangeKind::Invalidate,
568 entry_kind: None,
569 bytes: None,
570 allocated: None,
571 mtime_ns: None,
572 ignored: None,
573 clock,
574 }),
575 EffectiveChange::Reclassified { path, previous_ignored, current_ignored } => {
580 let name = path.file_name()?.to_string_lossy().into_owned();
581 let entry = facts.reclassified.get(path)?;
586 let admits = |ignored: bool| {
587 self.selection().admits(&crate::query::Candidate {
588 relative: path,
589 name: &name,
590 kind: entry.kind,
591 bytes: entry.bytes,
592 allocated: entry.allocated,
593 mtime_ns: entry.mtime_ns,
594 ignored,
595 })
596 };
597 match (admits(*previous_ignored), admits(*current_ignored)) {
598 (true, false) => Some(Change {
599 path: path.clone(),
600 kind: ChangeKind::Remove,
601 entry_kind: None,
602 bytes: None,
603 allocated: None,
604 mtime_ns: None,
605 ignored: Some(*current_ignored),
606 clock,
607 }),
608 (_, true) => Some(Change {
609 path: path.clone(),
610 kind: ChangeKind::Upsert,
611 entry_kind: Some(entry.kind),
612 bytes: Some(entry.bytes),
613 allocated: Some(entry.allocated),
614 mtime_ns: Some(entry.mtime_ns),
615 ignored: Some(*current_ignored),
616 clock,
617 }),
618 _ => None,
619 }
620 }
621 EffectiveChange::ControlUpdated { .. }
625 | EffectiveChange::ControlRefusalUpdated { .. } => None,
626 }
627 }
628
629 fn admits_by_path(&self, path: &std::path::Path, name: &str) -> bool {
631 let selection = self.selection();
632 if selection.exclude.iter().any(|pattern| pattern.matches(path, name)) {
633 return false;
634 }
635 selection.include.is_empty()
636 || selection.include.iter().any(|pattern| pattern.matches(path, name))
637 }
638
639 fn selection(&self) -> &Selection {
640 &self.request.query.selection
641 }
642}
643
644fn drain_initial_capture(
645 watcher: &Watcher,
646 index: &IndexHandle,
647 scan: &ScanConfig,
648) -> Result<bool> {
649 let mut dirty = false;
650 for _ in 0..2 {
651 watcher.flush_capture()?;
652 let mut drained = false;
653 for _ in 0..=watcher.capture_backlog_bound() {
654 if watcher
655 .apply_next(index, scan, Duration::ZERO, &mut |commit| {
656 dirty |= !commit.changes.is_empty();
657 })?
658 .is_none()
659 {
660 drained = true;
661 break;
662 }
663 }
664 if !drained {
665 return Err(Error::ObservationHandoffIncomplete);
666 }
667 }
668 Ok(dirty)
669}
670
671#[cfg(test)]
672mod tests {
673 use super::*;
674
675 #[test]
682 fn a_retained_session_rejects_a_request_for_another_root_before_binding() {
683 let a = tempfile::tempdir().expect("root a");
684 let b = tempfile::tempdir().expect("root b");
685 let basis = Basis {
686 root: a.path().into(),
687 scope: crate::query::Scope::default(),
688 content: crate::content::AnalysisSet::NONE,
689 };
690 let delivery = Delivery::new(crate::CachePolicy::Off, None);
691 let (index, _) = crate::open(&basis, &delivery).expect("open a");
692 let handle = IndexHandle::new(index);
693 let before = handle.clock().expect("clock");
694 let request = Request::new(
695 Basis { root: b.path().into(), ..basis },
696 Query::default(),
697 std::time::SystemTime::now(),
698 );
699 assert!(matches!(
700 Session::new(handle.clone(), request, &delivery, WatchConfig::default()),
701 Err(Error::InvalidRequest(crate::query::RequestError::RootMismatch { .. }))
702 ));
703 assert_eq!(handle.clock().expect("clock"), before);
704 }
705
706 #[test]
707 fn a_save_is_due_only_when_a_change_is_pending_and_the_throttle_has_elapsed() {
708 let interval = Duration::from_secs(1);
709 let cases = [
710 (true, Duration::from_secs(2), true, "pending and past the interval"),
712 (true, interval, true, "pending, exactly at the interval: inclusive"),
713 (true, Duration::from_millis(1), false, "pending but throttled"),
716 (false, Duration::from_secs(60), false, "nothing pending, however long it has been"),
717 (false, Duration::ZERO, false, "nothing pending and just saved"),
718 ];
719
720 for (pending, since, want, case) in cases {
721 assert_eq!(save_is_due(pending, since, interval), want, "{case}");
722 }
723 }
724
725 #[test]
727 fn only_a_completed_write_clears_the_pending_change() {
728 assert!(!pending_after(&SaveOutcome::Written), "a completed write persists the change");
731 assert!(
732 pending_after(&SaveOutcome::Skipped),
733 "a skipped save wrote nothing, so the change is still owed to disk",
734 );
735 assert!(
736 pending_after(&SaveOutcome::Failed(Error::Snapshot("failed".into()))),
737 "a failed save must be retried, not forgotten"
738 );
739 }
740
741 #[test]
743 fn a_burst_then_a_quiet_tree_still_persists() {
744 let interval = Duration::from_secs(1);
745
746 let mut pending = true;
748 assert!(!save_is_due(pending, Duration::from_millis(50), interval));
749 assert!(pending, "the throttle must not consume the change");
750
751 assert!(save_is_due(pending, Duration::from_secs(3), interval));
754
755 pending = pending_after(&SaveOutcome::Skipped);
758 assert!(pending);
759 pending = pending_after(&SaveOutcome::Written);
760 assert!(!pending, "once written, the loop stops rewriting an unchanged index");
761 }
762
763 #[test]
764 fn skips_and_failures_retry_only_after_another_interval() {
765 let start = Instant::now();
766 let interval = Duration::from_secs(2);
767 let mut persistence = Persistence { pending: true, last_attempt: start };
768 assert!(matches!(
769 persistence.persist_due(start + interval / 2, interval, || panic!("throttled")),
770 SaveOutcome::Skipped
771 ));
772 assert!(matches!(
773 persistence.persist_due(start + interval, interval, || Ok(false)),
774 SaveOutcome::Skipped
775 ));
776 assert!(persistence.pending);
777 assert!(matches!(
778 persistence.persist_due(start + interval, interval, || panic!("skip was throttled")),
779 SaveOutcome::Skipped
780 ));
781 assert!(matches!(
782 persistence.persist_due(start + interval * 2, interval, || {
783 Err(Error::Snapshot("disk unavailable".into()))
784 }),
785 SaveOutcome::Failed(_)
786 ));
787 assert!(persistence.pending);
788 assert!(matches!(
789 persistence
790 .persist_due(start + interval * 2, interval, || panic!("failure was throttled")),
791 SaveOutcome::Skipped
792 ));
793 assert!(matches!(
794 persistence.persist_due(start + interval * 3, interval, || Ok(true)),
795 SaveOutcome::Written
796 ));
797 assert!(!persistence.pending);
798 assert!(matches!(
799 persistence.persist_due(start + interval * 4, interval, || panic!("already persisted")),
800 SaveOutcome::Skipped
801 ));
802 }
803
804 #[test]
805 fn handoff_changes_are_persisted_after_the_tree_goes_quiet() {
806 let root = tempfile::tempdir().expect("root");
807 let cache = tempfile::tempdir().expect("cache");
808 let cache_path = cache.path().join("snapshot");
809 let scan = ScanConfig::default();
810 std::fs::write(root.path().join("before.txt"), b"before").expect("before");
811 let (index, _) = crate::scan::scan_into_index(root.path(), &scan).expect("scan");
812 let request = Request::new(
813 Basis {
814 root: root.path().to_path_buf(),
815 scope: scan.clone().into(),
816 content: crate::content::AnalysisSet::NONE,
817 },
818 Query::default(),
819 std::time::SystemTime::now(),
820 );
821 let interval = Duration::from_secs(2);
822 let delivery = Delivery {
823 cache: crate::CachePolicy::Auto,
824 cache_path: Some(cache_path.clone()),
825 accept_partial: false,
826 watch: Some(WatchDelivery { interval }),
827 workers: crate::query::Workers::default(),
828 batch_size: ScanConfig::default().batch_size,
829 order: crate::scan::ScanOrder::default(),
830 };
831 let script = tempfile::NamedTempFile::new().expect("script");
832 let (watcher, _sender) =
833 Watcher::scripted(root.path(), WatchConfig::default(), script.path()).expect("watcher");
834 std::fs::write(root.path().join("during-handoff.txt"), b"handoff").expect("handoff change");
835 let mut session = Session::finish_initial_handoff(
836 IndexHandle::new(index),
837 request,
838 &delivery,
839 watcher,
840 scan,
841 None,
842 )
843 .expect("handoff");
844 let started = session.persistence.last_attempt;
845 assert!(session.persistence.pending, "handoff changes need persistence too");
846 assert!(matches!(session.persist_due(started), SaveOutcome::Skipped));
847 assert!(!cache_path.exists(), "throttle delays the write");
848 assert!(matches!(session.persist_due(started + interval), SaveOutcome::Written));
849 let restored = crate::snapshot::load(&cache_path).expect("load").expect("saved");
850 assert!(matches!(
851 restored.path_state(std::path::Path::new("during-handoff.txt")),
852 crate::PathState::Present { .. }
853 ));
854 assert!(matches!(session.persist_due(started + interval * 2), SaveOutcome::Skipped));
855 }
856
857 #[test]
858 fn startup_save_failure_keeps_the_session_live_and_retries() {
859 let root = tempfile::tempdir().expect("root");
860 let cache = tempfile::tempdir().expect("cache");
861 let cache_path = cache.path().join("blocked-snapshot");
862 std::fs::create_dir(&cache_path).expect("directory blocks snapshot rename");
863 std::fs::write(root.path().join("file.txt"), b"content").expect("file");
864 let interval = Duration::from_secs(2);
865 let request = Request::new(
866 Basis {
867 root: root.path().to_path_buf(),
868 scope: ScanConfig::default().into(),
869 content: crate::content::AnalysisSet::NONE,
870 },
871 Query::default(),
872 std::time::SystemTime::now(),
873 );
874 let delivery = Delivery {
875 cache: crate::CachePolicy::Refresh,
876 cache_path: Some(cache_path.clone()),
877 accept_partial: false,
878 watch: Some(WatchDelivery { interval }),
879 workers: crate::query::Workers::default(),
880 batch_size: ScanConfig::default().batch_size,
881 order: crate::scan::ScanOrder::default(),
882 };
883 let mut session = Session::start(request, delivery).expect("save failure is nonfatal");
884 assert!(session.report(std::time::SystemTime::now()).expect("live report").status.complete);
885 let now = Instant::now();
886 assert!(matches!(session.persist_due(now), SaveOutcome::Failed(_)));
887 std::fs::remove_dir(&cache_path).expect("restore writable destination");
888 assert!(matches!(session.persist_due(now), SaveOutcome::Skipped));
889 assert!(matches!(session.persist_due(now + interval), SaveOutcome::Written));
890 assert!(crate::snapshot::load(&cache_path).expect("read snapshot").is_some());
891 }
892
893 #[test]
894 fn an_update_that_leaves_attribute_selection_emits_remove() {
895 let root = tempfile::tempdir().expect("tempdir");
896 std::fs::write(root.path().join("file.txt"), b"12345678").expect("fixture");
897 let scan = ScanConfig::default();
898 let (index, report) = crate::scan::scan_into_index(root.path(), &scan).expect("scan");
899 assert!(report.is_complete());
900 let request = Request::new(
901 Basis {
902 root: root.path().to_path_buf(),
903 scope: scan.clone().into(),
904 content: crate::content::AnalysisSet::NONE,
905 },
906 Query {
907 selection: Selection { min_size: Some(4), ..Selection::default() },
908 ..Query::default()
909 },
910 std::time::SystemTime::now(),
911 );
912 let delivery = Delivery {
913 cache: crate::CachePolicy::Off,
914 cache_path: None,
915 accept_partial: false,
916 watch: Some(WatchDelivery { interval: Duration::from_millis(50) }),
917 workers: crate::query::Workers::default(),
918 batch_size: ScanConfig::default().batch_size,
919 order: crate::scan::ScanOrder::default(),
920 };
921 let session =
922 Session::new(IndexHandle::new(index), request, &delivery, WatchConfig::default())
923 .expect("session");
924 let path = PathBuf::from("file.txt");
925 let change = session
926 .change_for(
927 &EffectiveChange::Updated {
928 path: path.clone(),
929 kind: EntryKind::File,
930 previous: crate::Attrs { size: 8, allocated: 8, ..crate::Attrs::default() },
931 current: crate::Attrs { size: 1, allocated: 1, ..crate::Attrs::default() },
932 },
933 1,
934 &BatchFacts {
935 ignored: Some(BTreeMap::from([(path, false)])),
936 reclassified: BTreeMap::new(),
937 },
938 )
939 .expect("membership transition");
940 assert_eq!(change.kind, ChangeKind::Remove);
941 }
942
943 #[test]
944 fn initial_handoff_drains_a_sticky_overflow_after_a_full_intent_queue() {
945 let root = tempfile::tempdir().expect("root");
946 let script = tempfile::NamedTempFile::new().expect("script");
947 let scan = ScanConfig::default();
948 let (index, report) = crate::scan::scan_into_index(root.path(), &scan).expect("scan");
949 assert!(report.is_complete());
950 let handle = IndexHandle::new(index);
951 let config = WatchConfig {
952 settle: Duration::from_millis(1),
953 max_hold: Duration::from_millis(2),
954 event_capacity: 8,
955 batch_path_capacity: 1,
956 intent_capacity: 1,
957 ..WatchConfig::default()
958 };
959 let (watcher, sender) =
960 Watcher::scripted(root.path(), config, script.path()).expect("scripted watcher");
961 std::fs::write(root.path().join("a.txt"), b"a").expect("a");
962 sender.send("create\ta.txt\n").expect("first event");
963 watcher.flush_capture().expect("first barrier fills the intent queue");
964 std::fs::write(root.path().join("b.txt"), b"b").expect("b");
965 sender.send("create\tb.txt\n").expect("second event");
966 watcher.flush_capture().expect("second barrier retains sticky overflow");
967
968 drain_initial_capture(&watcher, &handle, &scan).expect("bounded handoff");
969
970 assert!(
971 handle.snapshot().expect("snapshot").lookup(std::path::Path::new("a.txt")).is_some()
972 );
973 assert!(
974 handle.snapshot().expect("snapshot").lookup(std::path::Path::new("b.txt")).is_some()
975 );
976 assert!(
977 watcher
978 .apply_next(&handle, &scan, Duration::ZERO, &mut |_| {})
979 .expect("proof poll")
980 .is_none(),
981 "no queued or sticky pre-handoff work remains"
982 );
983 }
984
985 #[test]
986 fn initial_handoff_enforces_partial_acceptance_after_its_reconciliation() {
987 let root = tempfile::tempdir().expect("root");
988 std::fs::write(root.path().join("kept.txt"), b"kept").expect("fixture");
989 let scan = ScanConfig::default();
990 let (index, report) = crate::scan::scan_into_index(root.path(), &scan).expect("scan");
991 assert!(report.is_complete());
992 let request = Request::new(
993 Basis {
994 root: root.path().to_path_buf(),
995 scope: scan.clone().into(),
996 content: crate::content::AnalysisSet::NONE,
997 },
998 Query::default(),
999 std::time::SystemTime::now(),
1000 );
1001 let _fault = crate::scan::install_walk_hook(root.path(), |_| {
1002 Some(std::io::Error::new(
1003 std::io::ErrorKind::PermissionDenied,
1004 "deterministic handoff refusal",
1005 ))
1006 });
1007 let delivery = Delivery {
1008 cache: crate::CachePolicy::Off,
1009 cache_path: None,
1010 accept_partial: false,
1011 watch: Some(WatchDelivery { interval: Duration::from_millis(50) }),
1012 workers: crate::query::Workers::default(),
1013 batch_size: ScanConfig::default().batch_size,
1014 order: crate::scan::ScanOrder::default(),
1015 };
1016
1017 let Err(error) = Session::new(
1018 IndexHandle::new(index.clone()),
1019 request.clone(),
1020 &delivery,
1021 WatchConfig::default(),
1022 ) else {
1023 panic!("a partial handoff is refused");
1024 };
1025 assert!(matches!(error, Error::ObservationHandoffIncomplete));
1026
1027 let accepted = Session::new(
1028 IndexHandle::new(index),
1029 request,
1030 &Delivery { accept_partial: true, ..delivery },
1031 WatchConfig::default(),
1032 )
1033 .expect("the caller explicitly accepts a partial handoff");
1034 assert!(
1035 !accepted.report(std::time::SystemTime::now()).expect("partial report").status.complete
1036 );
1037 }
1038
1039 #[test]
1040 fn initial_handoff_rechecks_partial_acceptance_after_draining_capture() {
1041 use std::sync::Arc;
1042 use std::sync::atomic::{AtomicUsize, Ordering};
1043
1044 let root = tempfile::tempdir().expect("root");
1045 std::fs::write(root.path().join("kept.txt"), b"kept").expect("fixture");
1046 let scan = ScanConfig::default();
1047 let (index, report) = crate::scan::scan_into_index(root.path(), &scan).expect("scan");
1048 assert!(report.is_complete());
1049 let request = Request::new(
1050 Basis {
1051 root: root.path().to_path_buf(),
1052 scope: scan.clone().into(),
1053 content: crate::content::AnalysisSet::NONE,
1054 },
1055 Query::default(),
1056 std::time::SystemTime::now(),
1057 );
1058 let delivery = Delivery {
1059 cache: crate::CachePolicy::Off,
1060 cache_path: None,
1061 accept_partial: false,
1062 watch: Some(WatchDelivery { interval: Duration::from_millis(50) }),
1063 workers: crate::query::Workers::default(),
1064 batch_size: ScanConfig::default().batch_size,
1065 order: crate::scan::ScanOrder::default(),
1066 };
1067 let script = tempfile::NamedTempFile::new().expect("script");
1068 let (watcher, sender) =
1069 Watcher::scripted(root.path(), WatchConfig::default(), script.path()).expect("watcher");
1070 sender.send("rescan\t.\n").expect("queue initial gap");
1071
1072 let attempts = Arc::new(AtomicUsize::new(0));
1075 let hook_attempts = Arc::clone(&attempts);
1076 let _fault = crate::scan::install_walk_hook(root.path(), move |_| {
1077 (hook_attempts.fetch_add(1, Ordering::SeqCst) > 0).then(|| {
1078 std::io::Error::new(
1079 std::io::ErrorKind::PermissionDenied,
1080 "deterministic drain-only refusal",
1081 )
1082 })
1083 });
1084
1085 let Err(error) = Session::finish_initial_handoff(
1086 IndexHandle::new(index),
1087 request,
1088 &delivery,
1089 watcher,
1090 scan,
1091 None,
1092 ) else {
1093 panic!("a partial state created while draining is refused");
1094 };
1095 assert!(matches!(error, Error::ObservationHandoffIncomplete));
1096 assert!(attempts.load(Ordering::SeqCst) > 1, "the drain ran after startup reconciliation");
1097 }
1098
1099 #[test]
1103 fn the_handoff_pass_restarts_the_counts() {
1104 let root = tempfile::tempdir().expect("root");
1105 let mut bytes = 0;
1106 for directory in 0..2 {
1107 let dir = root.path().join(format!("d{directory}"));
1108 std::fs::create_dir(&dir).expect("directory");
1109 for file in 0..3 {
1110 let size = directory * 3 + file + 1;
1111 std::fs::write(dir.join(format!("f{file}.txt")), vec![b'.'; size]).expect("file");
1112 bytes += size as u64;
1113 }
1114 }
1115 let scan = ScanConfig::default();
1116 let (index, report) = crate::scan::scan_into_index(root.path(), &scan).expect("scan");
1117 assert!(report.is_complete());
1118 let request = Request::new(
1119 Basis {
1120 root: root.path().to_path_buf(),
1121 scope: scan.clone().into(),
1122 content: crate::content::AnalysisSet::NONE,
1123 },
1124 Query::default(),
1125 std::time::SystemTime::now(),
1126 );
1127 let delivery = Delivery {
1128 cache: crate::CachePolicy::Off,
1129 cache_path: None,
1130 accept_partial: false,
1131 watch: Some(WatchDelivery { interval: Duration::from_millis(50) }),
1132 workers: crate::query::Workers::default(),
1133 batch_size: ScanConfig::default().batch_size,
1134 order: crate::scan::ScanOrder::default(),
1135 };
1136 let script = tempfile::NamedTempFile::new().expect("script");
1137 let (watcher, _sender) =
1138 Watcher::scripted(root.path(), WatchConfig::default(), script.path()).expect("watcher");
1139 let progress = crate::Progress::new();
1140 progress.add_walked(100, 100, 100, 100);
1141 progress.enter(crate::ProgressPhase::Saving);
1142
1143 Session::finish_initial_handoff(
1144 IndexHandle::new(index),
1145 request,
1146 &delivery,
1147 watcher,
1148 scan,
1149 Some(&progress),
1150 )
1151 .expect("handoff");
1152
1153 let snapshot = progress.snapshot();
1154 assert_eq!(snapshot.phase, crate::ProgressPhase::Revalidating);
1155 assert_eq!(
1156 (snapshot.directories, snapshot.files, snapshot.bytes),
1157 (3, 6, bytes),
1158 "one walk of the root and its two directories, the first pass not added in"
1159 );
1160 assert_eq!(snapshot.allocated, report.allocated_walked, "allocated restarts with them");
1161 }
1162
1163 #[test]
1171 fn a_record_claims_a_classification_only_for_an_entry_the_index_still_holds() {
1172 let unobserved = BatchFacts { ignored: None, reclassified: BTreeMap::new() };
1173 assert_eq!(unobserved.is_ignored(std::path::Path::new("any.txt")), None);
1174
1175 let observed = BatchFacts {
1176 ignored: Some(BTreeMap::from([
1177 (PathBuf::from("build/out.bin"), true),
1178 (PathBuf::from("src/main.rs"), false),
1179 ])),
1180 reclassified: BTreeMap::new(),
1181 };
1182 assert_eq!(observed.is_ignored(std::path::Path::new("build/out.bin")), Some(true));
1183 assert_eq!(observed.is_ignored(std::path::Path::new("src/main.rs")), Some(false));
1184 assert_eq!(
1185 observed.is_ignored(std::path::Path::new("gone.tmp")),
1186 None,
1187 "an entry the batch removed is in neither partition, not in the unignored one"
1188 );
1189 }
1190
1191 #[test]
1198 fn a_started_session_reports_its_second_pass_and_then_nothing() {
1199 let root = tempfile::tempdir().expect("root");
1200 let cache = tempfile::tempdir().expect("cache");
1201 let mut bytes = 0;
1202 for directory in 0..4 {
1203 let dir = root.path().join(format!("d{directory}"));
1204 std::fs::create_dir(&dir).expect("directory");
1205 for file in 0..3 {
1206 let size = directory * 3 + file + 1;
1207 std::fs::write(dir.join(format!("f{file}.txt")), vec![b'.'; size]).expect("file");
1208 bytes += size as u64;
1209 }
1210 }
1211 let request = || {
1212 Request::new(
1213 Basis {
1214 root: root.path().to_path_buf(),
1215 scope: ScanConfig::default().into(),
1216 content: crate::content::AnalysisSet::NONE,
1217 },
1218 Query::default(),
1219 std::time::UNIX_EPOCH,
1220 )
1221 };
1222 let delivery = Delivery {
1223 cache: crate::CachePolicy::Auto,
1224 cache_path: Some(cache.path().join("snapshot")),
1225 accept_partial: false,
1226 watch: Some(WatchDelivery { interval: Duration::from_secs(2) }),
1227 workers: crate::query::Workers::default(),
1228 batch_size: 4,
1229 order: crate::scan::ScanOrder::default(),
1230 };
1231
1232 let progress = crate::Progress::new();
1233 let session = Session::start_with_progress(request(), delivery.clone(), &progress)
1234 .expect("observed start");
1235 let after_start = progress.snapshot();
1236 assert_eq!(
1237 after_start.phase,
1238 crate::ProgressPhase::Revalidating,
1239 "the handoff revalidation follows the joined save"
1240 );
1241 assert!(after_start.directories >= 5, "the root and four children: {after_start:?}");
1245 assert!(after_start.files >= 12, "{after_start:?}");
1246 assert!(after_start.bytes >= bytes, "{after_start:?}");
1247 assert_eq!(after_start.analysis, None);
1248
1249 let plain = Session::start(request(), delivery).expect("plain start");
1250 let generated_at = std::time::SystemTime::now();
1251 let observed_report = session.report(generated_at).expect("observed report");
1252 let mut plain_report = plain.report(generated_at).expect("plain report");
1253 plain_report.provenance = observed_report.provenance.clone();
1254 let json = |report: &Report| {
1255 crate::report_format::render(report, crate::report_format::Format::Json, false)
1256 .expect("render")
1257 };
1258 assert_eq!(json(&plain_report), json(&observed_report));
1259 assert_eq!(progress.snapshot(), after_start, "the second start was not observed");
1260 }
1261}