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