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