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