1use std::collections::VecDeque;
7use std::path::{Path, PathBuf};
8use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
9use std::sync::{Arc, Condvar, Mutex, MutexGuard};
10use std::thread::{self, JoinHandle};
11
12#[cfg(test)]
13use crate::EntryKind;
14#[cfg(test)]
15use crate::Observation;
16use crate::index::{DiscoveryCommit, DiscoveryTransition};
17use crate::scan::ReconcileControl;
18use crate::{Error, Index, IndexHandle, ObservationOp, Op, Result, ScanConfig, SessionId};
19
20mod continuation;
21#[cfg(all(test, feature = "watch"))]
22mod golden_support;
23#[cfg(all(test, feature = "watch"))]
24mod golden_tests;
25mod journal;
26pub(crate) mod read;
27
28const FIRST_SESSION_ORDINAL: u64 = 1;
30const FNV_OFFSET_BASIS: u64 = 0xcbf2_9ce4_8422_2325;
32const FNV_PRIME: u64 = 0x0000_0100_0000_01b3;
34const FIRST_SESSION_ID: u64 = 1;
36pub const MAX_PRIORITY_PATHS: usize = 64;
38pub const MAX_REFRESH_PATHS: usize = 1_024;
40#[cfg(test)]
41const TEST_GATE_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(5);
43
44impl crate::SessionId {
50 fn mint() -> Result<Self> {
51 static NEXT: AtomicU64 = AtomicU64::new(FIRST_SESSION_ORDINAL);
52
53 let ordinal = NEXT
54 .fetch_update(Ordering::Relaxed, Ordering::Relaxed, |current| current.checked_add(1))
55 .map_err(|_| Error::OpenedIdentityExhausted)?;
56 let nanos = std::time::SystemTime::now()
57 .duration_since(std::time::UNIX_EPOCH)
58 .map_or(0, |elapsed| {
59 u64::try_from(elapsed.as_nanos() & u128::from(u64::MAX)).unwrap_or(0)
60 });
61 let process = u64::from(std::process::id());
62 let mut hash = FNV_OFFSET_BASIS;
63 for byte in nanos
64 .to_le_bytes()
65 .iter()
66 .chain(process.to_le_bytes().iter())
67 .chain(ordinal.to_le_bytes().iter())
68 {
69 hash ^= u64::from(*byte);
70 hash = hash.wrapping_mul(FNV_PRIME);
71 }
72 Ok(Self(hash.max(FIRST_SESSION_ID)))
73 }
74}
75
76#[derive(Clone, Copy, PartialEq, Eq, Debug, Default)]
82pub struct DiscoveryBudget {
83 pub max_files: Option<u64>,
85}
86
87#[derive(Clone, Debug)]
94pub struct OpenOptions {
95 pub batch_size: usize,
97 pub follow_symlinks: bool,
99 pub one_filesystem: bool,
101 pub hidden: Option<Arc<crate::HiddenPolicy>>,
103 pub exclude_special: bool,
105 pub types: Option<Arc<crate::classify::TypeRegistry>>,
107 pub budget: DiscoveryBudget,
109 #[cfg(feature = "watch")]
111 pub observation: Option<crate::watch::WatchConfig>,
112 #[cfg(all(feature = "watch", test))]
114 #[doc(hidden)]
115 pub observation_script: Option<PathBuf>,
116 pub journal_capacity_bytes: usize,
122 pub control_limits: crate::control::ControlLimits,
126}
127
128impl Default for OpenOptions {
129 fn default() -> Self {
130 let scan = ScanConfig::default();
131 Self {
132 batch_size: scan.batch_size,
133 follow_symlinks: scan.follow_symlinks,
134 one_filesystem: scan.one_filesystem,
135 hidden: scan.hidden,
136 exclude_special: scan.exclude_special,
137 types: scan.types,
138 budget: DiscoveryBudget::default(),
139 #[cfg(feature = "watch")]
140 observation: None,
141 #[cfg(all(feature = "watch", test))]
142 observation_script: None,
143 journal_capacity_bytes: crate::DEFAULT_JOURNAL_CAPACITY_BYTES,
144 control_limits: scan.control_limits,
145 }
146 }
147}
148
149impl OpenOptions {
150 pub fn plan(&self, root: &Path, delivery: &crate::query::Delivery) -> Result<crate::Plan> {
152 let scan = self.clone().into_parts().0;
153 let basis = crate::query::Basis {
154 root: root.into(),
155 scope: scan.into(),
156 content: crate::content::AnalysisSet::NONE,
157 };
158 let request = crate::query::Request::new(
159 basis,
160 crate::query::Query::default(),
161 std::time::SystemTime::now(),
162 );
163 crate::plan(&request, delivery, crate::Route::Opened).map_err(Error::InvalidRequest)
164 }
165
166 fn into_parts(self) -> (ScanConfig, DiscoveryBudget, usize) {
167 let scan = ScanConfig {
168 max_depth: None,
169 batch_size: self.batch_size,
170 follow_symlinks: self.follow_symlinks,
171 one_filesystem: self.one_filesystem,
172 hidden: self.hidden,
173 exclude_special: self.exclude_special,
174 threads: Some(1),
180 order: crate::ScanOrder::BreadthFirst,
181 types: self.types,
182 read_controls: OpenedIndex::basis().scope.read_controls,
185 population: crate::query::IgnoredEntries::Include,
186 control_limits: self.control_limits,
187 progress: None,
190 };
191 (scan, self.budget, self.journal_capacity_bytes)
192 }
193}
194
195#[derive(Clone)]
201pub struct OpenedIndex {
202 state: Arc<OpenedState>,
203}
204
205impl std::fmt::Debug for OpenedIndex {
206 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
207 formatter
208 .debug_struct("OpenedIndex")
209 .field("session", &self.state.session)
210 .field("root", &self.state.root)
211 .finish_non_exhaustive()
212 }
213}
214
215#[cfg(test)]
216fn open_fixture(root: &Path, options: OpenOptions) -> Result<OpenedIndex> {
217 let mut delivery = crate::query::Delivery::new(crate::CachePolicy::Off, None);
218 delivery.batch_size = options.batch_size;
219 let plan = options.plan(root, &delivery)?;
220 OpenedIndex::open(&plan, options)
221}
222
223impl OpenedIndex {
224 pub fn basis() -> crate::query::Basis {
235 crate::query::Basis {
236 root: PathBuf::new(),
237 scope: crate::query::Scope { read_controls: true, ..crate::query::Scope::default() },
238 content: crate::content::AnalysisSet::NONE,
239 }
240 }
241
242 pub fn open(plan: &crate::Plan, mut options: OpenOptions) -> Result<Self> {
248 if plan.route() != crate::Route::Opened {
249 return Err(Error::InvalidRequest(crate::query::RequestError::DeliveryUnsupported {
250 route: "opened",
251 reason: "expected an opened-root execution plan",
252 }));
253 }
254 let basis = plan.basis();
255 options.clone().into_parts().0.validate_for_scope(basis.scope.scope())?;
256 options.batch_size = plan.delivery().batch_size;
257 let root = basis.root.as_path();
258 #[cfg(test)]
259 let opened = Self::open_inner(root, options, Arc::default());
260 #[cfg(not(test))]
261 let opened = Self::open_inner(root, options);
262 opened
263 }
264
265 #[cfg(not(test))]
266 fn open_inner(root: &Path, options: OpenOptions) -> Result<Self> {
267 Self::build(root, options)
268 }
269
270 #[cfg(test)]
271 fn open_inner(root: &Path, options: OpenOptions, controls: Arc<TestControls>) -> Result<Self> {
272 Self::build(root, options, controls)
273 }
274
275 #[cfg(not(test))]
276 fn build(root: &Path, options: OpenOptions) -> Result<Self> {
277 let state = OpenedState::new(root, options)?;
278 let opened = Self { state: Arc::new(state) };
279 opened.start_discovery()?;
280 #[cfg(feature = "watch")]
281 if let Err(error) = opened.start_observation() {
282 let _ = opened.close();
283 return Err(error);
284 }
285 Ok(opened)
286 }
287
288 #[cfg(test)]
289 fn build(root: &Path, options: OpenOptions, controls: Arc<TestControls>) -> Result<Self> {
290 let state = OpenedState::new(root, options, controls)?;
291 let opened = Self { state: Arc::new(state) };
292 if !opened.state.test_controls.discovery_disabled.load(Ordering::Acquire) {
293 opened.start_discovery()?;
294 }
295 #[cfg(feature = "watch")]
296 if let Err(error) = opened.start_observation() {
297 let _ = opened.close();
298 return Err(error);
299 }
300 Ok(opened)
301 }
302
303 pub fn close(&self) -> Result<()> {
309 self.state.shutdown()
310 }
311
312 pub fn prioritize(&self, paths: &[PathBuf]) -> Result<()> {
318 self.ensure_open()?;
319 if paths.len() > MAX_PRIORITY_PATHS {
320 return Err(Error::PriorityPathLimit {
321 attempted: paths.len(),
322 limit: MAX_PRIORITY_PATHS,
323 });
324 }
325 if self.state.index.state()?.phase == crate::LifecyclePhase::Stopped {
326 return Err(Error::OpenedIndexStopped);
327 }
328 let mut normalized = Vec::with_capacity(paths.len());
329 for path in paths {
330 normalized.push(crate::scan::normalize_subtree(path)?);
331 }
332 normalized.sort();
333 normalized.dedup();
334 self.state.frontier.prioritize(normalized);
335 Ok(())
336 }
337
338 pub fn read(&self, request: crate::ReadRequest) -> Result<crate::ReadResponse> {
357 self.ensure_open()?;
358 read::read(self, request)
359 }
360
361 pub fn changes(&self, request: crate::ChangeRequest) -> Result<crate::ChangePoll> {
372 self.ensure_open()?;
373 journal::poll(self, request)
374 }
375
376 pub fn refresh(&self, paths: &[PathBuf]) -> Result<crate::RefreshResult> {
383 let _active = self.state.begin_refresh()?;
384 if paths.len() > MAX_REFRESH_PATHS {
385 return Err(Error::RefreshPathLimit {
386 attempted: paths.len(),
387 limit: MAX_REFRESH_PATHS,
388 });
389 }
390 let control = OpenedReconcileControl { state: &self.state };
391 let (after, initial_state) = self.version_and_state()?;
392 let forbid_expansion = initial_state.phase == crate::LifecyclePhase::Stopped
396 && initial_state.coverage == crate::Coverage::Partial(crate::CoverageReason::Budget);
397 let report = crate::scan::reconcile_paths_handle_controlled(
398 &self.state.index,
399 paths,
400 &self.state.scan,
401 forbid_expansion,
402 &control,
403 &mut |_commit| self.state.journal.notify_commit(),
404 )?;
405 control.check_active()?;
406 let (version, state, impact) = self.state.index.read_with(|index| {
407 let scope = index.scope();
408 let since = index.since(after.sequence);
409 let version = crate::EngineVersion {
410 session: self.state.session,
411 sequence: since.clock,
412 scope: scope.entry_scope(),
413 semantics: scope.semantic_identity(),
414 };
415 let impact = journal::interval_impact(&since);
416 (version, since.state, impact)
417 })?;
418 let mut issues = Vec::new();
419 let mut omitted_issues = 0_u64;
420 for error in &report.reconciliation.scan.errors {
421 let issue = crate::Issue::from_error_under(&self.state.root, error);
422 if issues.len() < crate::MAX_RETAINED_ISSUES {
423 issues.push(issue);
424 } else {
425 omitted_issues = omitted_issues.saturating_add(1);
426 }
427 }
428 if report.reconciliation.apply.resource_refused > 0 {
429 let issue = crate::Issue::resource_budget(
430 self.state
431 .budget
432 .max_files
433 .expect("resource refusal requires a configured file limit"),
434 );
435 if issues.len() < crate::MAX_RETAINED_ISSUES {
436 issues.push(issue);
437 } else {
438 omitted_issues = omitted_issues.saturating_add(1);
439 }
440 }
441 if report.reconciliation.retry_required() {
445 let issue = crate::Issue::provider_failure(
446 None,
447 "reconciliation interrupted by newer verification; retry this refresh".to_string(),
448 );
449 if issues.len() < crate::MAX_RETAINED_ISSUES {
450 issues.push(issue);
451 } else {
452 omitted_issues = omitted_issues.saturating_add(1);
453 }
454 }
455 let work = crate::Work {
456 observations: report.reconciliation.observations,
457 unchanged: report.reconciliation.apply.unchanged,
458 stale: report.reconciliation.apply.stale,
459 resource_refused: report.reconciliation.apply.resource_refused,
460 directories_read: report.reconciliation.scan.dirs_read,
461 entries_visited: report.reconciliation.scan.entries,
462 files_visited: report.reconciliation.scan.files_walked,
463 bytes_visited: report.reconciliation.scan.bytes_walked,
464 ..crate::Work::default()
465 };
466 Ok(crate::RefreshResult {
467 after,
468 version,
469 state,
470 accepted: report.accepted,
471 rejected: report.rejected,
472 impact,
473 work,
474 issues,
475 omitted_issues,
476 })
477 }
478
479 fn version_and_state(&self) -> Result<(crate::EngineVersion, crate::IndexState)> {
480 self.state.index.read_with(|index| {
481 let scope = index.scope();
482 (
483 crate::EngineVersion {
484 session: self.state.session,
485 sequence: index.clock(),
486 scope: scope.entry_scope(),
487 semantics: scope.semantic_identity(),
488 },
489 index.state(),
490 )
491 })
492 }
493
494 fn start_discovery(&self) -> Result<()> {
495 publish_discovery_transition(
496 &self.state.index,
497 &self.state.journal,
498 DiscoveryTransition::Begin,
499 )?;
500 let root = self.state.root.clone();
501 let index = self.state.index.clone();
502 let journal = Arc::clone(&self.state.journal);
503 let scan = self.state.scan.clone();
504 let budget = self.state.budget;
505 let frontier = Arc::clone(&self.state.frontier);
506 #[cfg(feature = "watch")]
507 let baseline = Arc::clone(&self.state.baseline);
508 #[cfg(test)]
509 let controls = Arc::clone(&self.state.test_controls);
510 self.spawn_worker("discovery", move |cancellation| {
511 #[cfg(feature = "watch")]
512 let _baseline_finished = BaselineCompletion(baseline);
513 #[cfg(test)]
514 controls.reach(TestPoint::BeforeDiscovery);
515 #[cfg(not(test))]
516 let outcome =
517 { run_discovery(&root, &index, &journal, &scan, budget, &frontier, &cancellation) };
518 #[cfg(test)]
519 let outcome = {
520 run_discovery(
521 &root,
522 &index,
523 &journal,
524 &scan,
525 budget,
526 &frontier,
527 &cancellation,
528 &controls,
529 )
530 };
531 if let Err(error) = outcome {
532 let state = index.state()?;
533 if state.phase == crate::LifecyclePhase::Discovering {
534 publish_discovery_transition(
535 &index,
536 &journal,
537 DiscoveryTransition::Failed(crate::Issue::from_error_under(&root, &error)),
538 )?;
539 }
540 return Err(error);
541 }
542 Ok(())
543 })
544 }
545
546 #[cfg(feature = "watch")]
547 fn start_observation(&self) -> Result<()> {
548 let watcher =
549 self.state.observer.lock().map_err(|_| Error::OpenedLifecyclePoisoned)?.take();
550 let Some(watcher) = watcher else {
551 return Ok(());
552 };
553 let root = self.state.root.clone();
554 let index = self.state.index.clone();
555 let journal = Arc::clone(&self.state.journal);
556 let scan = self.state.scan.clone();
557 let budget = self.state.budget;
558 let baseline = Arc::clone(&self.state.baseline);
559 #[cfg(test)]
560 let controls = Arc::clone(&self.state.test_controls);
561 self.spawn_worker("observation", move |cancellation| {
562 let outcome = run_observation(
563 watcher,
564 &root,
565 &index,
566 &journal,
567 &scan,
568 budget,
569 &baseline,
570 &cancellation,
571 #[cfg(test)]
572 &controls,
573 );
574 if matches!(outcome, Err(Error::OpenedIndexClosed)) && cancellation.is_cancelled() {
575 return Ok(());
576 }
577 if let Err(error) = &outcome {
578 if !cancellation.is_cancelled() {
579 publish_observation_transition(
580 &index,
581 &journal,
582 crate::index::ObservationTransition::Failed(
583 crate::Issue::from_error_under(&root, error),
584 ),
585 )?;
586 }
587 }
588 outcome
589 })
590 }
591
592 #[allow(dead_code)]
593 fn ensure_open(&self) -> Result<()> {
594 self.state.ensure_open()
595 }
596
597 #[allow(dead_code)]
603 fn spawn_worker<F>(&self, name: &'static str, run: F) -> Result<()>
604 where
605 F: FnOnce(Arc<Cancellation>) -> Result<()> + Send + 'static,
606 {
607 let locked = self.state.lock_lifecycle();
608 if locked.poisoned {
609 return Err(Error::OpenedLifecyclePoisoned);
610 }
611 let mut lifecycle = locked.guard;
612 if lifecycle.phase != OwnerPhase::Open {
613 return Err(Error::OpenedIndexClosed);
614 }
615
616 let cancellation = Arc::clone(&self.state.cancellation);
617 let journal = Arc::clone(&self.state.journal);
618 let failures = Arc::clone(&self.state.failures);
619 #[cfg(test)]
620 let controls = Arc::clone(&self.state.test_controls);
621 let worker = thread::Builder::new()
622 .name(format!("fdu-{name}"))
623 .spawn(move || {
624 let outcome =
631 std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| run(cancellation)));
632 #[cfg(test)]
633 {
634 if name != "discovery" {
635 controls.reach(TestPoint::BeforeWorkerExit);
636 }
637 }
638 match outcome {
639 Ok(Ok(())) => {}
640 Ok(Err(error)) => {
641 failures.record(CloseOutcome::WorkerFailed {
642 worker: name,
643 source: Arc::new(error),
644 });
645 journal.wake();
646 }
647 Err(payload) => {
648 failures.record(CloseOutcome::WorkerPanicked { worker: name });
649 journal.wake();
650 std::panic::resume_unwind(payload);
651 }
652 }
653 })
654 .map_err(|source| Error::OpenedWorkerSpawn { worker: name, source })?;
655 lifecycle.workers.push(Worker { name, handle: worker });
656 Ok(())
657 }
658
659 #[cfg(test)]
660 fn open_for_test(
661 root: &Path,
662 options: OpenOptions,
663 controls: Arc<TestControls>,
664 ) -> Result<Self> {
665 Self::open_inner(root, options, controls)
666 }
667}
668
669struct OpenedState {
670 session: SessionId,
671 root: std::path::PathBuf,
672 index: IndexHandle,
673 #[allow(dead_code)]
676 scan: ScanConfig,
677 budget: DiscoveryBudget,
678 frontier: Arc<DiscoveryFrontier>,
679 continuations: Mutex<continuation::ContinuationTable>,
680 journal: Arc<journal::JournalWait>,
681 failures: Arc<WorkerFailures>,
683 cancellation: Arc<Cancellation>,
684 #[cfg(feature = "watch")]
685 baseline: Arc<BaselineLatch>,
686 #[cfg(feature = "watch")]
687 observer: Mutex<Option<crate::watch::Watcher>>,
688 lifecycle: Mutex<Lifecycle>,
689 lifecycle_changed: Condvar,
690 #[cfg(test)]
691 test_controls: Arc<TestControls>,
692}
693
694impl OpenedState {
695 #[cfg(not(test))]
696 fn new(root: &Path, options: OpenOptions) -> Result<Self> {
697 Self::build(root, options)
698 }
699
700 #[cfg(test)]
701 fn new(root: &Path, options: OpenOptions, controls: Arc<TestControls>) -> Result<Self> {
702 Self::build(root, options, controls)
703 }
704
705 #[cfg(not(test))]
706 fn build(root: &Path, options: OpenOptions) -> Result<Self> {
707 #[cfg(feature = "watch")]
708 let observation = options.observation;
709 let (root, index, scan, budget) = bind_root(root, options)?;
710 #[cfg(feature = "watch")]
711 let observer = if let Some(config) = observation {
712 scan.validate_for_watch_scope(index.scope()?)?;
713 Some(crate::watch::Watcher::new(&root, config)?)
714 } else {
715 None
716 };
717 Ok(Self {
718 session: SessionId::mint()?,
719 root,
720 index,
721 scan,
722 budget,
723 frontier: Arc::new(DiscoveryFrontier::new()),
724 continuations: Mutex::new(continuation::ContinuationTable::default()),
725 journal: Arc::new(journal::JournalWait::new()),
726 failures: Arc::new(WorkerFailures::default()),
727 cancellation: Arc::new(Cancellation::default()),
728 #[cfg(feature = "watch")]
729 baseline: Arc::new(BaselineLatch::default()),
730 #[cfg(feature = "watch")]
731 observer: Mutex::new(observer),
732 lifecycle: Mutex::new(Lifecycle::default()),
733 lifecycle_changed: Condvar::new(),
734 })
735 }
736
737 #[cfg(test)]
738 fn build(root: &Path, options: OpenOptions, controls: Arc<TestControls>) -> Result<Self> {
739 #[cfg(feature = "watch")]
740 let observation = options.observation;
741 #[cfg(feature = "watch")]
742 let observation_script = options.observation_script.clone();
743 let (root, index, scan, budget) = bind_root(root, options)?;
744 #[cfg(feature = "watch")]
745 if observation.is_some() {
746 scan.validate_for_watch_scope(index.scope()?)?;
747 }
748 #[cfg(feature = "watch")]
749 let (observer, scripted_sender) = match (observation, observation_script) {
750 (Some(config), Some(events)) => {
751 let (watcher, sender) = crate::watch::Watcher::scripted(&root, config, &events)?;
752 (Some(watcher), Some(sender))
753 }
754 (Some(config), None) => (Some(crate::watch::Watcher::new(&root, config)?), None),
755 (None, _) => (None, None),
756 };
757 #[cfg(feature = "watch")]
758 {
759 *controls.scripted_observer.lock().unwrap_or_else(std::sync::PoisonError::into_inner) =
760 scripted_sender;
761 }
762 Ok(Self {
763 session: SessionId::mint()?,
764 root,
765 index,
766 scan,
767 budget,
768 frontier: Arc::new(DiscoveryFrontier::new()),
769 continuations: Mutex::new(continuation::ContinuationTable::default()),
770 journal: Arc::new(journal::JournalWait::new()),
771 failures: Arc::new(WorkerFailures::default()),
772 cancellation: Arc::new(Cancellation::default()),
773 #[cfg(feature = "watch")]
774 baseline: Arc::new(BaselineLatch::default()),
775 #[cfg(feature = "watch")]
776 observer: Mutex::new(observer),
777 lifecycle: Mutex::new(Lifecycle::default()),
778 lifecycle_changed: Condvar::new(),
779 test_controls: controls,
780 })
781 }
782
783 #[allow(dead_code)]
784 fn ensure_open(&self) -> Result<()> {
785 let locked = self.lock_lifecycle();
786 if locked.poisoned {
787 return Err(Error::OpenedLifecyclePoisoned);
788 }
789 if locked.guard.phase != OwnerPhase::Open {
790 return Err(Error::OpenedIndexClosed);
791 }
792 Ok(())
793 }
794
795 fn lock_lifecycle(&self) -> LockedLifecycle<'_> {
796 match self.lifecycle.lock() {
797 Ok(guard) => LockedLifecycle { guard, poisoned: false },
798 Err(poisoned) => LockedLifecycle { guard: poisoned.into_inner(), poisoned: true },
799 }
800 }
801
802 fn begin_refresh(&self) -> Result<ActiveRefresh<'_>> {
803 let locked = self.lock_lifecycle();
804 if locked.poisoned {
805 return Err(Error::OpenedLifecyclePoisoned);
806 }
807 let mut lifecycle = locked.guard;
808 if lifecycle.phase != OwnerPhase::Open {
809 return Err(Error::OpenedIndexClosed);
810 }
811 lifecycle.active_refreshes = lifecycle.active_refreshes.saturating_add(1);
812 Ok(ActiveRefresh { state: self })
813 }
814
815 fn shutdown(&self) -> Result<()> {
816 let mut saw_poison = false;
817 let workers = loop {
818 let locked = self.lock_lifecycle();
819 saw_poison |= locked.poisoned;
820 let mut lifecycle = locked.guard;
821 match lifecycle.phase {
822 OwnerPhase::Open => {
823 lifecycle.phase = OwnerPhase::Closing;
824 let workers = std::mem::take(&mut lifecycle.workers);
825 self.cancellation.cancel();
826 #[cfg(feature = "watch")]
827 self.baseline.wake();
828 drop(lifecycle);
829 self.journal.close();
830 match self.continuations.lock() {
831 Ok(mut continuations) => continuations.close(),
832 Err(poisoned) => poisoned.into_inner().close(),
833 }
834 self.lifecycle_changed.notify_all();
835 break workers;
836 }
837 OwnerPhase::Closing => {
838 #[cfg(test)]
839 self.test_controls.reach(TestPoint::BeforeCloseWait);
840 let waited = self.lifecycle_changed.wait(lifecycle);
841 match waited {
842 Ok(_) => {}
843 Err(poisoned) => {
844 saw_poison = true;
845 drop(poisoned.into_inner());
846 }
847 }
848 }
849 OwnerPhase::Closed => {
850 return lifecycle
851 .terminal
852 .as_ref()
853 .expect("closed lifecycle stores one outcome")
854 .to_result();
855 }
856 }
857 };
858
859 let worker_outcome = join_workers(workers, &self.failures);
860 let locked = self.lock_lifecycle();
861 saw_poison |= locked.poisoned;
862 let mut lifecycle = locked.guard;
863 while lifecycle.active_refreshes > 0 {
864 match self.lifecycle_changed.wait(lifecycle) {
865 Ok(next) => lifecycle = next,
866 Err(poisoned) => {
867 saw_poison = true;
868 lifecycle = poisoned.into_inner();
869 }
870 }
871 }
872 drop(lifecycle);
873 let index_poisoned = self.index.clock().is_err();
874 let mut outcome = if saw_poison {
875 CloseOutcome::LifecyclePoisoned
876 } else if let Some(outcome) = worker_outcome {
877 outcome
878 } else if index_poisoned {
879 CloseOutcome::IndexPoisoned
880 } else {
881 CloseOutcome::Success
882 };
883
884 let locked = self.lock_lifecycle();
885 if locked.poisoned {
886 outcome = CloseOutcome::LifecyclePoisoned;
887 }
888 let mut lifecycle = locked.guard;
889 lifecycle.phase = OwnerPhase::Closed;
890 lifecycle.terminal = Some(outcome.clone());
891 self.lifecycle_changed.notify_all();
892 outcome.to_result()
893 }
894}
895
896impl Drop for OpenedState {
897 fn drop(&mut self) {
898 let _ = self.shutdown();
899 }
900}
901
902#[cfg(all(test, feature = "watch"))]
908pub(super) struct RetainedOwnership {
909 session: SessionId,
910 joined: bool,
912 workers: usize,
913 waiters: usize,
914 continuations: usize,
915 close: Option<Result<()>>,
917}
918
919#[cfg(all(test, feature = "watch"))]
920impl RetainedOwnership {
921 pub(super) fn is_released(&self) -> bool {
922 self.joined && self.workers == 0 && self.waiters == 0 && self.continuations == 0
923 }
924}
925
926#[cfg(all(test, feature = "watch"))]
927impl std::fmt::Display for RetainedOwnership {
928 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
929 write!(
930 formatter,
931 "session={:?} joined={} workers={} waiters={} continuations={} close=",
932 self.session, self.joined, self.workers, self.waiters, self.continuations
933 )?;
934 match &self.close {
935 Some(outcome) => write!(formatter, "{outcome:?}"),
936 None => formatter.write_str("none"),
937 }
938 }
939}
940
941#[cfg(all(test, feature = "watch"))]
942impl OpenedState {
943 pub(super) fn retained_ownership(&self) -> RetainedOwnership {
944 let (joined, workers, close) = {
945 let lifecycle = self.lock_lifecycle().guard;
946 (
947 lifecycle.phase == OwnerPhase::Closed,
948 lifecycle.workers.len(),
949 lifecycle.terminal.as_ref().map(CloseOutcome::to_result),
950 )
951 };
952 let continuations =
953 self.continuations.lock().unwrap_or_else(std::sync::PoisonError::into_inner).len();
954 RetainedOwnership {
955 session: self.session,
956 joined,
957 workers,
958 waiters: self.journal.waiters(),
959 continuations,
960 close,
961 }
962 }
963}
964
965fn bind_root(
966 root: &Path,
967 options: OpenOptions,
968) -> Result<(std::path::PathBuf, IndexHandle, ScanConfig, DiscoveryBudget)> {
969 let (scan, budget, journal_capacity_bytes) = options.into_parts();
970 scan.validate()?;
971 if budget.max_files == Some(0) {
972 return Err(Error::UnsupportedScanConfig(
973 "max_files must be nonzero; omit it for an unlimited discovery",
974 ));
975 }
976 if journal_capacity_bytes < crate::MIN_JOURNAL_CAPACITY_BYTES {
977 return Err(Error::JournalCapacityTooSmall {
978 requested: journal_capacity_bytes,
979 minimum: crate::MIN_JOURNAL_CAPACITY_BYTES,
980 });
981 }
982 let root = root.canonicalize().map_err(|source| Error::io(root, source))?;
983 let metadata = std::fs::symlink_metadata(&root).map_err(|source| Error::io(&root, source))?;
984 if !metadata.is_dir() {
985 return Err(Error::io(
986 &root,
987 std::io::Error::new(
988 std::io::ErrorKind::NotADirectory,
989 "opened-index root is not a directory",
990 ),
991 ));
992 }
993
994 let scope = scan.scope();
995 let types = scan.types_shared();
996 let mut index = Index::new_opened_with_scope_types_and_journal_capacity_bytes(
997 &root,
998 scope,
999 types,
1000 journal_capacity_bytes,
1001 );
1002 index.set_control_limits(scan.control_limits);
1003 let index = IndexHandle::new(index);
1004 Ok((root, index, scan, budget))
1005}
1006
1007#[derive(Clone, Debug)]
1008struct PendingDirectory {
1009 path: PathBuf,
1010 depth: usize,
1011}
1012
1013struct DiscoveryFrontier {
1014 state: Mutex<FrontierState>,
1015}
1016
1017struct FrontierState {
1026 unprioritized: VecDeque<QueuedDirectory>,
1028 serving: Vec<VecDeque<QueuedDirectory>>,
1031 priorities: Vec<PathBuf>,
1032 next_sequence: u64,
1034 stopped: bool,
1035}
1036
1037struct QueuedDirectory {
1038 sequence: u64,
1039 directory: PendingDirectory,
1040}
1041
1042impl FrontierState {
1043 fn enqueue(&mut self, queued: QueuedDirectory) {
1044 let path = &queued.directory.path;
1045 match self
1046 .priorities
1047 .iter()
1048 .position(|priority| priority.starts_with(path) || path.starts_with(priority))
1049 {
1050 Some(priority) => self.serving[priority].push_back(queued),
1051 None => self.unprioritized.push_back(queued),
1052 }
1053 }
1054}
1055
1056impl DiscoveryFrontier {
1057 fn new() -> Self {
1058 Self {
1059 state: Mutex::new(FrontierState {
1060 unprioritized: VecDeque::from([QueuedDirectory {
1061 sequence: 0,
1062 directory: PendingDirectory { path: PathBuf::new(), depth: 0 },
1063 }]),
1064 serving: Vec::new(),
1065 priorities: Vec::new(),
1066 next_sequence: 1,
1067 stopped: false,
1068 }),
1069 }
1070 }
1071
1072 fn pop(&self) -> Option<PendingDirectory> {
1073 let mut guard = self.state.lock().unwrap_or_else(std::sync::PoisonError::into_inner);
1074 let state = &mut *guard;
1075 if state.stopped {
1076 return None;
1077 }
1078 if let Some(queue) = state.serving.iter_mut().find(|queue| !queue.is_empty()) {
1079 return queue.pop_front().map(|queued| queued.directory);
1080 }
1081 state.unprioritized.pop_front().map(|queued| queued.directory)
1082 }
1083
1084 fn extend(&self, directories: impl IntoIterator<Item = PendingDirectory>) {
1085 let mut state = self.state.lock().unwrap_or_else(std::sync::PoisonError::into_inner);
1086 if state.stopped {
1087 return;
1088 }
1089 for directory in directories {
1090 let sequence = state.next_sequence;
1091 state.next_sequence = sequence.wrapping_add(1);
1092 state.enqueue(QueuedDirectory { sequence, directory });
1093 }
1094 }
1095
1096 fn prioritize(&self, priorities: Vec<PathBuf>) {
1097 let mut guard = self.state.lock().unwrap_or_else(std::sync::PoisonError::into_inner);
1098 let state = &mut *guard;
1099 if state.stopped {
1100 return;
1101 }
1102 let mut queued: Vec<_> =
1103 state.unprioritized.drain(..).chain(state.serving.drain(..).flatten()).collect();
1104 queued.sort_unstable_by_key(|queued| queued.sequence);
1105 state.serving = priorities.iter().map(|_| VecDeque::new()).collect();
1106 state.priorities = priorities;
1107 for directory in queued {
1108 state.enqueue(directory);
1109 }
1110 }
1111
1112 fn stop(&self) {
1113 let mut state = self.state.lock().unwrap_or_else(std::sync::PoisonError::into_inner);
1114 state.stopped = true;
1115 state.unprioritized.clear();
1116 state.serving.clear();
1117 state.priorities.clear();
1118 }
1119}
1120
1121#[allow(clippy::too_many_arguments)]
1122fn run_discovery(
1123 root: &Path,
1124 index: &IndexHandle,
1125 journal: &journal::JournalWait,
1126 scan: &ScanConfig,
1127 budget: DiscoveryBudget,
1128 frontier: &DiscoveryFrontier,
1129 cancellation: &Cancellation,
1130 #[cfg(test)] controls: &TestControls,
1131) -> Result<()> {
1132 let root_metadata =
1133 std::fs::symlink_metadata(root).map_err(|source| Error::io(root, source))?;
1134 let root_dev =
1135 crate::scan::root_device(root, &root_metadata).map_err(|source| Error::io(root, source))?;
1136 let mut readers = crate::scan::Readers::new();
1137
1138 while let Some(directory) = frontier.pop() {
1139 if cancellation.is_cancelled() {
1140 frontier.stop();
1141 publish_discovery_transition(index, journal, DiscoveryTransition::Cancelled)?;
1142 return Ok(());
1143 }
1144 match discover_directory(
1145 root,
1146 index,
1147 journal,
1148 scan,
1149 budget,
1150 root_dev,
1151 &directory,
1152 frontier,
1153 cancellation,
1154 &mut readers,
1155 #[cfg(test)]
1156 controls.deterministic_discovery_order.load(Ordering::Acquire),
1157 )? {
1158 DiscoveryStep::Continue => {}
1159 DiscoveryStep::Stopped => return Ok(()),
1160 }
1161 #[cfg(test)]
1162 if directory.path.as_os_str().is_empty() {
1163 controls.reach(TestPoint::AfterRootDirectory);
1164 }
1165 }
1166 publish_discovery_transition(index, journal, DiscoveryTransition::Finish)?;
1167 Ok(())
1168}
1169
1170#[derive(Clone, Copy, PartialEq, Eq)]
1171enum DiscoveryStep {
1172 Continue,
1173 Stopped,
1174}
1175
1176#[derive(Clone, Copy, Debug)]
1183enum DiscoveryAnswer {
1184 Accepted,
1186 Stopped,
1189 Stale,
1194}
1195
1196fn discovery_rejection(error: Error) -> Result<DiscoveryAnswer> {
1201 match error {
1202 Error::OpenedIndexStopped => Ok(DiscoveryAnswer::Stopped),
1203 Error::InvalidDirectoryCompletion(_) | Error::UnknownAncestry { .. } => {
1204 Ok(DiscoveryAnswer::Stale)
1205 }
1206 error => Err(error),
1207 }
1208}
1209
1210fn abandon_directory(frontier: &DiscoveryFrontier, answer: DiscoveryAnswer) -> DiscoveryStep {
1220 match answer {
1221 DiscoveryAnswer::Accepted | DiscoveryAnswer::Stale => DiscoveryStep::Continue,
1222 DiscoveryAnswer::Stopped => {
1223 frontier.stop();
1224 DiscoveryStep::Stopped
1225 }
1226 }
1227}
1228
1229#[allow(clippy::too_many_arguments)]
1230fn discover_directory(
1231 root: &Path,
1232 index: &IndexHandle,
1233 journal: &journal::JournalWait,
1234 scan: &ScanConfig,
1235 budget: DiscoveryBudget,
1236 root_dev: u64,
1237 directory: &PendingDirectory,
1238 frontier: &DiscoveryFrontier,
1239 cancellation: &Cancellation,
1240 readers: &mut crate::scan::Readers,
1241 #[cfg(test)] deterministic_discovery_order: bool,
1242) -> Result<DiscoveryStep> {
1243 let absolute = root.join(&directory.path);
1244 let policy = crate::scan::ListingPolicy::every_child_stated(scan.one_filesystem);
1245 let listing = match crate::scan::list_directory(readers, &absolute, policy, None) {
1246 Ok(listing) => listing,
1247 Err(source)
1253 if !directory.path.as_os_str().is_empty()
1254 && matches!(
1255 source.kind(),
1256 std::io::ErrorKind::NotFound | std::io::ErrorKind::NotADirectory
1257 ) =>
1258 {
1259 return Ok(DiscoveryStep::Continue);
1260 }
1261 Err(source) => {
1262 let error = Error::io(&absolute, source);
1263 publish_discovery_transition(
1264 index,
1265 journal,
1266 DiscoveryTransition::Inaccessible {
1267 issues: vec![crate::Issue::from_error_under(root, &error)],
1268 omitted: 0,
1269 },
1270 )?;
1271 return Ok(DiscoveryStep::Continue);
1272 }
1273 };
1274 #[cfg(test)]
1275 let listing = test_directory_listing(listing, deterministic_discovery_order);
1276 let mut batch = Vec::with_capacity(scan.batch_size);
1277 let mut discovered = Vec::new();
1279 let mut issues = Vec::new();
1280 let mut omitted_issues = 0_u64;
1281
1282 for item in listing {
1283 if cancellation.is_cancelled() {
1284 let answer = commit_discovery_batch(
1285 index,
1286 journal,
1287 &mut batch,
1288 None,
1289 Some(DiscoveryTransition::Cancelled),
1290 budget.max_files,
1291 )?;
1292 if matches!(answer, DiscoveryAnswer::Stale) {
1293 publish_discovery_transition(index, journal, DiscoveryTransition::Cancelled)?;
1295 }
1296 frontier.stop();
1297 return Ok(DiscoveryStep::Stopped);
1298 }
1299 let (name, observed) = match item {
1300 crate::scan::Listed::Child { name, observed } => (name, observed),
1301 crate::scan::Listed::Failed(source) => {
1302 retain_local_issue(
1303 &mut issues,
1304 &mut omitted_issues,
1305 crate::Issue::from_io_under(root, &absolute, &source),
1306 );
1307 continue;
1308 }
1309 };
1310 let (kind, attrs) = match observed {
1311 Ok(Some(observed)) => observed,
1312 Ok(None) => continue,
1313 Err(source) => {
1314 retain_local_issue(
1315 &mut issues,
1316 &mut omitted_issues,
1317 crate::Issue::from_io_under(root, &absolute.join(&name), &source),
1318 );
1319 continue;
1320 }
1321 };
1322 let Some(prepared) = crate::scan::prepare_walk_entry(
1323 root,
1324 &directory.path,
1325 directory.depth,
1326 &name,
1327 kind,
1328 attrs,
1329 root_dev,
1330 scan,
1331 ) else {
1332 continue;
1333 };
1334 let crate::scan::PreparedWalkEntry {
1335 path,
1336 kind,
1337 attrs,
1338 retained,
1339 control,
1340 descend,
1341 control_error,
1342 } = prepared;
1343 if let Some(error) = control_error {
1344 retain_local_issue(
1345 &mut issues,
1346 &mut omitted_issues,
1347 crate::Issue::from_error_under(root, &error),
1348 );
1349 }
1350 if !retained {
1351 if let Some(control) = control {
1352 let answer = push_discovery_op(
1353 index,
1354 journal,
1355 scan.batch_size,
1356 &mut batch,
1357 control,
1358 budget.max_files,
1359 )?;
1360 if !matches!(answer, DiscoveryAnswer::Accepted) {
1361 return Ok(abandon_directory(frontier, answer));
1362 }
1363 }
1364 continue;
1365 }
1366
1367 if descend {
1370 discovered.push(PendingDirectory {
1371 path: path.clone(),
1372 depth: directory.depth.saturating_add(1),
1373 });
1374 }
1375 let answer = push_discovery_op(
1376 index,
1377 journal,
1378 scan.batch_size,
1379 &mut batch,
1380 Op::Upsert { path, kind, attrs },
1381 budget.max_files,
1382 )?;
1383 if !matches!(answer, DiscoveryAnswer::Accepted) {
1384 return Ok(abandon_directory(frontier, answer));
1385 }
1386 if let Some(control) = control {
1387 let answer = push_discovery_op(
1388 index,
1389 journal,
1390 scan.batch_size,
1391 &mut batch,
1392 control,
1393 budget.max_files,
1394 )?;
1395 if !matches!(answer, DiscoveryAnswer::Accepted) {
1396 return Ok(abandon_directory(frontier, answer));
1397 }
1398 }
1399 }
1400
1401 let incomplete = !issues.is_empty() || omitted_issues > 0;
1402 let transition = incomplete.then(|| DiscoveryTransition::Inaccessible {
1403 issues: issues.clone(),
1404 omitted: omitted_issues,
1405 });
1406 let complete = (!incomplete).then(|| directory.path.clone());
1407 let answer =
1408 commit_discovery_batch(index, journal, &mut batch, complete, transition, budget.max_files)?;
1409 if !matches!(answer, DiscoveryAnswer::Accepted) {
1410 return Ok(abandon_directory(frontier, answer));
1411 }
1412 frontier.extend(discovered);
1413 Ok(DiscoveryStep::Continue)
1414}
1415
1416#[cfg(test)]
1417fn test_directory_listing<'r>(
1418 listing: crate::scan::Listing<'r>,
1419 deterministic: bool,
1420) -> Box<dyn Iterator<Item = crate::scan::Listed<'r>> + 'r> {
1421 use crate::scan::Listed;
1422 if !deterministic {
1423 return Box::new(listing);
1424 }
1425
1426 let mut entries = listing.collect::<Vec<_>>();
1429 entries.sort_by(|left, right| match (left, right) {
1430 (Listed::Child { name: left, .. }, Listed::Child { name: right, .. }) => left.cmp(right),
1431 (Listed::Child { .. }, Listed::Failed(_)) => std::cmp::Ordering::Less,
1432 (Listed::Failed(_), Listed::Child { .. }) => std::cmp::Ordering::Greater,
1433 (Listed::Failed(_), Listed::Failed(_)) => std::cmp::Ordering::Equal,
1434 });
1435 Box::new(entries.into_iter())
1436}
1437
1438fn retain_local_issue(issues: &mut Vec<crate::Issue>, omitted: &mut u64, issue: crate::Issue) {
1439 if issues.len() < crate::MAX_RETAINED_ISSUES {
1440 issues.push(issue);
1441 } else {
1442 *omitted = omitted.saturating_add(1);
1443 }
1444}
1445
1446fn push_discovery_op(
1447 index: &IndexHandle,
1448 journal: &journal::JournalWait,
1449 batch_size: usize,
1450 batch: &mut Vec<ObservationOp>,
1451 op: Op,
1452 max_files: Option<u64>,
1453) -> Result<DiscoveryAnswer> {
1454 batch.push(ObservationOp::unconditional(op));
1455 if batch.len() >= batch_size {
1456 return commit_discovery_batch(index, journal, batch, None, None, max_files);
1457 }
1458 Ok(DiscoveryAnswer::Accepted)
1459}
1460
1461fn commit_discovery_batch(
1463 index: &IndexHandle,
1464 journal: &journal::JournalWait,
1465 batch: &mut Vec<ObservationOp>,
1466 directory_complete: Option<PathBuf>,
1467 transition: Option<DiscoveryTransition>,
1468 max_files: Option<u64>,
1469) -> Result<DiscoveryAnswer> {
1470 let scanner_batch = crate::scan::ScannerBatch::new(std::mem::take(batch));
1471 let outcome = match index.apply_scanner_discovery_bounded(
1472 scanner_batch,
1473 DiscoveryCommit { directory_complete, transition },
1474 max_files,
1475 ) {
1476 Ok(outcome) => outcome,
1477 Err(error) => return discovery_rejection(error),
1478 };
1479 if outcome.commit.is_some() {
1480 journal.notify_commit();
1481 }
1482 Ok(if outcome.stats.resource_refused > 0 {
1483 DiscoveryAnswer::Stopped
1484 } else {
1485 DiscoveryAnswer::Accepted
1486 })
1487}
1488
1489fn publish_discovery_transition(
1495 index: &IndexHandle,
1496 journal: &journal::JournalWait,
1497 transition: DiscoveryTransition,
1498) -> Result<()> {
1499 let outcome = match index.transition_discovery(transition) {
1500 Ok(outcome) => outcome,
1501 Err(Error::OpenedIndexStopped) => return Ok(()),
1502 Err(error) => return Err(error),
1503 };
1504 if outcome.commit.is_some() {
1505 journal.notify_commit();
1506 }
1507 Ok(())
1508}
1509
1510#[cfg(feature = "watch")]
1511const OBSERVATION_POLL_INTERVAL: std::time::Duration = std::time::Duration::from_millis(100);
1513#[cfg(feature = "watch")]
1514const MAX_HANDOFF_RECONCILIATION_ATTEMPTS: usize = 3;
1516
1517#[cfg(feature = "watch")]
1518#[allow(clippy::too_many_arguments)]
1519#[allow(clippy::needless_pass_by_value)] fn run_observation(
1521 watcher: crate::watch::Watcher,
1522 root: &Path,
1523 index: &IndexHandle,
1524 journal: &journal::JournalWait,
1525 scan: &ScanConfig,
1526 budget: DiscoveryBudget,
1527 baseline: &BaselineLatch,
1528 cancellation: &Cancellation,
1529 #[cfg(test)] controls: &TestControls,
1530) -> Result<()> {
1531 if !baseline.wait(cancellation) {
1532 return Ok(());
1533 }
1534 let state = index.state()?;
1535 if state.phase != crate::LifecyclePhase::Ready {
1536 return Ok(());
1537 }
1538
1539 publish_observation_transition(
1540 index,
1541 journal,
1542 crate::index::ObservationTransition::Reconciling,
1543 )?;
1544 let state = index.state()?;
1545 if state.phase == crate::LifecyclePhase::Stopped {
1546 return Ok(());
1547 }
1548 if state.phase != crate::LifecyclePhase::Reconciling {
1549 return Err(Error::ObservationHandoffIncomplete);
1550 }
1551 #[cfg(test)]
1552 controls.reach(TestPoint::BeforeObservationHandoff);
1553 let control = ObservationReconcileControl {
1554 cancellation,
1555 max_files: budget.max_files,
1556 #[cfg(test)]
1557 controls,
1558 };
1559 let handoff_intents = watcher.capture_backlog_bound();
1560 let mut handoff_evidence = HandoffEvidence::default();
1561
1562 watcher.flush_capture()?;
1563 let _ = drain_observation_hints(
1564 &watcher,
1565 root,
1566 index,
1567 journal,
1568 scan,
1569 &control,
1570 handoff_intents,
1571 &mut handoff_evidence,
1572 )?;
1573 if index.state()?.phase == crate::LifecyclePhase::Stopped {
1574 return Ok(());
1575 }
1576 watcher.flush_capture()?;
1577 let mut handoff_complete = false;
1578 for _ in 0..MAX_HANDOFF_RECONCILIATION_ATTEMPTS {
1579 let final_pass = crate::scan::reconcile_paths_handle_controlled(
1580 index,
1581 &[PathBuf::new()],
1582 scan,
1583 false,
1584 &control,
1585 &mut |_commit| journal.notify_commit(),
1586 )?;
1587 handoff_evidence.retain(root, &final_pass.reconciliation);
1588 watcher.flush_capture()?;
1589 let drained = drain_observation_hints(
1590 &watcher,
1591 root,
1592 index,
1593 journal,
1594 scan,
1595 &control,
1596 handoff_intents,
1597 &mut handoff_evidence,
1598 )?;
1599 let final_settled = final_pass.reconciliation.apply.resource_refused == 0
1600 && final_pass.reconciliation.apply.stale == 0;
1601 if final_settled && drained.complete {
1602 handoff_complete = true;
1603 break;
1604 }
1605 if index.state()?.phase == crate::LifecyclePhase::Stopped {
1606 return Ok(());
1607 }
1608 let final_retryable =
1609 final_settled || reconcile_conflict_is_retryable(&final_pass.reconciliation);
1610 let retryable_conflict =
1611 final_retryable && (drained.complete || drained.retryable_conflict);
1612 if !retryable_conflict {
1613 break;
1614 }
1615 }
1616
1617 control.check_active()?;
1618 let state = index.state()?;
1619 if state.phase == crate::LifecyclePhase::Stopped {
1620 return Ok(());
1621 }
1622 if !handoff_complete {
1623 return Err(Error::ObservationHandoffIncomplete);
1624 }
1625 #[cfg(test)]
1626 controls.reach(TestPoint::BeforeObservationWatching);
1627 publish_observation_transition(
1628 index,
1629 journal,
1630 crate::index::ObservationTransition::Watching {
1631 issues: handoff_evidence.issues,
1632 omitted: handoff_evidence.omitted,
1633 },
1634 )?;
1635 let state = index.state()?;
1636 if state.phase == crate::LifecyclePhase::Stopped {
1637 return Ok(());
1638 }
1639 if state.phase != crate::LifecyclePhase::Watching {
1640 return Err(Error::ObservationHandoffIncomplete);
1641 }
1642
1643 loop {
1644 control.check_active()?;
1645 #[cfg(test)]
1646 controls.reach(TestPoint::BeforeObservationPoll);
1647 match watcher.apply_next_controlled(
1648 index,
1649 scan,
1650 OBSERVATION_POLL_INTERVAL,
1651 &control,
1652 &mut |_commit| journal.notify_commit(),
1653 ) {
1654 Ok(Some(report)) => {
1655 let mut unreadable = HandoffEvidence::default();
1659 unreadable.retain(root, &report.reconciliation);
1660 if !unreadable.issues.is_empty() || unreadable.omitted > 0 {
1661 publish_observation_transition(
1662 index,
1663 journal,
1664 crate::index::ObservationTransition::Unreadable {
1665 issues: unreadable.issues,
1666 omitted: unreadable.omitted,
1667 },
1668 )?;
1669 }
1670 }
1671 Ok(None) => {}
1672 Err(Error::OpenedIndexClosed) if cancellation.is_cancelled() => return Ok(()),
1673 Err(error) => return Err(error),
1674 }
1675 if index.state()?.phase == crate::LifecyclePhase::Stopped {
1676 return Ok(());
1677 }
1678 }
1679}
1680
1681#[cfg(feature = "watch")]
1682#[allow(clippy::too_many_arguments)]
1683fn drain_observation_hints(
1684 watcher: &crate::watch::Watcher,
1685 root: &Path,
1686 index: &IndexHandle,
1687 journal: &journal::JournalWait,
1688 scan: &ScanConfig,
1689 control: &dyn ReconcileControl,
1690 limit: usize,
1691 evidence: &mut HandoffEvidence,
1692) -> Result<HandoffDrain> {
1693 let mut drained = HandoffDrain { complete: true, retryable_conflict: true };
1694 for _ in 0..limit {
1695 let Some(report) = watcher.apply_next_controlled(
1696 index,
1697 scan,
1698 std::time::Duration::ZERO,
1699 control,
1700 &mut |_commit| journal.notify_commit(),
1701 )?
1702 else {
1703 break;
1704 };
1705 evidence.retain(root, &report.reconciliation);
1706 let report_complete = report.apply.resource_refused == 0
1707 && report.apply.stale == 0
1708 && report.reconciliation.apply.resource_refused == 0
1709 && report.reconciliation.apply.stale == 0;
1710 if !report_complete {
1711 drained.complete = false;
1712 drained.retryable_conflict &= report.apply.resource_refused == 0
1713 && report.reconciliation.apply.resource_refused == 0
1714 && (report.apply.stale > 0 || report.reconciliation.apply.stale > 0);
1715 }
1716 }
1717 Ok(drained)
1718}
1719
1720#[cfg(feature = "watch")]
1721struct HandoffDrain {
1722 complete: bool,
1723 retryable_conflict: bool,
1724}
1725
1726#[cfg(feature = "watch")]
1727fn reconcile_conflict_is_retryable(report: &crate::scan::ReconcileReport) -> bool {
1728 report.apply.resource_refused == 0 && report.apply.stale > 0
1729}
1730
1731#[cfg(feature = "watch")]
1732#[derive(Default)]
1733struct HandoffEvidence {
1734 issues: Vec<crate::Issue>,
1735 omitted: u64,
1736}
1737
1738#[cfg(feature = "watch")]
1739impl HandoffEvidence {
1740 fn retain(&mut self, root: &Path, report: &crate::scan::ReconcileReport) {
1741 for error in &report.scan.errors {
1742 retain_local_issue(
1743 &mut self.issues,
1744 &mut self.omitted,
1745 crate::Issue::from_error_under(root, error),
1746 );
1747 }
1748 }
1749}
1750
1751#[cfg(feature = "watch")]
1752fn publish_observation_transition(
1753 index: &IndexHandle,
1754 journal: &journal::JournalWait,
1755 transition: crate::index::ObservationTransition,
1756) -> Result<()> {
1757 let outcome = index.transition_observation(transition)?;
1758 if outcome.commit.is_some() {
1759 journal.notify_commit();
1760 }
1761 Ok(())
1762}
1763
1764#[cfg(feature = "watch")]
1765struct ObservationReconcileControl<'a> {
1766 cancellation: &'a Cancellation,
1767 max_files: Option<u64>,
1768 #[cfg(test)]
1769 controls: &'a TestControls,
1770}
1771
1772#[cfg(feature = "watch")]
1773impl ReconcileControl for ObservationReconcileControl<'_> {
1774 fn check_active(&self) -> Result<()> {
1775 if self.cancellation.is_cancelled() { Err(Error::OpenedIndexClosed) } else { Ok(()) }
1776 }
1777
1778 fn before_conditional_commit(&self) -> Result<()> {
1779 self.check_active()?;
1780 #[cfg(test)]
1781 self.controls.reach(TestPoint::AfterObservationVerification);
1782 self.check_active()
1783 }
1784
1785 fn max_files(&self) -> Option<u64> {
1786 self.max_files
1787 }
1788}
1789
1790struct LockedLifecycle<'a> {
1791 guard: MutexGuard<'a, Lifecycle>,
1792 poisoned: bool,
1793}
1794
1795#[derive(Default)]
1796struct Lifecycle {
1797 phase: OwnerPhase,
1798 workers: Vec<Worker>,
1799 active_refreshes: usize,
1800 terminal: Option<CloseOutcome>,
1801}
1802
1803struct ActiveRefresh<'a> {
1804 state: &'a OpenedState,
1805}
1806
1807impl Drop for ActiveRefresh<'_> {
1808 fn drop(&mut self) {
1809 let mut lifecycle = self.state.lock_lifecycle().guard;
1810 lifecycle.active_refreshes = lifecycle.active_refreshes.saturating_sub(1);
1811 self.state.lifecycle_changed.notify_all();
1812 }
1813}
1814
1815struct OpenedReconcileControl<'a> {
1816 state: &'a OpenedState,
1817}
1818
1819impl crate::scan::ReconcileControl for OpenedReconcileControl<'_> {
1820 fn check_active(&self) -> Result<()> {
1821 if self.state.cancellation.is_cancelled() { Err(Error::OpenedIndexClosed) } else { Ok(()) }
1822 }
1823
1824 fn before_conditional_commit(&self) -> Result<()> {
1825 self.check_active()?;
1826 #[cfg(test)]
1827 self.state.test_controls.reach(TestPoint::AfterRefreshVerification);
1828 self.check_active()
1829 }
1830
1831 fn max_files(&self) -> Option<u64> {
1832 self.state.budget.max_files
1833 }
1834}
1835
1836#[derive(Clone, Copy, PartialEq, Eq, Default)]
1837enum OwnerPhase {
1838 #[default]
1839 Open,
1840 Closing,
1841 Closed,
1842}
1843
1844struct Worker {
1845 name: &'static str,
1846 handle: JoinHandle<()>,
1847}
1848
1849#[derive(Default)]
1855struct WorkerFailures {
1856 recorded: Mutex<Vec<CloseOutcome>>,
1857}
1858
1859impl WorkerFailures {
1860 fn record(&self, outcome: CloseOutcome) {
1861 self.recorded.lock().unwrap_or_else(std::sync::PoisonError::into_inner).push(outcome);
1862 }
1863
1864 fn panicked(&self) -> Option<&'static str> {
1866 self.recorded.lock().unwrap_or_else(std::sync::PoisonError::into_inner).iter().find_map(
1867 |outcome| match outcome {
1868 CloseOutcome::WorkerPanicked { worker } => Some(*worker),
1869 _ => None,
1870 },
1871 )
1872 }
1873
1874 fn first(&self) -> Option<CloseOutcome> {
1884 let recorded = self.recorded.lock().unwrap_or_else(std::sync::PoisonError::into_inner);
1885 let panicked =
1886 recorded.iter().any(|outcome| matches!(outcome, CloseOutcome::WorkerPanicked { .. }));
1887 recorded.iter().find(|outcome| !(panicked && outcome.is_poison_trace())).cloned()
1888 }
1889}
1890
1891#[derive(Clone)]
1892enum CloseOutcome {
1893 Success,
1894 LifecyclePoisoned,
1895 IndexPoisoned,
1896 WorkerPanicked { worker: &'static str },
1897 WorkerFailed { worker: &'static str, source: Arc<Error> },
1898}
1899
1900impl CloseOutcome {
1901 fn is_poison_trace(&self) -> bool {
1903 match self {
1904 Self::WorkerFailed { source, .. } => matches!(
1905 **source,
1906 Error::IndexLockPoisoned
1907 | Error::OpenedLifecyclePoisoned
1908 | Error::OpenedJournalPoisoned
1909 ),
1910 Self::LifecyclePoisoned | Self::IndexPoisoned => true,
1911 Self::Success | Self::WorkerPanicked { .. } => false,
1912 }
1913 }
1914
1915 fn to_result(&self) -> Result<()> {
1916 match self {
1917 Self::Success => Ok(()),
1918 Self::LifecyclePoisoned => Err(Error::OpenedLifecyclePoisoned),
1919 Self::IndexPoisoned => Err(Error::IndexLockPoisoned),
1920 Self::WorkerPanicked { worker } => Err(Error::OpenedWorkerPanicked { worker }),
1921 Self::WorkerFailed { worker, source } => {
1922 Err(Error::OpenedWorkerFailed { worker, source: Arc::clone(source) })
1923 }
1924 }
1925 }
1926}
1927
1928fn join_workers(workers: Vec<Worker>, failures: &WorkerFailures) -> Option<CloseOutcome> {
1933 let mut first_unrecorded = None;
1934 for worker in workers {
1935 if worker.handle.join().is_err() && first_unrecorded.is_none() {
1936 first_unrecorded = Some(CloseOutcome::WorkerPanicked { worker: worker.name });
1937 }
1938 }
1939 failures.first().or(first_unrecorded)
1940}
1941
1942#[cfg(feature = "watch")]
1943#[derive(Default)]
1944struct BaselineLatch {
1945 finished: Mutex<bool>,
1946 changed: Condvar,
1947}
1948
1949#[cfg(feature = "watch")]
1950impl BaselineLatch {
1951 fn finish(&self) {
1952 let mut finished = self.finished.lock().unwrap_or_else(std::sync::PoisonError::into_inner);
1953 *finished = true;
1954 self.changed.notify_all();
1955 }
1956
1957 fn wake(&self) {
1958 let guard = self.finished.lock().unwrap_or_else(std::sync::PoisonError::into_inner);
1959 self.changed.notify_all();
1960 drop(guard);
1961 }
1962
1963 fn wait(&self, cancellation: &Cancellation) -> bool {
1964 let mut finished = self.finished.lock().unwrap_or_else(std::sync::PoisonError::into_inner);
1965 while !*finished && !cancellation.is_cancelled() {
1966 finished =
1967 self.changed.wait(finished).unwrap_or_else(std::sync::PoisonError::into_inner);
1968 }
1969 *finished && !cancellation.is_cancelled()
1970 }
1971}
1972
1973#[cfg(feature = "watch")]
1974struct BaselineCompletion(Arc<BaselineLatch>);
1975
1976#[cfg(feature = "watch")]
1977impl Drop for BaselineCompletion {
1978 fn drop(&mut self) {
1979 self.0.finish();
1980 }
1981}
1982
1983#[derive(Default)]
1984struct Cancellation {
1985 cancelled: AtomicBool,
1986 wait_lock: Mutex<()>,
1987 changed: Condvar,
1988}
1989
1990impl Cancellation {
1991 fn cancel(&self) {
1992 let guard = self.wait_lock.lock().unwrap_or_else(std::sync::PoisonError::into_inner);
1993 self.cancelled.store(true, Ordering::Release);
1994 self.changed.notify_all();
1995 drop(guard);
1996 }
1997
1998 #[allow(dead_code)]
1999 fn is_cancelled(&self) -> bool {
2000 self.cancelled.load(Ordering::Acquire)
2001 }
2002
2003 #[allow(dead_code)]
2004 fn wait_cancelled(&self) {
2005 let mut guard = self.wait_lock.lock().unwrap_or_else(std::sync::PoisonError::into_inner);
2006 while !self.cancelled.load(Ordering::Acquire) {
2007 guard = self.changed.wait(guard).unwrap_or_else(std::sync::PoisonError::into_inner);
2008 }
2009 }
2010}
2011
2012#[cfg(test)]
2013#[derive(Clone, Copy)]
2014enum TestPoint {
2015 BeforeWorkerExit,
2016 BeforeCloseWait,
2017 BeforeDiscovery,
2018 AfterRootDirectory,
2019 BeforeJournalWait,
2020 AfterRefreshVerification,
2021 DuringTreeProjection,
2022 #[cfg(feature = "watch")]
2023 BeforeObservationHandoff,
2024 #[cfg(feature = "watch")]
2025 BeforeObservationWatching,
2026 #[cfg(feature = "watch")]
2027 BeforeObservationPoll,
2028 #[cfg(feature = "watch")]
2029 AfterObservationVerification,
2030}
2031
2032#[cfg(test)]
2033#[derive(Default)]
2034struct TestControls {
2035 before_worker_exit: TestGate,
2036 before_close_wait: TestGate,
2037 before_discovery: TestGate,
2038 after_root_directory: TestGate,
2039 before_journal_wait: TestGate,
2040 after_refresh_verification: TestGate,
2041 during_tree_projection: TestGate,
2042 #[cfg(feature = "watch")]
2043 before_observation_handoff: TestGate,
2044 #[cfg(feature = "watch")]
2045 before_observation_watching: TestGate,
2046 #[cfg(feature = "watch")]
2047 before_observation_poll: TestGate,
2048 #[cfg(feature = "watch")]
2049 after_observation_verification: TestGate,
2050 #[cfg(feature = "watch")]
2051 scripted_observer: Mutex<Option<crate::watch::ScriptedSender>>,
2052 discovery_disabled: AtomicBool,
2053 deterministic_discovery_order: AtomicBool,
2054}
2055
2056#[cfg(test)]
2057impl TestControls {
2058 fn use_deterministic_discovery_order(&self) {
2059 self.deterministic_discovery_order.store(true, Ordering::Release);
2060 }
2061
2062 fn gate(&self, point: TestPoint) -> &TestGate {
2063 match point {
2064 TestPoint::BeforeWorkerExit => &self.before_worker_exit,
2065 TestPoint::BeforeCloseWait => &self.before_close_wait,
2066 TestPoint::BeforeDiscovery => &self.before_discovery,
2067 TestPoint::AfterRootDirectory => &self.after_root_directory,
2068 TestPoint::BeforeJournalWait => &self.before_journal_wait,
2069 TestPoint::AfterRefreshVerification => &self.after_refresh_verification,
2070 TestPoint::DuringTreeProjection => &self.during_tree_projection,
2071 #[cfg(feature = "watch")]
2072 TestPoint::BeforeObservationHandoff => &self.before_observation_handoff,
2073 #[cfg(feature = "watch")]
2074 TestPoint::BeforeObservationWatching => &self.before_observation_watching,
2075 #[cfg(feature = "watch")]
2076 TestPoint::BeforeObservationPoll => &self.before_observation_poll,
2077 #[cfg(feature = "watch")]
2078 TestPoint::AfterObservationVerification => &self.after_observation_verification,
2079 }
2080 }
2081
2082 fn reach(&self, point: TestPoint) {
2083 self.gate(point).reach();
2084 }
2085
2086 #[cfg(feature = "watch")]
2087 fn send_observation_hints(&self, source: &str) {
2088 self.scripted_observer
2089 .lock()
2090 .expect("scripted observer lock")
2091 .as_ref()
2092 .expect("scripted observer installed")
2093 .send(source)
2094 .expect("valid scripted hints");
2095 }
2096}
2097
2098#[cfg(test)]
2099#[derive(Default)]
2100struct TestGate {
2101 state: Mutex<TestGateState>,
2102 changed: Condvar,
2103}
2104
2105#[cfg(test)]
2106#[derive(Default)]
2107struct TestGateState {
2108 armed: bool,
2109 reached: bool,
2110 released: bool,
2111}
2112
2113#[cfg(test)]
2114impl TestGate {
2115 fn arm(&self) {
2116 let mut state = self.state.lock().expect("test gate lock");
2117 *state = TestGateState { armed: true, reached: false, released: false };
2118 }
2119
2120 fn reach(&self) {
2121 let mut state = self.state.lock().expect("test gate lock");
2122 if !state.armed {
2123 return;
2124 }
2125 state.reached = true;
2126 self.changed.notify_all();
2127 while !state.released {
2128 state = self.changed.wait(state).expect("test gate wait");
2129 }
2130 state.armed = false;
2131 }
2132
2133 fn wait_reached(&self) {
2134 let deadline = std::time::Instant::now() + TEST_GATE_TIMEOUT;
2135 let mut state = self.state.lock().expect("test gate lock");
2136 while !state.reached {
2137 let remaining = deadline.saturating_duration_since(std::time::Instant::now());
2138 assert!(!remaining.is_zero(), "test gate was not reached");
2139 let (next, timed_out) =
2140 self.changed.wait_timeout(state, remaining).expect("test gate wait");
2141 state = next;
2142 assert!(!timed_out.timed_out() || state.reached, "test gate was not reached");
2143 }
2144 }
2145
2146 fn release(&self) {
2147 let mut state = self.state.lock().expect("test gate lock");
2148 state.released = true;
2149 self.changed.notify_all();
2150 }
2151}
2152
2153#[cfg(test)]
2154mod tests {
2155 use super::*;
2156 use std::ffi::OsString;
2157 use std::sync::atomic::{AtomicUsize, Ordering};
2158
2159 fn wait_until_settled(opened: &OpenedIndex) -> crate::IndexState {
2160 let deadline = std::time::Instant::now() + TEST_GATE_TIMEOUT;
2161 loop {
2162 let state = opened.state.index.state().expect("read state");
2163 if state.phase != crate::LifecyclePhase::Discovering {
2164 return state;
2165 }
2166 assert!(std::time::Instant::now() < deadline, "discovery did not settle");
2167 std::thread::yield_now();
2168 }
2169 }
2170
2171 fn wait_for_worker_exit(opened: &OpenedIndex, name: &str) {
2176 let deadline = std::time::Instant::now() + TEST_GATE_TIMEOUT;
2177 loop {
2178 let (registered, finished) = {
2179 let lifecycle = opened.state.lock_lifecycle();
2180 let named = lifecycle.guard.workers.iter().filter(|worker| worker.name == name);
2181 named.fold((0, 0), |(registered, finished), worker| {
2182 (registered + 1, finished + usize::from(worker.handle.is_finished()))
2183 })
2184 };
2185 assert!(registered > 0, "no worker named {name}");
2186 if registered == finished {
2187 return;
2188 }
2189 assert!(std::time::Instant::now() < deadline, "worker {name} did not exit");
2190 std::thread::yield_now();
2191 }
2192 }
2193
2194 #[cfg(feature = "watch")]
2195 fn wait_until_phase(
2196 opened: &OpenedIndex,
2197 expected: crate::LifecyclePhase,
2198 ) -> crate::IndexState {
2199 let deadline = std::time::Instant::now() + TEST_GATE_TIMEOUT;
2200 loop {
2201 let state = opened.state.index.state().expect("read state");
2202 if state.phase == expected {
2203 return state;
2204 }
2205 assert!(std::time::Instant::now() < deadline, "phase did not become {expected:?}");
2206 std::thread::yield_now();
2207 }
2208 }
2209
2210 #[cfg(feature = "watch")]
2211 fn scripted_options(script: &Path) -> OpenOptions {
2212 OpenOptions {
2213 observation: Some(crate::watch::WatchConfig {
2214 settle: std::time::Duration::from_millis(1),
2215 max_hold: std::time::Duration::from_millis(10),
2216 ..crate::watch::WatchConfig::default()
2217 }),
2218 observation_script: Some(script.to_path_buf()),
2219 ..OpenOptions::default()
2220 }
2221 }
2222
2223 fn child_facts(index: &Index, path: &Path) -> Vec<(OsString, EntryKind, crate::Attrs)> {
2224 index
2225 .children(path)
2226 .expect("known directory")
2227 .map(|(name, id)| {
2228 (
2229 name.to_os_string(),
2230 index.kind(&index.path_of(id).expect("path")).expect("kind"),
2231 *index.attrs(&index.path_of(id).expect("path")).expect("attrs"),
2232 )
2233 })
2234 .collect()
2235 }
2236
2237 fn opened(controls: Arc<TestControls>) -> (tempfile::TempDir, OpenedIndex) {
2238 controls.discovery_disabled.store(true, Ordering::Release);
2239 let root = tempfile::tempdir().expect("temp root");
2240 let opened = OpenedIndex::open_for_test(root.path(), OpenOptions::default(), controls)
2241 .expect("open live root");
2242 (root, opened)
2243 }
2244
2245 fn current_version(opened: &OpenedIndex) -> crate::EngineVersion {
2246 opened.read(crate::ReadRequest::default()).expect("read version").version
2247 }
2248
2249 fn apply_and_notify(opened: &OpenedIndex, observation: &Observation) -> crate::ApplyOutcome {
2250 let outcome = opened.state.index.apply(observation).expect("apply observation");
2251 if outcome.commit.is_some() {
2252 opened.state.journal.notify_commit();
2253 }
2254 outcome
2255 }
2256
2257 #[test]
2261 fn an_opened_root_takes_a_case_variant_control_as_a_detached_scan_does() {
2262 use crate::test_support::CaseLookups;
2263
2264 let probe = tempfile::tempdir().expect("temp root");
2265 for (lookups, governs) in CaseLookups::on_this_host(probe.path()) {
2266 let root = tempfile::tempdir().expect("temp root");
2267 let _lookups = lookups.install(root.path());
2268 std::fs::create_dir(root.path().join("up")).expect("directory");
2269 std::fs::write(root.path().join("up/.GITIGNORE"), b"*.tmp\n").expect("variant");
2270 std::fs::write(root.path().join("up/x.tmp"), b"governed").expect("fixture");
2271 std::fs::write(root.path().join("up/notes.txt"), b"never").expect("fixture");
2272 let facts = |index: &Index| {
2273 let sources: Vec<(PathBuf, Vec<u8>)> = index
2274 .controls()
2275 .expect("observed")
2276 .sources()
2277 .map(|(path, source)| (path, source.to_vec()))
2278 .collect();
2279 (index.is_ignored(Path::new("up/x.tmp")).expect("observed"), sources)
2280 };
2281 let cold = || {
2282 let (cold, _) =
2283 crate::scan::scan_into_index(root.path(), &crate::ScanConfig::default())
2284 .expect("cold scan");
2285 facts(&cold)
2286 };
2287
2288 let options = OpenOptions { batch_size: 1, ..OpenOptions::default() };
2289 let opened = open_fixture(root.path(), options).expect("opened root");
2290 wait_until_settled(&opened);
2291 let discovered = opened.state.index.read_with(facts).expect("read");
2292 assert_eq!(discovered, cold(), "{lookups:?}: discovery");
2293 assert_eq!(discovered.0, Some(governs), "{lookups:?}");
2294
2295 std::fs::remove_file(root.path().join("up/.GITIGNORE")).expect("remove the variant");
2296 opened.refresh(&[PathBuf::from("up/.GITIGNORE")]).expect("refresh");
2297 let refreshed = opened.state.index.read_with(facts).expect("read");
2298 assert_eq!(refreshed, cold(), "{lookups:?}: refresh of the removed variant");
2299 assert_eq!(refreshed.0, Some(false), "{lookups:?}");
2300 opened.close().expect("close");
2301 }
2302 }
2303
2304 #[test]
2305 fn associated_and_free_open_contracts_coexist() {
2306 let root = tempfile::tempdir().expect("temp root");
2307 std::fs::write(root.path().join("file.txt"), b"one").expect("fixture");
2308
2309 let opened = open_fixture(root.path(), OpenOptions::default()).expect("opened root");
2310 let (detached, _) = crate::open_fixture(
2311 root.path(),
2312 &crate::OpenFixture {
2313 cache_path: None,
2314 policy: crate::CachePolicy::Off,
2315 ..crate::OpenFixture::default()
2316 },
2317 )
2318 .expect("blocking open");
2319
2320 assert_eq!(detached.total().files, 1);
2321 opened.close().expect("close");
2322 }
2323
2324 #[test]
2325 fn one_clone_closes_the_shared_authority_and_close_is_idempotent() {
2326 let (_root, opened) = opened(Arc::default());
2327 let clone = opened.clone();
2328 assert_eq!(opened.state.session, clone.state.session);
2329 assert!(Arc::ptr_eq(&opened.state, &clone.state));
2330
2331 clone.close().expect("first close");
2332 assert!(matches!(opened.ensure_open(), Err(Error::OpenedIndexClosed)));
2333 opened.close().expect("repeated close");
2334 }
2335
2336 #[test]
2337 fn shutdown_refuses_new_owned_work() {
2338 let (_root, opened) = opened(Arc::default());
2339 opened.close().expect("close");
2340 let ran = Arc::new(AtomicBool::new(false));
2341 let worker_ran = Arc::clone(&ran);
2342
2343 assert!(matches!(
2344 opened.spawn_worker("late", move |_cancellation| {
2345 worker_ran.store(true, Ordering::SeqCst);
2346 Ok(())
2347 }),
2348 Err(Error::OpenedIndexClosed)
2349 ));
2350 assert!(!ran.load(Ordering::SeqCst));
2351 }
2352
2353 #[test]
2354 fn concurrent_close_waits_for_one_stored_worker_failure() {
2355 let controls = Arc::new(TestControls::default());
2356 controls.gate(TestPoint::BeforeWorkerExit).arm();
2357 controls.gate(TestPoint::BeforeCloseWait).arm();
2358 let (_root, opened) = opened(Arc::clone(&controls));
2359 opened
2360 .spawn_worker("failure", |cancellation| {
2361 cancellation.wait_cancelled();
2362 Err(Error::CommitRejected("injected opened worker failure"))
2363 })
2364 .expect("spawn worker");
2365
2366 let first = opened.clone();
2367 let first_close = thread::spawn(move || first.close());
2368 controls.gate(TestPoint::BeforeWorkerExit).wait_reached();
2369 let second = opened.clone();
2370 let second_close = thread::spawn(move || second.close());
2371 controls.gate(TestPoint::BeforeCloseWait).wait_reached();
2372
2373 controls.gate(TestPoint::BeforeCloseWait).release();
2374 controls.gate(TestPoint::BeforeWorkerExit).release();
2375 let first_error = first_close.join().expect("first close thread").expect_err("failure");
2376 let second_error = second_close.join().expect("second close thread").expect_err("failure");
2377 assert_eq!(first_error.to_string(), second_error.to_string());
2378 assert!(matches!(first_error, Error::OpenedWorkerFailed { worker: "failure", .. }));
2379 assert_eq!(
2380 opened.close().expect_err("stored failure").to_string(),
2381 second_error.to_string()
2382 );
2383 }
2384
2385 #[test]
2386 fn a_panicking_worker_is_joined_and_reported() {
2387 let (_root, opened) = opened(Arc::default());
2388 opened
2389 .spawn_worker("panic", |_cancellation| panic!("injected worker panic"))
2390 .expect("spawn worker");
2391
2392 let error = opened.close().expect_err("panic is terminal");
2393 assert!(matches!(error, Error::OpenedWorkerPanicked { worker: "panic" }));
2394 assert!(matches!(opened.close(), Err(Error::OpenedWorkerPanicked { worker: "panic" })));
2395 }
2396
2397 #[test]
2402 fn a_worker_panic_wakes_a_blocked_change_poll_with_its_typed_failure() {
2403 for holds_index_lock in [false, true] {
2404 let controls = Arc::new(TestControls::default());
2405 controls.gate(TestPoint::BeforeJournalWait).arm();
2406 let (_root, opened) = opened(Arc::clone(&controls));
2407 let cursor = current_version(&opened);
2408 let poller = opened.clone();
2409 let poll = thread::spawn(move || {
2410 poller.changes(crate::ChangeRequest {
2411 after: cursor,
2412 timeout: std::time::Duration::from_secs(60),
2413 })
2414 });
2415 controls.gate(TestPoint::BeforeJournalWait).wait_reached();
2416 let index = opened.state.index.clone();
2417 opened
2418 .spawn_worker("panic", move |_cancellation| {
2419 if holds_index_lock {
2420 index.panic_holding_the_write_lock_for_test();
2421 }
2422 panic!("injected worker panic")
2423 })
2424 .expect("spawn worker");
2425 controls.gate(TestPoint::BeforeJournalWait).release();
2426
2427 let outcome = poll.join().expect("poll thread");
2428 assert!(
2429 matches!(outcome, Err(Error::OpenedWorkerPanicked { worker: "panic" })),
2430 "holds index lock {holds_index_lock}: {outcome:?}"
2431 );
2432 assert!(matches!(opened.close(), Err(Error::OpenedWorkerPanicked { worker: "panic" })));
2433 }
2434 }
2435
2436 #[test]
2438 fn close_reports_the_failure_that_happened_first_not_the_worker_spawned_first() {
2439 let (_root, opened) = opened(Arc::default());
2440 opened
2441 .spawn_worker("slow", |cancellation| {
2442 cancellation.wait_cancelled();
2443 Err(Error::CommitRejected("slow worker failed at close"))
2444 })
2445 .expect("spawn slow worker");
2446 opened
2447 .spawn_worker("fast", |_cancellation| {
2448 Err(Error::CommitRejected("fast worker failed first"))
2449 })
2450 .expect("spawn fast worker");
2451 wait_for_worker_exit(&opened, "fast");
2452
2453 assert!(matches!(opened.close(), Err(Error::OpenedWorkerFailed { worker: "fast", .. })));
2454 }
2455
2456 #[test]
2460 fn close_reports_a_panic_before_the_poisoning_it_left_behind() {
2461 let (_root, opened) = opened(Arc::default());
2462 opened
2463 .spawn_worker("tripped", |_cancellation| Err(Error::IndexLockPoisoned))
2464 .expect("spawn tripped worker");
2465 wait_for_worker_exit(&opened, "tripped");
2466 opened
2467 .spawn_worker("panicked", |_cancellation| panic!("injected worker panic"))
2468 .expect("spawn panicking worker");
2469 wait_for_worker_exit(&opened, "panicked");
2470
2471 assert!(matches!(opened.close(), Err(Error::OpenedWorkerPanicked { worker: "panicked" })));
2472 }
2473
2474 #[test]
2478 fn close_reports_a_poisoning_no_worker_panic_explains_when_it_came_first() {
2479 let (_root, opened) = opened(Arc::default());
2480 opened.state.index.poison_for_test();
2481 let index = opened.state.index.clone();
2482 opened
2483 .spawn_worker("tripped", move |_cancellation| index.clock().map(|_| ()))
2484 .expect("spawn tripped worker");
2485 wait_for_worker_exit(&opened, "tripped");
2486 opened
2487 .spawn_worker("later", |_cancellation| {
2488 Err(Error::CommitRejected("unrelated later failure"))
2489 })
2490 .expect("spawn later worker");
2491 wait_for_worker_exit(&opened, "later");
2492
2493 let closed = opened.close();
2494 assert!(
2495 matches!(
2496 &closed,
2497 Err(Error::OpenedWorkerFailed { worker: "tripped", source })
2498 if matches!(**source, Error::IndexLockPoisoned)
2499 ),
2500 "{closed:?}"
2501 );
2502 }
2503
2504 #[test]
2505 fn dropping_the_last_reference_cancels_and_joins() {
2506 let active = Arc::new(AtomicUsize::new(0));
2507 let (started_sender, started_receiver) = std::sync::mpsc::sync_channel(0);
2508 let (_root, opened) = opened(Arc::default());
2509 let worker_active = Arc::clone(&active);
2510 opened
2511 .spawn_worker("drop", move |cancellation| {
2512 worker_active.fetch_add(1, Ordering::SeqCst);
2513 started_sender.send(()).expect("report worker start");
2514 cancellation.wait_cancelled();
2515 worker_active.fetch_sub(1, Ordering::SeqCst);
2516 Ok(())
2517 })
2518 .expect("spawn worker");
2519
2520 started_receiver.recv_timeout(std::time::Duration::from_secs(5)).expect("worker started");
2521 drop(opened);
2522 assert_eq!(active.load(Ordering::SeqCst), 0, "drop returned before worker join");
2523 }
2524
2525 #[test]
2526 fn poisoned_lifecycle_still_joins_before_returning_its_typed_failure() {
2527 let active = Arc::new(AtomicUsize::new(0));
2528 let (started_sender, started_receiver) = std::sync::mpsc::sync_channel(0);
2529 let (_root, opened) = opened(Arc::default());
2530 let worker_active = Arc::clone(&active);
2531 opened
2532 .spawn_worker("poison", move |cancellation| {
2533 worker_active.fetch_add(1, Ordering::SeqCst);
2534 started_sender.send(()).expect("report worker start");
2535 cancellation.wait_cancelled();
2536 worker_active.fetch_sub(1, Ordering::SeqCst);
2537 Ok(())
2538 })
2539 .expect("spawn worker");
2540 started_receiver.recv_timeout(std::time::Duration::from_secs(5)).expect("worker started");
2541
2542 let state = Arc::clone(&opened.state);
2543 thread::spawn(move || {
2544 let _guard = state.lifecycle.lock().expect("lifecycle lock");
2545 panic!("inject lifecycle poison");
2546 })
2547 .join()
2548 .expect_err("injected panic");
2549
2550 assert!(matches!(opened.close(), Err(Error::OpenedLifecyclePoisoned)));
2551 assert_eq!(active.load(Ordering::SeqCst), 0, "poison bypassed worker join");
2552 assert!(matches!(opened.close(), Err(Error::OpenedLifecyclePoisoned)));
2553 }
2554
2555 #[test]
2556 fn poisoned_index_still_joins_and_replays_the_typed_failure() {
2557 let (_root, opened) = opened(Arc::default());
2558 opened.state.index.poison_for_test();
2559
2560 assert!(matches!(opened.close(), Err(Error::IndexLockPoisoned)));
2561 assert!(matches!(opened.close(), Err(Error::IndexLockPoisoned)));
2562 }
2563
2564 #[test]
2565 fn opened_root_rejects_invalid_scan_policy_and_nondirectories() {
2566 let root = tempfile::tempdir().expect("temp root");
2567 let options = OpenOptions { batch_size: 0, ..OpenOptions::default() };
2568 assert!(matches!(open_fixture(root.path(), options), Err(Error::UnsupportedScanConfig(_))));
2569
2570 let file = root.path().join("file");
2571 std::fs::write(&file, b"x").expect("fixture");
2572 assert!(matches!(open_fixture(&file, OpenOptions::default()), Err(Error::Io { .. })));
2573
2574 let zero_budget = OpenOptions {
2575 budget: DiscoveryBudget { max_files: Some(0) },
2576 ..OpenOptions::default()
2577 };
2578 assert!(matches!(
2579 open_fixture(root.path(), zero_budget),
2580 Err(Error::UnsupportedScanConfig(_))
2581 ));
2582
2583 let minimum = crate::MIN_JOURNAL_CAPACITY_BYTES;
2584 let below_minimum =
2585 OpenOptions { journal_capacity_bytes: minimum - 1, ..OpenOptions::default() };
2586 let error = open_fixture(root.path(), below_minimum).expect_err("refused");
2587 assert!(
2588 matches!(error, Error::JournalCapacityTooSmall { requested, minimum: stated }
2589 if requested == minimum - 1 && stated == minimum),
2590 "{error:?}"
2591 );
2592 let message = error.to_string();
2593 assert!(message.contains(&format!("at least {minimum} bytes")), "{message}");
2594
2595 let at_minimum = OpenOptions { journal_capacity_bytes: minimum, ..OpenOptions::default() };
2596 open_fixture(root.path(), at_minimum).expect("accepted").close().expect("close");
2597 }
2598
2599 #[test]
2600 fn distinct_opens_have_distinct_live_identity() {
2601 let root = tempfile::tempdir().expect("temp root");
2602 let first = open_fixture(root.path(), OpenOptions::default()).expect("first");
2603 let second = open_fixture(root.path(), OpenOptions::default()).expect("second");
2604 assert_ne!(first.state.session, second.state.session);
2605 first.close().expect("first close");
2606 second.close().expect("second close");
2607 }
2608
2609 #[test]
2610 fn terminal_discovery_failure_retains_bounded_typed_evidence() {
2611 let root = tempfile::tempdir().expect("temp root");
2612 let controls = Arc::new(TestControls::default());
2613 controls.gate(TestPoint::BeforeDiscovery).arm();
2614 let opened =
2615 OpenedIndex::open_for_test(root.path(), OpenOptions::default(), Arc::clone(&controls))
2616 .expect("opened root");
2617 std::fs::remove_dir(root.path()).expect("remove empty fixture root");
2618 controls.gate(TestPoint::BeforeDiscovery).release();
2619
2620 let state = wait_until_settled(&opened);
2621 assert_eq!(state.phase, crate::LifecyclePhase::Failed);
2622 assert_eq!(state.coverage, crate::Coverage::Partial(crate::CoverageReason::Failed));
2623 assert_eq!(state.issues.retained, 1);
2624 let issues = opened.state.index.issues().expect("issues");
2625 assert_eq!(issues.len(), 1);
2626 assert_eq!(issues[0].kind, crate::IssueKind::Disappeared);
2627 assert!(matches!(opened.close(), Err(Error::OpenedWorkerFailed { .. })));
2628 }
2629
2630 #[test]
2631 fn progressive_discovery_settles_to_the_one_shot_tree() {
2632 let root = tempfile::tempdir().expect("temp root");
2633 std::fs::create_dir_all(root.path().join("alpha/deep")).expect("fixture directories");
2634 std::fs::write(root.path().join("root.txt"), b"root").expect("root fixture");
2635 std::fs::write(root.path().join("alpha/child.bin"), b"child").expect("child fixture");
2636 std::fs::write(root.path().join("alpha/deep/leaf.rs"), b"leaf").expect("leaf fixture");
2637 crate::test_support::settle_allocations(root.path());
2638 let options = OpenOptions { batch_size: 2, ..OpenOptions::default() };
2639
2640 let opened = open_fixture(root.path(), options.clone()).expect("opened root");
2641 let state = wait_until_settled(&opened);
2642 assert_eq!(state.phase, crate::LifecyclePhase::Ready);
2643 assert_eq!(state.coverage, crate::Coverage::Complete);
2644 assert_eq!(state.progress.files_retained, 3);
2645
2646 let live = opened.state.index.snapshot().expect("live snapshot");
2647 let (one_shot, report) =
2648 crate::scan::scan_into_index(root.path(), &options.clone().into_parts().0)
2649 .expect("one-shot scan");
2650 assert!(report.is_complete());
2651 assert_eq!(live.total(), one_shot.total());
2652 assert_eq!(live.len(), one_shot.len());
2653 for path in [Path::new(""), Path::new("alpha"), Path::new("alpha/deep")] {
2654 assert_eq!(child_facts(&live, path), child_facts(&one_shot, path));
2655 assert_eq!(live.directory_complete(path), Some(true));
2656 }
2657 opened.close().expect("close");
2658 }
2659
2660 #[test]
2661 fn parent_listing_commits_before_prioritized_child_work_without_a_clock_change() {
2662 let root = tempfile::tempdir().expect("temp root");
2663 for directory in ["alpha", "target"] {
2664 std::fs::create_dir(root.path().join(directory)).expect("fixture directory");
2665 std::fs::write(root.path().join(directory).join("leaf"), directory)
2666 .expect("fixture file");
2667 }
2668 let controls = Arc::new(TestControls::default());
2669 controls.gate(TestPoint::AfterRootDirectory).arm();
2670 let opened = OpenedIndex::open_for_test(
2671 root.path(),
2672 OpenOptions { batch_size: 64, ..OpenOptions::default() },
2673 Arc::clone(&controls),
2674 )
2675 .expect("opened root");
2676 controls.gate(TestPoint::AfterRootDirectory).wait_reached();
2677
2678 let before_priority = opened.state.index.clock().expect("clock");
2679 assert_eq!(
2680 opened.state.index.directory_complete(Path::new("")).expect("root completeness"),
2681 Some(true)
2682 );
2683 assert_eq!(
2684 opened
2685 .state
2686 .index
2687 .directory_complete(Path::new("target"))
2688 .expect("target completeness"),
2689 Some(false)
2690 );
2691 opened.prioritize(&[PathBuf::from("target")]).expect("prioritize");
2692 assert_eq!(opened.state.index.clock().expect("clock"), before_priority);
2693 controls.gate(TestPoint::AfterRootDirectory).release();
2694 let state = wait_until_settled(&opened);
2695 assert_eq!(state.coverage, crate::Coverage::Complete);
2696
2697 let since = opened.state.index.since(before_priority).expect("commits after root");
2698 let first_file = since
2699 .commits
2700 .iter()
2701 .flat_map(|commit| &commit.changes)
2702 .find_map(|change| match change {
2703 crate::EffectiveChange::Inserted { path, kind: EntryKind::File, .. } => {
2704 Some(path.clone())
2705 }
2706 _ => None,
2707 })
2708 .expect("child file commit");
2709 assert_eq!(first_file, PathBuf::from("target/leaf"));
2710 opened.close().expect("close");
2711 }
2712
2713 #[test]
2720 fn the_grouped_frontier_pops_in_the_order_the_scan_chose() {
2721 fn reference_pop(
2722 pending: &mut VecDeque<PendingDirectory>,
2723 priorities: &[PathBuf],
2724 ) -> Option<PendingDirectory> {
2725 let selected = pending
2726 .iter()
2727 .enumerate()
2728 .filter_map(|(position, entry)| {
2729 priorities
2730 .iter()
2731 .position(|priority| {
2732 priority.starts_with(&entry.path) || entry.path.starts_with(priority)
2733 })
2734 .map(|priority| (priority, position))
2735 })
2736 .min()
2737 .map_or(0, |(_, position)| position);
2738 pending.remove(selected)
2739 }
2740
2741 const PATHS: [&str; 7] = ["a", "b", "a/b", "b/a", "a/b/c", "c", "c/a/b"];
2742 let mut seed = 0x2545_f491_4f6c_dd1d_u64;
2743 let mut next = move |bound: usize| {
2744 seed ^= seed << 13;
2745 seed ^= seed >> 7;
2746 seed ^= seed << 17;
2747 usize::try_from(seed % u64::try_from(bound).expect("small bound")).expect("fits")
2748 };
2749 let frontier = DiscoveryFrontier::new();
2750 let mut reference = VecDeque::from([PendingDirectory { path: PathBuf::new(), depth: 0 }]);
2751 let mut priorities = Vec::new();
2752 for step in 0..4_000 {
2753 match next(10) {
2754 0..=4 => {
2755 let expected = reference_pop(&mut reference, &priorities);
2756 let actual = frontier.pop();
2757 assert_eq!(
2758 actual.as_ref().map(|directory| &directory.path),
2759 expected.as_ref().map(|directory| &directory.path),
2760 "step {step}"
2761 );
2762 if let Some(parent) = expected {
2763 let children: Vec<_> = (0..next(3))
2764 .map(|child| PendingDirectory {
2765 path: parent.path.join(["a", "b", "c"][child]),
2766 depth: parent.depth + 1,
2767 })
2768 .collect();
2769 reference.extend(children.iter().cloned());
2770 frontier.extend(children);
2771 }
2772 }
2773 5..=7 => {
2774 let directory =
2775 PendingDirectory { path: PathBuf::from(PATHS[next(7)]), depth: 1 };
2776 reference.push_back(directory.clone());
2777 frontier.extend([directory]);
2778 }
2779 _ => {
2780 let mut chosen: Vec<_> =
2781 (0..next(4)).map(|_| PathBuf::from(PATHS[next(7)])).collect();
2782 chosen.sort();
2783 chosen.dedup();
2784 frontier.prioritize(chosen.clone());
2785 priorities = chosen;
2786 }
2787 }
2788 }
2789 while let Some(expected) = reference_pop(&mut reference, &priorities) {
2790 assert_eq!(frontier.pop().map(|directory| directory.path), Some(expected.path));
2791 }
2792 assert!(frontier.pop().is_none());
2793 }
2794
2795 #[test]
2796 fn reaching_a_file_limit_without_refusal_remains_complete() {
2797 let root = tempfile::tempdir().expect("temp root");
2798 std::fs::write(root.path().join("one"), b"1").expect("fixture");
2799 std::fs::write(root.path().join("two"), b"2").expect("fixture");
2800 let opened = open_fixture(
2801 root.path(),
2802 OpenOptions {
2803 budget: DiscoveryBudget { max_files: Some(2) },
2804 ..OpenOptions::default()
2805 },
2806 )
2807 .expect("opened root");
2808
2809 let state = wait_until_settled(&opened);
2810 assert_eq!(state.coverage, crate::Coverage::Complete);
2811 assert_eq!(state.progress.files_retained, 2);
2812 assert_eq!(
2813 opened.state.index.directory_complete(Path::new("")).expect("root completeness"),
2814 Some(true)
2815 );
2816 opened.close().expect("close");
2817 }
2818
2819 #[test]
2820 fn one_read_returns_lookup_state_and_version_from_one_boundary() {
2821 let (_root, opened) = opened(Arc::new(TestControls::default()));
2822 opened
2823 .state
2824 .index
2825 .apply(&Observation::new(vec![Op::Upsert {
2826 path: PathBuf::from("note.txt"),
2827 kind: EntryKind::File,
2828 attrs: crate::Attrs { size: 7, ..crate::Attrs::default() },
2829 }]))
2830 .expect("seed entry");
2831
2832 let response = opened
2833 .read(crate::ReadRequest {
2834 projections: vec![
2835 crate::ReadProjection::Lookup { path: PathBuf::from("note.txt") },
2836 crate::ReadProjection::Lookup { path: PathBuf::from("missing.txt") },
2837 ],
2838 ..crate::ReadRequest::default()
2839 })
2840 .expect("coherent read");
2841
2842 assert_eq!(response.version.sequence, opened.state.index.clock().expect("clock"));
2843 assert_eq!(response.state, opened.state.index.state().expect("state"));
2844 assert_eq!(response.results.len(), 2);
2845 assert!(matches!(
2846 &response.results[0],
2847 crate::ProjectionResult::Lookup(crate::Knowledge::Present(entry))
2848 if entry.path == Path::new("note.txt") && entry.attrs.size == 7
2849 ));
2850 assert!(matches!(
2851 response.results[1],
2852 crate::ProjectionResult::Lookup(crate::Knowledge::Absent)
2853 ));
2854 opened.close().expect("close");
2855 }
2856
2857 #[test]
2858 fn change_poll_returns_the_detached_exact_range_and_terminal_state() {
2859 let (_root, opened) = opened(Arc::new(TestControls::default()));
2860 let after = current_version(&opened);
2861 apply_and_notify(
2862 &opened,
2863 &Observation::new(vec![Op::Upsert {
2864 path: PathBuf::from("note.txt"),
2865 kind: EntryKind::File,
2866 attrs: crate::Attrs { size: 7, ..crate::Attrs::default() },
2867 }]),
2868 );
2869 let state_only = opened
2870 .state
2871 .index
2872 .transition_discovery(DiscoveryTransition::Begin)
2873 .expect("state-only commit");
2874 assert!(state_only.commit.as_ref().is_some_and(|commit| commit.changes.is_empty()));
2875 opened.state.journal.notify_commit();
2876
2877 let poll = opened
2878 .changes(crate::ChangeRequest { after, timeout: std::time::Duration::ZERO })
2879 .expect("immediate changes");
2880 let crate::ChangeOutcome::Changes { commits, impact } = &poll.outcome else {
2881 panic!("expected changes");
2882 };
2883 let detached = opened.state.index.since(after.sequence).expect("detached range");
2884 assert_eq!(commits, &detached.commits);
2885 assert_eq!(commits.len(), 2);
2886 assert!(commits[1].changes.is_empty(), "state-only commit remains observable");
2887 assert!(impact.domains.contains(&crate::ImpactDomain::State));
2888 assert_eq!(poll.cursor, poll.version);
2889 assert_eq!(poll.version.sequence, detached.clock);
2890 assert_eq!(poll.state, detached.state);
2891 assert_eq!(poll.work.commits_visited, 2);
2892 assert_eq!(poll.work.commits_returned, 2);
2893 opened.close().expect("close");
2894 }
2895
2896 #[test]
2897 fn idle_change_poll_waits_without_advancing_its_cursor() {
2898 let (_root, opened) = opened(Arc::new(TestControls::default()));
2899 let after = current_version(&opened);
2900 let poll = opened
2901 .changes(crate::ChangeRequest { after, timeout: std::time::Duration::from_millis(5) })
2902 .expect("idle poll");
2903
2904 assert!(matches!(poll.outcome, crate::ChangeOutcome::Idle));
2905 assert_eq!(poll.cursor, after);
2906 assert_eq!(poll.version, after);
2907 assert_eq!(poll.work, crate::Work::default());
2908 opened.close().expect("close");
2909 }
2910
2911 #[test]
2912 fn a_commit_at_the_wait_boundary_cannot_lose_its_wakeup() {
2913 let controls = Arc::new(TestControls::default());
2914 controls.gate(TestPoint::BeforeJournalWait).arm();
2915 let (_root, opened) = opened(Arc::clone(&controls));
2916 let after = current_version(&opened);
2917 let poller = opened.clone();
2918 let poll = thread::spawn(move || {
2919 poller
2920 .changes(crate::ChangeRequest { after, timeout: std::time::Duration::from_secs(5) })
2921 });
2922 controls.gate(TestPoint::BeforeJournalWait).wait_reached();
2923
2924 let (applied_sender, applied_receiver) = std::sync::mpsc::sync_channel(0);
2925 let committer = opened.clone();
2926 let commit = thread::spawn(move || {
2927 let outcome = committer
2928 .state
2929 .index
2930 .apply(&Observation::new(vec![Op::Upsert {
2931 path: PathBuf::from("arrived"),
2932 kind: EntryKind::File,
2933 attrs: crate::Attrs::default(),
2934 }]))
2935 .expect("commit at wait boundary");
2936 applied_sender.send(()).expect("report applied commit");
2937 if outcome.commit.is_some() {
2938 committer.state.journal.notify_commit();
2939 }
2940 });
2941 applied_receiver.recv_timeout(TEST_GATE_TIMEOUT).expect("commit applied");
2942 controls.gate(TestPoint::BeforeJournalWait).release();
2943
2944 commit.join().expect("committer");
2945 let poll = poll.join().expect("poller").expect("change poll");
2946 assert!(matches!(
2947 poll.outcome,
2948 crate::ChangeOutcome::Changes { ref commits, .. } if commits.len() == 1
2949 ));
2950 opened.close().expect("close");
2951 }
2952
2953 #[test]
2954 fn progressive_discovery_notifies_the_same_change_poll() {
2955 let root = tempfile::tempdir().expect("temp root");
2956 std::fs::write(root.path().join("discovered"), b"data").expect("fixture");
2957 let controls = Arc::new(TestControls::default());
2958 controls.gate(TestPoint::BeforeDiscovery).arm();
2959 let opened =
2960 OpenedIndex::open_for_test(root.path(), OpenOptions::default(), Arc::clone(&controls))
2961 .expect("opened root");
2962 let after = current_version(&opened);
2963 let poller = opened.clone();
2964 let poll = thread::spawn(move || {
2965 poller
2966 .changes(crate::ChangeRequest { after, timeout: std::time::Duration::from_secs(5) })
2967 });
2968
2969 controls.gate(TestPoint::BeforeDiscovery).release();
2970 let poll = poll.join().expect("poller").expect("discovery changes");
2971 assert!(matches!(
2972 poll.outcome,
2973 crate::ChangeOutcome::Changes { ref commits, .. } if !commits.is_empty()
2974 ));
2975 assert!(poll.version.sequence > after.sequence);
2976 wait_until_settled(&opened);
2977 opened.close().expect("close");
2978 }
2979
2980 #[test]
2981 fn terminal_state_only_discovery_commit_wakes_a_blocked_poll() {
2982 let root = tempfile::tempdir().expect("temp root");
2983 let controls = Arc::new(TestControls::default());
2984 controls.gate(TestPoint::AfterRootDirectory).arm();
2985 controls.gate(TestPoint::BeforeJournalWait).arm();
2986 let opened =
2987 OpenedIndex::open_for_test(root.path(), OpenOptions::default(), Arc::clone(&controls))
2988 .expect("opened root");
2989 controls.gate(TestPoint::AfterRootDirectory).wait_reached();
2990 let after = current_version(&opened);
2991 let poller = opened.clone();
2992 let poll = thread::spawn(move || {
2993 poller
2994 .changes(crate::ChangeRequest { after, timeout: std::time::Duration::from_secs(5) })
2995 });
2996 controls.gate(TestPoint::BeforeJournalWait).wait_reached();
2997 controls.gate(TestPoint::BeforeJournalWait).release();
2998 controls.gate(TestPoint::AfterRootDirectory).release();
2999
3000 let poll = poll.join().expect("poller").expect("terminal change");
3001 let crate::ChangeOutcome::Changes { commits, .. } = poll.outcome else {
3002 panic!("expected terminal change");
3003 };
3004 assert_eq!(commits.len(), 1);
3005 assert!(commits[0].changes.is_empty());
3006 assert!(commits[0].state.iter().any(|transition| matches!(
3007 transition,
3008 crate::StateTransition::IndexState {
3009 current: crate::IndexState { phase: crate::LifecyclePhase::Ready, .. },
3010 ..
3011 }
3012 )));
3013 assert_eq!(poll.state.phase, crate::LifecyclePhase::Ready);
3014 opened.close().expect("close");
3015 }
3016
3017 #[test]
3018 fn change_cursors_reject_foreign_identity_and_future_sequences() {
3019 let (_first_root, first) = opened(Arc::new(TestControls::default()));
3020 let (_second_root, second) = opened(Arc::new(TestControls::default()));
3021 let first_version = current_version(&first);
3022 let second_version = current_version(&second);
3023
3024 assert!(matches!(
3025 first.changes(crate::ChangeRequest {
3026 after: second_version,
3027 timeout: std::time::Duration::ZERO,
3028 }),
3029 Err(Error::ChangeCursorUnavailable { .. })
3030 ));
3031 let future = crate::EngineVersion {
3032 sequence: first_version.sequence.checked_next().expect("future sequence"),
3033 ..first_version
3034 };
3035 assert!(matches!(
3036 first.changes(crate::ChangeRequest {
3037 after: future,
3038 timeout: std::time::Duration::ZERO,
3039 }),
3040 Err(Error::ChangeCursorUnavailable { .. })
3041 ));
3042 first.close().expect("first close");
3043 second.close().expect("second close");
3044 }
3045
3046 fn opened_at_the_minimum_journal_budget() -> (tempfile::TempDir, OpenedIndex) {
3048 let controls = Arc::new(TestControls::default());
3049 controls.discovery_disabled.store(true, Ordering::Release);
3050 let root = tempfile::tempdir().expect("temp root");
3051 let opened = OpenedIndex::open_for_test(
3052 root.path(),
3053 OpenOptions {
3054 journal_capacity_bytes: crate::MIN_JOURNAL_CAPACITY_BYTES,
3055 ..OpenOptions::default()
3056 },
3057 controls,
3058 )
3059 .expect("opened root");
3060 (root, opened)
3061 }
3062
3063 #[test]
3066 fn the_minimum_journal_budget_delivers_a_burst_of_single_file_commits() {
3067 const BURST: usize = 64;
3068 let (_root, opened) = opened_at_the_minimum_journal_budget();
3069 let after = current_version(&opened);
3070 for index in 0..BURST {
3071 apply_and_notify(
3072 &opened,
3073 &Observation::new(vec![Op::Upsert {
3074 path: PathBuf::from(format!("file-{index:02}")),
3075 kind: EntryKind::File,
3076 attrs: crate::Attrs::default(),
3077 }]),
3078 );
3079 }
3080
3081 let poll = opened
3082 .changes(crate::ChangeRequest { after, timeout: std::time::Duration::ZERO })
3083 .expect("burst");
3084 let crate::ChangeOutcome::Changes { commits, .. } = &poll.outcome else {
3085 panic!("expected the burst's commits: {:?}", poll.outcome);
3086 };
3087 assert_eq!(commits.len(), BURST);
3088 opened.close().expect("close");
3089 }
3090
3091 #[test]
3092 fn a_slow_consumer_gets_one_coherent_all_dirty_reset() {
3093 let (_root, opened) = opened_at_the_minimum_journal_budget();
3094 let after = current_version(&opened);
3095 let outcome = apply_and_notify(
3097 &opened,
3098 &Observation::new(
3099 (0..=crate::MAX_DIRTY_PATHS)
3100 .map(|index| Op::Upsert {
3101 path: PathBuf::from(format!("{index:0>200}")),
3102 kind: EntryKind::File,
3103 attrs: crate::Attrs::default(),
3104 })
3105 .collect(),
3106 ),
3107 );
3108 let commit = outcome.commit.expect("effective commit");
3109 assert!(commit.retained_cost() > crate::MIN_JOURNAL_CAPACITY_BYTES);
3110
3111 let poll = opened
3112 .changes(crate::ChangeRequest { after, timeout: std::time::Duration::ZERO })
3113 .expect("consumer reset");
3114 let crate::ChangeOutcome::Reset { impact } = &poll.outcome else {
3115 panic!("expected reset");
3116 };
3117 assert!(impact.all_dirty);
3118 assert!(impact.dirty_paths.is_empty());
3119 assert_eq!(impact.domains.len(), 6);
3120 assert_eq!(poll.cursor, poll.version);
3121 assert_eq!(poll.state, opened.state.index.state().expect("terminal state"));
3122 assert!(opened.state.index.since(after.sequence).expect("history").truncated);
3123 opened.close().expect("close");
3124 }
3125
3126 #[test]
3127 fn close_wakes_a_blocked_change_poll() {
3128 let controls = Arc::new(TestControls::default());
3129 controls.gate(TestPoint::BeforeJournalWait).arm();
3130 let (_root, opened) = opened(Arc::clone(&controls));
3131 let after = current_version(&opened);
3132 let poller = opened.clone();
3133 let poll = thread::spawn(move || {
3134 poller.changes(crate::ChangeRequest {
3135 after,
3136 timeout: std::time::Duration::from_secs(60),
3137 })
3138 });
3139 controls.gate(TestPoint::BeforeJournalWait).wait_reached();
3140 controls.gate(TestPoint::BeforeJournalWait).release();
3141 opened.close().expect("close");
3142
3143 assert!(matches!(poll.join().expect("poller"), Err(Error::OpenedIndexClosed)));
3144 }
3145
3146 #[test]
3147 fn change_invalidations_fail_closed_at_the_existing_path_bound() {
3148 let (_root, opened) = opened(Arc::new(TestControls::default()));
3149 let after = current_version(&opened);
3150 let ops = (0..=crate::MAX_DIRTY_PATHS)
3151 .map(|index| Op::Upsert {
3152 path: PathBuf::from(format!("entry-{index}")),
3153 kind: EntryKind::File,
3154 attrs: crate::Attrs::default(),
3155 })
3156 .collect();
3157 apply_and_notify(&opened, &Observation::new(ops));
3158
3159 let poll = opened
3160 .changes(crate::ChangeRequest { after, timeout: std::time::Duration::ZERO })
3161 .expect("bounded invalidation");
3162 assert!(matches!(
3163 poll.outcome,
3164 crate::ChangeOutcome::Changes {
3165 ref commits,
3166 impact: crate::Impact { all_dirty: true, ref dirty_paths, .. },
3167 } if commits.len() == 1 && dirty_paths.is_empty()
3168 ));
3169 opened.close().expect("close");
3170 }
3171
3172 #[test]
3173 fn lookup_uses_directory_completeness_before_global_discovery_settles() {
3174 let (_root, opened) = opened(Arc::new(TestControls::default()));
3175 opened
3176 .state
3177 .index
3178 .transition_discovery(DiscoveryTransition::Begin)
3179 .expect("begin discovery");
3180 opened
3181 .state
3182 .index
3183 .apply(&Observation::new(vec![Op::Upsert {
3184 path: PathBuf::from("known"),
3185 kind: EntryKind::Dir,
3186 attrs: crate::Attrs::default(),
3187 }]))
3188 .expect("seed directory");
3189
3190 let lookup = || {
3191 opened
3192 .read(crate::ReadRequest {
3193 projections: vec![crate::ReadProjection::Lookup {
3194 path: PathBuf::from("known/missing"),
3195 }],
3196 ..crate::ReadRequest::default()
3197 })
3198 .expect("lookup")
3199 .results
3200 .into_iter()
3201 .next()
3202 .expect("lookup result")
3203 };
3204 assert!(matches!(
3205 lookup(),
3206 crate::ProjectionResult::Lookup(crate::Knowledge::Unknown {
3207 reason: crate::CoverageReason::Building
3208 })
3209 ));
3210
3211 opened
3212 .state
3213 .index
3214 .apply_discovery(
3215 &Observation::new(Vec::new()),
3216 DiscoveryCommit {
3217 directory_complete: Some(PathBuf::from("known")),
3218 transition: None,
3219 },
3220 )
3221 .expect("complete directory");
3222 assert!(matches!(lookup(), crate::ProjectionResult::Lookup(crate::Knowledge::Absent)));
3223 assert!(matches!(
3224 opened.state.index.state().expect("state").coverage,
3225 crate::Coverage::Partial(_)
3226 ));
3227 opened.close().expect("close");
3228 }
3229
3230 #[test]
3233 fn directory_completion_publishes_the_canonical_relative_path() {
3234 let (_root, opened) = opened(Arc::new(TestControls::default()));
3235 opened
3236 .state
3237 .index
3238 .transition_discovery(DiscoveryTransition::Begin)
3239 .expect("begin discovery");
3240 opened
3241 .state
3242 .index
3243 .apply(&Observation::new(vec![Op::Upsert {
3244 path: PathBuf::from("known"),
3245 kind: EntryKind::Dir,
3246 attrs: crate::Attrs::default(),
3247 }]))
3248 .expect("seed directory");
3249
3250 let outcome = opened
3251 .state
3252 .index
3253 .apply_discovery(
3254 &Observation::new(Vec::new()),
3255 DiscoveryCommit {
3256 directory_complete: Some(PathBuf::from("./known")),
3257 transition: None,
3258 },
3259 )
3260 .expect("complete directory");
3261
3262 let commit = outcome.commit.expect("completion commit");
3263 assert!(
3264 commit.state.contains(&crate::StateTransition::DirectoryComplete {
3265 path: PathBuf::from("known"),
3266 }),
3267 "{:?}",
3268 commit.state
3269 );
3270 assert_eq!(
3271 opened.state.index.directory_complete(Path::new("known")).expect("lookup"),
3272 Some(true)
3273 );
3274 opened.close().expect("close");
3275 }
3276
3277 #[test]
3286 fn a_path_of_the_wrong_kind_refuses_its_projection_and_the_read_still_answers() {
3287 let (_root, opened) = opened(Arc::new(TestControls::default()));
3288 let file = |path: &str| Op::Upsert {
3289 path: PathBuf::from(path),
3290 kind: EntryKind::File,
3291 attrs: crate::Attrs { size: 5, ..crate::Attrs::default() },
3292 };
3293 opened
3294 .state
3295 .index
3296 .apply(&Observation::new(vec![
3297 file("a"),
3298 file("README.md"),
3299 Op::Upsert {
3300 path: PathBuf::from("dir"),
3301 kind: EntryKind::Dir,
3302 attrs: crate::Attrs::default(),
3303 },
3304 ]))
3305 .expect("seed tree");
3306 let page = crate::PageRequest { limit: 16, max_work: 64 };
3307 let tree = |path: &str| crate::ReadProjection::Tree {
3308 path: PathBuf::from(path),
3309 depth: crate::query::Bound::Limit(1),
3310 include_ignored: true,
3311 page,
3312 };
3313 let before = opened
3315 .read(crate::ReadRequest { projections: vec![tree("dir")], ..Default::default() })
3316 .expect("tree of a directory");
3317 assert!(matches!(
3318 before.results.as_slice(),
3319 [crate::ProjectionResult::Tree(crate::Knowledge::Present(_))]
3320 ));
3321
3322 opened.state.index.apply(&Observation::new(vec![file("dir")])).expect("dir became a file");
3323 assert_eq!(opened.state.index.state().expect("state").coverage, crate::Coverage::Complete);
3324 let response = opened
3325 .read(crate::ReadRequest {
3326 projections: vec![
3327 crate::ReadProjection::Lookup { path: PathBuf::from("a") },
3328 tree("dir"),
3329 crate::ReadProjection::RollUp { path: PathBuf::from("README.md") },
3330 crate::ReadProjection::Lookup { path: PathBuf::from("dir") },
3331 ],
3332 ..crate::ReadRequest::default()
3333 })
3334 .expect("a refusal does not fail the read");
3335 match response.results.as_slice() {
3336 [
3337 crate::ProjectionResult::Lookup(crate::Knowledge::Present(a)),
3338 crate::ProjectionResult::Refused(crate::ProjectionRefusal::NotADirectory {
3339 path: tree_path,
3340 }),
3341 crate::ProjectionResult::Refused(crate::ProjectionRefusal::NotADirectory {
3342 path: rollup_path,
3343 }),
3344 crate::ProjectionResult::Lookup(crate::Knowledge::Present(dir)),
3345 ] => {
3346 assert_eq!(a.path, Path::new("a"));
3347 assert_eq!(tree_path, Path::new("dir"));
3348 assert_eq!(rollup_path, Path::new("README.md"));
3349 assert_eq!(dir.kind, EntryKind::File);
3350 }
3351 other => panic!("each projection answers for itself: {other:?}"),
3352 }
3353 assert_eq!(response.work.rows_returned, 2, "a refusal returns no rows");
3354
3355 assert!(matches!(
3357 opened
3358 .read(crate::ReadRequest {
3359 projections: vec![crate::ReadProjection::RollUp {
3360 path: PathBuf::from("README.md/inner"),
3361 }],
3362 ..crate::ReadRequest::default()
3363 })
3364 .expect("rollup below a file")
3365 .results[0],
3366 crate::ProjectionResult::RollUp(crate::Knowledge::Absent)
3367 ));
3368 opened.close().expect("close");
3369 }
3370
3371 #[cfg(unix)]
3373 fn opened_with_escaped_names() -> (tempfile::TempDir, OpenedIndex) {
3374 use std::os::unix::ffi::OsStrExt;
3375 let (root, opened) = opened(Arc::new(TestControls::default()));
3376 let native = |bytes: &[u8]| PathBuf::from(std::ffi::OsStr::from_bytes(bytes));
3377 let file = |path: PathBuf| Op::Upsert {
3378 path,
3379 kind: EntryKind::File,
3380 attrs: crate::Attrs { size: 3, ..crate::Attrs::default() },
3381 };
3382 let dir = |path: PathBuf| Op::Upsert {
3383 path,
3384 kind: EntryKind::Dir,
3385 attrs: crate::Attrs::default(),
3386 };
3387 opened
3388 .state
3389 .index
3390 .apply(&Observation::new(vec![
3391 dir(native(b"x\xff")),
3392 file(native(b"x\xff/inner.txt")),
3393 file(PathBuf::from("100%.txt")),
3394 dir(PathBuf::from("src")),
3395 file(PathBuf::from("src/lib.rs")),
3396 ]))
3397 .expect("seed escaped names");
3398 (root, opened)
3399 }
3400
3401 #[cfg(unix)]
3403 fn admitted_files(result: &crate::ProjectionResult) -> std::collections::BTreeSet<String> {
3404 match result {
3405 crate::ProjectionResult::Flat(page) => {
3406 assert!(page.next.is_none(), "one page holds the fixture");
3407 page.rows
3408 .iter()
3409 .filter(|row| row.kind == EntryKind::File)
3410 .map(|row| row.portable_path.as_str().to_owned())
3411 .collect()
3412 }
3413 crate::ProjectionResult::Report(report) => match report.sections.as_slice() {
3414 [crate::query::Section::Files { rows, .. }] => rows
3415 .iter()
3416 .filter(|row| row.kind == EntryKind::File)
3417 .map(|row| read::portable_path(&row.path).as_str().to_owned())
3418 .collect(),
3419 other => panic!("a files report: {other:?}"),
3420 },
3421 other => panic!("a page or a report: {other:?}"),
3422 }
3423 }
3424
3425 #[cfg(unix)]
3433 #[test]
3434 fn every_projection_filters_by_the_portable_identity_a_page_returns() {
3435 let (_root, opened) = opened_with_escaped_names();
3436 let glob = |source: &str| crate::query::Pattern::parse(source).expect("pattern");
3437 let names = |values: &[&str]| values.iter().map(ToString::to_string).collect::<Vec<_>>();
3438 let cases: Vec<(&str, crate::query::EntrySelection, &[&str])> = vec![
3439 (
3440 "an exact name, escaped",
3441 crate::query::EntrySelection {
3442 exact_names: names(&["100%25.txt"]),
3443 ..Default::default()
3444 },
3445 &["100%25.txt"],
3446 ),
3447 (
3448 "an exact name in its native spelling",
3449 crate::query::EntrySelection {
3450 exact_names: names(&["100%.txt"]),
3451 ..Default::default()
3452 },
3453 &[],
3454 ),
3455 (
3456 "a non-UTF-8 ancestor",
3457 crate::query::EntrySelection {
3458 ancestor_names: names(&["x%FF"]),
3459 ..Default::default()
3460 },
3461 &["x%FF/inner.txt"],
3462 ),
3463 (
3464 "a terminal suffix below that ancestor",
3465 crate::query::EntrySelection {
3466 terminal_extensions: names(&[".txt"]),
3467 ancestor_names: names(&["x%FF"]),
3468 ..Default::default()
3469 },
3470 &["x%FF/inner.txt"],
3471 ),
3472 (
3473 "an anchored glob through the escaped directory",
3474 crate::query::EntrySelection {
3475 query: crate::query::Selection {
3476 include: vec![glob("x%FF/*")],
3477 ..Default::default()
3478 },
3479 ..Default::default()
3480 },
3481 &["x%FF/inner.txt"],
3482 ),
3483 (
3484 "an unanchored glob on an escaped name",
3485 crate::query::EntrySelection {
3486 query: crate::query::Selection {
3487 include: vec![glob("100%25.txt")],
3488 ..Default::default()
3489 },
3490 ..Default::default()
3491 },
3492 &["100%25.txt"],
3493 ),
3494 (
3495 "an unanchored glob in the native spelling",
3496 crate::query::EntrySelection {
3497 query: crate::query::Selection {
3498 include: vec![glob("100%.txt")],
3499 ..Default::default()
3500 },
3501 ..Default::default()
3502 },
3503 &[],
3504 ),
3505 (
3506 "an exclusion by escaped directory",
3507 crate::query::EntrySelection {
3508 query: crate::query::Selection {
3509 exclude: vec![glob("x%FF/**")],
3510 ..Default::default()
3511 },
3512 ..Default::default()
3513 },
3514 &["100%25.txt", "src/lib.rs"],
3515 ),
3516 ];
3517 for (case, selection, expected) in cases {
3518 let expected: std::collections::BTreeSet<String> =
3519 expected.iter().map(ToString::to_string).collect();
3520 let report = crate::ReadProjection::Report(crate::ReportRequest {
3521 query: crate::query::Query {
3522 views: vec![crate::query::ViewSpec::Files],
3523 selection: selection.query.clone(),
3524 ..crate::query::Query::default()
3525 },
3526 now: std::time::SystemTime::UNIX_EPOCH,
3527 max_work: 1_000,
3528 });
3529 let response = opened
3530 .read(crate::ReadRequest {
3531 projections: vec![
3532 crate::ReadProjection::Flat {
3533 selection: selection.clone(),
3534 shape: crate::RowShape::Compact,
3535 page: crate::PageRequest { limit: 64, max_work: 1_000 },
3536 },
3537 crate::ReadProjection::Aggregate {
3538 selection: crate::query::EntrySelection {
3539 query: crate::query::Selection {
3540 kinds: vec![EntryKind::File],
3541 ..selection.query.clone()
3542 },
3543 ..selection.clone()
3544 },
3545 count_cap: 64,
3546 max_work: 1_000,
3547 },
3548 report,
3549 ],
3550 ..crate::ReadRequest::default()
3551 })
3552 .expect(case);
3553 assert_eq!(admitted_files(&response.results[0]), expected, "flat: {case}");
3554 assert!(
3555 matches!(
3556 response.results[1],
3557 crate::ProjectionResult::Aggregate(crate::CountResult::Exact(count))
3558 if count == expected.len() as u64
3559 ),
3560 "aggregate: {case}: {:?}",
3561 response.results[1]
3562 );
3563 if selection.exact_names.is_empty()
3566 && selection.ancestor_names.is_empty()
3567 && selection.terminal_extensions.is_empty()
3568 {
3569 assert_eq!(admitted_files(&response.results[2]), expected, "report: {case}");
3570 }
3571 }
3572 opened.close().expect("close");
3573 }
3574
3575 #[cfg(unix)]
3577 #[test]
3578 fn a_path_from_a_page_passes_back_into_a_filter_unchanged() {
3579 let (_root, opened) = opened_with_escaped_names();
3580 let flat = |selection: crate::query::EntrySelection| {
3581 let response = opened
3582 .read(crate::ReadRequest {
3583 projections: vec![crate::ReadProjection::Flat {
3584 selection,
3585 shape: crate::RowShape::Compact,
3586 page: crate::PageRequest { limit: 64, max_work: 1_000 },
3587 }],
3588 ..crate::ReadRequest::default()
3589 })
3590 .expect("flat page");
3591 admitted_files(&response.results[0])
3592 };
3593 let every = flat(crate::query::EntrySelection::default());
3594 assert_eq!(
3595 every,
3596 ["100%25.txt", "src/lib.rs", "x%FF/inner.txt"].map(String::from).into(),
3597 "the page names every file by its portable path"
3598 );
3599 for shown in &every {
3600 let only: std::collections::BTreeSet<String> = [shown.clone()].into();
3601 let (ancestors, name) = match shown.rsplit_once('/') {
3602 Some((ancestors, name)) => (Some(ancestors), name),
3603 None => (None, shown.as_str()),
3604 };
3605 let anchored = crate::query::Pattern::parse(&format!("**/{shown}")).expect("glob");
3606 assert_eq!(
3607 flat(crate::query::EntrySelection {
3608 query: crate::query::Selection {
3609 include: vec![anchored],
3610 ..Default::default()
3611 },
3612 ..Default::default()
3613 }),
3614 only,
3615 "the whole path as a glob: {shown}"
3616 );
3617 assert_eq!(
3618 flat(crate::query::EntrySelection {
3619 exact_names: vec![name.to_owned()],
3620 ..Default::default()
3621 }),
3622 only,
3623 "its name as an exact name: {shown}"
3624 );
3625 if let Some(ancestors) = ancestors {
3626 let mut selection = crate::query::EntrySelection::default();
3627 selection.admit_ancestor_name(ancestors).expect("a page component is a valid name");
3628 assert_eq!(flat(selection), only, "its parent as an ancestor name: {shown}");
3629 }
3630 }
3631 opened.close().expect("close");
3632 }
3633
3634 #[test]
3636 fn a_read_refuses_a_hand_written_selection_the_constructors_would_refuse() {
3637 let (_root, opened) = opened(Arc::new(TestControls::default()));
3638 for selection in [
3639 crate::query::EntrySelection {
3640 terminal_extensions: vec!["rs".to_string()],
3641 ..Default::default()
3642 },
3643 crate::query::EntrySelection {
3644 ancestor_names: vec!["..".to_string()],
3645 ..Default::default()
3646 },
3647 ] {
3648 let read = opened.read(crate::ReadRequest {
3649 projections: vec![
3650 crate::ReadProjection::Lookup { path: PathBuf::new() },
3651 crate::ReadProjection::Aggregate { selection, count_cap: 8, max_work: 64 },
3652 ],
3653 ..crate::ReadRequest::default()
3654 });
3655 assert!(matches!(read, Err(Error::InvalidValue { .. })), "{read:?}");
3656 }
3657 opened.close().expect("close");
3658 }
3659
3660 #[test]
3666 fn a_page_whose_continuation_cannot_be_retained_refuses_alone() {
3667 let (_root, opened) = opened(Arc::new(TestControls::default()));
3668 let file = |path: &str| Op::Upsert {
3669 path: PathBuf::from(path),
3670 kind: EntryKind::File,
3671 attrs: crate::Attrs { size: 1, ..crate::Attrs::default() },
3672 };
3673 opened
3674 .state
3675 .index
3676 .apply(&Observation::new(vec![file("a.txt"), file("b.txt")]))
3677 .expect("seed tree");
3678 let mut exact_names = vec!["a.txt".to_string(), "b.txt".to_string()];
3681 exact_names.extend((0..4_000).map(|number| format!("unused-{number:05}.txt")));
3682 let selection = crate::query::EntrySelection { exact_names, ..Default::default() };
3683 let flat = crate::ReadProjection::Flat {
3684 selection,
3685 shape: crate::RowShape::Compact,
3686 page: crate::PageRequest { limit: 1, max_work: 64 },
3687 };
3688 let response = opened
3689 .read(crate::ReadRequest {
3690 projections: vec![
3691 flat,
3692 crate::ReadProjection::Lookup { path: PathBuf::from("b.txt") },
3693 ],
3694 ..crate::ReadRequest::default()
3695 })
3696 .expect("a refused page does not fail the read");
3697 match response.results.as_slice() {
3698 [
3699 crate::ProjectionResult::Refused(
3700 crate::ProjectionRefusal::ContinuationRecordLimit { attempted, limit },
3701 ),
3702 crate::ProjectionResult::Lookup(crate::Knowledge::Present(_)),
3703 ] => {
3704 assert_eq!(*limit, crate::MAX_CONTINUATION_RECORD_BYTES);
3705 assert!(attempted > limit, "{attempted} > {limit}");
3706 }
3707 other => panic!("the page refuses and the lookup answers: {other:?}"),
3708 }
3709 assert_eq!(response.work.rows_returned, 1, "the refused page returns no rows");
3710 assert_eq!(
3711 opened.state.continuations.lock().expect("table").len(),
3712 0,
3713 "a refused page retains nothing"
3714 );
3715 opened.close().expect("close");
3716 }
3717
3718 #[test]
3719 fn mixed_read_preserves_projection_order_and_uses_maintained_rollups() {
3720 let (_root, opened) = opened(Arc::new(TestControls::default()));
3721 opened
3722 .state
3723 .index
3724 .apply(&Observation::new(vec![Op::Upsert {
3725 path: PathBuf::from("note.txt"),
3726 kind: EntryKind::File,
3727 attrs: crate::Attrs { size: 7, ..crate::Attrs::default() },
3728 }]))
3729 .expect("seed entry");
3730
3731 let response = opened
3732 .read(crate::ReadRequest {
3733 projections: vec![
3734 crate::ReadProjection::Diagnostics,
3735 crate::ReadProjection::RollUp { path: PathBuf::new() },
3736 ],
3737 ..crate::ReadRequest::default()
3738 })
3739 .expect("coherent read");
3740
3741 assert!(matches!(
3742 &response.results[0],
3743 crate::ProjectionResult::Diagnostics(diagnostics)
3744 if diagnostics.root == opened.state.root && diagnostics.entries == 2
3745 ));
3746 assert!(matches!(
3747 &response.results[1],
3748 crate::ProjectionResult::RollUp(crate::Knowledge::Present(rollup))
3749 if rollup.all.files == 1 && rollup.all.bytes == 7
3750 ));
3751 assert_eq!(response.work.maintained_index_work, 1);
3752 opened.close().expect("close");
3753 }
3754
3755 #[test]
3756 fn tree_pages_are_directory_first_and_resume_at_the_same_version() {
3757 let (_root, opened) = opened(Arc::new(TestControls::default()));
3758 opened
3759 .state
3760 .index
3761 .apply(&Observation::new(vec![
3762 Op::Upsert {
3763 path: PathBuf::from("z-dir"),
3764 kind: EntryKind::Dir,
3765 attrs: crate::Attrs::default(),
3766 },
3767 Op::Upsert {
3768 path: PathBuf::from("a.txt"),
3769 kind: EntryKind::File,
3770 attrs: crate::Attrs { size: 1, ..crate::Attrs::default() },
3771 },
3772 Op::Upsert {
3773 path: PathBuf::from("b.txt"),
3774 kind: EntryKind::File,
3775 attrs: crate::Attrs { size: 2, ..crate::Attrs::default() },
3776 },
3777 ]))
3778 .expect("seed entries");
3779
3780 let first = opened
3781 .read(crate::ReadRequest {
3782 projections: vec![crate::ReadProjection::Tree {
3783 path: PathBuf::new(),
3784 depth: crate::query::Bound::Limit(1),
3785 include_ignored: true,
3786 page: crate::PageRequest { limit: 2, max_work: 4 },
3787 }],
3788 ..crate::ReadRequest::default()
3789 })
3790 .expect("first page");
3791 let crate::ProjectionResult::Tree(crate::Knowledge::Present(first_page)) =
3792 &first.results[0]
3793 else {
3794 panic!("tree page");
3795 };
3796 assert_eq!(
3797 first_page.rows.iter().map(|row| row.path.as_path()).collect::<Vec<_>>(),
3798 vec![Path::new("z-dir"), Path::new("a.txt")]
3799 );
3800 let continuation = first_page.next.expect("more rows");
3801
3802 let second = opened
3803 .read(crate::ReadRequest {
3804 projections: vec![crate::ReadProjection::Continue {
3805 continuation,
3806 page: crate::PageRequest { limit: 2, max_work: 2 },
3807 }],
3808 expected: Some(first.version),
3809 })
3810 .expect("second page");
3811 let crate::ProjectionResult::Tree(crate::Knowledge::Present(second_page)) =
3812 &second.results[0]
3813 else {
3814 panic!("continued tree page");
3815 };
3816 assert_eq!(
3817 second_page.rows.iter().map(|row| row.path.as_path()).collect::<Vec<_>>(),
3818 vec![Path::new("b.txt")]
3819 );
3820 assert!(second_page.next.is_none());
3821 assert_eq!(first.version, second.version);
3822 opened.close().expect("close");
3823 }
3824
3825 #[test]
3826 fn flat_pages_follow_complete_portable_path_order_without_rescanning() {
3827 let (_root, opened) = opened(Arc::new(TestControls::default()));
3828 opened
3829 .state
3830 .index
3831 .apply(&Observation::new(vec![
3832 Op::Upsert {
3833 path: PathBuf::from("c.txt"),
3834 kind: EntryKind::File,
3835 attrs: crate::Attrs::default(),
3836 },
3837 Op::Upsert {
3838 path: PathBuf::from("a.txt"),
3839 kind: EntryKind::File,
3840 attrs: crate::Attrs::default(),
3841 },
3842 Op::Upsert {
3843 path: PathBuf::from("b.txt"),
3844 kind: EntryKind::File,
3845 attrs: crate::Attrs::default(),
3846 },
3847 ]))
3848 .expect("seed entries");
3849
3850 let first = opened
3851 .read(crate::ReadRequest {
3852 projections: vec![crate::ReadProjection::Flat {
3853 selection: crate::query::EntrySelection::default(),
3854 shape: crate::RowShape::Compact,
3855 page: crate::PageRequest { limit: 2, max_work: 3 },
3856 }],
3857 ..crate::ReadRequest::default()
3858 })
3859 .expect("first page");
3860 let crate::ProjectionResult::Flat(first_page) = &first.results[0] else {
3861 panic!("flat page");
3862 };
3863 assert_eq!(
3864 first_page.rows.iter().map(|row| row.portable_path.as_str()).collect::<Vec<_>>(),
3865 vec!["a.txt", "b.txt"]
3866 );
3867 let continuation = first_page.next.expect("more rows");
3868
3869 let second = opened
3870 .read(crate::ReadRequest {
3871 projections: vec![crate::ReadProjection::Continue {
3872 continuation,
3873 page: crate::PageRequest { limit: 2, max_work: 2 },
3874 }],
3875 expected: Some(first.version),
3876 })
3877 .expect("second page");
3878 let crate::ProjectionResult::Flat(second_page) = &second.results[0] else {
3879 panic!("continued flat page");
3880 };
3881 assert_eq!(
3882 second_page.rows.iter().map(|row| row.portable_path.as_str()).collect::<Vec<_>>(),
3883 vec!["c.txt"]
3884 );
3885 assert!(second_page.next.is_none());
3886 assert!(second.work.rows_visited <= 2, "continuation resumed from retained position");
3887 opened.close().expect("close");
3888 }
3889
3890 #[test]
3899 fn a_full_flat_page_survives_any_budget_that_filled_it() {
3900 let (_root, opened) = opened(Arc::new(TestControls::default()));
3901 let file = |name: &str| Op::Upsert {
3902 path: PathBuf::from(name),
3903 kind: EntryKind::File,
3904 attrs: crate::Attrs::default(),
3905 };
3906 let mut ops = vec![file("a0.rs"), file("a1.rs")];
3907 ops.extend((0..5).map(|i| file(&format!("b{i}.txt"))));
3908 ops.push(file("c.rs"));
3909 opened.state.index.apply(&Observation::new(ops)).expect("seed entries");
3910 let selection = crate::query::EntrySelection {
3911 query: crate::query::Selection {
3912 include: vec![crate::query::Pattern::parse("*.rs").expect("pattern")],
3913 ..crate::query::Selection::default()
3914 },
3915 ..crate::query::EntrySelection::default()
3916 };
3917 let rows = |page: &crate::FlatPage| {
3918 page.rows.iter().map(|row| row.portable_path.as_str().to_string()).collect::<Vec<_>>()
3919 };
3920
3921 for max_work in 1..=12 {
3922 let first = opened
3923 .read(crate::ReadRequest {
3924 projections: vec![crate::ReadProjection::Flat {
3925 selection: selection.clone(),
3926 shape: crate::RowShape::Compact,
3927 page: crate::PageRequest { limit: 2, max_work },
3928 }],
3929 ..crate::ReadRequest::default()
3930 })
3931 .expect("first page");
3932 if max_work < 2 {
3933 assert!(
3935 matches!(first.results[0], crate::ProjectionResult::Limit(_)),
3936 "max_work {max_work}: {:?}",
3937 first.results[0]
3938 );
3939 continue;
3940 }
3941 let crate::ProjectionResult::Flat(first_page) = &first.results[0] else {
3942 panic!("max_work {max_work}: a full page was refused: {:?}", first.results[0]);
3943 };
3944 assert_eq!(rows(first_page), ["a0.rs", "a1.rs"], "max_work {max_work}");
3945 assert!(first.work.rows_visited <= max_work, "max_work {max_work}: {:?}", first.work);
3946 let continuation = first_page.next.expect("the page stopped with rows left");
3947
3948 let second = opened
3949 .read(crate::ReadRequest {
3950 projections: vec![crate::ReadProjection::Continue {
3951 continuation,
3952 page: crate::PageRequest { limit: 2, max_work: crate::MAX_PAGE_WORK },
3953 }],
3954 expected: Some(first.version),
3955 })
3956 .expect("second page");
3957 let crate::ProjectionResult::Flat(second_page) = &second.results[0] else {
3958 panic!("max_work {max_work}: continued flat page: {:?}", second.results[0]);
3959 };
3960 assert_eq!(rows(second_page), ["c.rs"], "max_work {max_work}");
3961 assert!(second_page.next.is_none(), "max_work {max_work}");
3962 }
3963 opened.close().expect("close");
3964 }
3965
3966 #[test]
3972 fn a_read_proceeds_while_another_read_is_projecting() {
3973 let controls = Arc::new(TestControls::default());
3974 let (_root, opened) = opened(Arc::clone(&controls));
3975 opened
3976 .state
3977 .index
3978 .apply(&Observation::new(vec![Op::Upsert {
3979 path: PathBuf::from("a.txt"),
3980 kind: EntryKind::File,
3981 attrs: crate::Attrs::default(),
3982 }]))
3983 .expect("seed entry");
3984 controls.gate(TestPoint::DuringTreeProjection).arm();
3985 let tree_reader = opened.clone();
3986 let tree = thread::spawn(move || {
3987 tree_reader.read(crate::ReadRequest {
3988 projections: vec![crate::ReadProjection::Tree {
3989 path: PathBuf::new(),
3990 depth: crate::query::Bound::Limit(1),
3991 include_ignored: true,
3992 page: crate::PageRequest { limit: 16, max_work: 64 },
3993 }],
3994 ..crate::ReadRequest::default()
3995 })
3996 });
3997 controls.gate(TestPoint::DuringTreeProjection).wait_reached();
3998
3999 let (sender, receiver) = std::sync::mpsc::channel();
4000 let lookup_reader = opened.clone();
4001 let lookup = thread::spawn(move || {
4002 let _ = sender.send(lookup_reader.read(crate::ReadRequest {
4003 projections: vec![crate::ReadProjection::Lookup { path: PathBuf::from("a.txt") }],
4004 ..crate::ReadRequest::default()
4005 }));
4006 });
4007 let concurrent = receiver.recv_timeout(TEST_GATE_TIMEOUT);
4008 controls.gate(TestPoint::DuringTreeProjection).release();
4009 let concurrent = concurrent.expect("a second read finished while the first projected");
4010 assert!(matches!(
4011 concurrent.expect("lookup").results[0],
4012 crate::ProjectionResult::Lookup(crate::Knowledge::Present(_))
4013 ));
4014 lookup.join().expect("lookup thread");
4015 tree.join().expect("tree thread").expect("tree page");
4016 opened.close().expect("close");
4017 }
4018
4019 #[test]
4021 fn a_read_racing_close_leaves_no_continuation_behind() {
4022 let controls = Arc::new(TestControls::default());
4023 let (_root, opened) = opened(Arc::clone(&controls));
4024 opened
4025 .state
4026 .index
4027 .apply(&Observation::new(vec![
4028 Op::Upsert {
4029 path: PathBuf::from("a.txt"),
4030 kind: EntryKind::File,
4031 attrs: crate::Attrs::default(),
4032 },
4033 Op::Upsert {
4034 path: PathBuf::from("b.txt"),
4035 kind: EntryKind::File,
4036 attrs: crate::Attrs::default(),
4037 },
4038 ]))
4039 .expect("seed entries");
4040 controls.gate(TestPoint::DuringTreeProjection).arm();
4041 let reader = opened.clone();
4042 let page = thread::spawn(move || {
4043 reader.read(crate::ReadRequest {
4044 projections: vec![crate::ReadProjection::Tree {
4045 path: PathBuf::new(),
4046 depth: crate::query::Bound::Limit(1),
4047 include_ignored: true,
4048 page: crate::PageRequest { limit: 1, max_work: 64 },
4049 }],
4050 ..crate::ReadRequest::default()
4051 })
4052 });
4053 controls.gate(TestPoint::DuringTreeProjection).wait_reached();
4054
4055 let (sender, receiver) = std::sync::mpsc::channel();
4056 let closer = opened.clone();
4057 let close = thread::spawn(move || {
4058 let _ = sender.send(closer.close());
4059 });
4060 let shutdown = receiver.recv_timeout(TEST_GATE_TIMEOUT);
4061 controls.gate(TestPoint::DuringTreeProjection).release();
4062 shutdown.expect("close did not wait for a read in progress").expect("close");
4063 close.join().expect("close thread");
4064
4065 assert!(matches!(page.join().expect("page thread"), Err(Error::OpenedIndexClosed)));
4066 assert_eq!(opened.state.continuations.lock().expect("continuations").len(), 0);
4067 }
4068
4069 #[test]
4070 fn flat_continuation_retains_its_normalized_native_query() {
4071 let (_root, opened) = opened(Arc::new(TestControls::default()));
4072 opened
4073 .state
4074 .index
4075 .apply(&Observation::new(vec![
4076 Op::Upsert {
4077 path: PathBuf::from("a.rs"),
4078 kind: EntryKind::File,
4079 attrs: crate::Attrs::default(),
4080 },
4081 Op::Upsert {
4082 path: PathBuf::from("b.txt"),
4083 kind: EntryKind::File,
4084 attrs: crate::Attrs::default(),
4085 },
4086 Op::Upsert {
4087 path: PathBuf::from("c.rs"),
4088 kind: EntryKind::File,
4089 attrs: crate::Attrs::default(),
4090 },
4091 ]))
4092 .expect("seed entries");
4093 let selection = crate::query::EntrySelection {
4094 query: crate::query::Selection {
4095 include: vec![crate::query::Pattern::parse("*.rs").expect("pattern")],
4096 ..crate::query::Selection::default()
4097 },
4098 ..crate::query::EntrySelection::default()
4099 };
4100
4101 let first = opened
4102 .read(crate::ReadRequest {
4103 projections: vec![crate::ReadProjection::Flat {
4104 selection,
4105 shape: crate::RowShape::Compact,
4106 page: crate::PageRequest { limit: 1, max_work: 3 },
4107 }],
4108 ..crate::ReadRequest::default()
4109 })
4110 .expect("first page");
4111 let crate::ProjectionResult::Flat(first_page) = &first.results[0] else {
4112 panic!("flat page");
4113 };
4114 assert_eq!(first_page.rows[0].portable_path.as_str(), "a.rs");
4115
4116 let second = opened
4117 .read(crate::ReadRequest {
4118 projections: vec![crate::ReadProjection::Continue {
4119 continuation: first_page.next.expect("continuation"),
4120 page: crate::PageRequest { limit: 1, max_work: 2 },
4121 }],
4122 expected: Some(first.version),
4123 })
4124 .expect("continued page");
4125 let crate::ProjectionResult::Flat(second_page) = &second.results[0] else {
4126 panic!("continued flat page");
4127 };
4128 assert_eq!(second_page.rows[0].portable_path.as_str(), "c.rs");
4129 assert!(second_page.next.is_none());
4130 assert_eq!(second.work.rows_visited, 2);
4134 opened.close().expect("close");
4135 }
4136
4137 #[test]
4140 fn continuations_are_single_use_version_pinned_handle_local_and_bounded() {
4141 let (_root, opened) = opened(Arc::new(TestControls::default()));
4142 opened
4143 .state
4144 .index
4145 .apply(&Observation::new(vec![
4146 Op::Upsert {
4147 path: PathBuf::from("a"),
4148 kind: EntryKind::File,
4149 attrs: crate::Attrs::default(),
4150 },
4151 Op::Upsert {
4152 path: PathBuf::from("b"),
4153 kind: EntryKind::File,
4154 attrs: crate::Attrs::default(),
4155 },
4156 ]))
4157 .expect("seed entries");
4158 let page = crate::PageRequest { limit: 1, max_work: 2 };
4159 let new_token = || {
4160 let response = opened
4161 .read(crate::ReadRequest {
4162 projections: vec![crate::ReadProjection::Flat {
4163 selection: crate::query::EntrySelection::default(),
4164 shape: crate::RowShape::Compact,
4165 page,
4166 }],
4167 ..crate::ReadRequest::default()
4168 })
4169 .expect("first page");
4170 let crate::ProjectionResult::Flat(result) = &response.results[0] else {
4171 panic!("flat page");
4172 };
4173 result.next.expect("continuation")
4174 };
4175
4176 let refuses_beside_a_lookup = |continuation| {
4179 let response = opened
4180 .read(crate::ReadRequest {
4181 projections: vec![
4182 crate::ReadProjection::Continue { continuation, page },
4183 crate::ReadProjection::Lookup { path: PathBuf::from("a") },
4184 ],
4185 ..crate::ReadRequest::default()
4186 })
4187 .expect("a refused continuation does not fail the read");
4188 matches!(
4189 response.results.as_slice(),
4190 [
4191 crate::ProjectionResult::Refused(
4192 crate::ProjectionRefusal::ContinuationUnavailable
4193 ),
4194 crate::ProjectionResult::Lookup(crate::Knowledge::Present(_)),
4195 ]
4196 )
4197 };
4198
4199 let replay = new_token();
4200 opened
4201 .read(crate::ReadRequest {
4202 projections: vec![crate::ReadProjection::Continue { continuation: replay, page }],
4203 ..crate::ReadRequest::default()
4204 })
4205 .expect("first continuation use");
4206 assert!(refuses_beside_a_lookup(replay), "a consumed token refuses its page");
4207
4208 let stale = new_token();
4209 opened
4210 .state
4211 .index
4212 .apply(&Observation::new(vec![Op::Upsert {
4213 path: PathBuf::from("c"),
4214 kind: EntryKind::File,
4215 attrs: crate::Attrs::default(),
4216 }]))
4217 .expect("advance version");
4218 assert!(matches!(
4219 opened.read(crate::ReadRequest {
4220 projections: vec![crate::ReadProjection::Continue { continuation: stale, page }],
4221 ..crate::ReadRequest::default()
4222 }),
4223 Err(Error::ContinuationStale { .. })
4224 ));
4225
4226 let retryable = new_token();
4227 let limited = opened
4230 .read(crate::ReadRequest {
4231 projections: vec![crate::ReadProjection::Continue {
4232 continuation: retryable,
4233 page: crate::PageRequest { limit: 2, max_work: 1 },
4234 }],
4235 ..crate::ReadRequest::default()
4236 })
4237 .expect("bounded continuation");
4238 assert!(matches!(limited.results[0], crate::ProjectionResult::Limit(_)));
4239 assert!(matches!(
4240 opened
4241 .read(crate::ReadRequest {
4242 projections: vec![crate::ReadProjection::Continue {
4243 continuation: retryable,
4244 page: crate::PageRequest { limit: 1, max_work: 2 },
4245 }],
4246 ..crate::ReadRequest::default()
4247 })
4248 .expect("retry continuation")
4249 .results[0],
4250 crate::ProjectionResult::Flat(_)
4251 ));
4252
4253 let foreign = new_token();
4254 let (_other_root, other) = self::opened(Arc::new(TestControls::default()));
4255 assert!(matches!(
4256 other.read(crate::ReadRequest {
4257 projections: vec![crate::ReadProjection::Continue { continuation: foreign, page }],
4258 ..crate::ReadRequest::default()
4259 }),
4260 Err(Error::ContinuationUnavailable)
4261 ));
4262
4263 let oldest = new_token();
4264 for _ in 0..super::continuation::MAX_CONTINUATIONS {
4265 let _ = new_token();
4266 }
4267 assert!(refuses_beside_a_lookup(oldest), "an evicted token refuses its page");
4268 other.close().expect("close other");
4269 opened.close().expect("close");
4270 }
4271
4272 #[test]
4288 fn tree_pages_are_breadth_first_across_levels() {
4289 let (_root, opened) = opened(Arc::new(TestControls::default()));
4290 opened
4291 .state
4292 .index
4293 .apply(&Observation::new(vec![
4294 Op::Upsert {
4295 path: PathBuf::from("a"),
4296 kind: EntryKind::Dir,
4297 attrs: crate::Attrs::default(),
4298 },
4299 Op::Upsert {
4300 path: PathBuf::from("b"),
4301 kind: EntryKind::Dir,
4302 attrs: crate::Attrs::default(),
4303 },
4304 Op::Upsert {
4305 path: PathBuf::from("z.txt"),
4306 kind: EntryKind::File,
4307 attrs: crate::Attrs { size: 1, ..crate::Attrs::default() },
4308 },
4309 Op::Upsert {
4310 path: PathBuf::from("a/a1"),
4311 kind: EntryKind::Dir,
4312 attrs: crate::Attrs::default(),
4313 },
4314 Op::Upsert {
4315 path: PathBuf::from("b/b1"),
4316 kind: EntryKind::Dir,
4317 attrs: crate::Attrs::default(),
4318 },
4319 Op::Upsert {
4320 path: PathBuf::from("a/a1/deep.txt"),
4321 kind: EntryKind::File,
4322 attrs: crate::Attrs { size: 2, ..crate::Attrs::default() },
4323 },
4324 ]))
4325 .expect("seed tree");
4326
4327 let rows = |depth: crate::query::Bound| -> Vec<String> {
4328 let response = opened
4329 .read(crate::ReadRequest {
4330 projections: vec![crate::ReadProjection::Tree {
4331 path: PathBuf::new(),
4332 depth,
4333 include_ignored: true,
4334 page: crate::PageRequest {
4335 limit: crate::MAX_PAGE_ROWS,
4336 max_work: crate::MAX_PAGE_WORK,
4337 },
4338 }],
4339 ..crate::ReadRequest::default()
4340 })
4341 .expect("tree read");
4342 let crate::ProjectionResult::Tree(crate::Knowledge::Present(page)) =
4343 &response.results[0]
4344 else {
4345 panic!("tree page");
4346 };
4347 page.rows.iter().map(|row| row.portable_path.as_str().to_owned()).collect()
4348 };
4349
4350 assert_eq!(rows(crate::query::Bound::Limit(1)), vec!["a", "b", "z.txt"]);
4353
4354 assert_eq!(rows(crate::query::Bound::Limit(2)), vec!["a", "b", "z.txt", "a/a1", "b/b1"]);
4356
4357 assert_eq!(
4360 rows(crate::query::Bound::All),
4361 vec!["a", "b", "z.txt", "a/a1", "b/b1", "a/a1/deep.txt"]
4362 );
4363
4364 opened.close().expect("close");
4365 }
4366
4367 #[test]
4374 fn tree_pages_resume_across_level_boundaries() {
4375 let (_root, opened) = opened(Arc::new(TestControls::default()));
4376 opened
4377 .state
4378 .index
4379 .apply(&Observation::new(vec![
4380 Op::Upsert {
4381 path: PathBuf::from("a"),
4382 kind: EntryKind::Dir,
4383 attrs: crate::Attrs::default(),
4384 },
4385 Op::Upsert {
4386 path: PathBuf::from("b"),
4387 kind: EntryKind::Dir,
4388 attrs: crate::Attrs::default(),
4389 },
4390 Op::Upsert {
4391 path: PathBuf::from("z.txt"),
4392 kind: EntryKind::File,
4393 attrs: crate::Attrs { size: 1, ..crate::Attrs::default() },
4394 },
4395 Op::Upsert {
4396 path: PathBuf::from("a/a1"),
4397 kind: EntryKind::Dir,
4398 attrs: crate::Attrs::default(),
4399 },
4400 Op::Upsert {
4401 path: PathBuf::from("b/b1"),
4402 kind: EntryKind::Dir,
4403 attrs: crate::Attrs::default(),
4404 },
4405 Op::Upsert {
4406 path: PathBuf::from("a/a1/deep.txt"),
4407 kind: EntryKind::File,
4408 attrs: crate::Attrs { size: 2, ..crate::Attrs::default() },
4409 },
4410 ]))
4411 .expect("seed tree");
4412
4413 let whole = vec!["a", "b", "z.txt", "a/a1", "b/b1", "a/a1/deep.txt"];
4414 let page = crate::PageRequest { limit: 1, max_work: crate::MAX_PAGE_WORK };
4415
4416 let mut seen: Vec<String> = Vec::new();
4417 let response = opened
4418 .read(crate::ReadRequest {
4419 projections: vec![crate::ReadProjection::Tree {
4420 path: PathBuf::new(),
4421 depth: crate::query::Bound::All,
4422 include_ignored: true,
4423 page,
4424 }],
4425 ..crate::ReadRequest::default()
4426 })
4427 .expect("first page");
4428 let crate::ProjectionResult::Tree(crate::Knowledge::Present(first)) = &response.results[0]
4429 else {
4430 panic!("tree page");
4431 };
4432 seen.extend(first.rows.iter().map(|row| row.portable_path.as_str().to_owned()));
4433 let mut continuation = first.next;
4434
4435 while let Some(token) = continuation {
4436 let response = opened
4437 .read(crate::ReadRequest {
4438 projections: vec![crate::ReadProjection::Continue {
4439 continuation: token,
4440 page,
4441 }],
4442 ..crate::ReadRequest::default()
4443 })
4444 .expect("resumed page");
4445 let crate::ProjectionResult::Tree(crate::Knowledge::Present(next)) =
4446 &response.results[0]
4447 else {
4448 panic!("tree page");
4449 };
4450 seen.extend(next.rows.iter().map(|row| row.portable_path.as_str().to_owned()));
4451 continuation = next.next;
4452 assert!(seen.len() <= whole.len(), "paging must terminate, saw {seen:?}");
4453 }
4454
4455 assert_eq!(seen, whole, "one row at a time reassembles the single-page order");
4456 opened.close().expect("close");
4457 }
4458
4459 #[test]
4472 fn a_tree_page_stopped_by_the_work_budget_is_resumable() {
4473 let (_root, opened) = opened(Arc::new(TestControls::default()));
4474 opened
4475 .state
4476 .index
4477 .apply(&Observation::new(vec![
4478 Op::Upsert {
4479 path: PathBuf::from("a"),
4480 kind: EntryKind::Dir,
4481 attrs: crate::Attrs::default(),
4482 },
4483 Op::Upsert {
4484 path: PathBuf::from("b"),
4485 kind: EntryKind::Dir,
4486 attrs: crate::Attrs::default(),
4487 },
4488 Op::Upsert {
4489 path: PathBuf::from("z.txt"),
4490 kind: EntryKind::File,
4491 attrs: crate::Attrs { size: 1, ..crate::Attrs::default() },
4492 },
4493 Op::Upsert {
4494 path: PathBuf::from("a/a1"),
4495 kind: EntryKind::Dir,
4496 attrs: crate::Attrs::default(),
4497 },
4498 Op::Upsert {
4499 path: PathBuf::from("b/b1"),
4500 kind: EntryKind::Dir,
4501 attrs: crate::Attrs::default(),
4502 },
4503 Op::Upsert {
4504 path: PathBuf::from("a/a1/deep.txt"),
4505 kind: EntryKind::File,
4506 attrs: crate::Attrs { size: 2, ..crate::Attrs::default() },
4507 },
4508 ]))
4509 .expect("seed tree");
4510
4511 let whole = ["a", "b", "z.txt", "a/a1", "b/b1", "a/a1/deep.txt"];
4512
4513 for max_work in 2..=14_u64 {
4516 let page = crate::PageRequest { limit: crate::MAX_PAGE_ROWS, max_work };
4517 let mut seen: Vec<String> = Vec::new();
4518 let mut continuation = None;
4519 let mut pages = 0;
4520
4521 loop {
4522 let projection = match continuation {
4523 None => crate::ReadProjection::Tree {
4524 path: PathBuf::new(),
4525 depth: crate::query::Bound::All,
4526 include_ignored: true,
4527 page,
4528 },
4529 Some(token) => crate::ReadProjection::Continue { continuation: token, page },
4530 };
4531 let response = opened
4532 .read(crate::ReadRequest {
4533 projections: vec![projection],
4534 ..crate::ReadRequest::default()
4535 })
4536 .expect("tree read");
4537 let current = match &response.results[0] {
4538 crate::ProjectionResult::Tree(crate::Knowledge::Present(page)) => page,
4539 other => panic!("unexpected result at budget {max_work}: {other:?}"),
4543 };
4544 seen.extend(current.rows.iter().map(|row| row.portable_path.as_str().to_owned()));
4545 continuation = current.next;
4546 pages += 1;
4547 assert!(
4548 pages <= whole.len() * 4 + 16,
4549 "budget {max_work} never finished paging, saw {seen:?}"
4550 );
4551 if continuation.is_none() {
4552 break;
4553 }
4554 }
4555
4556 let mut sorted = seen.clone();
4557 sorted.sort();
4558 let mut expected: Vec<String> = whole.iter().map(|row| (*row).to_owned()).collect();
4559 expected.sort();
4560 assert_eq!(
4561 sorted, expected,
4562 "budget {max_work} finished with next: None while missing rows; saw {seen:?}"
4563 );
4564 }
4565
4566 opened.close().expect("close");
4567 }
4568
4569 #[test]
4582 fn a_page_moves_even_when_the_path_walk_spends_the_budget() {
4583 let (_root, opened) = opened(Arc::new(TestControls::default()));
4584 opened
4585 .state
4586 .index
4587 .apply(&Observation::new(vec![
4588 Op::Upsert {
4589 path: PathBuf::from("a"),
4590 kind: EntryKind::Dir,
4591 attrs: crate::Attrs::default(),
4592 },
4593 Op::Upsert {
4594 path: PathBuf::from("a/b"),
4595 kind: EntryKind::Dir,
4596 attrs: crate::Attrs::default(),
4597 },
4598 Op::Upsert {
4599 path: PathBuf::from("a/b/x.txt"),
4600 kind: EntryKind::File,
4601 attrs: crate::Attrs { size: 1, ..crate::Attrs::default() },
4602 },
4603 Op::Upsert {
4604 path: PathBuf::from("a/b/y.txt"),
4605 kind: EntryKind::File,
4606 attrs: crate::Attrs { size: 2, ..crate::Attrs::default() },
4607 },
4608 ]))
4609 .expect("seed tree");
4610
4611 for max_work in 3..=12_u64 {
4615 let page = crate::PageRequest { limit: crate::MAX_PAGE_ROWS, max_work };
4616 let mut seen: Vec<String> = Vec::new();
4617 let mut continuation = None;
4618 let mut pages = 0;
4619
4620 loop {
4621 let projection = match continuation {
4622 None => crate::ReadProjection::Tree {
4623 path: PathBuf::from("a/b"),
4624 depth: crate::query::Bound::All,
4625 include_ignored: true,
4626 page,
4627 },
4628 Some(token) => crate::ReadProjection::Continue { continuation: token, page },
4629 };
4630 let response = opened
4631 .read(crate::ReadRequest {
4632 projections: vec![projection],
4633 ..crate::ReadRequest::default()
4634 })
4635 .expect("tree read");
4636 let crate::ProjectionResult::Tree(crate::Knowledge::Present(current)) =
4637 &response.results[0]
4638 else {
4639 panic!("tree page");
4640 };
4641 seen.extend(current.rows.iter().map(|row| row.portable_path.as_str().to_owned()));
4642 continuation = current.next;
4643 pages += 1;
4644 assert!(
4645 pages <= 16,
4646 "budget {max_work} never terminated; after {pages} pages saw {seen:?}"
4647 );
4648 if continuation.is_none() {
4649 break;
4650 }
4651 }
4652
4653 assert_eq!(seen, vec!["a/b/x.txt", "a/b/y.txt"], "budget {max_work} lost rows");
4654 }
4655 opened.close().expect("close");
4656 }
4657
4658 #[test]
4666 fn descending_costs_nothing_on_a_level_of_leaves() {
4667 let (_root, opened) = opened(Arc::new(TestControls::default()));
4668 let mut ops = Vec::new();
4669 for index in 0..60 {
4670 ops.push(Op::Upsert {
4671 path: PathBuf::from(format!("d{index:03}")),
4672 kind: EntryKind::Dir,
4673 attrs: crate::Attrs::default(),
4674 });
4675 ops.push(Op::Upsert {
4676 path: PathBuf::from(format!("d{index:03}/leaf.txt")),
4677 kind: EntryKind::File,
4678 attrs: crate::Attrs { size: 1, ..crate::Attrs::default() },
4679 });
4680 }
4681 opened.state.index.apply(&Observation::new(ops)).expect("seed wide leaf level");
4682
4683 let response = opened
4684 .read(crate::ReadRequest {
4685 projections: vec![crate::ReadProjection::Tree {
4686 path: PathBuf::new(),
4687 depth: crate::query::Bound::All,
4688 include_ignored: true,
4689 page: crate::PageRequest {
4690 limit: crate::MAX_PAGE_ROWS,
4691 max_work: crate::MAX_PAGE_WORK,
4692 },
4693 }],
4694 ..crate::ReadRequest::default()
4695 })
4696 .expect("tree read");
4697 let crate::ProjectionResult::Tree(crate::Knowledge::Present(page)) = &response.results[0]
4698 else {
4699 panic!("tree page");
4700 };
4701 assert_eq!(page.rows.len(), 120, "60 directories and their 60 files");
4702 assert!(page.next.is_none(), "one page holds the whole tree");
4703 let width = 60;
4712 assert!(
4713 response.work.rows_visited <= 3 * width + 10,
4714 "descent scanned the level it was leaving: {} steps for {} rows",
4715 response.work.rows_visited,
4716 response.work.rows_returned
4717 );
4718 opened.close().expect("close");
4719 }
4720
4721 fn tree_rows(opened: &OpenedIndex, include_ignored: bool) -> Vec<String> {
4723 let response = opened
4724 .read(crate::ReadRequest {
4725 projections: vec![crate::ReadProjection::Tree {
4726 path: PathBuf::new(),
4727 depth: crate::query::Bound::All,
4728 include_ignored,
4729 page: crate::PageRequest {
4730 limit: crate::MAX_PAGE_ROWS,
4731 max_work: crate::MAX_PAGE_WORK,
4732 },
4733 }],
4734 ..crate::ReadRequest::default()
4735 })
4736 .expect("tree read");
4737 let crate::ProjectionResult::Tree(crate::Knowledge::Present(page)) = &response.results[0]
4738 else {
4739 panic!("tree page");
4740 };
4741 page.rows.iter().map(|row| row.portable_path.as_str().to_owned()).collect()
4742 }
4743
4744 #[test]
4749 fn an_opened_tree_read_excluding_ignored_entries_relies_on_an_observing_root() {
4750 let (_root, opened) = opened(Arc::new(TestControls::default()));
4751 opened
4752 .state
4753 .index
4754 .apply(&Observation::new(vec![
4755 Op::Upsert {
4756 path: PathBuf::from("src"),
4757 kind: EntryKind::Dir,
4758 attrs: crate::Attrs::default(),
4759 },
4760 Op::Upsert {
4761 path: PathBuf::from("src/main.rs"),
4762 kind: EntryKind::File,
4763 attrs: crate::Attrs { size: 1, ..crate::Attrs::default() },
4764 },
4765 ]))
4766 .expect("seed tree");
4767
4768 let image = opened.state.index.snapshot().expect("snapshot");
4769 assert!(image.observes_controls());
4770 let excluded = tree_rows(&opened, false);
4771 assert_eq!(excluded, ["src", "src/main.rs"]);
4772 assert_eq!(excluded, tree_rows(&opened, true));
4773
4774 opened.close().expect("close");
4775 }
4776
4777 #[test]
4783 fn excluding_ignored_prunes_the_subtree() {
4784 let (_root, opened) = opened(Arc::new(TestControls::default()));
4785 opened
4786 .state
4787 .index
4788 .apply(&Observation::new(vec![
4789 Op::Upsert {
4790 path: PathBuf::from("src"),
4791 kind: EntryKind::Dir,
4792 attrs: crate::Attrs::default(),
4793 },
4794 Op::Upsert {
4795 path: PathBuf::from("src/main.rs"),
4796 kind: EntryKind::File,
4797 attrs: crate::Attrs { size: 1, ..crate::Attrs::default() },
4798 },
4799 Op::ControlUpsert {
4800 path: PathBuf::from(".gitignore"),
4801 source: b"vendor/\n".to_vec(),
4802 },
4803 Op::Upsert {
4804 path: PathBuf::from("vendor"),
4805 kind: EntryKind::Dir,
4806 attrs: crate::Attrs::default(),
4807 },
4808 Op::Upsert {
4809 path: PathBuf::from("vendor/keep.txt"),
4810 kind: EntryKind::File,
4811 attrs: crate::Attrs { size: 2, ..crate::Attrs::default() },
4812 },
4813 ]))
4814 .expect("seed tree");
4815
4816 let included = tree_rows(&opened, true);
4817 assert!(included.iter().any(|row| row == "vendor"));
4818 assert!(included.iter().any(|row| row == "vendor/keep.txt"));
4819
4820 let excluded = tree_rows(&opened, false);
4821 assert!(!excluded.iter().any(|row| row == "vendor"), "the row is gone");
4822 assert!(
4823 !excluded.iter().any(|row| row == "vendor/keep.txt"),
4824 "and so is everything beneath it, which is what pruning means"
4825 );
4826 assert!(excluded.iter().any(|row| row == "src/main.rs"), "unignored work is untouched");
4827
4828 opened.close().expect("close");
4829 }
4830
4831 #[test]
4847 fn the_remembered_descent_skips_a_pruned_first_child() {
4848 let (_root, opened) = opened(Arc::new(TestControls::default()));
4849 opened
4850 .state
4851 .index
4852 .apply(&Observation::new(vec![
4853 Op::ControlUpsert {
4854 path: PathBuf::from(".gitignore"),
4855 source: b"aaa_vendor/\n".to_vec(),
4856 },
4857 Op::Upsert {
4858 path: PathBuf::from("src"),
4859 kind: EntryKind::Dir,
4860 attrs: crate::Attrs::default(),
4861 },
4862 Op::Upsert {
4863 path: PathBuf::from("src/aaa_vendor"),
4864 kind: EntryKind::Dir,
4865 attrs: crate::Attrs::default(),
4866 },
4867 Op::Upsert {
4868 path: PathBuf::from("src/bbb_keep"),
4869 kind: EntryKind::Dir,
4870 attrs: crate::Attrs::default(),
4871 },
4872 Op::Upsert {
4873 path: PathBuf::from("src/bbb_keep/kept.txt"),
4874 kind: EntryKind::File,
4875 attrs: crate::Attrs { size: 2, ..crate::Attrs::default() },
4876 },
4877 ]))
4878 .expect("seed tree");
4879 let buried: Vec<Op> = (0..40)
4880 .map(|index| Op::Upsert {
4881 path: PathBuf::from(format!("src/aaa_vendor/hidden{index:03}.txt")),
4882 kind: EntryKind::File,
4883 attrs: crate::Attrs { size: 1, ..crate::Attrs::default() },
4884 })
4885 .collect();
4886 opened.state.index.apply(&Observation::new(buried)).expect("seed the pruned subtree");
4887
4888 let response = opened
4889 .read(crate::ReadRequest {
4890 projections: vec![crate::ReadProjection::Tree {
4891 path: PathBuf::new(),
4892 depth: crate::query::Bound::All,
4893 include_ignored: false,
4894 page: crate::PageRequest {
4895 limit: crate::MAX_PAGE_ROWS,
4896 max_work: crate::MAX_PAGE_WORK,
4897 },
4898 }],
4899 ..crate::ReadRequest::default()
4900 })
4901 .expect("tree read");
4902 let crate::ProjectionResult::Tree(crate::Knowledge::Present(page)) = &response.results[0]
4903 else {
4904 panic!("tree page");
4905 };
4906 let rows: Vec<String> =
4907 page.rows.iter().map(|row| row.portable_path.as_str().to_owned()).collect();
4908
4909 assert!(rows.iter().any(|row| row == "src/bbb_keep"), "the kept directory is listed");
4910 assert!(
4911 rows.iter().any(|row| row == "src/bbb_keep/kept.txt"),
4912 "and the descent reached the level below it"
4913 );
4914 assert!(
4915 !rows.iter().any(|row| row == "src/aaa_vendor"),
4916 "the pruned directory is not a row"
4917 );
4918 assert!(
4919 !rows.iter().any(|row| row.starts_with("src/aaa_vendor/")),
4920 "and nothing beneath it is listed"
4921 );
4922
4923 assert!(
4928 response.work.rows_visited < 40,
4929 "the pruned subtree was expanded: {} steps for {} rows",
4930 response.work.rows_visited,
4931 response.work.rows_returned
4932 );
4933
4934 opened.close().expect("close");
4935 }
4936
4937 #[cfg(unix)]
4951 #[test]
4952 fn non_utf8_names_are_escaped_into_pages_rather_than_omitted() {
4953 use std::os::unix::ffi::OsStringExt;
4954
4955 let (_root, opened) = opened(Arc::new(TestControls::default()));
4956 let invalid = PathBuf::from(OsString::from_vec(vec![b'x', 0xff]));
4957 let literal_percent = PathBuf::from("y%FE");
4960 opened
4961 .state
4962 .index
4963 .apply(&Observation::new(vec![
4964 Op::Upsert {
4965 path: invalid.clone(),
4966 kind: EntryKind::File,
4967 attrs: crate::Attrs { size: 9, ..crate::Attrs::default() },
4968 },
4969 Op::Upsert {
4970 path: literal_percent.clone(),
4971 kind: EntryKind::File,
4972 attrs: crate::Attrs { size: 1, ..crate::Attrs::default() },
4973 },
4974 ]))
4975 .expect("seed non-utf8 and literal-percent entries");
4976
4977 let response = opened
4978 .read(crate::ReadRequest {
4979 projections: vec![
4980 crate::ReadProjection::Lookup { path: PathBuf::from("missing") },
4981 crate::ReadProjection::Tree {
4982 path: PathBuf::new(),
4983 depth: crate::query::Bound::Limit(1),
4984 include_ignored: true,
4985 page: crate::PageRequest {
4986 limit: crate::MAX_PAGE_ROWS,
4987 max_work: crate::MAX_PAGE_WORK,
4988 },
4989 },
4990 crate::ReadProjection::Flat {
4991 selection: crate::query::EntrySelection::default(),
4992 shape: crate::RowShape::Compact,
4993 page: crate::PageRequest {
4994 limit: crate::MAX_PAGE_ROWS,
4995 max_work: crate::MAX_PAGE_WORK,
4996 },
4997 },
4998 crate::ReadProjection::RollUp { path: PathBuf::new() },
4999 ],
5000 ..crate::ReadRequest::default()
5001 })
5002 .expect("portable read");
5003
5004 assert!(matches!(
5006 response.results[0],
5007 crate::ProjectionResult::Lookup(crate::Knowledge::Absent)
5008 ));
5009
5010 let crate::ProjectionResult::Tree(crate::Knowledge::Present(tree)) = &response.results[1]
5011 else {
5012 panic!("tree page");
5013 };
5014 assert!(tree.complete);
5015 let names: Vec<_> =
5016 tree.rows.iter().map(|row| row.portable_path.as_str().to_owned()).collect();
5017 assert_eq!(names, vec!["x%FF".to_owned(), "y%25FE".to_owned()]);
5018
5019 let crate::ProjectionResult::Flat(flat) = &response.results[2] else {
5020 panic!("flat page");
5021 };
5022 let flat_names: Vec<_> =
5023 flat.rows.iter().map(|row| row.portable_path.as_str().to_owned()).collect();
5024 assert_eq!(flat_names, names, "ordered pages agree on one population");
5025
5026 assert!(matches!(
5028 &response.results[3],
5029 crate::ProjectionResult::RollUp(crate::Knowledge::Present(rollup))
5030 if rollup.all.files == 2
5031 && rollup.all.files
5032 == u64::try_from(flat.rows.len()).expect("row count fits u64")
5033 ));
5034
5035 opened
5036 .state
5037 .index
5038 .apply(&Observation::new(vec![
5039 Op::Remove { path: invalid },
5040 Op::Remove { path: literal_percent },
5041 ]))
5042 .expect("remove escaped entries");
5043 let after = opened
5044 .read(crate::ReadRequest {
5045 projections: vec![crate::ReadProjection::Tree {
5046 path: PathBuf::new(),
5047 depth: crate::query::Bound::Limit(1),
5048 include_ignored: true,
5049 page: crate::PageRequest {
5050 limit: crate::MAX_PAGE_ROWS,
5051 max_work: crate::MAX_PAGE_WORK,
5052 },
5053 }],
5054 ..crate::ReadRequest::default()
5055 })
5056 .expect("read after removal");
5057 let crate::ProjectionResult::Tree(crate::Knowledge::Present(tree)) = &after.results[0]
5058 else {
5059 panic!("tree page");
5060 };
5061 assert!(tree.rows.is_empty(), "removal is symmetric for escaped names too");
5062 opened.close().expect("close");
5063 }
5064
5065 #[test]
5066 fn read_bounds_are_validated_before_any_continuation_is_consumed() {
5067 let (_root, opened) = opened(Arc::new(TestControls::default()));
5068 opened
5069 .state
5070 .index
5071 .apply(&Observation::new(vec![
5072 Op::Upsert {
5073 path: PathBuf::from("a"),
5074 kind: EntryKind::File,
5075 attrs: crate::Attrs::default(),
5076 },
5077 Op::Upsert {
5078 path: PathBuf::from("b"),
5079 kind: EntryKind::File,
5080 attrs: crate::Attrs::default(),
5081 },
5082 ]))
5083 .expect("seed entries");
5084 let first = opened
5085 .read(crate::ReadRequest {
5086 projections: vec![crate::ReadProjection::Flat {
5087 selection: crate::query::EntrySelection::default(),
5088 shape: crate::RowShape::Compact,
5089 page: crate::PageRequest { limit: 1, max_work: 2 },
5090 }],
5091 ..crate::ReadRequest::default()
5092 })
5093 .expect("first page");
5094 let crate::ProjectionResult::Flat(page) = &first.results[0] else {
5095 panic!("flat page");
5096 };
5097 let continuation = page.next.expect("continuation");
5098
5099 assert!(matches!(
5100 opened.read(crate::ReadRequest {
5101 projections: vec![
5102 crate::ReadProjection::Continue {
5103 continuation,
5104 page: crate::PageRequest { limit: 1, max_work: 2 },
5105 },
5106 crate::ReadProjection::Tree {
5107 path: PathBuf::new(),
5108 depth: crate::query::Bound::Limit(1),
5109 include_ignored: true,
5110 page: crate::PageRequest { limit: 0, max_work: 1 },
5111 },
5112 ],
5113 ..crate::ReadRequest::default()
5114 }),
5115 Err(Error::PageRowLimit { attempted: 0, .. })
5116 ));
5117 opened
5118 .read(crate::ReadRequest {
5119 projections: vec![crate::ReadProjection::Continue {
5120 continuation,
5121 page: crate::PageRequest { limit: 1, max_work: 2 },
5122 }],
5123 ..crate::ReadRequest::default()
5124 })
5125 .expect("validation preserved continuation");
5126
5127 assert!(matches!(
5128 opened.read(crate::ReadRequest {
5129 projections: vec![
5130 crate::ReadProjection::Diagnostics;
5131 crate::MAX_READ_PROJECTIONS + 1
5132 ],
5133 ..crate::ReadRequest::default()
5134 }),
5135 Err(Error::ReadProjectionLimit { .. })
5136 ));
5137 assert!(matches!(
5138 opened.read(crate::ReadRequest {
5139 projections: vec![crate::ReadProjection::Aggregate {
5140 selection: crate::query::EntrySelection::default(),
5141 count_cap: 0,
5142 max_work: 1,
5143 }],
5144 ..crate::ReadRequest::default()
5145 }),
5146 Err(Error::CountCapLimit { attempted: 0, .. })
5147 ));
5148 assert!(matches!(
5152 opened.read(crate::ReadRequest {
5153 projections: vec![crate::ReadProjection::Tree {
5154 path: PathBuf::new(),
5155 depth: crate::query::Bound::Limit(0),
5156 include_ignored: true,
5157 page: crate::PageRequest { limit: 1, max_work: 1 },
5158 }],
5159 ..crate::ReadRequest::default()
5160 }),
5161 Err(Error::TreeDepthZero)
5162 ));
5163 let bounded_path = opened
5164 .read(crate::ReadRequest {
5165 projections: vec![crate::ReadProjection::Tree {
5166 path: PathBuf::from("missing/deep"),
5167 depth: crate::query::Bound::Limit(1),
5168 include_ignored: true,
5169 page: crate::PageRequest { limit: 1, max_work: 1 },
5170 }],
5171 ..crate::ReadRequest::default()
5172 })
5173 .expect("bounded path traversal");
5174 assert!(matches!(
5175 bounded_path.results[0],
5176 crate::ProjectionResult::Limit(crate::QueryLimit {
5177 projection: crate::LimitedProjection::Tree,
5178 rows_visited: 1,
5179 ..
5180 })
5181 ));
5182 assert!(matches!(
5183 opened.read(crate::ReadRequest {
5184 projections: vec![crate::ReadProjection::Report(crate::ReportRequest {
5185 query: crate::query::Query {
5186 views: vec![crate::query::ViewSpec::Summary; crate::MAX_REPORT_VIEWS + 1],
5187 ..crate::query::Query::default()
5188 },
5189 now: std::time::UNIX_EPOCH,
5190 max_work: crate::MAX_PAGE_WORK,
5191 })],
5192 ..crate::ReadRequest::default()
5193 }),
5194 Err(Error::ReportViewLimit { .. })
5195 ));
5196 opened.close().expect("close");
5197 }
5198
5199 #[test]
5200 fn a_coherent_read_cannot_straddle_a_commit() {
5201 let (_root, opened) = opened(Arc::new(TestControls::default()));
5202 let stop = Arc::new(AtomicBool::new(false));
5203 let writer_stop = Arc::clone(&stop);
5204 let writer_index = opened.state.index.clone();
5205 let writer = std::thread::spawn(move || {
5206 for round in 0..400_u64 {
5207 if writer_stop.load(Ordering::Relaxed) {
5208 break;
5209 }
5210 writer_index
5211 .apply(&Observation::new(vec![Op::Upsert {
5212 path: PathBuf::from(format!("file-{round}")),
5213 kind: EntryKind::File,
5214 attrs: crate::Attrs { size: 1, ..crate::Attrs::default() },
5215 }]))
5216 .expect("writer commit");
5217 }
5218 });
5219
5220 for _ in 0..500 {
5221 let response = opened
5222 .read(crate::ReadRequest {
5223 projections: vec![
5224 crate::ReadProjection::Tree {
5225 path: PathBuf::new(),
5226 depth: crate::query::Bound::Limit(1),
5227 include_ignored: true,
5228 page: crate::PageRequest {
5229 limit: crate::MAX_PAGE_ROWS,
5230 max_work: crate::MAX_PAGE_WORK,
5231 },
5232 },
5233 crate::ReadProjection::RollUp { path: PathBuf::new() },
5234 ],
5235 ..crate::ReadRequest::default()
5236 })
5237 .expect("coherent read");
5238 let crate::ProjectionResult::Tree(crate::Knowledge::Present(tree)) =
5239 &response.results[0]
5240 else {
5241 panic!("tree page");
5242 };
5243 let crate::ProjectionResult::RollUp(crate::Knowledge::Present(rollup)) =
5244 &response.results[1]
5245 else {
5246 panic!("root roll-up");
5247 };
5248 assert!(tree.next.is_none());
5249 assert_eq!(
5250 tree.rows.iter().filter(|row| row.kind == EntryKind::File).count() as u64,
5251 rollup.all.files
5252 );
5253 }
5254 stop.store(true, Ordering::Relaxed);
5255 writer.join().expect("writer");
5256 opened.close().expect("close");
5257 }
5258
5259 #[test]
5260 fn aggregates_distinguish_maintained_exact_totals_from_capped_counts() {
5261 let (_root, opened) = opened(Arc::new(TestControls::default()));
5262 opened
5263 .state
5264 .index
5265 .apply(&Observation::new(
5266 ["a", "b", "c"]
5267 .into_iter()
5268 .map(|path| Op::Upsert {
5269 path: PathBuf::from(path),
5270 kind: EntryKind::File,
5271 attrs: crate::Attrs::default(),
5272 })
5273 .collect(),
5274 ))
5275 .expect("seed entries");
5276
5277 let response = opened
5278 .read(crate::ReadRequest {
5279 projections: vec![
5280 crate::ReadProjection::Aggregate {
5281 selection: crate::query::EntrySelection::default(),
5282 count_cap: 1,
5283 max_work: 1,
5284 },
5285 crate::ReadProjection::Aggregate {
5286 selection: crate::query::EntrySelection {
5287 query: crate::query::Selection {
5288 kinds: vec![EntryKind::File],
5289 ..crate::query::Selection::default()
5290 },
5291 ..crate::query::EntrySelection::default()
5292 },
5293 count_cap: 2,
5294 max_work: 3,
5295 },
5296 ],
5297 ..crate::ReadRequest::default()
5298 })
5299 .expect("aggregate read");
5300
5301 assert!(matches!(
5302 response.results[0],
5303 crate::ProjectionResult::Aggregate(crate::CountResult::Exact(3))
5304 ));
5305 assert!(matches!(
5306 response.results[1],
5307 crate::ProjectionResult::Aggregate(crate::CountResult::AtLeast(2))
5308 ));
5309 opened.close().expect("close");
5310 }
5311
5312 #[test]
5313 fn report_projection_matches_the_existing_query_and_fails_closed_at_its_work_bound() {
5314 let (_root, opened) = opened(Arc::new(TestControls::default()));
5315 opened
5316 .state
5317 .index
5318 .apply(&Observation::new(vec![
5319 Op::Upsert {
5320 path: PathBuf::from("a"),
5321 kind: EntryKind::File,
5322 attrs: crate::Attrs { size: 3, ..crate::Attrs::default() },
5323 },
5324 Op::Upsert {
5325 path: PathBuf::from("b"),
5326 kind: EntryKind::File,
5327 attrs: crate::Attrs { size: 5, ..crate::Attrs::default() },
5328 },
5329 ]))
5330 .expect("seed entries");
5331 let query = crate::query::Query {
5332 views: vec![crate::query::ViewSpec::Summary],
5333 ..crate::query::Query::default()
5334 };
5335
5336 let response = opened
5337 .read(crate::ReadRequest {
5338 projections: vec![crate::ReadProjection::Report(crate::ReportRequest {
5339 query: query.clone(),
5340 now: std::time::UNIX_EPOCH,
5341 max_work: 1,
5342 })],
5343 ..crate::ReadRequest::default()
5344 })
5345 .expect("maintained report");
5346 let crate::ProjectionResult::Report(report) = &response.results[0] else {
5347 panic!("report projection");
5348 };
5349 let crate::query::Section::Summary(summary) = &report.sections[0] else {
5350 panic!("summary section");
5351 };
5352 assert_eq!((summary.files, summary.bytes), (2, 8));
5353
5354 let limited = opened
5355 .read(crate::ReadRequest {
5356 projections: vec![crate::ReadProjection::Report(crate::ReportRequest {
5357 query: crate::query::Query {
5358 selection: crate::query::Selection {
5359 kinds: vec![EntryKind::File],
5360 ..crate::query::Selection::default()
5361 },
5362 views: vec![crate::query::ViewSpec::Summary],
5363 ..crate::query::Query::default()
5364 },
5365 now: std::time::UNIX_EPOCH,
5366 max_work: 1,
5367 })],
5368 ..crate::ReadRequest::default()
5369 })
5370 .expect("bounded report");
5371 assert!(matches!(
5372 limited.results[0],
5373 crate::ProjectionResult::Limit(crate::QueryLimit {
5374 projection: crate::LimitedProjection::Report,
5375 ..
5376 })
5377 ));
5378 opened.close().expect("close");
5379 }
5380
5381 #[test]
5382 fn first_refusal_stops_expansion_and_commits_the_budget_state_with_prior_facts() {
5383 let root = tempfile::tempdir().expect("temp root");
5384 std::fs::create_dir(root.path().join("nested")).expect("fixture directory");
5385 std::fs::write(root.path().join("one"), b"1").expect("fixture");
5386 std::fs::write(root.path().join("two"), b"2").expect("fixture");
5387 std::fs::write(root.path().join("nested/deep"), b"deep").expect("deep fixture");
5388 let opened = open_fixture(
5389 root.path(),
5390 OpenOptions {
5391 batch_size: 64,
5392 budget: DiscoveryBudget { max_files: Some(1) },
5393 ..OpenOptions::default()
5394 },
5395 )
5396 .expect("opened root");
5397
5398 let state = wait_until_settled(&opened);
5399 assert_eq!(state.phase, crate::LifecyclePhase::Stopped);
5400 assert_eq!(state.coverage, crate::Coverage::Partial(crate::CoverageReason::Budget));
5401 assert_eq!(state.progress.files_retained, 1);
5402 assert_eq!(opened.state.index.total().expect("total").files, 1);
5403 assert_eq!(opened.state.index.kind(Path::new("nested/deep")).expect("deep lookup"), None);
5404 assert_eq!(
5405 opened.state.index.directory_complete(Path::new("")).expect("root completeness"),
5406 Some(false)
5407 );
5408 let partial = opened.state.index.snapshot().expect("partial snapshot image");
5409 assert!(matches!(
5410 crate::snapshot::save(&partial, &root.path().join("partial.fdu")),
5411 Err(Error::Snapshot(_))
5412 ));
5413 assert!(matches!(
5414 opened.prioritize(&[PathBuf::from("nested")]),
5415 Err(Error::OpenedIndexStopped)
5416 ));
5417
5418 let terminal = opened
5419 .state
5420 .index
5421 .since(crate::Clock::ZERO)
5422 .expect("journal")
5423 .commits
5424 .into_iter()
5425 .find(|commit| {
5426 commit.state.iter().any(|transition| {
5427 matches!(
5428 transition,
5429 crate::StateTransition::IndexState {
5430 current: crate::IndexState {
5431 coverage: crate::Coverage::Partial(crate::CoverageReason::Budget),
5432 ..
5433 },
5434 ..
5435 }
5436 )
5437 })
5438 })
5439 .expect("budget commit");
5440 assert!(terminal.changes.iter().any(|change| matches!(
5441 change,
5442 crate::EffectiveChange::Inserted { kind: EntryKind::File, .. }
5443 )));
5444 opened.close().expect("close");
5445 }
5446
5447 #[test]
5448 fn refresh_receipt_counts_verified_no_op_work_without_a_fact_commit() {
5449 let root = tempfile::tempdir().expect("temp root");
5450 std::fs::write(root.path().join("stable.txt"), b"stable").expect("fixture");
5451 let opened = open_fixture(root.path(), OpenOptions::default()).expect("opened root");
5452 assert_eq!(wait_until_settled(&opened).phase, crate::LifecyclePhase::Ready);
5453
5454 let result = opened.refresh(&[PathBuf::from("stable.txt")]).expect("refresh");
5455
5456 assert_eq!(result.accepted, vec![PathBuf::from("stable.txt")]);
5457 assert!(result.rejected.is_empty());
5458 assert_eq!(result.work.observations, 1, "the verified observation is still work");
5459 assert_eq!(result.work.unchanged, 1, "the matching fact is reported as unchanged");
5460 assert_eq!(result.work.stale, 0);
5461 opened.close().expect("close");
5462 }
5463
5464 #[test]
5465 fn refresh_can_fill_remaining_file_budget_without_exceeding_it() {
5466 let root = tempfile::tempdir().expect("temp root");
5467 std::fs::write(root.path().join("one"), b"1").expect("fixture");
5468 let opened = open_fixture(
5469 root.path(),
5470 OpenOptions {
5471 budget: DiscoveryBudget { max_files: Some(2) },
5472 ..OpenOptions::default()
5473 },
5474 )
5475 .expect("opened root");
5476 assert_eq!(wait_until_settled(&opened).progress.files_retained, 1);
5477 std::fs::write(root.path().join("two"), b"2").expect("new file");
5478
5479 let result = opened.refresh(&[PathBuf::from("two")]).expect("refresh");
5480
5481 assert_eq!(result.accepted, vec![PathBuf::from("two")]);
5482 assert!(result.rejected.is_empty());
5483 assert_eq!(opened.state.index.total().expect("total").files, 2);
5484 assert_eq!(result.state.phase, crate::LifecyclePhase::Ready);
5485 assert_eq!(result.state.progress.files_retained, 2);
5486 opened.close().expect("close");
5487 }
5488
5489 #[test]
5490 fn refresh_classifies_paths_and_collapses_overlapping_walks() {
5491 let root = tempfile::tempdir().expect("temp root");
5492 std::fs::create_dir_all(root.path().join("visible/nested")).expect("fixture directories");
5493 std::fs::write(root.path().join("visible/nested/leaf"), b"leaf").expect("fixture");
5494 std::fs::create_dir(root.path().join(".hidden")).expect("hidden directory");
5495 std::fs::write(root.path().join(".hidden/leaf"), b"hidden").expect("hidden fixture");
5496 let opened = open_fixture(
5497 root.path(),
5498 OpenOptions {
5499 hidden: Some(Arc::new(crate::HiddenPolicy::prune_hidden::<[&str; 0], &str>([]))),
5500 ..OpenOptions::default()
5501 },
5502 )
5503 .expect("opened root");
5504 assert_eq!(wait_until_settled(&opened).phase, crate::LifecyclePhase::Ready);
5505
5506 let result = opened
5507 .refresh(&[
5508 PathBuf::from("visible/nested"),
5509 PathBuf::from("visible"),
5510 PathBuf::from("visible/nested"),
5511 PathBuf::from("../escape"),
5512 PathBuf::from(".hidden/leaf"),
5513 ])
5514 .expect("refresh");
5515
5516 assert_eq!(
5517 result.accepted,
5518 vec![PathBuf::from("visible"), PathBuf::from("visible/nested")]
5519 );
5520 assert_eq!(
5521 result.rejected,
5522 vec![
5523 crate::RejectedRefreshPath {
5524 path: PathBuf::from("../escape"),
5525 reason: crate::RefreshRejection::OutsideRoot,
5526 },
5527 crate::RejectedRefreshPath {
5528 path: PathBuf::from(".hidden/leaf"),
5529 reason: crate::RefreshRejection::NotAdmitted,
5530 },
5531 ]
5532 );
5533 assert_eq!(result.work.directories_read, 2, "the descendant was not walked twice");
5534 opened.close().expect("close");
5535 }
5536
5537 #[test]
5538 fn refresh_widens_through_a_replaced_ancestor_and_reports_exact_commits() {
5539 let root = tempfile::tempdir().expect("temp root");
5540 std::fs::create_dir_all(root.path().join("parent/child")).expect("fixture directories");
5541 std::fs::write(root.path().join("parent/child/leaf"), b"leaf").expect("fixture");
5542 let opened = open_fixture(root.path(), OpenOptions::default()).expect("opened root");
5543 assert_eq!(wait_until_settled(&opened).phase, crate::LifecyclePhase::Ready);
5544 std::fs::remove_dir_all(root.path().join("parent")).expect("remove old subtree");
5545 std::fs::write(root.path().join("parent"), b"replacement").expect("replacement file");
5546 let before = current_version(&opened);
5547
5548 let result =
5549 opened.refresh(&[PathBuf::from("parent/child/leaf")]).expect("refresh widened path");
5550 let poll = opened
5551 .changes(crate::ChangeRequest { after: before, timeout: std::time::Duration::ZERO })
5552 .expect("refresh commits");
5553
5554 assert_eq!(result.after, before);
5555 assert_eq!(result.version, poll.version);
5556 assert_eq!(result.accepted, vec![PathBuf::from("parent/child/leaf")]);
5557 assert_eq!(
5558 opened.state.index.kind(Path::new("parent")).expect("kind"),
5559 Some(EntryKind::File)
5560 );
5561 assert_eq!(
5562 opened.state.index.kind(Path::new("parent/child/leaf")).expect("removed child"),
5563 None
5564 );
5565 let crate::ChangeOutcome::Changes { commits, impact } = poll.outcome else {
5566 panic!("refresh must advance the journal");
5567 };
5568 assert!(!commits.is_empty());
5569 assert_eq!(impact, result.impact);
5570 assert!(commits.iter().all(|commit| {
5571 commit.clock.0 > result.after.sequence.0 && commit.clock.0 <= result.version.sequence.0
5572 }));
5573 opened.close().expect("close");
5574 }
5575
5576 #[cfg(unix)]
5577 #[test]
5578 fn refresh_rejects_symlink_shadowed_ancestry_without_aborting_other_paths() {
5579 use std::os::unix::fs::symlink;
5580
5581 let root = tempfile::tempdir().expect("temp root");
5582 let outside = tempfile::tempdir().expect("outside root");
5583 std::fs::create_dir_all(root.path().join("shadow/child")).expect("baseline ancestry");
5584 std::fs::write(root.path().join("good.txt"), b"before").expect("baseline file");
5585 let opened = open_fixture(root.path(), OpenOptions::default()).expect("opened root");
5586 assert_eq!(wait_until_settled(&opened).phase, crate::LifecyclePhase::Ready);
5587
5588 std::fs::remove_dir_all(root.path().join("shadow")).expect("remove ancestry");
5589 symlink(outside.path(), root.path().join("shadow")).expect("shadow with symlink");
5590 std::fs::write(root.path().join("good.txt"), b"after and larger").expect("mutate file");
5591
5592 let result = opened
5593 .refresh(&[PathBuf::from("shadow/child/leaf"), PathBuf::from("good.txt")])
5594 .expect("one unsafe path is a rejection, not a batch error");
5595
5596 assert_eq!(result.accepted, vec![PathBuf::from("good.txt")]);
5597 assert_eq!(result.rejected.len(), 1);
5598 assert_eq!(result.rejected[0].path, Path::new("shadow/child/leaf"));
5599 assert_eq!(result.rejected[0].reason, crate::RefreshRejection::UnsafeAncestry);
5600 assert_eq!(
5601 opened.state.index.attrs(Path::new("good.txt")).expect("attrs").expect("retained").size,
5602 16
5603 );
5604 opened.close().expect("close");
5605 }
5606
5607 #[test]
5608 fn refresh_refusal_is_atomic_with_the_shared_file_budget() {
5609 let root = tempfile::tempdir().expect("temp root");
5610 std::fs::write(root.path().join("one"), b"1").expect("fixture");
5611 let opened = open_fixture(
5612 root.path(),
5613 OpenOptions {
5614 budget: DiscoveryBudget { max_files: Some(2) },
5615 ..OpenOptions::default()
5616 },
5617 )
5618 .expect("opened root");
5619 assert_eq!(wait_until_settled(&opened).progress.files_retained, 1);
5620 std::fs::write(root.path().join("two"), b"2").expect("new file");
5621 std::fs::write(root.path().join("three"), b"3").expect("new file");
5622
5623 let result = opened
5624 .refresh(&[PathBuf::from("two"), PathBuf::from("three")])
5625 .expect("bounded refresh");
5626
5627 assert_eq!(result.accepted.len(), 2);
5628 assert!(result.rejected.is_empty());
5629 assert_eq!(result.work.observations, 2);
5630 assert_eq!(result.work.resource_refused, 1);
5631 assert_eq!(opened.state.index.total().expect("total").files, 2);
5632 assert_eq!(result.state.progress.files_retained, 2);
5633 assert_eq!(result.state.phase, crate::LifecyclePhase::Stopped);
5634 assert_eq!(result.state.coverage, crate::Coverage::Partial(crate::CoverageReason::Budget));
5635 assert_eq!(result.issues.len(), 1);
5636 assert_eq!(result.issues[0].kind, crate::IssueKind::ResourceBudget);
5637
5638 std::fs::write(root.path().join("four"), b"4").expect("later file");
5639 let stopped = opened.refresh(&[PathBuf::from("four")]).expect("stopped refresh receipt");
5640 assert!(stopped.accepted.is_empty());
5641 assert_eq!(
5642 stopped.rejected,
5643 vec![crate::RejectedRefreshPath {
5644 path: PathBuf::from("four"),
5645 reason: crate::RefreshRejection::ResourceBudget,
5646 }]
5647 );
5648 assert_eq!(stopped.work.entries_visited, 1, "the refusal reports its probe");
5649 assert_eq!(stopped.work.files_visited, 1);
5650 assert_eq!(stopped.work.bytes_visited, 1);
5651 assert_eq!(opened.state.index.total().expect("bounded total").files, 2);
5652 opened.close().expect("close");
5653 }
5654
5655 #[test]
5656 fn concurrent_discovery_and_refresh_share_one_atomic_file_budget() {
5657 let root = tempfile::tempdir().expect("temp root");
5658 std::fs::write(root.path().join("from-discovery"), b"discovery").expect("fixture");
5659 let controls = Arc::new(TestControls::default());
5660 controls.gate(TestPoint::BeforeDiscovery).arm();
5661 let opened = OpenedIndex::open_for_test(
5662 root.path(),
5663 OpenOptions {
5664 budget: DiscoveryBudget { max_files: Some(1) },
5665 ..OpenOptions::default()
5666 },
5667 Arc::clone(&controls),
5668 )
5669 .expect("opened root");
5670 controls.gate(TestPoint::BeforeDiscovery).wait_reached();
5671 std::fs::write(root.path().join("from-refresh"), b"refresh").expect("new file");
5672
5673 let refreshed =
5674 opened.refresh(&[PathBuf::from("from-refresh")]).expect("refresh during discovery");
5675 assert_eq!(refreshed.accepted, vec![PathBuf::from("from-refresh")]);
5676 assert_eq!(opened.state.index.total().expect("after refresh").files, 1);
5677 controls.gate(TestPoint::BeforeDiscovery).release();
5678
5679 let state = wait_until_settled(&opened);
5680 assert_eq!(opened.state.index.total().expect("bounded total").files, 1);
5681 assert_eq!(state.progress.files_retained, 1);
5682 assert_eq!(state.phase, crate::LifecyclePhase::Stopped);
5683 assert_eq!(state.coverage, crate::Coverage::Partial(crate::CoverageReason::Budget));
5684 opened.close().expect("close");
5685 }
5686
5687 #[test]
5697 fn a_refresh_racing_discovery_leaves_the_queued_directory_as_stale_work() {
5698 for batch_size in [1, OpenOptions::default().batch_size] {
5699 let controls = Arc::new(TestControls::default());
5700 controls.gate(TestPoint::AfterRootDirectory).arm();
5701 let root = tempfile::tempdir().expect("temp root");
5702 std::fs::create_dir(root.path().join("sub")).expect("sub");
5703 std::fs::write(root.path().join("sub/inner.txt"), b"x").expect("fixture");
5704 let opened = OpenedIndex::open_for_test(
5705 root.path(),
5706 OpenOptions { batch_size, ..OpenOptions::default() },
5707 Arc::clone(&controls),
5708 )
5709 .expect("open");
5710 controls.gate(TestPoint::AfterRootDirectory).wait_reached();
5711 assert_eq!(
5712 opened.state.index.kind(Path::new("sub")).expect("lookup"),
5713 Some(EntryKind::Dir),
5714 "the root listing queued `sub`"
5715 );
5716
5717 std::fs::remove_dir_all(root.path().join("sub")).expect("remove sub");
5718 let refreshed = opened.refresh(&[PathBuf::from("sub")]).expect("refresh");
5719 assert_eq!(refreshed.accepted, vec![PathBuf::from("sub")]);
5720 std::fs::create_dir(root.path().join("sub")).expect("recreate sub");
5721 std::fs::write(root.path().join("sub/again.txt"), b"y").expect("fixture");
5722 controls.gate(TestPoint::AfterRootDirectory).release();
5723
5724 let state = wait_until_settled(&opened);
5725 assert_eq!(state.phase, crate::LifecyclePhase::Ready, "batch size {batch_size}");
5726 assert_eq!(state.coverage, crate::Coverage::Complete, "batch size {batch_size}");
5727 assert_eq!(state.issues.retained, 0, "batch size {batch_size}");
5728 assert_eq!(opened.state.index.kind(Path::new("sub")).expect("lookup"), None);
5731 opened.close().unwrap_or_else(|error| panic!("batch size {batch_size}: {error}"));
5732 }
5733 }
5734
5735 #[test]
5738 fn a_directory_that_vanishes_during_discovery_is_not_inaccessible() {
5739 for replace_with_file in [false, true] {
5740 let controls = Arc::new(TestControls::default());
5741 controls.gate(TestPoint::AfterRootDirectory).arm();
5742 let root = tempfile::tempdir().expect("temp root");
5743 std::fs::create_dir(root.path().join("sub")).expect("sub");
5744 std::fs::write(root.path().join("sub/inner.txt"), b"x").expect("fixture");
5745 std::fs::write(root.path().join("keep.txt"), b"k").expect("fixture");
5746 let opened = OpenedIndex::open_for_test(
5747 root.path(),
5748 OpenOptions::default(),
5749 Arc::clone(&controls),
5750 )
5751 .expect("open");
5752 controls.gate(TestPoint::AfterRootDirectory).wait_reached();
5753 std::fs::remove_dir_all(root.path().join("sub")).expect("remove sub");
5754 if replace_with_file {
5755 std::fs::write(root.path().join("sub"), b"now a file").expect("replacement");
5756 }
5757 controls.gate(TestPoint::AfterRootDirectory).release();
5758
5759 let state = wait_until_settled(&opened);
5760 let case = if replace_with_file { "replaced by a file" } else { "removed" };
5761 assert_eq!(state.phase, crate::LifecyclePhase::Ready, "{case}");
5762 assert_eq!(state.coverage, crate::Coverage::Complete, "{case}");
5763 assert_eq!(state.freshness, crate::Freshness::Fresh, "{case}");
5764 assert_eq!(state.issues.retained, 0, "{case}");
5765 opened.close().unwrap_or_else(|error| panic!("{case}: {error}"));
5766 }
5767 }
5768
5769 #[test]
5776 fn a_budget_stop_during_discovery_stays_terminal() {
5777 let controls = Arc::new(TestControls::default());
5778 controls.gate(TestPoint::AfterRootDirectory).arm();
5779 let root = tempfile::tempdir().expect("temp root");
5780 std::fs::write(root.path().join("a.txt"), b"a").expect("fixture");
5781 std::fs::create_dir(root.path().join("emptydir")).expect("fixture");
5782 let opened = OpenedIndex::open_for_test(
5783 root.path(),
5784 OpenOptions {
5785 budget: DiscoveryBudget { max_files: Some(1) },
5786 ..OpenOptions::default()
5787 },
5788 Arc::clone(&controls),
5789 )
5790 .expect("open");
5791 controls.gate(TestPoint::AfterRootDirectory).wait_reached();
5792 std::fs::write(root.path().join("b.txt"), b"b").expect("over budget");
5793 let refreshed = opened.refresh(&[PathBuf::from("b.txt")]).expect("refresh");
5794 assert_eq!(refreshed.work.resource_refused, 1);
5795 assert_eq!(refreshed.state.phase, crate::LifecyclePhase::Stopped);
5796 controls.gate(TestPoint::AfterRootDirectory).release();
5797
5798 wait_for_worker_exit(&opened, "discovery");
5799 let state = opened.state.index.state().expect("state");
5800 assert_eq!(state.phase, crate::LifecyclePhase::Stopped);
5801 assert_eq!(state.coverage, crate::Coverage::Partial(crate::CoverageReason::Budget));
5802 assert_eq!(
5803 opened.state.index.directory_complete(Path::new("emptydir")).expect("lookup"),
5804 Some(false),
5805 "a listing that arrives after the stop must not land"
5806 );
5807 assert!(matches!(
5808 opened.prioritize(&[PathBuf::from("emptydir")]),
5809 Err(Error::OpenedIndexStopped)
5810 ));
5811 opened.close().expect("close");
5812 }
5813
5814 fn tree_with_a_guarded_control() -> tempfile::TempDir {
5817 let root = tempfile::tempdir().expect("temp root");
5818 let a = root.path().join("a");
5819 std::fs::create_dir_all(a.join("-early")).expect("fixture");
5820 std::fs::write(a.join("-early").join("leaf.txt"), b"l").expect("fixture");
5821 let mut line = b"*.txt\n".to_vec();
5822 line.extend(std::iter::repeat_n(b'x', crate::control::DEFAULT_CONTROL_LINE_LIMIT + 1));
5823 std::fs::write(a.join(crate::control::CONTROL_FILE_NAME), &line).expect("control");
5824 std::fs::write(a.join("zzz.txt"), b"z").expect("fixture");
5825 std::fs::create_dir(root.path().join("b")).expect("fixture");
5826 std::fs::write(root.path().join("b/kept.txt"), b"k").expect("fixture");
5827 root
5828 }
5829
5830 fn diagnostics(opened: &OpenedIndex) -> crate::ReadDiagnostics {
5831 let response = opened
5832 .read(crate::ReadRequest {
5833 projections: vec![crate::ReadProjection::Diagnostics],
5834 ..crate::ReadRequest::default()
5835 })
5836 .expect("read");
5837 let [crate::ProjectionResult::Diagnostics(diagnostics)] = response.results.as_slice()
5838 else {
5839 panic!("one diagnostics result: {:?}", response.results);
5840 };
5841 diagnostics.clone()
5842 }
5843
5844 #[test]
5852 fn a_control_over_a_bound_is_refused_while_its_listing_commits() {
5853 let controls = Arc::new(TestControls::default());
5854 controls.use_deterministic_discovery_order();
5855 let root = tree_with_a_guarded_control();
5856 let opened = OpenedIndex::open_for_test(
5857 root.path(),
5858 OpenOptions { batch_size: 1, ..OpenOptions::default() },
5859 Arc::clone(&controls),
5860 )
5861 .expect("open");
5862 let state = wait_until_settled(&opened);
5863 assert_eq!(state.phase, crate::LifecyclePhase::Ready);
5864 assert_eq!(state.coverage, crate::Coverage::Complete);
5865 assert_eq!(state.issues, crate::IssueSummary::default());
5866
5867 let index = &opened.state.index;
5868 assert_eq!(index.directory_complete(Path::new("a")).expect("lookup"), Some(true));
5869 assert_eq!(index.directory_complete(Path::new("a/-early")).expect("lookup"), Some(true));
5870 assert_eq!(index.kind(Path::new("a/zzz.txt")).expect("lookup"), Some(EntryKind::File));
5871 assert_eq!(index.kind(Path::new("b/kept.txt")).expect("lookup"), Some(EntryKind::File));
5872 assert_eq!(
5873 index
5874 .snapshot()
5875 .expect("snapshot")
5876 .is_ignored(Path::new("a/zzz.txt"))
5877 .expect("observed"),
5878 None,
5879 "the refused file could have governed this entry"
5880 );
5881 assert_eq!(
5882 diagnostics(&opened).controls,
5883 crate::control::ControlObservation {
5884 limits: crate::control::ControlLimits::default(),
5885 applied: 0,
5886 rules: 0,
5887 refused: 1,
5888 refusals: vec![crate::control::RefusedControl {
5889 path: PathBuf::from("a/.gitignore"),
5890 reason: crate::control::ControlRefusalReason::LineLimit,
5891 }],
5892 }
5893 );
5894 opened.close().expect("close");
5895 }
5896
5897 #[test]
5900 fn an_opened_root_without_a_line_limit_applies_what_the_default_limit_refuses() {
5901 let root = tree_with_a_guarded_control();
5902 let limits = crate::control::ControlLimits {
5903 line_limit: None,
5904 ..crate::control::ControlLimits::default()
5905 };
5906 let options = OpenOptions { control_limits: limits, ..OpenOptions::default() };
5907 let opened = open_fixture(root.path(), options).expect("open");
5908 wait_until_settled(&opened);
5909
5910 let diagnostics = diagnostics(&opened);
5911 assert_eq!((diagnostics.controls.limits, diagnostics.controls.applied), (limits, 1));
5912 assert_eq!(diagnostics.controls.refused, 0);
5913 assert_eq!(
5914 diagnostics.scope,
5915 ScanConfig { control_limits: limits, ..ScanConfig::default() }.scope()
5916 );
5917 let snapshot = opened.state.index.snapshot().expect("snapshot");
5918 assert_eq!(snapshot.is_ignored(Path::new("a/zzz.txt")).expect("observed"), Some(true));
5919 opened.close().expect("close");
5920 }
5921
5922 #[cfg(feature = "watch")]
5925 fn observe_then_settle(
5926 opened: &OpenedIndex,
5927 controls: &TestControls,
5928 root: &Path,
5929 hints: &str,
5930 ) {
5931 static MARKERS: AtomicUsize = AtomicUsize::new(0);
5932 let marker = format!("settled-{}.txt", MARKERS.fetch_add(1, Ordering::Relaxed));
5933 std::fs::write(root.join(&marker), b"marker").expect("marker");
5934 controls.send_observation_hints(&format!("{hints}create\t{marker}\n"));
5935 let deadline = std::time::Instant::now() + TEST_GATE_TIMEOUT;
5936 while opened.state.index.kind(Path::new(&marker)).expect("lookup") != Some(EntryKind::File)
5937 {
5938 assert!(std::time::Instant::now() < deadline, "{marker} was not applied");
5939 std::thread::yield_now();
5940 }
5941 }
5942
5943 #[cfg(feature = "watch")]
5951 #[test]
5952 fn a_watched_root_over_a_control_bound_keeps_watching_through_events() {
5953 let root = tree_with_a_guarded_control();
5954 let scripts = tempfile::tempdir().expect("script root");
5955 let script = scripts.path().join("events.script");
5956 std::fs::write(&script, b"").expect("script");
5957 let controls = Arc::new(TestControls::default());
5958 let opened = OpenedIndex::open_for_test(
5959 root.path(),
5960 scripted_options(&script),
5961 Arc::clone(&controls),
5962 )
5963 .expect("open");
5964 let state = wait_until_phase(&opened, crate::LifecyclePhase::Watching);
5965 assert_eq!(state.coverage, crate::Coverage::Complete);
5966 assert_eq!(diagnostics(&opened).controls.refused, 1);
5967
5968 let control = root.path().join("a/.gitignore");
5970 let mut grown = std::fs::read(&control).expect("control");
5971 grown.extend_from_slice(b"\n*.md\n");
5972 std::fs::write(&control, grown).expect("grown control");
5973 std::fs::write(root.path().join("a/new.txt"), b"new").expect("fixture");
5974 let hints = "modify\ta/.gitignore\ncreate\ta/new.txt\n";
5975 observe_then_settle(&opened, &controls, root.path(), hints);
5976 let watching = crate::LifecyclePhase::Watching;
5977 assert_eq!(opened.state.index.state().expect("state").phase, watching);
5978 assert_eq!(
5979 opened.state.index.kind(Path::new("a/new.txt")).expect("lookup"),
5980 Some(EntryKind::File)
5981 );
5982 assert_eq!(diagnostics(&opened).controls.refused, 1);
5983
5984 let refreshed = opened.refresh(&[PathBuf::from("a")]).expect("refresh over the refusal");
5986 assert_eq!(refreshed.state.phase, watching);
5987
5988 std::fs::write(&control, b"*.txt\n").expect("fitting control");
5990 observe_then_settle(&opened, &controls, root.path(), "modify\ta/.gitignore\n");
5991 let coverage = diagnostics(&opened).controls;
5992 assert_eq!((coverage.applied, coverage.refused), (1, 0));
5993 let snapshot = opened.state.index.snapshot().expect("snapshot");
5994 assert_eq!(snapshot.is_ignored(Path::new("a/new.txt")).expect("observed"), Some(true));
5995 assert_eq!(opened.state.index.state().expect("state").phase, watching);
5996 opened.close().expect("close");
5997 }
5998
5999 #[test]
6001 fn a_directory_created_after_discovery_is_complete_once_refreshed() {
6002 let root = tempfile::tempdir().expect("temp root");
6003 std::fs::write(root.path().join("before.txt"), b"b").expect("fixture");
6004 let opened = open_fixture(root.path(), OpenOptions::default()).expect("open");
6005 let settled = wait_until_settled(&opened);
6006 assert_eq!(settled.coverage, crate::Coverage::Complete);
6007
6008 std::fs::create_dir_all(root.path().join("later/deeper")).expect("fixture");
6009 std::fs::write(root.path().join("later/deeper/inner.txt"), b"i").expect("fixture");
6010 let receipt = opened.refresh(&[PathBuf::from("later")]).expect("refresh");
6011 assert!(receipt.issues.is_empty(), "{:?}", receipt.issues);
6012
6013 let index = &opened.state.index;
6014 assert_eq!(index.directory_complete(Path::new("later")).expect("lookup"), Some(true));
6015 assert_eq!(
6016 index.directory_complete(Path::new("later/deeper")).expect("lookup"),
6017 Some(true)
6018 );
6019 let response = opened
6020 .read(crate::ReadRequest {
6021 projections: vec![
6022 crate::ReadProjection::Lookup { path: PathBuf::from("later/missing") },
6023 crate::ReadProjection::Lookup { path: PathBuf::from("later/deeper/missing") },
6024 ],
6025 ..crate::ReadRequest::default()
6026 })
6027 .expect("read");
6028 assert_eq!(response.state.coverage, crate::Coverage::Complete);
6029 for result in &response.results {
6030 assert!(
6031 matches!(result, crate::ProjectionResult::Lookup(crate::Knowledge::Absent)),
6032 "{result:?}"
6033 );
6034 }
6035 let completed: Vec<_> = index
6036 .since(receipt.after.sequence)
6037 .expect("journal")
6038 .commits
6039 .iter()
6040 .flat_map(|commit| commit.state.iter())
6041 .filter_map(|transition| match transition {
6042 crate::StateTransition::DirectoryComplete { path } => Some(path.clone()),
6043 _ => None,
6044 })
6045 .collect();
6046 assert_eq!(completed, [PathBuf::from("later"), PathBuf::from("later/deeper")]);
6047 assert_eq!(
6048 receipt.state.progress.directories_complete,
6049 settled.progress.directories_complete + 2
6050 );
6051 opened.close().expect("close");
6052 }
6053
6054 #[test]
6055 fn refresh_rejects_an_unbounded_input_before_filesystem_work() {
6056 let root = tempfile::tempdir().expect("temp root");
6057 std::fs::write(root.path().join("same"), b"same").expect("fixture");
6058 let opened = open_fixture(root.path(), OpenOptions::default()).expect("opened root");
6059 assert_eq!(wait_until_settled(&opened).phase, crate::LifecyclePhase::Ready);
6060 let paths = vec![PathBuf::from("same"); MAX_REFRESH_PATHS + 2];
6061
6062 assert!(matches!(
6063 opened.refresh(&paths),
6064 Err(Error::RefreshPathLimit {
6065 attempted,
6066 limit: MAX_REFRESH_PATHS,
6067 }) if attempted == MAX_REFRESH_PATHS + 2
6068 ));
6069 opened.close().expect("close");
6070 }
6071
6072 #[test]
6073 fn refresh_rejects_stale_preparation_and_counts_the_lost_race() {
6074 let controls = Arc::new(TestControls::default());
6075 let (root, opened) = opened(Arc::clone(&controls));
6076 std::fs::write(root.path().join("race"), b"filesystem").expect("fixture");
6077 controls.gate(TestPoint::AfterRefreshVerification).arm();
6078 let refresher = opened.clone();
6079 let refresh = thread::spawn(move || refresher.refresh(&[PathBuf::from("race")]));
6080 controls.gate(TestPoint::AfterRefreshVerification).wait_reached();
6081 let concurrent = Observation::new(vec![
6082 Op::Upsert {
6083 path: PathBuf::from("race"),
6084 kind: EntryKind::File,
6085 attrs: crate::Attrs { size: 99, ..crate::Attrs::default() },
6086 },
6087 Op::Upsert {
6088 path: PathBuf::from("other"),
6089 kind: EntryKind::File,
6090 attrs: crate::Attrs { size: 5, ..crate::Attrs::default() },
6091 },
6092 ]);
6093 apply_and_notify(&opened, &concurrent);
6094 controls.gate(TestPoint::AfterRefreshVerification).release();
6095
6096 let result = refresh.join().expect("refresh thread").expect("refresh receipt");
6097 assert_eq!(result.work.observations, 1);
6098 assert_eq!(result.work.stale, 1);
6099 assert!(
6100 result.impact.dirty_paths.contains(&PathBuf::from("other")),
6101 "advancing to the receipt version must cover a concurrent producer"
6102 );
6103 assert_eq!(
6104 opened.state.index.attrs(Path::new("race")).expect("attrs").expect("retained").size,
6105 99
6106 );
6107 opened.close().expect("close");
6108 }
6109
6110 #[test]
6111 fn refresh_receipt_names_a_pass_retired_by_newer_verification() {
6112 let controls = Arc::new(TestControls::default());
6113 let (root, opened) = opened(Arc::clone(&controls));
6114 std::fs::write(root.path().join("stable.txt"), b"stable").expect("fixture");
6115 controls.gate(TestPoint::AfterRefreshVerification).arm();
6116 let refresher = opened.clone();
6117 let refresh = thread::spawn(move || refresher.refresh(&[PathBuf::from("")]));
6118 controls.gate(TestPoint::AfterRefreshVerification).wait_reached();
6119 let budget = opened.state.index.len().expect("entry count");
6123 for number in 0..=budget {
6124 let path = PathBuf::from(format!("missing-{number}"));
6125 let (epoch, _) = opened.state.index.begin_reconcile(&path).expect("begin newer");
6126 opened
6127 .state
6128 .index
6129 .finish_reconcile(
6130 &path,
6131 epoch,
6132 true,
6133 &[],
6134 &[],
6135 crate::index::ReconcileErrors {
6136 errors: &[],
6137 terminal: None,
6138 disproves_old: true,
6139 },
6140 )
6141 .expect("finish newer");
6142 }
6143 controls.gate(TestPoint::AfterRefreshVerification).release();
6144
6145 let result = refresh.join().expect("refresh thread").expect("refresh receipt");
6146
6147 assert_eq!(
6148 result.state.coverage,
6149 crate::Coverage::Partial(crate::CoverageReason::Inaccessible),
6150 "the retired pass publishes its scope partial"
6151 );
6152 assert!(
6153 result.issues.iter().any(|issue| issue.message.contains("retry")),
6154 "the receipt must name the retry the retired pass earned: {:?}",
6155 result.issues
6156 );
6157 opened.close().expect("close");
6158 }
6159
6160 #[cfg(unix)]
6164 #[test]
6165 fn multi_path_refresh_closes_each_subtree_on_its_own_walk() {
6166 use std::os::unix::fs::PermissionsExt;
6167
6168 if !crate::test_support::require_permission_bits() {
6169 return;
6170 }
6171
6172 let root = tempfile::tempdir().expect("temp root");
6173 std::fs::create_dir(root.path().join("readable")).expect("readable directory");
6174 std::fs::write(root.path().join("readable/file"), b"ok").expect("readable fixture");
6175 let blocked = root.path().join("blocked");
6176 std::fs::create_dir(&blocked).expect("blocked directory");
6177 std::fs::write(blocked.join("secret"), b"secret").expect("blocked fixture");
6178 let opened = open_fixture(root.path(), OpenOptions::default()).expect("open");
6179 let settled = wait_until_settled(&opened);
6180 assert_eq!(settled.phase, crate::LifecyclePhase::Ready);
6181 std::fs::set_permissions(&blocked, std::fs::Permissions::from_mode(0o000))
6182 .expect("make directory unreadable");
6183
6184 let receipt = opened
6185 .refresh(&[PathBuf::from("readable"), PathBuf::from("blocked")])
6186 .expect("refresh");
6187 std::fs::set_permissions(&blocked, std::fs::Permissions::from_mode(0o700))
6188 .expect("restore directory permissions");
6189
6190 assert_eq!(receipt.issues.len(), 1, "{:?}", receipt.issues);
6191 let index = &opened.state.index;
6192 assert_eq!(
6193 index.freshness_at(Path::new("readable")).expect("freshness"),
6194 crate::Freshness::Fresh
6195 );
6196 assert_eq!(
6197 index.freshness_at(Path::new("blocked")).expect("freshness"),
6198 crate::Freshness::Partial
6199 );
6200 let since = index.since(receipt.after.sequence).expect("journal");
6201 assert!(
6202 since.commits.iter().flat_map(|commit| commit.state.iter()).any(|transition| {
6203 matches!(
6204 transition,
6205 crate::StateTransition::Verified { path } if path == Path::new("readable")
6206 )
6207 }),
6208 "the readable subtree was not verified"
6209 );
6210 opened.close().expect("close");
6211 }
6212
6213 #[test]
6216 fn discovery_omits_a_child_deleted_between_listing_and_stat() {
6217 let root = tempfile::tempdir().expect("temp root");
6218 std::fs::write(root.path().join("kept"), b"kept").expect("kept fixture");
6219 std::fs::write(root.path().join("gone"), b"gone").expect("gone fixture");
6220 let hook = crate::scan::install_child_metadata_hook(root.path(), |path| {
6221 if path.file_name() == Some(std::ffi::OsStr::new("gone")) {
6222 std::fs::remove_file(path).expect("delete between listing and stat");
6223 }
6224 None
6225 });
6226 let opened = open_fixture(root.path(), OpenOptions::default()).expect("open");
6227 let settled = wait_until_settled(&opened);
6228 drop(hook);
6229
6230 assert_eq!(settled.coverage, crate::Coverage::Complete);
6231 assert_eq!(settled.issues.retained, 0);
6232 let index = &opened.state.index;
6233 assert!(index.kind(Path::new("gone")).expect("lookup").is_none());
6234 assert!(index.kind(Path::new("kept")).expect("lookup").is_some());
6235 opened.close().expect("close");
6236 }
6237
6238 #[test]
6243 fn refresh_records_completeness_per_listed_directory_despite_a_child_error() {
6244 let root = tempfile::tempdir().expect("temp root");
6245 std::fs::create_dir(root.path().join("steady")).expect("steady directory");
6246 std::fs::write(root.path().join("steady/kept"), b"kept").expect("steady fixture");
6247 let opened = open_fixture(root.path(), OpenOptions::default()).expect("open");
6248 let settled = wait_until_settled(&opened);
6249 assert_eq!(settled.coverage, crate::Coverage::Complete);
6250 std::fs::create_dir(root.path().join("fresh")).expect("fresh directory");
6251 std::fs::write(root.path().join("fresh/new"), b"new").expect("fresh fixture");
6252
6253 let hook = crate::scan::install_child_metadata_hook(root.path(), |path| {
6254 (path.file_name() == Some(std::ffi::OsStr::new("kept"))).then(|| {
6255 std::io::Error::new(std::io::ErrorKind::PermissionDenied, "injected child error")
6256 })
6257 });
6258 let receipt = opened.refresh(&[PathBuf::new()]);
6259 drop(hook);
6260 let receipt = receipt.expect("refresh");
6261
6262 assert_eq!(receipt.issues.len(), 1, "{:?}", receipt.issues);
6263 let index = &opened.state.index;
6264 assert_eq!(
6265 index.freshness_at(Path::new("")).expect("freshness"),
6266 crate::Freshness::Partial
6267 );
6268 assert_eq!(index.directory_complete(Path::new("fresh")).expect("lookup"), Some(true));
6269 assert_eq!(index.directory_complete(Path::new("steady")).expect("lookup"), Some(false));
6271 let lookup = opened
6272 .read(crate::ReadRequest {
6273 projections: vec![crate::ReadProjection::Lookup {
6274 path: PathBuf::from("fresh/missing"),
6275 }],
6276 ..crate::ReadRequest::default()
6277 })
6278 .expect("lookup")
6279 .results
6280 .into_iter()
6281 .next()
6282 .expect("lookup result");
6283 assert!(
6284 matches!(lookup, crate::ProjectionResult::Lookup(crate::Knowledge::Absent)),
6285 "{lookup:?}"
6286 );
6287 opened.close().expect("close");
6288 }
6289
6290 #[test]
6294 fn refresh_on_a_failed_root_keeps_the_issue_that_explains_it() {
6295 let (root, opened) = opened(Arc::default());
6296 std::fs::create_dir(root.path().join("sub")).expect("fixture directory");
6297 opened
6298 .state
6299 .index
6300 .transition_discovery(DiscoveryTransition::Begin)
6301 .expect("begin discovery");
6302 let failure = crate::Issue::from_error_under(
6303 root.path(),
6304 &Error::io(root.path().join("sub"), std::io::Error::other("provider failed here")),
6305 );
6306 opened
6307 .state
6308 .index
6309 .transition_discovery(DiscoveryTransition::Failed(failure))
6310 .expect("fail discovery");
6311 let failed = opened.state.index.state().expect("state");
6312 assert_eq!(failed.phase, crate::LifecyclePhase::Failed);
6313 assert_eq!(failed.issues.retained, 1);
6314
6315 let receipt = opened.refresh(&[PathBuf::from("sub")]).expect("refresh on a failed root");
6316 assert_eq!(receipt.work.stale, 0);
6317
6318 let after = opened.state.index.state().expect("state");
6319 assert_eq!(after.phase, crate::LifecyclePhase::Failed);
6320 let issues = opened.state.index.issues().expect("issues");
6321 assert_eq!(issues.len(), 1, "{issues:?}");
6322 assert_eq!(issues[0].path.as_deref(), Some(Path::new("sub")));
6323 assert_eq!(after.issues.retained, 1);
6324 opened.close().expect("close");
6325 }
6326
6327 #[test]
6328 fn close_cancels_verified_refresh_before_its_conditional_commit() {
6329 let controls = Arc::new(TestControls::default());
6330 let (root, opened) = opened(Arc::clone(&controls));
6331 std::fs::write(root.path().join("late"), b"late").expect("fixture");
6332 controls.gate(TestPoint::AfterRefreshVerification).arm();
6333 let refresher = opened.clone();
6334 let refresh = thread::spawn(move || refresher.refresh(&[PathBuf::from("late")]));
6335 controls.gate(TestPoint::AfterRefreshVerification).wait_reached();
6336 let closer = opened.clone();
6337 let close = thread::spawn(move || closer.close());
6338 let deadline = std::time::Instant::now() + TEST_GATE_TIMEOUT;
6339 while !opened.state.cancellation.is_cancelled() {
6340 assert!(std::time::Instant::now() < deadline, "close did not cancel refresh");
6341 thread::yield_now();
6342 }
6343 controls.gate(TestPoint::AfterRefreshVerification).release();
6344
6345 assert!(matches!(refresh.join().expect("refresh thread"), Err(Error::OpenedIndexClosed)));
6346 close.join().expect("close thread").expect("joined close");
6347 assert_eq!(opened.state.index.kind(Path::new("late")).expect("lookup"), None);
6348 opened.close().expect("repeat close");
6349 }
6350
6351 #[test]
6352 fn refresh_tracks_hidden_control_creation_edit_and_deletion() {
6353 let root = tempfile::tempdir().expect("temp root");
6354 std::fs::write(root.path().join("debug.log"), b"log").expect("fixture");
6355 std::fs::write(root.path().join("keep.rs"), b"keep").expect("fixture");
6356 let opened = open_fixture(
6357 root.path(),
6358 OpenOptions {
6359 hidden: Some(Arc::new(crate::HiddenPolicy::prune_hidden::<[&str; 0], &str>([]))),
6360 ..OpenOptions::default()
6361 },
6362 )
6363 .expect("opened root");
6364 assert_eq!(wait_until_settled(&opened).phase, crate::LifecyclePhase::Ready);
6365
6366 std::fs::write(root.path().join(".gitignore"), b"*.log\n").expect("create control");
6367 let created = opened.refresh(&[PathBuf::from(".gitignore")]).expect("create refresh");
6368 assert_eq!(created.accepted, vec![PathBuf::from(".gitignore")]);
6369 let image = opened.state.index.snapshot().expect("snapshot");
6370 assert!(
6371 image
6372 .controls()
6373 .expect("control state observed")
6374 .source_is(Path::new(".gitignore"), b"*.log\n")
6375 );
6376 assert_eq!(
6377 image.is_ignored(Path::new("debug.log")).expect("control state observed"),
6378 Some(true)
6379 );
6380
6381 std::fs::write(root.path().join(".gitignore"), b"*.tmp\n").expect("edit control");
6382 opened.refresh(&[PathBuf::from(".gitignore")]).expect("edit refresh");
6383 let image = opened.state.index.snapshot().expect("snapshot");
6384 assert!(
6385 image
6386 .controls()
6387 .expect("control state observed")
6388 .source_is(Path::new(".gitignore"), b"*.tmp\n")
6389 );
6390 assert_eq!(
6391 image.is_ignored(Path::new("debug.log")).expect("control state observed"),
6392 Some(false)
6393 );
6394 let unchanged =
6395 opened.refresh(&[PathBuf::from(".gitignore")]).expect("unchanged control refresh");
6396 assert_eq!(unchanged.work.observations, 1);
6397
6398 std::fs::remove_file(root.path().join(".gitignore")).expect("delete control");
6399 opened.refresh(&[PathBuf::from(".gitignore")]).expect("delete refresh");
6400 let image = opened.state.index.snapshot().expect("snapshot");
6401 assert!(image.controls().expect("control state observed").is_empty());
6402 let partitions = image.partition_total().expect("control state observed");
6403 assert_eq!(partitions.all, partitions.unignored);
6404 opened.close().expect("close");
6405 }
6406
6407 #[cfg(feature = "watch")]
6408 #[test]
6409 #[cfg(unix)]
6410 fn inaccessible_baseline_enters_watching_with_partial_coverage() {
6411 use std::os::unix::fs::PermissionsExt;
6412
6413 if !crate::test_support::require_permission_bits() {
6414 return;
6415 }
6416
6417 let root = tempfile::tempdir().expect("temp root");
6418 let scripts = tempfile::tempdir().expect("script root");
6419 let inaccessible = root.path().join("inaccessible");
6420 std::fs::create_dir(&inaccessible).expect("inaccessible directory");
6421 std::fs::write(inaccessible.join("secret"), b"secret").expect("fixture");
6422 std::fs::set_permissions(&inaccessible, std::fs::Permissions::from_mode(0o000))
6423 .expect("make directory inaccessible");
6424 let script = scripts.path().join("events.script");
6425 std::fs::write(&script, b"").expect("script");
6426
6427 let opened = open_fixture(root.path(), scripted_options(&script)).expect("open");
6428 let deadline = std::time::Instant::now() + TEST_GATE_TIMEOUT;
6429 let state = loop {
6430 let state = opened.state.index.state().expect("read state");
6431 if matches!(
6432 state.phase,
6433 crate::LifecyclePhase::Watching | crate::LifecyclePhase::Failed
6434 ) {
6435 break state;
6436 }
6437 assert!(std::time::Instant::now() < deadline, "observation handoff did not settle");
6438 std::thread::yield_now();
6439 };
6440 std::fs::set_permissions(&inaccessible, std::fs::Permissions::from_mode(0o700))
6441 .expect("restore directory permissions");
6442
6443 assert_eq!(state.phase, crate::LifecyclePhase::Watching);
6444 assert_eq!(state.coverage, crate::Coverage::Partial(crate::CoverageReason::Inaccessible));
6445 assert_eq!(state.freshness, crate::Freshness::Partial);
6446 assert!(state.issues.retained > 0);
6447 opened.close().expect("close");
6448 }
6449
6450 #[cfg(all(unix, feature = "watch"))]
6457 #[test]
6458 fn watching_after_a_clean_handoff_rederives_complete_coverage() {
6459 use std::os::unix::fs::PermissionsExt;
6460
6461 if !crate::test_support::require_permission_bits() {
6462 return;
6463 }
6464 let root = tempfile::tempdir().expect("temp root");
6465 let scripts = tempfile::tempdir().expect("script root");
6466 let script = scripts.path().join("events.script");
6467 std::fs::write(&script, b"").expect("script");
6468 let blocked = root.path().join("blocked");
6469 std::fs::create_dir(&blocked).expect("blocked directory");
6470 std::fs::write(blocked.join("secret"), b"secret").expect("fixture");
6471 std::fs::set_permissions(&blocked, std::fs::Permissions::from_mode(0o000))
6472 .expect("make directory inaccessible");
6473 let controls = Arc::new(TestControls::default());
6474 controls.gate(TestPoint::BeforeObservationHandoff).arm();
6475 let opened = OpenedIndex::open_for_test(
6476 root.path(),
6477 scripted_options(&script),
6478 Arc::clone(&controls),
6479 )
6480 .expect("open");
6481 controls.gate(TestPoint::BeforeObservationHandoff).wait_reached();
6482 let discovered = opened.state.index.state().expect("state after discovery");
6483 std::fs::set_permissions(&blocked, std::fs::Permissions::from_mode(0o700))
6484 .expect("restore directory permissions");
6485 controls.gate(TestPoint::BeforeObservationHandoff).release();
6486
6487 let state = wait_until_phase(&opened, crate::LifecyclePhase::Watching);
6488 assert_eq!(
6489 discovered.coverage,
6490 crate::Coverage::Partial(crate::CoverageReason::Inaccessible)
6491 );
6492 assert_eq!(state.coverage, crate::Coverage::Complete);
6493 assert_eq!(state.freshness, crate::Freshness::Fresh);
6494 assert_eq!(
6495 opened.state.index.kind(Path::new("blocked/secret")).expect("lookup"),
6496 Some(EntryKind::File),
6497 "the handoff read the formerly inaccessible directory"
6498 );
6499 assert_eq!(
6503 opened.state.index.directory_complete(Path::new("blocked")).expect("lookup"),
6504 Some(true)
6505 );
6506 let response = opened
6507 .read(crate::ReadRequest {
6508 projections: vec![crate::ReadProjection::Lookup {
6509 path: PathBuf::from("blocked/missing"),
6510 }],
6511 ..crate::ReadRequest::default()
6512 })
6513 .expect("read");
6514 assert_eq!(response.state.phase, crate::LifecyclePhase::Watching);
6515 assert_eq!(response.state.coverage, crate::Coverage::Complete);
6516 assert!(
6517 matches!(
6518 response.results[0],
6519 crate::ProjectionResult::Lookup(crate::Knowledge::Absent)
6520 ),
6521 "{:?}",
6522 response.results[0]
6523 );
6524 opened.close().expect("close");
6525 }
6526
6527 #[cfg(feature = "watch")]
6533 #[test]
6534 fn a_directory_the_observer_adds_is_complete_after_its_relist() {
6535 let root = tempfile::tempdir().expect("temp root");
6536 let scripts = tempfile::tempdir().expect("script root");
6537 let script = scripts.path().join("events.script");
6538 std::fs::write(&script, b"").expect("script");
6539 let controls = Arc::new(TestControls::default());
6540 let opened = OpenedIndex::open_for_test(
6541 root.path(),
6542 scripted_options(&script),
6543 Arc::clone(&controls),
6544 )
6545 .expect("open scripted observer");
6546 wait_until_phase(&opened, crate::LifecyclePhase::Watching);
6547 let start = current_version(&opened);
6548 let watching = opened.read(crate::ReadRequest::default()).expect("read state").state;
6549 assert_eq!(watching.coverage, crate::Coverage::Complete);
6550
6551 std::fs::create_dir_all(root.path().join("later/deeper")).expect("fixture");
6552 std::fs::write(root.path().join("later/deeper/inner.txt"), b"i").expect("fixture");
6553 controls.send_observation_hints("create-dir\tlater\n");
6554 wait_for_observed_walk(&opened, &controls, root.path(), start, Path::new("later"), 1);
6555
6556 let index = &opened.state.index;
6557 assert_eq!(
6558 index.kind(Path::new("later/deeper/inner.txt")).expect("lookup"),
6559 Some(EntryKind::File)
6560 );
6561 assert_eq!(index.directory_complete(Path::new("later")).expect("lookup"), Some(true));
6562 assert_eq!(
6563 index.directory_complete(Path::new("later/deeper")).expect("lookup"),
6564 Some(true)
6565 );
6566 let response = opened
6567 .read(crate::ReadRequest {
6568 projections: vec![
6569 crate::ReadProjection::Lookup { path: PathBuf::from("later/missing") },
6570 crate::ReadProjection::Lookup { path: PathBuf::from("later/deeper/missing") },
6571 ],
6572 ..crate::ReadRequest::default()
6573 })
6574 .expect("read");
6575 assert_eq!(response.state.coverage, crate::Coverage::Complete);
6576 for result in &response.results {
6577 assert!(
6578 matches!(result, crate::ProjectionResult::Lookup(crate::Knowledge::Absent)),
6579 "{result:?}"
6580 );
6581 }
6582 let since = index.since(start.sequence).expect("journal");
6583 let completed: Vec<_> = since
6584 .commits
6585 .iter()
6586 .flat_map(|commit| commit.state.iter())
6587 .filter_map(|transition| match transition {
6588 crate::StateTransition::DirectoryComplete { path } => Some(path.clone()),
6589 _ => None,
6590 })
6591 .collect();
6592 assert_eq!(completed, [PathBuf::from("later"), PathBuf::from("later/deeper")]);
6593 assert_eq!(
6594 response.state.progress.directories_complete,
6595 watching.progress.directories_complete + 2,
6596 "the progress count agrees with the transitions"
6597 );
6598 opened.close().expect("close");
6599 }
6600
6601 #[cfg(feature = "watch")]
6603 #[test]
6604 fn a_directory_that_vanishes_during_discovery_leaves_a_watched_root_complete() {
6605 let controls = Arc::new(TestControls::default());
6606 controls.gate(TestPoint::AfterRootDirectory).arm();
6607 let root = tempfile::tempdir().expect("temp root");
6608 let scripts = tempfile::tempdir().expect("script root");
6609 let script = scripts.path().join("events.script");
6610 std::fs::write(&script, b"").expect("script");
6611 std::fs::create_dir(root.path().join("sub")).expect("sub");
6612 std::fs::write(root.path().join("sub/inner.txt"), b"x").expect("fixture");
6613 std::fs::write(root.path().join("keep.txt"), b"k").expect("fixture");
6614 let opened = OpenedIndex::open_for_test(
6615 root.path(),
6616 scripted_options(&script),
6617 Arc::clone(&controls),
6618 )
6619 .expect("open");
6620 controls.gate(TestPoint::AfterRootDirectory).wait_reached();
6621 std::fs::remove_dir_all(root.path().join("sub")).expect("remove sub during discovery");
6622 controls.gate(TestPoint::AfterRootDirectory).release();
6623
6624 let state = wait_until_phase(&opened, crate::LifecyclePhase::Watching);
6625 assert_eq!(state.coverage, crate::Coverage::Complete);
6626 assert_eq!(state.freshness, crate::Freshness::Fresh);
6627 assert_eq!(state.issues.retained, 0);
6628 assert_eq!(
6629 opened.state.index.kind(Path::new("sub")).expect("lookup"),
6630 None,
6631 "the handoff pass removed the vanished directory"
6632 );
6633 assert_eq!(opened.state.index.directory_complete(Path::new("")).expect("root"), Some(true));
6634 opened.close().expect("close");
6635 }
6636
6637 #[cfg(feature = "watch")]
6638 #[test]
6639 fn observation_is_captured_before_baseline_and_closes_the_handoff_gap() {
6640 let root = tempfile::tempdir().expect("temp root");
6641 let scripts = tempfile::tempdir().expect("script root");
6642 let path = root.path().join("during.txt");
6643 std::fs::write(&path, b"before").expect("fixture");
6644 let script = scripts.path().join("events.script");
6645 std::fs::write(&script, b"modify\tduring.txt\n").expect("script");
6646 let controls = Arc::new(TestControls::default());
6647 controls.gate(TestPoint::BeforeDiscovery).arm();
6648
6649 let opened = OpenedIndex::open_for_test(
6650 root.path(),
6651 scripted_options(&script),
6652 Arc::clone(&controls),
6653 )
6654 .expect("opened observed root");
6655 controls.gate(TestPoint::BeforeDiscovery).wait_reached();
6656 std::fs::write(&path, b"changed-during-baseline").expect("mutate during handoff");
6657 controls.gate(TestPoint::BeforeDiscovery).release();
6658
6659 let state = wait_until_phase(&opened, crate::LifecyclePhase::Watching);
6660 assert_eq!(state.freshness, crate::Freshness::Fresh);
6661 assert_eq!(state.coverage, crate::Coverage::Complete);
6662 assert_eq!(
6663 opened
6664 .state
6665 .index
6666 .attrs(Path::new("during.txt"))
6667 .expect("attrs")
6668 .expect("retained")
6669 .size,
6670 23
6671 );
6672 let since = opened.state.index.since(crate::Clock::ZERO).expect("journal");
6673 assert!(since.commits.iter().any(|commit| {
6674 commit.state.iter().any(|transition| {
6675 matches!(
6676 transition,
6677 crate::StateTransition::IndexState { current, .. }
6678 if current.phase == crate::LifecyclePhase::Reconciling
6679 )
6680 })
6681 }));
6682 assert!(since.commits.iter().any(|commit| {
6683 commit.state.iter().any(|transition| {
6684 matches!(
6685 transition,
6686 crate::StateTransition::IndexState { current, .. }
6687 if current.phase == crate::LifecyclePhase::Watching
6688 )
6689 })
6690 }));
6691 opened.close().expect("joined close");
6692 }
6693
6694 #[cfg(feature = "watch")]
6695 #[test]
6696 fn scripted_overflow_is_provider_recovery_not_a_consumer_reset() {
6697 let root = tempfile::tempdir().expect("temp root");
6698 let scripts = tempfile::tempdir().expect("script root");
6699 std::fs::create_dir(root.path().join("src")).expect("directory");
6700 std::fs::write(root.path().join("src/present"), b"present").expect("fixture");
6701 let script = scripts.path().join("events.script");
6702 std::fs::write(&script, b"rescan\tsrc\n").expect("script");
6703 let opened = open_fixture(root.path(), scripted_options(&script)).expect("open");
6704
6705 let state = wait_until_phase(&opened, crate::LifecyclePhase::Watching);
6706 assert_eq!(state.freshness, crate::Freshness::Fresh);
6707 assert!(state.issues.retained > 0);
6708 let since = opened.state.index.since(crate::Clock::ZERO).expect("journal");
6709 assert!(!since.truncated);
6710 assert!(since.commits.iter().any(|commit| commit.changes.iter().any(|change| {
6711 matches!(
6712 change,
6713 crate::EffectiveChange::Invalidated {
6714 path,
6715 reason: crate::InvalidateReason::WatchOverflow,
6716 } if path == Path::new("src")
6717 )
6718 })));
6719 assert!(
6720 opened
6721 .state
6722 .index
6723 .issues()
6724 .expect("issues")
6725 .iter()
6726 .any(|issue| issue.kind == crate::IssueKind::ObservationGap)
6727 );
6728 opened.close().expect("close");
6729 }
6730
6731 #[cfg(feature = "watch")]
6732 #[test]
6733 fn scripted_directory_creation_closes_the_registration_gap() {
6734 let root = tempfile::tempdir().expect("temp root");
6735 let scripts = tempfile::tempdir().expect("script root");
6736 std::fs::create_dir(root.path().join("newdir")).expect("directory");
6737 std::fs::write(root.path().join("newdir/child"), b"child").expect("fixture");
6738 let script = scripts.path().join("events.script");
6739 std::fs::write(&script, b"create-dir\tnewdir\n").expect("script");
6740 let opened = open_fixture(root.path(), scripted_options(&script)).expect("open");
6741
6742 wait_until_phase(&opened, crate::LifecyclePhase::Watching);
6743 let since = opened.state.index.since(crate::Clock::ZERO).expect("journal");
6744 assert!(since.commits.iter().any(|commit| commit.changes.iter().any(|change| {
6745 matches!(
6746 change,
6747 crate::EffectiveChange::Invalidated {
6748 path,
6749 reason: crate::InvalidateReason::WatchSetupRace,
6750 } if path == Path::new("newdir")
6751 )
6752 })));
6753 assert_eq!(
6754 opened.state.index.kind(Path::new("newdir/child")).expect("lookup"),
6755 Some(EntryKind::File)
6756 );
6757 opened.close().expect("close");
6758 }
6759
6760 #[cfg(feature = "watch")]
6761 #[test]
6762 fn scripted_observation_keeps_the_opened_index_live_after_handoff() {
6763 let root = tempfile::tempdir().expect("temp root");
6764 let scripts = tempfile::tempdir().expect("script root");
6765 let script = scripts.path().join("events.script");
6766 std::fs::write(&script, b"").expect("script");
6767 let controls = Arc::new(TestControls::default());
6768 let opened = OpenedIndex::open_for_test(
6769 root.path(),
6770 scripted_options(&script),
6771 Arc::clone(&controls),
6772 )
6773 .expect("open scripted observer");
6774 wait_until_phase(&opened, crate::LifecyclePhase::Watching);
6775 let before = current_version(&opened);
6776
6777 std::fs::write(root.path().join("live.txt"), b"live").expect("live mutation");
6778 controls.send_observation_hints("create\tlive.txt\n");
6779 let deadline = std::time::Instant::now() + TEST_GATE_TIMEOUT;
6780 loop {
6781 if opened.state.index.kind(Path::new("live.txt")).expect("lookup")
6782 == Some(EntryKind::File)
6783 {
6784 break;
6785 }
6786 assert!(std::time::Instant::now() < deadline, "scripted hint was not applied");
6787 std::thread::yield_now();
6788 }
6789 let since = opened.state.index.since(before.sequence).expect("journal");
6790 assert!(since.commits.iter().any(|commit| commit.changes.iter().any(|change| {
6791 matches!(
6792 change,
6793 crate::EffectiveChange::Inserted { path, .. }
6794 if path == Path::new("live.txt")
6795 )
6796 })));
6797 opened.close().expect("close");
6798 }
6799
6800 #[cfg(feature = "watch")]
6801 #[test]
6802 fn live_observation_gap_recovers_before_reporting_freshness() {
6803 let root = tempfile::tempdir().expect("temp root");
6804 let scripts = tempfile::tempdir().expect("script root");
6805 let script = scripts.path().join("events.script");
6806 std::fs::write(&script, b"").expect("script");
6807 let controls = Arc::new(TestControls::default());
6808 let opened = OpenedIndex::open_for_test(
6809 root.path(),
6810 scripted_options(&script),
6811 Arc::clone(&controls),
6812 )
6813 .expect("open scripted observer");
6814 wait_until_phase(&opened, crate::LifecyclePhase::Watching);
6815 let before = current_version(&opened);
6816
6817 std::fs::write(root.path().join("recovered.txt"), b"recovered").expect("missed mutation");
6818 controls.send_observation_hints("rescan\t.\n");
6819 let deadline = std::time::Instant::now() + TEST_GATE_TIMEOUT;
6820 loop {
6821 let state = opened.state.index.state().expect("read state");
6822 let recovered = opened.state.index.kind(Path::new("recovered.txt")).expect("lookup")
6823 == Some(EntryKind::File);
6824 if recovered && state.freshness == crate::Freshness::Fresh {
6825 break;
6826 }
6827 assert!(std::time::Instant::now() < deadline, "gap recovery did not finish");
6828 std::thread::yield_now();
6829 }
6830
6831 let since = opened.state.index.since(before.sequence).expect("journal");
6832 assert!(since.commits.iter().any(|commit| commit.changes.iter().any(|change| {
6833 matches!(
6834 change,
6835 crate::EffectiveChange::Invalidated {
6836 path,
6837 reason: crate::InvalidateReason::WatchOverflow,
6838 } if path.as_os_str().is_empty()
6839 )
6840 })));
6841 assert!(since.commits.iter().any(|commit| commit.state.iter().any(|transition| {
6842 matches!(
6843 transition,
6844 crate::StateTransition::Freshness {
6845 current: crate::Freshness::Reconciling | crate::Freshness::Stale,
6846 ..
6847 }
6848 )
6849 })));
6850 assert!(since.commits.iter().any(|commit| commit.state.iter().any(|transition| {
6851 matches!(
6852 transition,
6853 crate::StateTransition::Freshness { current: crate::Freshness::Fresh, .. }
6854 )
6855 })));
6856 assert!(
6857 opened
6858 .state
6859 .index
6860 .issues()
6861 .expect("issues")
6862 .iter()
6863 .any(|issue| issue.kind == crate::IssueKind::ObservationGap)
6864 );
6865 opened.close().expect("close");
6866 }
6867
6868 #[cfg(feature = "watch")]
6869 fn reconciles_of(opened: &OpenedIndex, since: crate::EngineVersion, path: &Path) -> usize {
6870 opened
6871 .state
6872 .index
6873 .since(since.sequence)
6874 .expect("journal")
6875 .commits
6876 .iter()
6877 .flat_map(|commit| commit.state.iter())
6878 .filter(|transition| {
6879 matches!(
6880 transition,
6881 crate::StateTransition::Freshness { path: marked, current, .. }
6882 if marked == path && *current == crate::Freshness::Reconciling
6883 )
6884 })
6885 .count()
6886 }
6887
6888 #[cfg(feature = "watch")]
6896 fn wait_for_observed_walk(
6897 opened: &OpenedIndex,
6898 controls: &TestControls,
6899 root: &Path,
6900 start: crate::EngineVersion,
6901 path: &Path,
6902 walks: usize,
6903 ) {
6904 let deadline = std::time::Instant::now() + TEST_GATE_TIMEOUT;
6905 while reconciles_of(opened, start, path) < walks {
6906 assert!(std::time::Instant::now() < deadline, "walk {walks} of {path:?} did not begin");
6907 std::thread::yield_now();
6908 }
6909 let marker = format!("marker-{}-{walks}.txt", path.display());
6910 std::fs::write(root.join(&marker), &marker).expect("marker");
6911 controls.send_observation_hints(&format!("create\t{marker}\n"));
6912 let deadline = std::time::Instant::now() + TEST_GATE_TIMEOUT;
6913 while opened.state.index.kind(Path::new(&marker)).expect("lookup") != Some(EntryKind::File)
6914 {
6915 assert!(std::time::Instant::now() < deadline, "{marker} was not applied");
6916 std::thread::yield_now();
6917 }
6918 }
6919
6920 #[cfg(all(unix, feature = "watch"))]
6928 #[test]
6929 fn an_unreadable_gap_is_walked_once_and_explains_itself() {
6930 use std::os::unix::fs::PermissionsExt;
6931
6932 if !crate::test_support::require_permission_bits() {
6933 return;
6934 }
6935 let root = tempfile::tempdir().expect("temp root");
6936 let scripts = tempfile::tempdir().expect("script root");
6937 let script = scripts.path().join("events.script");
6938 std::fs::write(&script, b"").expect("script");
6939 let blocked = root.path().join("blocked");
6940 std::fs::create_dir(&blocked).expect("blocked");
6941 std::fs::write(blocked.join("secret"), b"s").expect("fixture");
6942 let controls = Arc::new(TestControls::default());
6943 let opened = OpenedIndex::open_for_test(
6944 root.path(),
6945 scripted_options(&script),
6946 Arc::clone(&controls),
6947 )
6948 .expect("open scripted observer");
6949 let watching = wait_until_phase(&opened, crate::LifecyclePhase::Watching);
6950 assert_eq!(watching.freshness, crate::Freshness::Fresh);
6951 let start = current_version(&opened);
6952
6953 std::fs::set_permissions(&blocked, std::fs::Permissions::from_mode(0o000))
6954 .expect("make directory inaccessible");
6955 controls.send_observation_hints("rescan\tblocked\n");
6956 let deadline = std::time::Instant::now() + TEST_GATE_TIMEOUT;
6957 while reconciles_of(&opened, start, Path::new("blocked")) < 1 {
6958 assert!(std::time::Instant::now() < deadline, "the gap was not reconciled");
6959 std::thread::yield_now();
6960 }
6961
6962 for name in ["live.txt", "marker.txt"] {
6965 std::fs::write(root.path().join(name), name).expect("unrelated mutation");
6966 controls.send_observation_hints(&format!("create\t{name}\n"));
6967 let deadline = std::time::Instant::now() + TEST_GATE_TIMEOUT;
6968 while opened.state.index.kind(Path::new(name)).expect("lookup") != Some(EntryKind::File)
6969 {
6970 assert!(std::time::Instant::now() < deadline, "{name} was not applied");
6971 std::thread::yield_now();
6972 }
6973 }
6974
6975 let walks = reconciles_of(&opened, start, Path::new("blocked"));
6976 let state = opened.state.index.state().expect("state");
6977 let issues = opened.state.index.issues().expect("issues");
6978 std::fs::set_permissions(&blocked, std::fs::Permissions::from_mode(0o700))
6979 .expect("restore directory permissions");
6980 assert_eq!(walks, 1, "an unreadable subtree must not be re-walked per unrelated event");
6981 assert_eq!(state.phase, crate::LifecyclePhase::Watching);
6982 assert_eq!(state.freshness, crate::Freshness::Partial);
6983 assert!(
6984 issues.iter().any(|issue| issue.kind == crate::IssueKind::Permission
6985 && issue.path.as_deref().is_some_and(|path| path.ends_with("blocked"))),
6986 "{issues:?}"
6987 );
6988 opened.close().expect("close");
6989 }
6990
6991 #[cfg(all(unix, feature = "watch"))]
6999 #[test]
7000 fn repeated_unreadable_reconciles_retain_one_issue_per_boundary() {
7001 use std::os::unix::fs::PermissionsExt;
7002
7003 if !crate::test_support::require_permission_bits() {
7004 return;
7005 }
7006 let root = tempfile::tempdir().expect("temp root");
7007 let scripts = tempfile::tempdir().expect("script root");
7008 let script = scripts.path().join("events.script");
7009 std::fs::write(&script, b"").expect("script");
7010 let blocked = root.path().join("blocked");
7011 std::fs::create_dir(&blocked).expect("blocked");
7012 std::fs::write(blocked.join("secret"), b"s").expect("fixture");
7013 let controls = Arc::new(TestControls::default());
7014 let opened = OpenedIndex::open_for_test(
7015 root.path(),
7016 scripted_options(&script),
7017 Arc::clone(&controls),
7018 )
7019 .expect("open scripted observer");
7020 let watching = wait_until_phase(&opened, crate::LifecyclePhase::Watching);
7021 assert_eq!(watching.issues.retained, 0);
7022 let start = current_version(&opened);
7023 let walks = std::cell::Cell::new(0);
7024 let rescan_then_wait = || {
7025 controls.send_observation_hints("rescan\tblocked\n");
7026 walks.set(walks.get() + 1);
7027 wait_for_observed_walk(
7028 &opened,
7029 &controls,
7030 root.path(),
7031 start,
7032 Path::new("blocked"),
7033 walks.get(),
7034 );
7035 };
7036 let issues_of = |kind: crate::IssueKind| {
7037 opened
7038 .state
7039 .index
7040 .issues()
7041 .expect("issues")
7042 .into_iter()
7043 .filter(|issue| issue.kind == kind)
7044 .collect::<Vec<_>>()
7045 };
7046
7047 std::fs::set_permissions(&blocked, std::fs::Permissions::from_mode(0o000))
7048 .expect("make directory inaccessible");
7049 for _ in 0..5 {
7050 rescan_then_wait();
7051 }
7052 let state = opened.state.index.state().expect("state");
7053 let permission = issues_of(crate::IssueKind::Permission);
7054 let gaps = issues_of(crate::IssueKind::ObservationGap);
7055 std::fs::set_permissions(&blocked, std::fs::Permissions::from_mode(0o700))
7056 .expect("restore directory permissions");
7057 assert_eq!(state.phase, crate::LifecyclePhase::Watching);
7058 assert_eq!(permission.len(), 1, "{permission:?}");
7059 assert_eq!(permission[0].path.as_deref(), Some(Path::new("blocked")), "{permission:?}");
7060 assert_eq!(gaps.len(), 1, "{gaps:?}");
7061 assert_eq!(gaps[0].path.as_deref(), Some(Path::new("blocked")), "{gaps:?}");
7062 assert_eq!(state.issues, crate::IssueSummary { retained: 2, omitted: 0 });
7063
7064 rescan_then_wait();
7067 let state = opened.state.index.state().expect("state");
7068 assert_eq!(issues_of(crate::IssueKind::Permission), []);
7069 assert_eq!(issues_of(crate::IssueKind::ObservationGap).len(), 1);
7070 assert_eq!(state.issues, crate::IssueSummary { retained: 1, omitted: 0 });
7071 assert_eq!(state.freshness, crate::Freshness::Fresh);
7072 assert_eq!(
7073 opened.state.index.kind(Path::new("blocked/secret")).expect("lookup"),
7074 Some(EntryKind::File)
7075 );
7076 opened.close().expect("close");
7077 }
7078
7079 #[cfg(feature = "watch")]
7080 #[test]
7081 fn live_observation_shares_the_exact_opened_root_file_budget() {
7082 let root = tempfile::tempdir().expect("temp root");
7083 let scripts = tempfile::tempdir().expect("script root");
7084 std::fs::write(root.path().join("baseline.txt"), b"baseline").expect("fixture");
7085 let script = scripts.path().join("events.script");
7086 std::fs::write(&script, b"").expect("script");
7087 let controls = Arc::new(TestControls::default());
7088 let mut options = scripted_options(&script);
7089 options.budget.max_files = Some(1);
7090 let opened = OpenedIndex::open_for_test(root.path(), options, Arc::clone(&controls))
7091 .expect("open scripted observer");
7092 wait_until_phase(&opened, crate::LifecyclePhase::Watching);
7093
7094 std::fs::write(root.path().join("over-budget.txt"), b"refused")
7095 .expect("over-budget mutation");
7096 controls.send_observation_hints("create\tover-budget.txt\n");
7097 let state = wait_until_phase(&opened, crate::LifecyclePhase::Stopped);
7098
7099 assert_eq!(state.coverage, crate::Coverage::Partial(crate::CoverageReason::Budget));
7100 assert_eq!(opened.state.index.kind(Path::new("over-budget.txt")).expect("lookup"), None);
7101 assert_eq!(opened.state.index.total().expect("total").files, 1);
7102 assert!(
7103 opened
7104 .state
7105 .index
7106 .issues()
7107 .expect("issues")
7108 .iter()
7109 .any(|issue| issue.kind == crate::IssueKind::ResourceBudget)
7110 );
7111 opened.close().expect("close");
7112 }
7113
7114 #[cfg(feature = "watch")]
7115 #[test]
7116 fn close_after_observation_verification_prevents_publication() {
7117 let root = tempfile::tempdir().expect("temp root");
7118 let scripts = tempfile::tempdir().expect("script root");
7119 let script = scripts.path().join("events.script");
7120 std::fs::write(&script, b"").expect("script");
7121 let controls = Arc::new(TestControls::default());
7122 let opened = OpenedIndex::open_for_test(
7123 root.path(),
7124 scripted_options(&script),
7125 Arc::clone(&controls),
7126 )
7127 .expect("open scripted observer");
7128 wait_until_phase(&opened, crate::LifecyclePhase::Watching);
7129
7130 controls.gate(TestPoint::AfterObservationVerification).arm();
7131 std::fs::write(root.path().join("too-late.txt"), b"verified").expect("late mutation");
7132 controls.send_observation_hints("create\ttoo-late.txt\n");
7133 controls.gate(TestPoint::AfterObservationVerification).wait_reached();
7134
7135 let closer = opened.clone();
7136 let close = thread::spawn(move || closer.close());
7137 let deadline = std::time::Instant::now() + TEST_GATE_TIMEOUT;
7138 while !opened.state.cancellation.is_cancelled() {
7139 assert!(std::time::Instant::now() < deadline, "close did not cancel observation");
7140 thread::yield_now();
7141 }
7142 assert!(!close.is_finished(), "close returned before the commit boundary released");
7143 controls.gate(TestPoint::AfterObservationVerification).release();
7144 close.join().expect("close thread").expect("joined close");
7145
7146 assert_eq!(opened.state.index.kind(Path::new("too-late.txt")).expect("lookup"), None);
7147 opened.close().expect("repeat close");
7148 }
7149
7150 #[cfg(feature = "watch")]
7151 #[test]
7152 fn stopped_discovery_never_claims_to_be_watching() {
7153 let root = tempfile::tempdir().expect("temp root");
7154 let scripts = tempfile::tempdir().expect("script root");
7155 std::fs::write(root.path().join("one"), b"one").expect("fixture");
7156 std::fs::write(root.path().join("two"), b"two").expect("fixture");
7157 let script = scripts.path().join("events.script");
7158 std::fs::write(&script, b"").expect("script");
7159 let mut options = scripted_options(&script);
7160 options.budget.max_files = Some(1);
7161 let opened = open_fixture(root.path(), options).expect("open");
7162
7163 wait_until_phase(&opened, crate::LifecyclePhase::Stopped);
7164 let since = opened.state.index.since(crate::Clock::ZERO).expect("journal");
7165 assert!(!since.commits.iter().any(|commit| commit.state.iter().any(|transition| {
7166 matches!(
7167 transition,
7168 crate::StateTransition::IndexState { current, .. }
7169 if current.phase == crate::LifecyclePhase::Watching
7170 )
7171 })));
7172 opened.close().expect("close");
7173 }
7174
7175 #[cfg(feature = "watch")]
7176 #[test]
7177 fn close_joins_an_observation_worker_blocked_at_a_named_boundary() {
7178 let root = tempfile::tempdir().expect("temp root");
7179 let scripts = tempfile::tempdir().expect("script root");
7180 let script = scripts.path().join("events.script");
7181 std::fs::write(&script, b"").expect("script");
7182 let controls = Arc::new(TestControls::default());
7183 controls.gate(TestPoint::BeforeObservationPoll).arm();
7184 let opened = OpenedIndex::open_for_test(
7185 root.path(),
7186 scripted_options(&script),
7187 Arc::clone(&controls),
7188 )
7189 .expect("open");
7190 controls.gate(TestPoint::BeforeObservationPoll).wait_reached();
7191
7192 let closer = opened.clone();
7193 let close = thread::spawn(move || closer.close());
7194 let deadline = std::time::Instant::now() + TEST_GATE_TIMEOUT;
7195 while !opened.state.cancellation.is_cancelled() {
7196 assert!(std::time::Instant::now() < deadline, "close did not cancel observation");
7197 thread::yield_now();
7198 }
7199 assert!(!close.is_finished(), "close returned before the owned worker was released");
7200 controls.gate(TestPoint::BeforeObservationPoll).release();
7201 close.join().expect("close thread").expect("joined close");
7202 opened.close().expect("repeat close");
7203 }
7204
7205 #[cfg(feature = "watch")]
7206 #[test]
7207 fn malformed_script_fails_before_discovery_starts() {
7208 let root = tempfile::tempdir().expect("temp root");
7209 let scripts = tempfile::tempdir().expect("script root");
7210 let script = scripts.path().join("events.script");
7211 std::fs::write(&script, b"teleport\tmissing\n").expect("script");
7212 let error = open_fixture(root.path(), scripted_options(&script))
7213 .expect_err("invalid observer configuration must fail open");
7214 assert!(matches!(error, Error::WatchScript(_)));
7215 }
7216
7217 #[cfg(feature = "watch")]
7218 #[test]
7219 fn observation_rejects_a_restricted_scope_before_open_returns() {
7220 let root = tempfile::tempdir().expect("temp root");
7221 let scripts = tempfile::tempdir().expect("script root");
7222 let script = scripts.path().join("events.script");
7223 std::fs::write(&script, b"").expect("script");
7224 let mut options = scripted_options(&script);
7225 options.one_filesystem = true;
7226
7227 let error = open_fixture(root.path(), options)
7228 .expect_err("unsupported observed scope must fail open");
7229 #[cfg(not(unix))]
7232 assert!(matches!(
7233 error,
7234 Error::InvalidRequest(crate::query::RequestError::ScopeUnsupported {
7235 axis: crate::query::ScopeAxis::OneFilesystem,
7236 ..
7237 })
7238 ));
7239 #[cfg(unix)]
7240 assert!(matches!(error, Error::UnsupportedScanConfig(_)));
7241 }
7242
7243 #[cfg(feature = "watch")]
7244 #[test]
7245 fn handoff_retries_a_benign_refresh_conflict_before_watching() {
7246 let root = tempfile::tempdir().expect("temp root");
7247 let scripts = tempfile::tempdir().expect("script root");
7248 let path = root.path().join("shared.txt");
7249 std::fs::write(&path, b"before").expect("fixture");
7250 let script = scripts.path().join("events.script");
7251 std::fs::write(&script, b"").expect("script");
7252 let controls = Arc::new(TestControls::default());
7253 controls.gate(TestPoint::AfterObservationVerification).arm();
7254 let opened = OpenedIndex::open_for_test(
7255 root.path(),
7256 scripted_options(&script),
7257 Arc::clone(&controls),
7258 )
7259 .expect("open scripted observer");
7260 controls.gate(TestPoint::AfterObservationVerification).wait_reached();
7261
7262 std::fs::write(&path, b"updated-by-refresh").expect("concurrent mutation");
7263 let refreshed =
7264 opened.refresh(&[PathBuf::from("shared.txt")]).expect("overlapping refresh succeeds");
7265 assert_eq!(refreshed.work.stale, 0);
7266 controls.gate(TestPoint::AfterObservationVerification).release();
7267
7268 wait_until_phase(&opened, crate::LifecyclePhase::Watching);
7269 assert_eq!(
7270 opened
7271 .state
7272 .index
7273 .attrs(Path::new("shared.txt"))
7274 .expect("attrs")
7275 .expect("retained")
7276 .size,
7277 18
7278 );
7279 opened.close().expect("close");
7280 }
7281
7282 #[cfg(feature = "watch")]
7289 #[test]
7290 fn handoff_settles_through_a_convergent_refresh_without_a_second_walk() {
7291 for name in ["shared.txt", ".gitignore"] {
7292 let root = tempfile::tempdir().expect("temp root");
7293 let scripts = tempfile::tempdir().expect("script root");
7294 let path = root.path().join(name);
7295 std::fs::write(&path, b"before").expect("fixture");
7296 let script = scripts.path().join("events.script");
7297 std::fs::write(&script, b"").expect("script");
7298 let controls = Arc::new(TestControls::default());
7299 controls.gate(TestPoint::BeforeObservationHandoff).arm();
7300 controls.gate(TestPoint::AfterObservationVerification).arm();
7301 let opened = OpenedIndex::open_for_test(
7302 root.path(),
7303 scripted_options(&script),
7304 Arc::clone(&controls),
7305 )
7306 .expect("open scripted observer");
7307 controls.gate(TestPoint::BeforeObservationHandoff).wait_reached();
7308 std::fs::write(&path, b"changed").expect("mutation before the handoff walk");
7310 controls.gate(TestPoint::BeforeObservationHandoff).release();
7311 controls.gate(TestPoint::AfterObservationVerification).wait_reached();
7312
7313 let refreshed = opened.refresh(&[PathBuf::from(name)]).expect("refresh");
7315 assert_eq!(refreshed.work.stale, 0, "{name}");
7316 let ops = if name == crate::control::CONTROL_FILE_NAME { 2 } else { 1 };
7317 assert_eq!(refreshed.work.observations, ops, "{name}");
7318 controls.gate(TestPoint::AfterObservationVerification).release();
7319
7320 let state = wait_until_phase(&opened, crate::LifecyclePhase::Watching);
7321 assert_eq!(state.coverage, crate::Coverage::Complete, "{name}");
7322 assert_eq!(state.freshness, crate::Freshness::Fresh, "{name}");
7323 let since = opened.state.index.since(crate::Clock::ZERO).expect("journal");
7324 let transitions: Vec<&crate::StateTransition> =
7325 since.commits.iter().flat_map(|commit| commit.state.iter()).collect();
7326 assert!(
7328 !transitions.iter().any(|transition| matches!(
7329 transition,
7330 crate::StateTransition::Freshness { current: crate::Freshness::Partial, .. }
7331 )),
7332 "{name}: the handoff's first pass was refused: {transitions:?}"
7333 );
7334 assert_eq!(
7335 transitions
7336 .iter()
7337 .filter(|transition| matches!(
7338 transition,
7339 crate::StateTransition::Verified { path } if path.as_os_str().is_empty()
7340 ))
7341 .count(),
7342 1,
7343 "{name}"
7344 );
7345 let index = &opened.state.index;
7346 assert_eq!(
7347 index.attrs(Path::new(name)).expect("attrs").expect("retained").size,
7348 7,
7349 "{name}"
7350 );
7351 if name == crate::control::CONTROL_FILE_NAME {
7352 assert!(
7353 index
7354 .read_with(|index| {
7355 index.controls().is_ok_and(|controls| {
7356 controls.source_is(Path::new(name), b"changed")
7357 })
7358 })
7359 .expect("controls"),
7360 "the refreshed rules are retained"
7361 );
7362 }
7363 opened.close().expect("close");
7364 }
7365 }
7366
7367 #[cfg(feature = "watch")]
7368 #[test]
7369 fn budget_stop_wins_a_race_with_the_transition_to_watching() {
7370 let root = tempfile::tempdir().expect("temp root");
7371 let scripts = tempfile::tempdir().expect("script root");
7372 std::fs::write(root.path().join("baseline.txt"), b"baseline").expect("fixture");
7373 let script = scripts.path().join("events.script");
7374 std::fs::write(&script, b"").expect("script");
7375 let controls = Arc::new(TestControls::default());
7376 controls.gate(TestPoint::BeforeObservationWatching).arm();
7377 let mut options = scripted_options(&script);
7378 options.budget.max_files = Some(1);
7379 let opened =
7380 OpenedIndex::open_for_test(root.path(), options, Arc::clone(&controls)).expect("open");
7381 controls.gate(TestPoint::BeforeObservationWatching).wait_reached();
7382
7383 std::fs::write(root.path().join("over-budget.txt"), b"refused")
7384 .expect("over-budget mutation");
7385 let refreshed = opened
7386 .refresh(&[PathBuf::from("over-budget.txt")])
7387 .expect("resource refusal is a typed result");
7388 assert_eq!(refreshed.work.resource_refused, 1);
7389 assert_eq!(refreshed.state.phase, crate::LifecyclePhase::Stopped);
7390 controls.gate(TestPoint::BeforeObservationWatching).release();
7391
7392 let state = wait_until_phase(&opened, crate::LifecyclePhase::Stopped);
7393 assert_eq!(state.coverage, crate::Coverage::Partial(crate::CoverageReason::Budget));
7394 let since = opened.state.index.since(crate::Clock::ZERO).expect("journal");
7395 assert!(!since.commits.iter().any(|commit| commit.state.iter().any(|transition| {
7396 matches!(
7397 transition,
7398 crate::StateTransition::IndexState { current, .. }
7399 if current.phase == crate::LifecyclePhase::Watching
7400 )
7401 })));
7402 opened.close().expect("close");
7403 }
7404}