Skip to main content

fdu_core/
opened.rs

1//! Ownership and joined shutdown for one long-lived opened root.
2//!
3//! [`OpenedIndex`] is the public behavior surface. Its private shared state contains
4//! data and synchronization only; it is deliberately not a second API-shaped service.
5
6use 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
28/// First ordinal reserved for a minted session; zero never identifies a live owner.
29const FIRST_SESSION_ORDINAL: u64 = 1;
30/// FNV-1a offset used to mix the process and open instance into an opaque identity.
31const FNV_OFFSET_BASIS: u64 = 0xcbf2_9ce4_8422_2325;
32/// FNV-1a prime used for the opened-root identity mix.
33const FNV_PRIME: u64 = 0x0000_0100_0000_01b3;
34/// Nonzero fallback for the reserved zero identity.
35const FIRST_SESSION_ID: u64 = 1;
36/// Maximum paths accepted by one best-effort priority request.
37pub const MAX_PRIORITY_PATHS: usize = 64;
38/// Maximum paths accepted by one refresh operation.
39pub const MAX_REFRESH_PATHS: usize = 1_024;
40#[cfg(test)]
41/// Deadline for a missing deterministic test barrier to fail instead of hanging.
42const TEST_GATE_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(5);
43
44/// Identity of one live opened-root lifetime.
45///
46/// The value is process-local, opaque, and never persisted. It prevents a future cursor
47/// or continuation from being accepted by another open whose sequence also began at
48/// zero; it is not a credential.
49impl 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/// Resource bounds applied to progressive discovery.
77///
78/// The first version deliberately has one measured resource: retained regular files.
79/// The value is execution policy, not semantic scan scope, and therefore is absent from
80/// [`crate::ScanScope`] and snapshot identity.
81#[derive(Clone, Copy, PartialEq, Eq, Debug, Default)]
82pub struct DiscoveryBudget {
83    /// Maximum regular files retained, or no limit when absent.
84    pub max_files: Option<u64>,
85}
86
87/// Configuration for a long-lived [`OpenedIndex`].
88///
89/// Scope and execution settings are flat here because each is one independent decision;
90/// display depth is absent by design and belongs to a read request. The options bind
91/// progressive discovery, exact-history bounds, and optional live observation without
92/// changing the one-shot API.
93#[derive(Clone, Debug)]
94pub struct OpenOptions {
95    /// Ops per committed discovery batch.
96    pub batch_size: usize,
97    /// Follow directory symlinks when the configured platform semantics support it.
98    pub follow_symlinks: bool,
99    /// Stay on the filesystem containing the opened root.
100    pub one_filesystem: bool,
101    /// Hidden-component admission, or `None` to retain every component.
102    pub hidden: Option<Arc<crate::HiddenPolicy>>,
103    /// Exclude objects other than files, directories, and symlinks.
104    pub exclude_special: bool,
105    /// File-type rules, or `None` for the rules compiled into fdu.
106    pub types: Option<Arc<crate::classify::TypeRegistry>>,
107    /// Resource policy for the cold progressive walk.
108    pub budget: DiscoveryBudget,
109    /// Optional filesystem observation captured before cold discovery begins.
110    #[cfg(feature = "watch")]
111    pub observation: Option<crate::watch::WatchConfig>,
112    /// Test-only replacement for the native event source.
113    #[cfg(all(feature = "watch", test))]
114    #[doc(hidden)]
115    pub observation_script: Option<PathBuf>,
116    /// Approximate bytes the exact commit journal may retain, as
117    /// [`crate::Commit::retained_cost`] estimates them; see
118    /// [`crate::DEFAULT_JOURNAL_CAPACITY_BYTES`] for the default and why there is no
119    /// unbounded setting. [`OpenedIndex::open`] refuses a budget below
120    /// [`crate::MIN_JOURNAL_CAPACITY_BYTES`] with [`Error::JournalCapacityTooSmall`].
121    pub journal_capacity_bytes: usize,
122    /// The budget and the line limit `.gitignore` files are applied under. See
123    /// [`ScanConfig::control_limits`]: a refused file ends nothing, and
124    /// [`crate::ReadDiagnostics::controls`] names it and the limit that fired.
125    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    /// Construct a validated cold-progressive plan from these options and a delivery.
151    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            // The first opened-root scheduler is intentionally one parent-first
175            // producer. Parallel I/O is an internal optimization, not a public
176            // semantic or tuning promise for this new API. `crate::plan` refuses a
177            // `Route::Opened` delivery that asks for another worker count or order, so
178            // nothing a caller asked for is dropped here.
179            threads: Some(1),
180            order: crate::ScanOrder::BreadthFirst,
181            types: self.types,
182            // Never optional here: the basis an opened root holds says control state is
183            // always observed, and this scan is what makes that true.
184            read_controls: OpenedIndex::basis().scope.read_controls,
185            population: crate::query::IgnoredEntries::Include,
186            control_limits: self.control_limits,
187            // An opened root reports through its own `DiscoveryProgress`, which is
188            // state the lifecycle publishes rather than a count of work done.
189            progress: None,
190        };
191        (scan, self.budget, self.journal_capacity_bytes)
192    }
193}
194
195/// A long-lived, synchronously controlled filesystem index.
196///
197/// Clones are cheap references to one authority. Calling [`Self::close`] through any
198/// clone cancels and joins every worker owned by that authority, and every concurrent
199/// caller receives the same stored terminal outcome.
200#[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    /// The basis every opened root holds, stated once for every route into one.
225    ///
226    /// No analyzers, because an opened root runs none, and control state always observed,
227    /// because its ignored and unignored partitions are part of what it serves --
228    /// [`OpenOptions`] has no switch for either, and reads this rather than restating it.
229    /// The root is empty because no rule in the request model compares roots and an opened
230    /// read names none: it asks the root it already holds.
231    ///
232    /// Stated rather than read back from the retained index, so a read is refused before
233    /// any stored state is touched, which is where every other surface refuses one.
234    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    /// Open one live root without changing the existing blocking [`crate::open`] API.
243    ///
244    /// This constructor validates and binds the root and semantic configuration, starts
245    /// cold progressive discovery, and captures optional observation before that
246    /// baseline begins. No cache image or second mutable index is created here.
247    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    /// Cancel and join all work owned by this opened root.
304    ///
305    /// The first caller performs shutdown. Concurrent and repeated callers wait for or
306    /// replay its stored terminal outcome; success is never reported while a worker is
307    /// still live.
308    pub fn close(&self) -> Result<()> {
309        self.state.shutdown()
310    }
311
312    /// Reorder pending discovery toward the supplied relative paths.
313    ///
314    /// This is a bounded best-effort scheduling hint. It does not change scope, facts,
315    /// lifecycle state, or the index clock, and paths that are already complete simply
316    /// have no effect.
317    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    /// Return requested projections from one committed version and state boundary.
339    ///
340    /// # Errors
341    ///
342    /// The whole read fails only when no projection in it can be trusted: the request's
343    /// shape is invalid (a bound out of range, too many projections, a path that escapes the
344    /// root, a continuation another root issued or this one never did), the root is closed,
345    /// or [`crate::ReadRequest::expected`] or a continuation pins a version the index no
346    /// longer holds. What one projection finds at the pinned version -- including a
347    /// continuation this root issued and has since consumed or evicted -- is that
348    /// projection's [`crate::ProjectionRefusal`] instead, returned in its position while
349    /// every other projection answers.
350    ///
351    /// The lifecycle lock guards only the phase check. The projection's coherence comes
352    /// from the index read boundary, which any number of readers share, so holding the
353    /// lifecycle lock across it would serialize every read with every other read and with
354    /// refresh, worker registration, and the start of close. A read that races close cannot
355    /// leave state behind: the continuation table refuses records once shutdown clears it.
356    pub fn read(&self, request: crate::ReadRequest) -> Result<crate::ReadResponse> {
357        self.ensure_open()?;
358        read::read(self, request)
359    }
360
361    /// Return exact commits after one version, waiting up to the supplied timeout.
362    ///
363    /// A poll waiting on an idle journal returns as soon as an owned worker panics, with
364    /// [`Error::OpenedWorkerPanicked`] naming the worker -- the cause [`Self::close`] will
365    /// report -- rather than at its timeout. Commits retained before the panic are still
366    /// returned first, unless the panic struck inside a commit: that poisons the index,
367    /// which leaves nothing to read. A poll parked on the journal then reports the panic.
368    /// A poll already reading when the commit panics can observe the poisoned index before
369    /// the worker has recorded its panic, and returns [`Error::IndexLockPoisoned`]; the
370    /// panic is recorded moments later, and [`Self::close`] names it either way.
371    pub fn changes(&self, request: crate::ChangeRequest) -> Result<crate::ChangePoll> {
372        self.ensure_open()?;
373        journal::poll(self, request)
374    }
375
376    /// Verify a bounded set of relative paths and conditionally commit exact changes.
377    ///
378    /// Inputs are classified as one set: duplicates are removed, descendants covered
379    /// by an accepted ancestor cost no second walk, and an empty set is a no-op. The
380    /// returned `(after, version]` interval is safe as the next journal boundary even
381    /// when another producer committed concurrently.
382    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        // A stopped partial root remains readable and may verify work that proves it
393        // cannot expand retained truth. Before that terminal state, the exact commit
394        // boundary arbitrates every file upsert against the shared budget.
395        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        // A pass retired by newer verification closed its scope partial without an
442        // error of its own: the receipt has to say so, or a caller sees a partial
443        // state with nothing to retry.
444        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    /// Register one worker with the shared owner before it can race with close.
598    ///
599    /// The worker receives only cancellation, not a strong owner reference. Discovery
600    /// and observation both use this boundary, which deterministic lifecycle tests
601    /// exercise directly.
602    #[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                // The worker records its own failure as it leaves, and a panic is caught
625                // for exactly that long before it resumes. Left to the join, a panic would
626                // be learned of only at close -- a change poll blocked on the journal would
627                // sleep to its timeout -- and close could report failures only in spawn
628                // order, letting discovery's poisoned-lock error stand in for the panic
629                // that poisoned it.
630                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    /// Retained for discovery and verified producers; semantic identity is also
674    /// fixed in `index`, so this cannot reinterpret already-retained facts.
675    #[allow(dead_code)]
676    scan: ScanConfig,
677    budget: DiscoveryBudget,
678    frontier: Arc<DiscoveryFrontier>,
679    continuations: Mutex<continuation::ContinuationTable>,
680    journal: Arc<journal::JournalWait>,
681    /// Worker failures in the order they happened; shared with the workers that record them.
682    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/// What one opened root still retains, read from its own state.
903///
904/// The session goldens' `final` record is derived from this rather than written as a
905/// literal: a regression that left a worker, a blocked poll, or a page record behind a
906/// closed root would otherwise print the same text as a clean shutdown.
907#[cfg(all(test, feature = "watch"))]
908pub(super) struct RetainedOwnership {
909    session: SessionId,
910    /// Shutdown finished: every worker was joined before its outcome was stored.
911    joined: bool,
912    workers: usize,
913    waiters: usize,
914    continuations: usize,
915    /// The outcome every later `close()` replays, once shutdown has stored one.
916    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
1017/// Pending directories, grouped by the first priority each one serves.
1018///
1019/// The next directory is the earliest queued one serving the earliest priority, and
1020/// otherwise the earliest queued one. Which priority a directory serves depends only on
1021/// its path and the priority list, so it is decided once -- when the directory is queued,
1022/// or when the priorities change -- rather than on every pop. Rescanning the whole queue
1023/// against every priority per pop made one breadth-first level of width `w` cost
1024/// `O(w² · priorities)`: 4,000 pops under 64 priorities took 40 seconds.
1025struct FrontierState {
1026    /// Directories serving no priority, in queue order.
1027    unprioritized: VecDeque<QueuedDirectory>,
1028    /// `serving[k]` holds the directories whose first matching priority is
1029    /// `priorities[k]`, in queue order.
1030    serving: Vec<VecDeque<QueuedDirectory>>,
1031    priorities: Vec<PathBuf>,
1032    /// Queue order, kept so that changing priorities regroups without reordering.
1033    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/// How the index answered one discovery commit.
1175///
1176/// Discovery walks a tree other producers are free to change underneath it, so a refused
1177/// commit is usually news about the world rather than a failure of the walk. Only an
1178/// engine failure -- a poisoned lock, an exhausted clock, a malformed observation -- is
1179/// returned as an error and ends discovery.
1180#[derive(Clone, Copy, Debug)]
1181enum DiscoveryAnswer {
1182    /// Committed, or verified to change nothing; the listing continues.
1183    Accepted,
1184    /// The root is terminal: this commit's upsert crossed the shared file budget, or an
1185    /// earlier commit had already stopped or failed the root.
1186    Stopped,
1187    /// The commit named a directory the index no longer holds, or a child whose ancestry
1188    /// it no longer holds. A refresh or the observer verified that part of the tree after
1189    /// this directory was queued, so the frontier entry is stale and the newer commit
1190    /// stands; whichever producer next verifies the path records what is there now.
1191    Stale,
1192}
1193
1194/// Classify a refused discovery commit, keeping engine failures fatal.
1195///
1196/// A control file the control budget or line limit cannot admit is not a refusal of the
1197/// commit: the index refuses that one source and commits the listing (fdu-1onj).
1198fn 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
1208/// Stop listing one directory whose commit the index did not accept.
1209///
1210/// A directory whose listing was cut short is never marked complete, so absence below it
1211/// stays unknown. The answer decides what happens next:
1212///
1213/// - `Stale`: the index no longer holds the directory or an ancestor, so the
1214///   subdirectories earlier batches committed left with it, and the producer that replaced
1215///   it records what is there.
1216/// - `Stopped`: the root is terminal and nothing more is queued.
1217fn 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        // Removed or replaced after its parent was listed -- a build cache, an editor's
1245        // temporary directory. That is stale frontier work, not an inaccessible boundary:
1246        // marking the root inaccessible for it would outlive every later verification
1247        // that the path is simply gone. The root itself is never stale work; losing it is
1248        // a failure of the whole walk.
1249        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    // Subdirectories to queue once the listing commits.
1275    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                // The stale batch carried the transition with it.
1291                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        // Queued before its upsert is pushed; the whole list is queued once the listing
1368        // commits.
1369        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    // Schedule real filesystem inputs before they enter production admission; do not
1426    // normalize or reorder the commits whose exact sequence the golden records.
1427    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
1460/// Commit one batch of `directory`'s listing.
1461fn 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
1488/// Publish a state-only discovery transition.
1489///
1490/// A root that has already stopped or failed refuses every discovery commit, and a
1491/// transition it refuses has nothing left to say: the terminal state it would have
1492/// replaced is the answer.
1493fn 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")]
1510/// Maximum time the owner waits for an idle observation before checking cancellation.
1511const OBSERVATION_POLL_INTERVAL: std::time::Duration = std::time::Duration::from_millis(100);
1512#[cfg(feature = "watch")]
1513/// Full-root handoff retries allowed after benign conditional conflicts with a producer.
1514const MAX_HANDOFF_RECONCILIATION_ATTEMPTS: usize = 3;
1515
1516#[cfg(feature = "watch")]
1517#[allow(clippy::too_many_arguments)]
1518#[allow(clippy::needless_pass_by_value)] // Ownership keeps the backend alive for this worker.
1519fn 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                // The reconciliation already left what it could not read partial, and that
1655                // boundary is settled rather than retried on every later event. Retaining
1656                // the causes is what lets a consumer see why.
1657                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/// Worker exits that ended in failure, in the order they happened.
1849///
1850/// Each worker records its own exit before its thread ends, so close can report the
1851/// failure that came first rather than the first in spawn order, and a change poll
1852/// blocked on the journal can be answered with a typed cause instead of its timeout.
1853#[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    /// The worker whose panic ended this root, if one did.
1864    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    /// The failure close reports.
1874    ///
1875    /// The earliest, unless the earliest is only the trace a recorded worker panic left. A
1876    /// poisoned lock is then a consequence: the guard is poisoned while its thread is still
1877    /// unwinding, before that thread can record its panic, so the worker that trips over
1878    /// the poison can record first, and discovery's `IndexLockPoisoned` would stand in for
1879    /// the observation worker's panic. With no panic recorded, nothing here explains the
1880    /// poisoning -- the caller's own thread may have panicked inside a commit -- so it is
1881    /// a failure like any other, and a later unrelated error does not displace it.
1882    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    /// Whether this failure only reports a lock some panic poisoned, rather than a cause.
1901    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
1927/// Join every worker, then report the failure the workers themselves recorded first.
1928///
1929/// The join only waits: each worker records its exit before its thread ends. The join's
1930/// own view is kept as the answer of last resort for a panic that escaped recording.
1931fn 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    /// Block until every worker registered under `name` has returned, without closing.
2171    ///
2172    /// A state that should stay put cannot be awaited by polling for a change, so this
2173    /// waits for the worker that might change it to finish instead.
2174    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    /// An opened root takes a control file spelled `.GITIGNORE` as a detached scan does:
2257    /// in discovery, and in a refresh of the variant's own path once it is gone, which no
2258    /// stat can verify and the directory's control lookup answers (fdu-0w1b).
2259    #[test]
2260    fn an_opened_root_takes_a_case_variant_control_as_a_detached_scan_does() {
2261        use crate::test_support::CaseLookups;
2262
2263        let probe = tempfile::tempdir().expect("temp root");
2264        for (lookups, governs) in CaseLookups::on_this_host(probe.path()) {
2265            let root = tempfile::tempdir().expect("temp root");
2266            let _lookups = lookups.install(root.path());
2267            std::fs::create_dir(root.path().join("up")).expect("directory");
2268            std::fs::write(root.path().join("up/.GITIGNORE"), b"*.tmp\n").expect("variant");
2269            std::fs::write(root.path().join("up/x.tmp"), b"governed").expect("fixture");
2270            std::fs::write(root.path().join("up/notes.txt"), b"never").expect("fixture");
2271            let facts = |index: &Index| {
2272                let sources: Vec<(PathBuf, Vec<u8>)> = index
2273                    .controls()
2274                    .expect("observed")
2275                    .sources()
2276                    .map(|(path, source)| (path, source.to_vec()))
2277                    .collect();
2278                (index.is_ignored(Path::new("up/x.tmp")).expect("observed"), sources)
2279            };
2280            let cold = || {
2281                let (cold, _) =
2282                    crate::scan::scan_into_index(root.path(), &crate::ScanConfig::default())
2283                        .expect("cold scan");
2284                facts(&cold)
2285            };
2286
2287            let options = OpenOptions { batch_size: 1, ..OpenOptions::default() };
2288            let opened = open_fixture(root.path(), options).expect("opened root");
2289            wait_until_settled(&opened);
2290            let discovered = opened.state.index.read_with(facts).expect("read");
2291            assert_eq!(discovered, cold(), "{lookups:?}: discovery");
2292            assert_eq!(discovered.0, Some(governs), "{lookups:?}");
2293
2294            std::fs::remove_file(root.path().join("up/.GITIGNORE")).expect("remove the variant");
2295            opened.refresh(&[PathBuf::from("up/.GITIGNORE")]).expect("refresh");
2296            let refreshed = opened.state.index.read_with(facts).expect("read");
2297            assert_eq!(refreshed, cold(), "{lookups:?}: refresh of the removed variant");
2298            assert_eq!(refreshed.0, Some(false), "{lookups:?}");
2299            opened.close().expect("close");
2300        }
2301    }
2302
2303    #[test]
2304    fn associated_and_free_open_contracts_coexist() {
2305        let root = tempfile::tempdir().expect("temp root");
2306        std::fs::write(root.path().join("file.txt"), b"one").expect("fixture");
2307
2308        let opened = open_fixture(root.path(), OpenOptions::default()).expect("opened root");
2309        let (detached, _) = crate::open_fixture(
2310            root.path(),
2311            &crate::OpenFixture {
2312                cache_path: None,
2313                policy: crate::CachePolicy::Off,
2314                ..crate::OpenFixture::default()
2315            },
2316        )
2317        .expect("blocking open");
2318
2319        assert_eq!(detached.total().files, 1);
2320        opened.close().expect("close");
2321    }
2322
2323    #[test]
2324    fn one_clone_closes_the_shared_authority_and_close_is_idempotent() {
2325        let (_root, opened) = opened(Arc::default());
2326        let clone = opened.clone();
2327        assert_eq!(opened.state.session, clone.state.session);
2328        assert!(Arc::ptr_eq(&opened.state, &clone.state));
2329
2330        clone.close().expect("first close");
2331        assert!(matches!(opened.ensure_open(), Err(Error::OpenedIndexClosed)));
2332        opened.close().expect("repeated close");
2333    }
2334
2335    #[test]
2336    fn shutdown_refuses_new_owned_work() {
2337        let (_root, opened) = opened(Arc::default());
2338        opened.close().expect("close");
2339        let ran = Arc::new(AtomicBool::new(false));
2340        let worker_ran = Arc::clone(&ran);
2341
2342        assert!(matches!(
2343            opened.spawn_worker("late", move |_cancellation| {
2344                worker_ran.store(true, Ordering::SeqCst);
2345                Ok(())
2346            }),
2347            Err(Error::OpenedIndexClosed)
2348        ));
2349        assert!(!ran.load(Ordering::SeqCst));
2350    }
2351
2352    #[test]
2353    fn concurrent_close_waits_for_one_stored_worker_failure() {
2354        let controls = Arc::new(TestControls::default());
2355        controls.gate(TestPoint::BeforeWorkerExit).arm();
2356        controls.gate(TestPoint::BeforeCloseWait).arm();
2357        let (_root, opened) = opened(Arc::clone(&controls));
2358        opened
2359            .spawn_worker("failure", |cancellation| {
2360                cancellation.wait_cancelled();
2361                Err(Error::CommitRejected("injected opened worker failure"))
2362            })
2363            .expect("spawn worker");
2364
2365        let first = opened.clone();
2366        let first_close = thread::spawn(move || first.close());
2367        controls.gate(TestPoint::BeforeWorkerExit).wait_reached();
2368        let second = opened.clone();
2369        let second_close = thread::spawn(move || second.close());
2370        controls.gate(TestPoint::BeforeCloseWait).wait_reached();
2371
2372        controls.gate(TestPoint::BeforeCloseWait).release();
2373        controls.gate(TestPoint::BeforeWorkerExit).release();
2374        let first_error = first_close.join().expect("first close thread").expect_err("failure");
2375        let second_error = second_close.join().expect("second close thread").expect_err("failure");
2376        assert_eq!(first_error.to_string(), second_error.to_string());
2377        assert!(matches!(first_error, Error::OpenedWorkerFailed { worker: "failure", .. }));
2378        assert_eq!(
2379            opened.close().expect_err("stored failure").to_string(),
2380            second_error.to_string()
2381        );
2382    }
2383
2384    #[test]
2385    fn a_panicking_worker_is_joined_and_reported() {
2386        let (_root, opened) = opened(Arc::default());
2387        opened
2388            .spawn_worker("panic", |_cancellation| panic!("injected worker panic"))
2389            .expect("spawn worker");
2390
2391        let error = opened.close().expect_err("panic is terminal");
2392        assert!(matches!(error, Error::OpenedWorkerPanicked { worker: "panic" }));
2393        assert!(matches!(opened.close(), Err(Error::OpenedWorkerPanicked { worker: "panic" })));
2394    }
2395
2396    /// A poll blocked on the journal learns of a worker panic when it happens, with the
2397    /// cause close will report, instead of sleeping to its timeout. That holds for a panic
2398    /// inside a commit too, the likeliest place for one: the unwinding guard poisons the
2399    /// index lock, and the poll reports the panic rather than the poisoning it left.
2400    #[test]
2401    fn a_worker_panic_wakes_a_blocked_change_poll_with_its_typed_failure() {
2402        for holds_index_lock in [false, true] {
2403            let controls = Arc::new(TestControls::default());
2404            controls.gate(TestPoint::BeforeJournalWait).arm();
2405            let (_root, opened) = opened(Arc::clone(&controls));
2406            let cursor = current_version(&opened);
2407            let poller = opened.clone();
2408            let poll = thread::spawn(move || {
2409                poller.changes(crate::ChangeRequest {
2410                    after: cursor,
2411                    timeout: std::time::Duration::from_secs(60),
2412                })
2413            });
2414            controls.gate(TestPoint::BeforeJournalWait).wait_reached();
2415            let index = opened.state.index.clone();
2416            opened
2417                .spawn_worker("panic", move |_cancellation| {
2418                    if holds_index_lock {
2419                        index.panic_holding_the_write_lock_for_test();
2420                    }
2421                    panic!("injected worker panic")
2422                })
2423                .expect("spawn worker");
2424            controls.gate(TestPoint::BeforeJournalWait).release();
2425
2426            let outcome = poll.join().expect("poll thread");
2427            assert!(
2428                matches!(outcome, Err(Error::OpenedWorkerPanicked { worker: "panic" })),
2429                "holds index lock {holds_index_lock}: {outcome:?}"
2430            );
2431            assert!(matches!(opened.close(), Err(Error::OpenedWorkerPanicked { worker: "panic" })));
2432        }
2433    }
2434
2435    /// Close reports the failure that happened first, not the worker that was spawned first.
2436    #[test]
2437    fn close_reports_the_failure_that_happened_first_not_the_worker_spawned_first() {
2438        let (_root, opened) = opened(Arc::default());
2439        opened
2440            .spawn_worker("slow", |cancellation| {
2441                cancellation.wait_cancelled();
2442                Err(Error::CommitRejected("slow worker failed at close"))
2443            })
2444            .expect("spawn slow worker");
2445        opened
2446            .spawn_worker("fast", |_cancellation| {
2447                Err(Error::CommitRejected("fast worker failed first"))
2448            })
2449            .expect("spawn fast worker");
2450        wait_for_worker_exit(&opened, "fast");
2451
2452        assert!(matches!(opened.close(), Err(Error::OpenedWorkerFailed { worker: "fast", .. })));
2453    }
2454
2455    /// A poisoned lock is what a panic leaves behind, and the worker that trips over it can
2456    /// record its error before the unwinding thread records the panic. The panic is the
2457    /// cause, so it is the failure close reports.
2458    #[test]
2459    fn close_reports_a_panic_before_the_poisoning_it_left_behind() {
2460        let (_root, opened) = opened(Arc::default());
2461        opened
2462            .spawn_worker("tripped", |_cancellation| Err(Error::IndexLockPoisoned))
2463            .expect("spawn tripped worker");
2464        wait_for_worker_exit(&opened, "tripped");
2465        opened
2466            .spawn_worker("panicked", |_cancellation| panic!("injected worker panic"))
2467            .expect("spawn panicking worker");
2468        wait_for_worker_exit(&opened, "panicked");
2469
2470        assert!(matches!(opened.close(), Err(Error::OpenedWorkerPanicked { worker: "panicked" })));
2471    }
2472
2473    /// Without a recorded panic, nothing explains a poisoning: the caller's own thread can
2474    /// poison the index by panicking inside a commit. A worker that trips over it first is
2475    /// then the earliest failure, and close reports it rather than a later, unrelated one.
2476    #[test]
2477    fn close_reports_a_poisoning_no_worker_panic_explains_when_it_came_first() {
2478        let (_root, opened) = opened(Arc::default());
2479        opened.state.index.poison_for_test();
2480        let index = opened.state.index.clone();
2481        opened
2482            .spawn_worker("tripped", move |_cancellation| index.clock().map(|_| ()))
2483            .expect("spawn tripped worker");
2484        wait_for_worker_exit(&opened, "tripped");
2485        opened
2486            .spawn_worker("later", |_cancellation| {
2487                Err(Error::CommitRejected("unrelated later failure"))
2488            })
2489            .expect("spawn later worker");
2490        wait_for_worker_exit(&opened, "later");
2491
2492        let closed = opened.close();
2493        assert!(
2494            matches!(
2495                &closed,
2496                Err(Error::OpenedWorkerFailed { worker: "tripped", source })
2497                    if matches!(**source, Error::IndexLockPoisoned)
2498            ),
2499            "{closed:?}"
2500        );
2501    }
2502
2503    #[test]
2504    fn dropping_the_last_reference_cancels_and_joins() {
2505        let active = Arc::new(AtomicUsize::new(0));
2506        let (started_sender, started_receiver) = std::sync::mpsc::sync_channel(0);
2507        let (_root, opened) = opened(Arc::default());
2508        let worker_active = Arc::clone(&active);
2509        opened
2510            .spawn_worker("drop", move |cancellation| {
2511                worker_active.fetch_add(1, Ordering::SeqCst);
2512                started_sender.send(()).expect("report worker start");
2513                cancellation.wait_cancelled();
2514                worker_active.fetch_sub(1, Ordering::SeqCst);
2515                Ok(())
2516            })
2517            .expect("spawn worker");
2518
2519        started_receiver.recv_timeout(std::time::Duration::from_secs(5)).expect("worker started");
2520        drop(opened);
2521        assert_eq!(active.load(Ordering::SeqCst), 0, "drop returned before worker join");
2522    }
2523
2524    #[test]
2525    fn poisoned_lifecycle_still_joins_before_returning_its_typed_failure() {
2526        let active = Arc::new(AtomicUsize::new(0));
2527        let (started_sender, started_receiver) = std::sync::mpsc::sync_channel(0);
2528        let (_root, opened) = opened(Arc::default());
2529        let worker_active = Arc::clone(&active);
2530        opened
2531            .spawn_worker("poison", move |cancellation| {
2532                worker_active.fetch_add(1, Ordering::SeqCst);
2533                started_sender.send(()).expect("report worker start");
2534                cancellation.wait_cancelled();
2535                worker_active.fetch_sub(1, Ordering::SeqCst);
2536                Ok(())
2537            })
2538            .expect("spawn worker");
2539        started_receiver.recv_timeout(std::time::Duration::from_secs(5)).expect("worker started");
2540
2541        let state = Arc::clone(&opened.state);
2542        thread::spawn(move || {
2543            let _guard = state.lifecycle.lock().expect("lifecycle lock");
2544            panic!("inject lifecycle poison");
2545        })
2546        .join()
2547        .expect_err("injected panic");
2548
2549        assert!(matches!(opened.close(), Err(Error::OpenedLifecyclePoisoned)));
2550        assert_eq!(active.load(Ordering::SeqCst), 0, "poison bypassed worker join");
2551        assert!(matches!(opened.close(), Err(Error::OpenedLifecyclePoisoned)));
2552    }
2553
2554    #[test]
2555    fn poisoned_index_still_joins_and_replays_the_typed_failure() {
2556        let (_root, opened) = opened(Arc::default());
2557        opened.state.index.poison_for_test();
2558
2559        assert!(matches!(opened.close(), Err(Error::IndexLockPoisoned)));
2560        assert!(matches!(opened.close(), Err(Error::IndexLockPoisoned)));
2561    }
2562
2563    #[test]
2564    fn opened_root_rejects_invalid_scan_policy_and_nondirectories() {
2565        let root = tempfile::tempdir().expect("temp root");
2566        let options = OpenOptions { batch_size: 0, ..OpenOptions::default() };
2567        assert!(matches!(open_fixture(root.path(), options), Err(Error::UnsupportedScanConfig(_))));
2568
2569        let file = root.path().join("file");
2570        std::fs::write(&file, b"x").expect("fixture");
2571        assert!(matches!(open_fixture(&file, OpenOptions::default()), Err(Error::Io { .. })));
2572
2573        let zero_budget = OpenOptions {
2574            budget: DiscoveryBudget { max_files: Some(0) },
2575            ..OpenOptions::default()
2576        };
2577        assert!(matches!(
2578            open_fixture(root.path(), zero_budget),
2579            Err(Error::UnsupportedScanConfig(_))
2580        ));
2581
2582        let minimum = crate::MIN_JOURNAL_CAPACITY_BYTES;
2583        let below_minimum =
2584            OpenOptions { journal_capacity_bytes: minimum - 1, ..OpenOptions::default() };
2585        let error = open_fixture(root.path(), below_minimum).expect_err("refused");
2586        assert!(
2587            matches!(error, Error::JournalCapacityTooSmall { requested, minimum: stated }
2588                if requested == minimum - 1 && stated == minimum),
2589            "{error:?}"
2590        );
2591        let message = error.to_string();
2592        assert!(message.contains(&format!("at least {minimum} bytes")), "{message}");
2593
2594        let at_minimum = OpenOptions { journal_capacity_bytes: minimum, ..OpenOptions::default() };
2595        open_fixture(root.path(), at_minimum).expect("accepted").close().expect("close");
2596    }
2597
2598    #[test]
2599    fn distinct_opens_have_distinct_live_identity() {
2600        let root = tempfile::tempdir().expect("temp root");
2601        let first = open_fixture(root.path(), OpenOptions::default()).expect("first");
2602        let second = open_fixture(root.path(), OpenOptions::default()).expect("second");
2603        assert_ne!(first.state.session, second.state.session);
2604        first.close().expect("first close");
2605        second.close().expect("second close");
2606    }
2607
2608    #[test]
2609    fn terminal_discovery_failure_retains_bounded_typed_evidence() {
2610        let root = tempfile::tempdir().expect("temp root");
2611        let controls = Arc::new(TestControls::default());
2612        controls.gate(TestPoint::BeforeDiscovery).arm();
2613        let opened =
2614            OpenedIndex::open_for_test(root.path(), OpenOptions::default(), Arc::clone(&controls))
2615                .expect("opened root");
2616        std::fs::remove_dir(root.path()).expect("remove empty fixture root");
2617        controls.gate(TestPoint::BeforeDiscovery).release();
2618
2619        let state = wait_until_settled(&opened);
2620        assert_eq!(state.phase, crate::LifecyclePhase::Failed);
2621        assert_eq!(state.coverage, crate::Coverage::Partial(crate::CoverageReason::Failed));
2622        assert_eq!(state.issues.retained, 1);
2623        let issues = opened.state.index.issues().expect("issues");
2624        assert_eq!(issues.len(), 1);
2625        assert_eq!(issues[0].kind, crate::IssueKind::Disappeared);
2626        assert!(matches!(opened.close(), Err(Error::OpenedWorkerFailed { .. })));
2627    }
2628
2629    #[test]
2630    fn progressive_discovery_settles_to_the_one_shot_tree() {
2631        let root = tempfile::tempdir().expect("temp root");
2632        std::fs::create_dir_all(root.path().join("alpha/deep")).expect("fixture directories");
2633        std::fs::write(root.path().join("root.txt"), b"root").expect("root fixture");
2634        std::fs::write(root.path().join("alpha/child.bin"), b"child").expect("child fixture");
2635        std::fs::write(root.path().join("alpha/deep/leaf.rs"), b"leaf").expect("leaf fixture");
2636        let options = OpenOptions { batch_size: 2, ..OpenOptions::default() };
2637
2638        let opened = open_fixture(root.path(), options.clone()).expect("opened root");
2639        let state = wait_until_settled(&opened);
2640        assert_eq!(state.phase, crate::LifecyclePhase::Ready);
2641        assert_eq!(state.coverage, crate::Coverage::Complete);
2642        assert_eq!(state.progress.files_retained, 3);
2643
2644        let live = opened.state.index.snapshot().expect("live snapshot");
2645        let (one_shot, report) =
2646            crate::scan::scan_into_index(root.path(), &options.clone().into_parts().0)
2647                .expect("one-shot scan");
2648        assert!(report.is_complete());
2649        assert_eq!(live.total(), one_shot.total());
2650        assert_eq!(live.len(), one_shot.len());
2651        for path in [Path::new(""), Path::new("alpha"), Path::new("alpha/deep")] {
2652            assert_eq!(child_facts(&live, path), child_facts(&one_shot, path));
2653            assert_eq!(live.directory_complete(path), Some(true));
2654        }
2655        opened.close().expect("close");
2656    }
2657
2658    #[test]
2659    fn parent_listing_commits_before_prioritized_child_work_without_a_clock_change() {
2660        let root = tempfile::tempdir().expect("temp root");
2661        for directory in ["alpha", "target"] {
2662            std::fs::create_dir(root.path().join(directory)).expect("fixture directory");
2663            std::fs::write(root.path().join(directory).join("leaf"), directory)
2664                .expect("fixture file");
2665        }
2666        let controls = Arc::new(TestControls::default());
2667        controls.gate(TestPoint::AfterRootDirectory).arm();
2668        let opened = OpenedIndex::open_for_test(
2669            root.path(),
2670            OpenOptions { batch_size: 64, ..OpenOptions::default() },
2671            Arc::clone(&controls),
2672        )
2673        .expect("opened root");
2674        controls.gate(TestPoint::AfterRootDirectory).wait_reached();
2675
2676        let before_priority = opened.state.index.clock().expect("clock");
2677        assert_eq!(
2678            opened.state.index.directory_complete(Path::new("")).expect("root completeness"),
2679            Some(true)
2680        );
2681        assert_eq!(
2682            opened
2683                .state
2684                .index
2685                .directory_complete(Path::new("target"))
2686                .expect("target completeness"),
2687            Some(false)
2688        );
2689        opened.prioritize(&[PathBuf::from("target")]).expect("prioritize");
2690        assert_eq!(opened.state.index.clock().expect("clock"), before_priority);
2691        controls.gate(TestPoint::AfterRootDirectory).release();
2692        let state = wait_until_settled(&opened);
2693        assert_eq!(state.coverage, crate::Coverage::Complete);
2694
2695        let since = opened.state.index.since(before_priority).expect("commits after root");
2696        let first_file = since
2697            .commits
2698            .iter()
2699            .flat_map(|commit| &commit.changes)
2700            .find_map(|change| match change {
2701                crate::EffectiveChange::Inserted { path, kind: EntryKind::File, .. } => {
2702                    Some(path.clone())
2703                }
2704                _ => None,
2705            })
2706            .expect("child file commit");
2707        assert_eq!(first_file, PathBuf::from("target/leaf"));
2708        opened.close().expect("close");
2709    }
2710
2711    /// Grouping by priority changes what a pop costs, never which directory it returns.
2712    ///
2713    /// The reference is the per-pop scan the frontier used to run: the pending directory
2714    /// whose first matching priority is earliest, ties broken by queue position, and
2715    /// otherwise the front of the queue. A deterministic mix of pops that queue children,
2716    /// arbitrary extends, and priority changes must pop the same sequence from both.
2717    #[test]
2718    fn the_grouped_frontier_pops_in_the_order_the_scan_chose() {
2719        fn reference_pop(
2720            pending: &mut VecDeque<PendingDirectory>,
2721            priorities: &[PathBuf],
2722        ) -> Option<PendingDirectory> {
2723            let selected = pending
2724                .iter()
2725                .enumerate()
2726                .filter_map(|(position, entry)| {
2727                    priorities
2728                        .iter()
2729                        .position(|priority| {
2730                            priority.starts_with(&entry.path) || entry.path.starts_with(priority)
2731                        })
2732                        .map(|priority| (priority, position))
2733                })
2734                .min()
2735                .map_or(0, |(_, position)| position);
2736            pending.remove(selected)
2737        }
2738
2739        const PATHS: [&str; 7] = ["a", "b", "a/b", "b/a", "a/b/c", "c", "c/a/b"];
2740        let mut seed = 0x2545_f491_4f6c_dd1d_u64;
2741        let mut next = move |bound: usize| {
2742            seed ^= seed << 13;
2743            seed ^= seed >> 7;
2744            seed ^= seed << 17;
2745            usize::try_from(seed % u64::try_from(bound).expect("small bound")).expect("fits")
2746        };
2747        let frontier = DiscoveryFrontier::new();
2748        let mut reference = VecDeque::from([PendingDirectory { path: PathBuf::new(), depth: 0 }]);
2749        let mut priorities = Vec::new();
2750        for step in 0..4_000 {
2751            match next(10) {
2752                0..=4 => {
2753                    let expected = reference_pop(&mut reference, &priorities);
2754                    let actual = frontier.pop();
2755                    assert_eq!(
2756                        actual.as_ref().map(|directory| &directory.path),
2757                        expected.as_ref().map(|directory| &directory.path),
2758                        "step {step}"
2759                    );
2760                    if let Some(parent) = expected {
2761                        let children: Vec<_> = (0..next(3))
2762                            .map(|child| PendingDirectory {
2763                                path: parent.path.join(["a", "b", "c"][child]),
2764                                depth: parent.depth + 1,
2765                            })
2766                            .collect();
2767                        reference.extend(children.iter().cloned());
2768                        frontier.extend(children);
2769                    }
2770                }
2771                5..=7 => {
2772                    let directory =
2773                        PendingDirectory { path: PathBuf::from(PATHS[next(7)]), depth: 1 };
2774                    reference.push_back(directory.clone());
2775                    frontier.extend([directory]);
2776                }
2777                _ => {
2778                    let mut chosen: Vec<_> =
2779                        (0..next(4)).map(|_| PathBuf::from(PATHS[next(7)])).collect();
2780                    chosen.sort();
2781                    chosen.dedup();
2782                    frontier.prioritize(chosen.clone());
2783                    priorities = chosen;
2784                }
2785            }
2786        }
2787        while let Some(expected) = reference_pop(&mut reference, &priorities) {
2788            assert_eq!(frontier.pop().map(|directory| directory.path), Some(expected.path));
2789        }
2790        assert!(frontier.pop().is_none());
2791    }
2792
2793    #[test]
2794    fn reaching_a_file_limit_without_refusal_remains_complete() {
2795        let root = tempfile::tempdir().expect("temp root");
2796        std::fs::write(root.path().join("one"), b"1").expect("fixture");
2797        std::fs::write(root.path().join("two"), b"2").expect("fixture");
2798        let opened = open_fixture(
2799            root.path(),
2800            OpenOptions {
2801                budget: DiscoveryBudget { max_files: Some(2) },
2802                ..OpenOptions::default()
2803            },
2804        )
2805        .expect("opened root");
2806
2807        let state = wait_until_settled(&opened);
2808        assert_eq!(state.coverage, crate::Coverage::Complete);
2809        assert_eq!(state.progress.files_retained, 2);
2810        assert_eq!(
2811            opened.state.index.directory_complete(Path::new("")).expect("root completeness"),
2812            Some(true)
2813        );
2814        opened.close().expect("close");
2815    }
2816
2817    #[test]
2818    fn one_read_returns_lookup_state_and_version_from_one_boundary() {
2819        let (_root, opened) = opened(Arc::new(TestControls::default()));
2820        opened
2821            .state
2822            .index
2823            .apply(&Observation::new(vec![Op::Upsert {
2824                path: PathBuf::from("note.txt"),
2825                kind: EntryKind::File,
2826                attrs: crate::Attrs { size: 7, ..crate::Attrs::default() },
2827            }]))
2828            .expect("seed entry");
2829
2830        let response = opened
2831            .read(crate::ReadRequest {
2832                projections: vec![
2833                    crate::ReadProjection::Lookup { path: PathBuf::from("note.txt") },
2834                    crate::ReadProjection::Lookup { path: PathBuf::from("missing.txt") },
2835                ],
2836                ..crate::ReadRequest::default()
2837            })
2838            .expect("coherent read");
2839
2840        assert_eq!(response.version.sequence, opened.state.index.clock().expect("clock"));
2841        assert_eq!(response.state, opened.state.index.state().expect("state"));
2842        assert_eq!(response.results.len(), 2);
2843        assert!(matches!(
2844            &response.results[0],
2845            crate::ProjectionResult::Lookup(crate::Knowledge::Present(entry))
2846                if entry.path == Path::new("note.txt") && entry.attrs.size == 7
2847        ));
2848        assert!(matches!(
2849            response.results[1],
2850            crate::ProjectionResult::Lookup(crate::Knowledge::Absent)
2851        ));
2852        opened.close().expect("close");
2853    }
2854
2855    #[test]
2856    fn change_poll_returns_the_detached_exact_range_and_terminal_state() {
2857        let (_root, opened) = opened(Arc::new(TestControls::default()));
2858        let after = current_version(&opened);
2859        apply_and_notify(
2860            &opened,
2861            &Observation::new(vec![Op::Upsert {
2862                path: PathBuf::from("note.txt"),
2863                kind: EntryKind::File,
2864                attrs: crate::Attrs { size: 7, ..crate::Attrs::default() },
2865            }]),
2866        );
2867        let state_only = opened
2868            .state
2869            .index
2870            .transition_discovery(DiscoveryTransition::Begin)
2871            .expect("state-only commit");
2872        assert!(state_only.commit.as_ref().is_some_and(|commit| commit.changes.is_empty()));
2873        opened.state.journal.notify_commit();
2874
2875        let poll = opened
2876            .changes(crate::ChangeRequest { after, timeout: std::time::Duration::ZERO })
2877            .expect("immediate changes");
2878        let crate::ChangeOutcome::Changes { commits, impact } = &poll.outcome else {
2879            panic!("expected changes");
2880        };
2881        let detached = opened.state.index.since(after.sequence).expect("detached range");
2882        assert_eq!(commits, &detached.commits);
2883        assert_eq!(commits.len(), 2);
2884        assert!(commits[1].changes.is_empty(), "state-only commit remains observable");
2885        assert!(impact.domains.contains(&crate::ImpactDomain::State));
2886        assert_eq!(poll.cursor, poll.version);
2887        assert_eq!(poll.version.sequence, detached.clock);
2888        assert_eq!(poll.state, detached.state);
2889        assert_eq!(poll.work.commits_visited, 2);
2890        assert_eq!(poll.work.commits_returned, 2);
2891        opened.close().expect("close");
2892    }
2893
2894    #[test]
2895    fn idle_change_poll_waits_without_advancing_its_cursor() {
2896        let (_root, opened) = opened(Arc::new(TestControls::default()));
2897        let after = current_version(&opened);
2898        let poll = opened
2899            .changes(crate::ChangeRequest { after, timeout: std::time::Duration::from_millis(5) })
2900            .expect("idle poll");
2901
2902        assert!(matches!(poll.outcome, crate::ChangeOutcome::Idle));
2903        assert_eq!(poll.cursor, after);
2904        assert_eq!(poll.version, after);
2905        assert_eq!(poll.work, crate::Work::default());
2906        opened.close().expect("close");
2907    }
2908
2909    #[test]
2910    fn a_commit_at_the_wait_boundary_cannot_lose_its_wakeup() {
2911        let controls = Arc::new(TestControls::default());
2912        controls.gate(TestPoint::BeforeJournalWait).arm();
2913        let (_root, opened) = opened(Arc::clone(&controls));
2914        let after = current_version(&opened);
2915        let poller = opened.clone();
2916        let poll = thread::spawn(move || {
2917            poller
2918                .changes(crate::ChangeRequest { after, timeout: std::time::Duration::from_secs(5) })
2919        });
2920        controls.gate(TestPoint::BeforeJournalWait).wait_reached();
2921
2922        let (applied_sender, applied_receiver) = std::sync::mpsc::sync_channel(0);
2923        let committer = opened.clone();
2924        let commit = thread::spawn(move || {
2925            let outcome = committer
2926                .state
2927                .index
2928                .apply(&Observation::new(vec![Op::Upsert {
2929                    path: PathBuf::from("arrived"),
2930                    kind: EntryKind::File,
2931                    attrs: crate::Attrs::default(),
2932                }]))
2933                .expect("commit at wait boundary");
2934            applied_sender.send(()).expect("report applied commit");
2935            if outcome.commit.is_some() {
2936                committer.state.journal.notify_commit();
2937            }
2938        });
2939        applied_receiver.recv_timeout(TEST_GATE_TIMEOUT).expect("commit applied");
2940        controls.gate(TestPoint::BeforeJournalWait).release();
2941
2942        commit.join().expect("committer");
2943        let poll = poll.join().expect("poller").expect("change poll");
2944        assert!(matches!(
2945            poll.outcome,
2946            crate::ChangeOutcome::Changes { ref commits, .. } if commits.len() == 1
2947        ));
2948        opened.close().expect("close");
2949    }
2950
2951    #[test]
2952    fn progressive_discovery_notifies_the_same_change_poll() {
2953        let root = tempfile::tempdir().expect("temp root");
2954        std::fs::write(root.path().join("discovered"), b"data").expect("fixture");
2955        let controls = Arc::new(TestControls::default());
2956        controls.gate(TestPoint::BeforeDiscovery).arm();
2957        let opened =
2958            OpenedIndex::open_for_test(root.path(), OpenOptions::default(), Arc::clone(&controls))
2959                .expect("opened root");
2960        let after = current_version(&opened);
2961        let poller = opened.clone();
2962        let poll = thread::spawn(move || {
2963            poller
2964                .changes(crate::ChangeRequest { after, timeout: std::time::Duration::from_secs(5) })
2965        });
2966
2967        controls.gate(TestPoint::BeforeDiscovery).release();
2968        let poll = poll.join().expect("poller").expect("discovery changes");
2969        assert!(matches!(
2970            poll.outcome,
2971            crate::ChangeOutcome::Changes { ref commits, .. } if !commits.is_empty()
2972        ));
2973        assert!(poll.version.sequence > after.sequence);
2974        wait_until_settled(&opened);
2975        opened.close().expect("close");
2976    }
2977
2978    #[test]
2979    fn terminal_state_only_discovery_commit_wakes_a_blocked_poll() {
2980        let root = tempfile::tempdir().expect("temp root");
2981        let controls = Arc::new(TestControls::default());
2982        controls.gate(TestPoint::AfterRootDirectory).arm();
2983        controls.gate(TestPoint::BeforeJournalWait).arm();
2984        let opened =
2985            OpenedIndex::open_for_test(root.path(), OpenOptions::default(), Arc::clone(&controls))
2986                .expect("opened root");
2987        controls.gate(TestPoint::AfterRootDirectory).wait_reached();
2988        let after = current_version(&opened);
2989        let poller = opened.clone();
2990        let poll = thread::spawn(move || {
2991            poller
2992                .changes(crate::ChangeRequest { after, timeout: std::time::Duration::from_secs(5) })
2993        });
2994        controls.gate(TestPoint::BeforeJournalWait).wait_reached();
2995        controls.gate(TestPoint::BeforeJournalWait).release();
2996        controls.gate(TestPoint::AfterRootDirectory).release();
2997
2998        let poll = poll.join().expect("poller").expect("terminal change");
2999        let crate::ChangeOutcome::Changes { commits, .. } = poll.outcome else {
3000            panic!("expected terminal change");
3001        };
3002        assert_eq!(commits.len(), 1);
3003        assert!(commits[0].changes.is_empty());
3004        assert!(commits[0].state.iter().any(|transition| matches!(
3005            transition,
3006            crate::StateTransition::IndexState {
3007                current: crate::IndexState { phase: crate::LifecyclePhase::Ready, .. },
3008                ..
3009            }
3010        )));
3011        assert_eq!(poll.state.phase, crate::LifecyclePhase::Ready);
3012        opened.close().expect("close");
3013    }
3014
3015    #[test]
3016    fn change_cursors_reject_foreign_identity_and_future_sequences() {
3017        let (_first_root, first) = opened(Arc::new(TestControls::default()));
3018        let (_second_root, second) = opened(Arc::new(TestControls::default()));
3019        let first_version = current_version(&first);
3020        let second_version = current_version(&second);
3021
3022        assert!(matches!(
3023            first.changes(crate::ChangeRequest {
3024                after: second_version,
3025                timeout: std::time::Duration::ZERO,
3026            }),
3027            Err(Error::ChangeCursorUnavailable { .. })
3028        ));
3029        let future = crate::EngineVersion {
3030            sequence: first_version.sequence.checked_next().expect("future sequence"),
3031            ..first_version
3032        };
3033        assert!(matches!(
3034            first.changes(crate::ChangeRequest {
3035                after: future,
3036                timeout: std::time::Duration::ZERO,
3037            }),
3038            Err(Error::ChangeCursorUnavailable { .. })
3039        ));
3040        first.close().expect("first close");
3041        second.close().expect("second close");
3042    }
3043
3044    /// Opens a root with no discovery, at the smallest journal budget it accepts.
3045    fn opened_at_the_minimum_journal_budget() -> (tempfile::TempDir, OpenedIndex) {
3046        let controls = Arc::new(TestControls::default());
3047        controls.discovery_disabled.store(true, Ordering::Release);
3048        let root = tempfile::tempdir().expect("temp root");
3049        let opened = OpenedIndex::open_for_test(
3050            root.path(),
3051            OpenOptions {
3052                journal_capacity_bytes: crate::MIN_JOURNAL_CAPACITY_BYTES,
3053                ..OpenOptions::default()
3054            },
3055            controls,
3056        )
3057        .expect("opened root");
3058        (root, opened)
3059    }
3060
3061    /// The least budget accepted holds history worth polling, not one commit: a consumer
3062    /// that falls behind by a burst of single-file commits still receives every one.
3063    #[test]
3064    fn the_minimum_journal_budget_delivers_a_burst_of_single_file_commits() {
3065        const BURST: usize = 64;
3066        let (_root, opened) = opened_at_the_minimum_journal_budget();
3067        let after = current_version(&opened);
3068        for index in 0..BURST {
3069            apply_and_notify(
3070                &opened,
3071                &Observation::new(vec![Op::Upsert {
3072                    path: PathBuf::from(format!("file-{index:02}")),
3073                    kind: EntryKind::File,
3074                    attrs: crate::Attrs::default(),
3075                }]),
3076            );
3077        }
3078
3079        let poll = opened
3080            .changes(crate::ChangeRequest { after, timeout: std::time::Duration::ZERO })
3081            .expect("burst");
3082        let crate::ChangeOutcome::Changes { commits, .. } = &poll.outcome else {
3083            panic!("expected the burst's commits: {:?}", poll.outcome);
3084        };
3085        assert_eq!(commits.len(), BURST);
3086        opened.close().expect("close");
3087    }
3088
3089    #[test]
3090    fn a_slow_consumer_gets_one_coherent_all_dirty_reset() {
3091        let (_root, opened) = opened_at_the_minimum_journal_budget();
3092        let after = current_version(&opened);
3093        // One commit that costs more than the whole budget: many long names at once.
3094        let outcome = apply_and_notify(
3095            &opened,
3096            &Observation::new(
3097                (0..=crate::MAX_DIRTY_PATHS)
3098                    .map(|index| Op::Upsert {
3099                        path: PathBuf::from(format!("{index:0>200}")),
3100                        kind: EntryKind::File,
3101                        attrs: crate::Attrs::default(),
3102                    })
3103                    .collect(),
3104            ),
3105        );
3106        let commit = outcome.commit.expect("effective commit");
3107        assert!(commit.retained_cost() > crate::MIN_JOURNAL_CAPACITY_BYTES);
3108
3109        let poll = opened
3110            .changes(crate::ChangeRequest { after, timeout: std::time::Duration::ZERO })
3111            .expect("consumer reset");
3112        let crate::ChangeOutcome::Reset { impact } = &poll.outcome else {
3113            panic!("expected reset");
3114        };
3115        assert!(impact.all_dirty);
3116        assert!(impact.dirty_paths.is_empty());
3117        assert_eq!(impact.domains.len(), 6);
3118        assert_eq!(poll.cursor, poll.version);
3119        assert_eq!(poll.state, opened.state.index.state().expect("terminal state"));
3120        assert!(opened.state.index.since(after.sequence).expect("history").truncated);
3121        opened.close().expect("close");
3122    }
3123
3124    #[test]
3125    fn close_wakes_a_blocked_change_poll() {
3126        let controls = Arc::new(TestControls::default());
3127        controls.gate(TestPoint::BeforeJournalWait).arm();
3128        let (_root, opened) = opened(Arc::clone(&controls));
3129        let after = current_version(&opened);
3130        let poller = opened.clone();
3131        let poll = thread::spawn(move || {
3132            poller.changes(crate::ChangeRequest {
3133                after,
3134                timeout: std::time::Duration::from_secs(60),
3135            })
3136        });
3137        controls.gate(TestPoint::BeforeJournalWait).wait_reached();
3138        controls.gate(TestPoint::BeforeJournalWait).release();
3139        opened.close().expect("close");
3140
3141        assert!(matches!(poll.join().expect("poller"), Err(Error::OpenedIndexClosed)));
3142    }
3143
3144    #[test]
3145    fn change_invalidations_fail_closed_at_the_existing_path_bound() {
3146        let (_root, opened) = opened(Arc::new(TestControls::default()));
3147        let after = current_version(&opened);
3148        let ops = (0..=crate::MAX_DIRTY_PATHS)
3149            .map(|index| Op::Upsert {
3150                path: PathBuf::from(format!("entry-{index}")),
3151                kind: EntryKind::File,
3152                attrs: crate::Attrs::default(),
3153            })
3154            .collect();
3155        apply_and_notify(&opened, &Observation::new(ops));
3156
3157        let poll = opened
3158            .changes(crate::ChangeRequest { after, timeout: std::time::Duration::ZERO })
3159            .expect("bounded invalidation");
3160        assert!(matches!(
3161            poll.outcome,
3162            crate::ChangeOutcome::Changes {
3163                ref commits,
3164                impact: crate::Impact { all_dirty: true, ref dirty_paths, .. },
3165            } if commits.len() == 1 && dirty_paths.is_empty()
3166        ));
3167        opened.close().expect("close");
3168    }
3169
3170    #[test]
3171    fn lookup_uses_directory_completeness_before_global_discovery_settles() {
3172        let (_root, opened) = opened(Arc::new(TestControls::default()));
3173        opened
3174            .state
3175            .index
3176            .transition_discovery(DiscoveryTransition::Begin)
3177            .expect("begin discovery");
3178        opened
3179            .state
3180            .index
3181            .apply(&Observation::new(vec![Op::Upsert {
3182                path: PathBuf::from("known"),
3183                kind: EntryKind::Dir,
3184                attrs: crate::Attrs::default(),
3185            }]))
3186            .expect("seed directory");
3187
3188        let lookup = || {
3189            opened
3190                .read(crate::ReadRequest {
3191                    projections: vec![crate::ReadProjection::Lookup {
3192                        path: PathBuf::from("known/missing"),
3193                    }],
3194                    ..crate::ReadRequest::default()
3195                })
3196                .expect("lookup")
3197                .results
3198                .into_iter()
3199                .next()
3200                .expect("lookup result")
3201        };
3202        assert!(matches!(
3203            lookup(),
3204            crate::ProjectionResult::Lookup(crate::Knowledge::Unknown {
3205                reason: crate::CoverageReason::Building
3206            })
3207        ));
3208
3209        opened
3210            .state
3211            .index
3212            .apply_discovery(
3213                &Observation::new(Vec::new()),
3214                DiscoveryCommit {
3215                    directory_complete: Some(PathBuf::from("known")),
3216                    transition: None,
3217                },
3218            )
3219            .expect("complete directory");
3220        assert!(matches!(lookup(), crate::ProjectionResult::Lookup(crate::Knowledge::Absent)));
3221        assert!(matches!(
3222            opened.state.index.state().expect("state").coverage,
3223            crate::Coverage::Partial(_)
3224        ));
3225        opened.close().expect("close");
3226    }
3227
3228    /// The completion transition carries the canonical relative path, whatever spelling the
3229    /// producer used. Discovery happened to build canonical paths; nothing else guaranteed it.
3230    #[test]
3231    fn directory_completion_publishes_the_canonical_relative_path() {
3232        let (_root, opened) = opened(Arc::new(TestControls::default()));
3233        opened
3234            .state
3235            .index
3236            .transition_discovery(DiscoveryTransition::Begin)
3237            .expect("begin discovery");
3238        opened
3239            .state
3240            .index
3241            .apply(&Observation::new(vec![Op::Upsert {
3242                path: PathBuf::from("known"),
3243                kind: EntryKind::Dir,
3244                attrs: crate::Attrs::default(),
3245            }]))
3246            .expect("seed directory");
3247
3248        let outcome = opened
3249            .state
3250            .index
3251            .apply_discovery(
3252                &Observation::new(Vec::new()),
3253                DiscoveryCommit {
3254                    directory_complete: Some(PathBuf::from("./known")),
3255                    transition: None,
3256                },
3257            )
3258            .expect("complete directory");
3259
3260        let commit = outcome.commit.expect("completion commit");
3261        assert!(
3262            commit.state.contains(&crate::StateTransition::DirectoryComplete {
3263                path: PathBuf::from("known"),
3264            }),
3265            "{:?}",
3266            commit.state
3267        );
3268        assert_eq!(
3269            opened.state.index.directory_complete(Path::new("known")).expect("lookup"),
3270            Some(true)
3271        );
3272        opened.close().expect("close");
3273    }
3274
3275    /// A tree page and a roll-up on a retained file refuse that projection, and only it.
3276    ///
3277    /// They used to contradict the lookup of the same path on a complete root: the tree
3278    /// answered `Absent`, which claims coverage proves the path missing, and the roll-up
3279    /// answered `Unknown { reason: Building }`, which a caller polling for an answer would
3280    /// wait on forever. Then they failed the whole read, so a mixed request lost the
3281    /// lookup beside them because one path had changed kind since an earlier page
3282    /// (`fdu-l89e`).
3283    #[test]
3284    fn a_path_of_the_wrong_kind_refuses_its_projection_and_the_read_still_answers() {
3285        let (_root, opened) = opened(Arc::new(TestControls::default()));
3286        let file = |path: &str| Op::Upsert {
3287            path: PathBuf::from(path),
3288            kind: EntryKind::File,
3289            attrs: crate::Attrs { size: 5, ..crate::Attrs::default() },
3290        };
3291        opened
3292            .state
3293            .index
3294            .apply(&Observation::new(vec![
3295                file("a"),
3296                file("README.md"),
3297                Op::Upsert {
3298                    path: PathBuf::from("dir"),
3299                    kind: EntryKind::Dir,
3300                    attrs: crate::Attrs::default(),
3301                },
3302            ]))
3303            .expect("seed tree");
3304        let page = crate::PageRequest { limit: 16, max_work: 64 };
3305        let tree = |path: &str| crate::ReadProjection::Tree {
3306            path: PathBuf::from(path),
3307            depth: crate::query::Bound::Limit(1),
3308            include_ignored: true,
3309            page,
3310        };
3311        // The directory answers while it is one: the path a caller holds is a good path.
3312        let before = opened
3313            .read(crate::ReadRequest { projections: vec![tree("dir")], ..Default::default() })
3314            .expect("tree of a directory");
3315        assert!(matches!(
3316            before.results.as_slice(),
3317            [crate::ProjectionResult::Tree(crate::Knowledge::Present(_))]
3318        ));
3319
3320        opened.state.index.apply(&Observation::new(vec![file("dir")])).expect("dir became a file");
3321        assert_eq!(opened.state.index.state().expect("state").coverage, crate::Coverage::Complete);
3322        let response = opened
3323            .read(crate::ReadRequest {
3324                projections: vec![
3325                    crate::ReadProjection::Lookup { path: PathBuf::from("a") },
3326                    tree("dir"),
3327                    crate::ReadProjection::RollUp { path: PathBuf::from("README.md") },
3328                    crate::ReadProjection::Lookup { path: PathBuf::from("dir") },
3329                ],
3330                ..crate::ReadRequest::default()
3331            })
3332            .expect("a refusal does not fail the read");
3333        match response.results.as_slice() {
3334            [
3335                crate::ProjectionResult::Lookup(crate::Knowledge::Present(a)),
3336                crate::ProjectionResult::Refused(crate::ProjectionRefusal::NotADirectory {
3337                    path: tree_path,
3338                }),
3339                crate::ProjectionResult::Refused(crate::ProjectionRefusal::NotADirectory {
3340                    path: rollup_path,
3341                }),
3342                crate::ProjectionResult::Lookup(crate::Knowledge::Present(dir)),
3343            ] => {
3344                assert_eq!(a.path, Path::new("a"));
3345                assert_eq!(tree_path, Path::new("dir"));
3346                assert_eq!(rollup_path, Path::new("README.md"));
3347                assert_eq!(dir.kind, EntryKind::File);
3348            }
3349            other => panic!("each projection answers for itself: {other:?}"),
3350        }
3351        assert_eq!(response.work.rows_returned, 2, "a refusal returns no rows");
3352
3353        // Below a file nothing can exist, and a complete root can say so.
3354        assert!(matches!(
3355            opened
3356                .read(crate::ReadRequest {
3357                    projections: vec![crate::ReadProjection::RollUp {
3358                        path: PathBuf::from("README.md/inner"),
3359                    }],
3360                    ..crate::ReadRequest::default()
3361                })
3362                .expect("rollup below a file")
3363                .results[0],
3364            crate::ProjectionResult::RollUp(crate::Knowledge::Absent)
3365        ));
3366        opened.close().expect("close");
3367    }
3368
3369    /// Seed an opened root whose names need escaping, without touching a filesystem.
3370    #[cfg(unix)]
3371    fn opened_with_escaped_names() -> (tempfile::TempDir, OpenedIndex) {
3372        use std::os::unix::ffi::OsStrExt;
3373        let (root, opened) = opened(Arc::new(TestControls::default()));
3374        let native = |bytes: &[u8]| PathBuf::from(std::ffi::OsStr::from_bytes(bytes));
3375        let file = |path: PathBuf| Op::Upsert {
3376            path,
3377            kind: EntryKind::File,
3378            attrs: crate::Attrs { size: 3, ..crate::Attrs::default() },
3379        };
3380        let dir = |path: PathBuf| Op::Upsert {
3381            path,
3382            kind: EntryKind::Dir,
3383            attrs: crate::Attrs::default(),
3384        };
3385        opened
3386            .state
3387            .index
3388            .apply(&Observation::new(vec![
3389                dir(native(b"x\xff")),
3390                file(native(b"x\xff/inner.txt")),
3391                file(PathBuf::from("100%.txt")),
3392                dir(PathBuf::from("src")),
3393                file(PathBuf::from("src/lib.rs")),
3394            ]))
3395            .expect("seed escaped names");
3396        (root, opened)
3397    }
3398
3399    /// The portable paths of the files one projection admits.
3400    #[cfg(unix)]
3401    fn admitted_files(result: &crate::ProjectionResult) -> std::collections::BTreeSet<String> {
3402        match result {
3403            crate::ProjectionResult::Flat(page) => {
3404                assert!(page.next.is_none(), "one page holds the fixture");
3405                page.rows
3406                    .iter()
3407                    .filter(|row| row.kind == EntryKind::File)
3408                    .map(|row| row.portable_path.as_str().to_owned())
3409                    .collect()
3410            }
3411            crate::ProjectionResult::Report(report) => match report.sections.as_slice() {
3412                [crate::query::Section::Files { rows, .. }] => rows
3413                    .iter()
3414                    .filter(|row| row.kind == EntryKind::File)
3415                    .map(|row| read::portable_path(&row.path).as_str().to_owned())
3416                    .collect(),
3417                other => panic!("a files report: {other:?}"),
3418            },
3419            other => panic!("a page or a report: {other:?}"),
3420        }
3421    }
3422
3423    /// Every projection in an opened read filters by one spelling: the portable path.
3424    ///
3425    /// Flat and aggregate took the name from the portable path and the relative path from
3426    /// the native one, ancestor names compared native components, and a report projection
3427    /// matched native names, so `exact_names: ["100%.txt"]`, an unanchored glob, and an
3428    /// anchored one each answered differently, and a non-UTF-8 ancestor could not be named
3429    /// at all (`fdu-8w5k`).
3430    #[cfg(unix)]
3431    #[test]
3432    fn every_projection_filters_by_the_portable_identity_a_page_returns() {
3433        let (_root, opened) = opened_with_escaped_names();
3434        let glob = |source: &str| crate::query::Pattern::parse(source).expect("pattern");
3435        let names = |values: &[&str]| values.iter().map(ToString::to_string).collect::<Vec<_>>();
3436        let cases: Vec<(&str, crate::query::EntrySelection, &[&str])> = vec![
3437            (
3438                "an exact name, escaped",
3439                crate::query::EntrySelection {
3440                    exact_names: names(&["100%25.txt"]),
3441                    ..Default::default()
3442                },
3443                &["100%25.txt"],
3444            ),
3445            (
3446                "an exact name in its native spelling",
3447                crate::query::EntrySelection {
3448                    exact_names: names(&["100%.txt"]),
3449                    ..Default::default()
3450                },
3451                &[],
3452            ),
3453            (
3454                "a non-UTF-8 ancestor",
3455                crate::query::EntrySelection {
3456                    ancestor_names: names(&["x%FF"]),
3457                    ..Default::default()
3458                },
3459                &["x%FF/inner.txt"],
3460            ),
3461            (
3462                "a terminal suffix below that ancestor",
3463                crate::query::EntrySelection {
3464                    terminal_extensions: names(&[".txt"]),
3465                    ancestor_names: names(&["x%FF"]),
3466                    ..Default::default()
3467                },
3468                &["x%FF/inner.txt"],
3469            ),
3470            (
3471                "an anchored glob through the escaped directory",
3472                crate::query::EntrySelection {
3473                    query: crate::query::Selection {
3474                        include: vec![glob("x%FF/*")],
3475                        ..Default::default()
3476                    },
3477                    ..Default::default()
3478                },
3479                &["x%FF/inner.txt"],
3480            ),
3481            (
3482                "an unanchored glob on an escaped name",
3483                crate::query::EntrySelection {
3484                    query: crate::query::Selection {
3485                        include: vec![glob("100%25.txt")],
3486                        ..Default::default()
3487                    },
3488                    ..Default::default()
3489                },
3490                &["100%25.txt"],
3491            ),
3492            (
3493                "an unanchored glob in the native spelling",
3494                crate::query::EntrySelection {
3495                    query: crate::query::Selection {
3496                        include: vec![glob("100%.txt")],
3497                        ..Default::default()
3498                    },
3499                    ..Default::default()
3500                },
3501                &[],
3502            ),
3503            (
3504                "an exclusion by escaped directory",
3505                crate::query::EntrySelection {
3506                    query: crate::query::Selection {
3507                        exclude: vec![glob("x%FF/**")],
3508                        ..Default::default()
3509                    },
3510                    ..Default::default()
3511                },
3512                &["100%25.txt", "src/lib.rs"],
3513            ),
3514        ];
3515        for (case, selection, expected) in cases {
3516            let expected: std::collections::BTreeSet<String> =
3517                expected.iter().map(ToString::to_string).collect();
3518            let report = crate::ReadProjection::Report(crate::ReportRequest {
3519                query: crate::query::Query {
3520                    views: vec![crate::query::ViewSpec::Files],
3521                    selection: selection.query.clone(),
3522                    ..crate::query::Query::default()
3523                },
3524                now: std::time::SystemTime::UNIX_EPOCH,
3525                max_work: 1_000,
3526            });
3527            let response = opened
3528                .read(crate::ReadRequest {
3529                    projections: vec![
3530                        crate::ReadProjection::Flat {
3531                            selection: selection.clone(),
3532                            shape: crate::RowShape::Compact,
3533                            page: crate::PageRequest { limit: 64, max_work: 1_000 },
3534                        },
3535                        crate::ReadProjection::Aggregate {
3536                            selection: crate::query::EntrySelection {
3537                                query: crate::query::Selection {
3538                                    kinds: vec![EntryKind::File],
3539                                    ..selection.query.clone()
3540                                },
3541                                ..selection.clone()
3542                            },
3543                            count_cap: 64,
3544                            max_work: 1_000,
3545                        },
3546                        report,
3547                    ],
3548                    ..crate::ReadRequest::default()
3549                })
3550                .expect(case);
3551            assert_eq!(admitted_files(&response.results[0]), expected, "flat: {case}");
3552            assert!(
3553                matches!(
3554                    response.results[1],
3555                    crate::ProjectionResult::Aggregate(crate::CountResult::Exact(count))
3556                        if count == expected.len() as u64
3557                ),
3558                "aggregate: {case}: {:?}",
3559                response.results[1]
3560            );
3561            // A report carries only the base selection, so it is compared where that is
3562            // the whole question.
3563            if selection.exact_names.is_empty()
3564                && selection.ancestor_names.is_empty()
3565                && selection.terminal_extensions.is_empty()
3566            {
3567                assert_eq!(admitted_files(&response.results[2]), expected, "report: {case}");
3568            }
3569        }
3570        opened.close().expect("close");
3571    }
3572
3573    /// A path a page returned is a filter a caller can write back, on every axis.
3574    #[cfg(unix)]
3575    #[test]
3576    fn a_path_from_a_page_passes_back_into_a_filter_unchanged() {
3577        let (_root, opened) = opened_with_escaped_names();
3578        let flat = |selection: crate::query::EntrySelection| {
3579            let response = opened
3580                .read(crate::ReadRequest {
3581                    projections: vec![crate::ReadProjection::Flat {
3582                        selection,
3583                        shape: crate::RowShape::Compact,
3584                        page: crate::PageRequest { limit: 64, max_work: 1_000 },
3585                    }],
3586                    ..crate::ReadRequest::default()
3587                })
3588                .expect("flat page");
3589            admitted_files(&response.results[0])
3590        };
3591        let every = flat(crate::query::EntrySelection::default());
3592        assert_eq!(
3593            every,
3594            ["100%25.txt", "src/lib.rs", "x%FF/inner.txt"].map(String::from).into(),
3595            "the page names every file by its portable path"
3596        );
3597        for shown in &every {
3598            let only: std::collections::BTreeSet<String> = [shown.clone()].into();
3599            let (ancestors, name) = match shown.rsplit_once('/') {
3600                Some((ancestors, name)) => (Some(ancestors), name),
3601                None => (None, shown.as_str()),
3602            };
3603            let anchored = crate::query::Pattern::parse(&format!("**/{shown}")).expect("glob");
3604            assert_eq!(
3605                flat(crate::query::EntrySelection {
3606                    query: crate::query::Selection {
3607                        include: vec![anchored],
3608                        ..Default::default()
3609                    },
3610                    ..Default::default()
3611                }),
3612                only,
3613                "the whole path as a glob: {shown}"
3614            );
3615            assert_eq!(
3616                flat(crate::query::EntrySelection {
3617                    exact_names: vec![name.to_owned()],
3618                    ..Default::default()
3619                }),
3620                only,
3621                "its name as an exact name: {shown}"
3622            );
3623            if let Some(ancestors) = ancestors {
3624                let mut selection = crate::query::EntrySelection::default();
3625                selection.admit_ancestor_name(ancestors).expect("a page component is a valid name");
3626                assert_eq!(flat(selection), only, "its parent as an ancestor name: {shown}");
3627            }
3628        }
3629        opened.close().expect("close");
3630    }
3631
3632    /// A selection a constructor would refuse is refused by a read too, before any answer.
3633    #[test]
3634    fn a_read_refuses_a_hand_written_selection_the_constructors_would_refuse() {
3635        let (_root, opened) = opened(Arc::new(TestControls::default()));
3636        for selection in [
3637            crate::query::EntrySelection {
3638                terminal_extensions: vec!["rs".to_string()],
3639                ..Default::default()
3640            },
3641            crate::query::EntrySelection {
3642                ancestor_names: vec!["..".to_string()],
3643                ..Default::default()
3644            },
3645        ] {
3646            let read = opened.read(crate::ReadRequest {
3647                projections: vec![
3648                    crate::ReadProjection::Lookup { path: PathBuf::new() },
3649                    crate::ReadProjection::Aggregate { selection, count_cap: 8, max_work: 64 },
3650                ],
3651                ..crate::ReadRequest::default()
3652            });
3653            assert!(matches!(read, Err(Error::InvalidValue { .. })), "{read:?}");
3654        }
3655        opened.close().expect("close");
3656    }
3657
3658    /// A page whose resume state is too large to retain refuses that page, and only it.
3659    ///
3660    /// The page has rows left, so returning them with no continuation would present a
3661    /// truncated page as a finished one. Failing the read instead discarded every other
3662    /// projection in it (READ-8).
3663    #[test]
3664    fn a_page_whose_continuation_cannot_be_retained_refuses_alone() {
3665        let (_root, opened) = opened(Arc::new(TestControls::default()));
3666        let file = |path: &str| Op::Upsert {
3667            path: PathBuf::from(path),
3668            kind: EntryKind::File,
3669            attrs: crate::Attrs { size: 1, ..crate::Attrs::default() },
3670        };
3671        opened
3672            .state
3673            .index
3674            .apply(&Observation::new(vec![file("a.txt"), file("b.txt")]))
3675            .expect("seed tree");
3676        // A selection this large is valid, and admits both files; it is only too large to
3677        // carry into a continuation record.
3678        let mut exact_names = vec!["a.txt".to_string(), "b.txt".to_string()];
3679        exact_names.extend((0..4_000).map(|number| format!("unused-{number:05}.txt")));
3680        let selection = crate::query::EntrySelection { exact_names, ..Default::default() };
3681        let flat = crate::ReadProjection::Flat {
3682            selection,
3683            shape: crate::RowShape::Compact,
3684            page: crate::PageRequest { limit: 1, max_work: 64 },
3685        };
3686        let response = opened
3687            .read(crate::ReadRequest {
3688                projections: vec![
3689                    flat,
3690                    crate::ReadProjection::Lookup { path: PathBuf::from("b.txt") },
3691                ],
3692                ..crate::ReadRequest::default()
3693            })
3694            .expect("a refused page does not fail the read");
3695        match response.results.as_slice() {
3696            [
3697                crate::ProjectionResult::Refused(
3698                    crate::ProjectionRefusal::ContinuationRecordLimit { attempted, limit },
3699                ),
3700                crate::ProjectionResult::Lookup(crate::Knowledge::Present(_)),
3701            ] => {
3702                assert_eq!(*limit, crate::MAX_CONTINUATION_RECORD_BYTES);
3703                assert!(attempted > limit, "{attempted} > {limit}");
3704            }
3705            other => panic!("the page refuses and the lookup answers: {other:?}"),
3706        }
3707        assert_eq!(response.work.rows_returned, 1, "the refused page returns no rows");
3708        assert_eq!(
3709            opened.state.continuations.lock().expect("table").len(),
3710            0,
3711            "a refused page retains nothing"
3712        );
3713        opened.close().expect("close");
3714    }
3715
3716    #[test]
3717    fn mixed_read_preserves_projection_order_and_uses_maintained_rollups() {
3718        let (_root, opened) = opened(Arc::new(TestControls::default()));
3719        opened
3720            .state
3721            .index
3722            .apply(&Observation::new(vec![Op::Upsert {
3723                path: PathBuf::from("note.txt"),
3724                kind: EntryKind::File,
3725                attrs: crate::Attrs { size: 7, ..crate::Attrs::default() },
3726            }]))
3727            .expect("seed entry");
3728
3729        let response = opened
3730            .read(crate::ReadRequest {
3731                projections: vec![
3732                    crate::ReadProjection::Diagnostics,
3733                    crate::ReadProjection::RollUp { path: PathBuf::new() },
3734                ],
3735                ..crate::ReadRequest::default()
3736            })
3737            .expect("coherent read");
3738
3739        assert!(matches!(
3740            &response.results[0],
3741            crate::ProjectionResult::Diagnostics(diagnostics)
3742                if diagnostics.root == opened.state.root && diagnostics.entries == 2
3743        ));
3744        assert!(matches!(
3745            &response.results[1],
3746            crate::ProjectionResult::RollUp(crate::Knowledge::Present(rollup))
3747                if rollup.all.files == 1 && rollup.all.bytes == 7
3748        ));
3749        assert_eq!(response.work.maintained_index_work, 1);
3750        opened.close().expect("close");
3751    }
3752
3753    #[test]
3754    fn tree_pages_are_directory_first_and_resume_at_the_same_version() {
3755        let (_root, opened) = opened(Arc::new(TestControls::default()));
3756        opened
3757            .state
3758            .index
3759            .apply(&Observation::new(vec![
3760                Op::Upsert {
3761                    path: PathBuf::from("z-dir"),
3762                    kind: EntryKind::Dir,
3763                    attrs: crate::Attrs::default(),
3764                },
3765                Op::Upsert {
3766                    path: PathBuf::from("a.txt"),
3767                    kind: EntryKind::File,
3768                    attrs: crate::Attrs { size: 1, ..crate::Attrs::default() },
3769                },
3770                Op::Upsert {
3771                    path: PathBuf::from("b.txt"),
3772                    kind: EntryKind::File,
3773                    attrs: crate::Attrs { size: 2, ..crate::Attrs::default() },
3774                },
3775            ]))
3776            .expect("seed entries");
3777
3778        let first = opened
3779            .read(crate::ReadRequest {
3780                projections: vec![crate::ReadProjection::Tree {
3781                    path: PathBuf::new(),
3782                    depth: crate::query::Bound::Limit(1),
3783                    include_ignored: true,
3784                    page: crate::PageRequest { limit: 2, max_work: 4 },
3785                }],
3786                ..crate::ReadRequest::default()
3787            })
3788            .expect("first page");
3789        let crate::ProjectionResult::Tree(crate::Knowledge::Present(first_page)) =
3790            &first.results[0]
3791        else {
3792            panic!("tree page");
3793        };
3794        assert_eq!(
3795            first_page.rows.iter().map(|row| row.path.as_path()).collect::<Vec<_>>(),
3796            vec![Path::new("z-dir"), Path::new("a.txt")]
3797        );
3798        let continuation = first_page.next.expect("more rows");
3799
3800        let second = opened
3801            .read(crate::ReadRequest {
3802                projections: vec![crate::ReadProjection::Continue {
3803                    continuation,
3804                    page: crate::PageRequest { limit: 2, max_work: 2 },
3805                }],
3806                expected: Some(first.version),
3807            })
3808            .expect("second page");
3809        let crate::ProjectionResult::Tree(crate::Knowledge::Present(second_page)) =
3810            &second.results[0]
3811        else {
3812            panic!("continued tree page");
3813        };
3814        assert_eq!(
3815            second_page.rows.iter().map(|row| row.path.as_path()).collect::<Vec<_>>(),
3816            vec![Path::new("b.txt")]
3817        );
3818        assert!(second_page.next.is_none());
3819        assert_eq!(first.version, second.version);
3820        opened.close().expect("close");
3821    }
3822
3823    #[test]
3824    fn flat_pages_follow_complete_portable_path_order_without_rescanning() {
3825        let (_root, opened) = opened(Arc::new(TestControls::default()));
3826        opened
3827            .state
3828            .index
3829            .apply(&Observation::new(vec![
3830                Op::Upsert {
3831                    path: PathBuf::from("c.txt"),
3832                    kind: EntryKind::File,
3833                    attrs: crate::Attrs::default(),
3834                },
3835                Op::Upsert {
3836                    path: PathBuf::from("a.txt"),
3837                    kind: EntryKind::File,
3838                    attrs: crate::Attrs::default(),
3839                },
3840                Op::Upsert {
3841                    path: PathBuf::from("b.txt"),
3842                    kind: EntryKind::File,
3843                    attrs: crate::Attrs::default(),
3844                },
3845            ]))
3846            .expect("seed entries");
3847
3848        let first = opened
3849            .read(crate::ReadRequest {
3850                projections: vec![crate::ReadProjection::Flat {
3851                    selection: crate::query::EntrySelection::default(),
3852                    shape: crate::RowShape::Compact,
3853                    page: crate::PageRequest { limit: 2, max_work: 3 },
3854                }],
3855                ..crate::ReadRequest::default()
3856            })
3857            .expect("first page");
3858        let crate::ProjectionResult::Flat(first_page) = &first.results[0] else {
3859            panic!("flat page");
3860        };
3861        assert_eq!(
3862            first_page.rows.iter().map(|row| row.portable_path.as_str()).collect::<Vec<_>>(),
3863            vec!["a.txt", "b.txt"]
3864        );
3865        let continuation = first_page.next.expect("more rows");
3866
3867        let second = opened
3868            .read(crate::ReadRequest {
3869                projections: vec![crate::ReadProjection::Continue {
3870                    continuation,
3871                    page: crate::PageRequest { limit: 2, max_work: 2 },
3872                }],
3873                expected: Some(first.version),
3874            })
3875            .expect("second page");
3876        let crate::ProjectionResult::Flat(second_page) = &second.results[0] else {
3877            panic!("continued flat page");
3878        };
3879        assert_eq!(
3880            second_page.rows.iter().map(|row| row.portable_path.as_str()).collect::<Vec<_>>(),
3881            vec!["c.txt"]
3882        );
3883        assert!(second_page.next.is_none());
3884        assert!(second.work.rows_visited <= 2, "continuation resumed from retained position");
3885        opened.close().expect("close");
3886    }
3887
3888    /// A full page is an answer, whatever the budget left over after filling it.
3889    ///
3890    /// The page used to keep scanning past its last row for the next *admitted* entry to
3891    /// name as its cursor, charging the budget as it went, and a budget that ran out during
3892    /// that look-ahead returned `Limit` and threw the finished page away. The same request at
3893    /// the same budget did the same thing forever, so a fixed-budget client could not page
3894    /// past a run of unselected entries. Swept rather than named, because the defect lived
3895    /// in a band of budgets rather than at one value.
3896    #[test]
3897    fn a_full_flat_page_survives_any_budget_that_filled_it() {
3898        let (_root, opened) = opened(Arc::new(TestControls::default()));
3899        let file = |name: &str| Op::Upsert {
3900            path: PathBuf::from(name),
3901            kind: EntryKind::File,
3902            attrs: crate::Attrs::default(),
3903        };
3904        let mut ops = vec![file("a0.rs"), file("a1.rs")];
3905        ops.extend((0..5).map(|i| file(&format!("b{i}.txt"))));
3906        ops.push(file("c.rs"));
3907        opened.state.index.apply(&Observation::new(ops)).expect("seed entries");
3908        let selection = crate::query::EntrySelection {
3909            query: crate::query::Selection {
3910                include: vec![crate::query::Pattern::parse("*.rs").expect("pattern")],
3911                ..crate::query::Selection::default()
3912            },
3913            ..crate::query::EntrySelection::default()
3914        };
3915        let rows = |page: &crate::FlatPage| {
3916            page.rows.iter().map(|row| row.portable_path.as_str().to_string()).collect::<Vec<_>>()
3917        };
3918
3919        for max_work in 1..=12 {
3920            let first = opened
3921                .read(crate::ReadRequest {
3922                    projections: vec![crate::ReadProjection::Flat {
3923                        selection: selection.clone(),
3924                        shape: crate::RowShape::Compact,
3925                        page: crate::PageRequest { limit: 2, max_work },
3926                    }],
3927                    ..crate::ReadRequest::default()
3928                })
3929                .expect("first page");
3930            if max_work < 2 {
3931                // Too small to fill the page: no position to name, so a typed limit.
3932                assert!(
3933                    matches!(first.results[0], crate::ProjectionResult::Limit(_)),
3934                    "max_work {max_work}: {:?}",
3935                    first.results[0]
3936                );
3937                continue;
3938            }
3939            let crate::ProjectionResult::Flat(first_page) = &first.results[0] else {
3940                panic!("max_work {max_work}: a full page was refused: {:?}", first.results[0]);
3941            };
3942            assert_eq!(rows(first_page), ["a0.rs", "a1.rs"], "max_work {max_work}");
3943            assert!(first.work.rows_visited <= max_work, "max_work {max_work}: {:?}", first.work);
3944            let continuation = first_page.next.expect("the page stopped with rows left");
3945
3946            let second = opened
3947                .read(crate::ReadRequest {
3948                    projections: vec![crate::ReadProjection::Continue {
3949                        continuation,
3950                        page: crate::PageRequest { limit: 2, max_work: crate::MAX_PAGE_WORK },
3951                    }],
3952                    expected: Some(first.version),
3953                })
3954                .expect("second page");
3955            let crate::ProjectionResult::Flat(second_page) = &second.results[0] else {
3956                panic!("max_work {max_work}: continued flat page: {:?}", second.results[0]);
3957            };
3958            assert_eq!(rows(second_page), ["c.rs"], "max_work {max_work}");
3959            assert!(second_page.next.is_none(), "max_work {max_work}");
3960        }
3961        opened.close().expect("close");
3962    }
3963
3964    /// Two reads proceed together: a page in progress does not hold the lifecycle lock.
3965    ///
3966    /// `read()` used to bind the lifecycle guard and project as its tail expression, so the
3967    /// guard lived until the page returned. Every read then excluded every other read, and
3968    /// refresh, worker registration, and the start of close, for up to a full page of work.
3969    #[test]
3970    fn a_read_proceeds_while_another_read_is_projecting() {
3971        let controls = Arc::new(TestControls::default());
3972        let (_root, opened) = opened(Arc::clone(&controls));
3973        opened
3974            .state
3975            .index
3976            .apply(&Observation::new(vec![Op::Upsert {
3977                path: PathBuf::from("a.txt"),
3978                kind: EntryKind::File,
3979                attrs: crate::Attrs::default(),
3980            }]))
3981            .expect("seed entry");
3982        controls.gate(TestPoint::DuringTreeProjection).arm();
3983        let tree_reader = opened.clone();
3984        let tree = thread::spawn(move || {
3985            tree_reader.read(crate::ReadRequest {
3986                projections: vec![crate::ReadProjection::Tree {
3987                    path: PathBuf::new(),
3988                    depth: crate::query::Bound::Limit(1),
3989                    include_ignored: true,
3990                    page: crate::PageRequest { limit: 16, max_work: 64 },
3991                }],
3992                ..crate::ReadRequest::default()
3993            })
3994        });
3995        controls.gate(TestPoint::DuringTreeProjection).wait_reached();
3996
3997        let (sender, receiver) = std::sync::mpsc::channel();
3998        let lookup_reader = opened.clone();
3999        let lookup = thread::spawn(move || {
4000            let _ = sender.send(lookup_reader.read(crate::ReadRequest {
4001                projections: vec![crate::ReadProjection::Lookup { path: PathBuf::from("a.txt") }],
4002                ..crate::ReadRequest::default()
4003            }));
4004        });
4005        let concurrent = receiver.recv_timeout(TEST_GATE_TIMEOUT);
4006        controls.gate(TestPoint::DuringTreeProjection).release();
4007        let concurrent = concurrent.expect("a second read finished while the first projected");
4008        assert!(matches!(
4009            concurrent.expect("lookup").results[0],
4010            crate::ProjectionResult::Lookup(crate::Knowledge::Present(_))
4011        ));
4012        lookup.join().expect("lookup thread");
4013        tree.join().expect("tree thread").expect("tree page");
4014        opened.close().expect("close");
4015    }
4016
4017    /// A page that finishes after close began cannot leave a continuation in a closed root.
4018    #[test]
4019    fn a_read_racing_close_leaves_no_continuation_behind() {
4020        let controls = Arc::new(TestControls::default());
4021        let (_root, opened) = opened(Arc::clone(&controls));
4022        opened
4023            .state
4024            .index
4025            .apply(&Observation::new(vec![
4026                Op::Upsert {
4027                    path: PathBuf::from("a.txt"),
4028                    kind: EntryKind::File,
4029                    attrs: crate::Attrs::default(),
4030                },
4031                Op::Upsert {
4032                    path: PathBuf::from("b.txt"),
4033                    kind: EntryKind::File,
4034                    attrs: crate::Attrs::default(),
4035                },
4036            ]))
4037            .expect("seed entries");
4038        controls.gate(TestPoint::DuringTreeProjection).arm();
4039        let reader = opened.clone();
4040        let page = thread::spawn(move || {
4041            reader.read(crate::ReadRequest {
4042                projections: vec![crate::ReadProjection::Tree {
4043                    path: PathBuf::new(),
4044                    depth: crate::query::Bound::Limit(1),
4045                    include_ignored: true,
4046                    page: crate::PageRequest { limit: 1, max_work: 64 },
4047                }],
4048                ..crate::ReadRequest::default()
4049            })
4050        });
4051        controls.gate(TestPoint::DuringTreeProjection).wait_reached();
4052
4053        let (sender, receiver) = std::sync::mpsc::channel();
4054        let closer = opened.clone();
4055        let close = thread::spawn(move || {
4056            let _ = sender.send(closer.close());
4057        });
4058        let shutdown = receiver.recv_timeout(TEST_GATE_TIMEOUT);
4059        controls.gate(TestPoint::DuringTreeProjection).release();
4060        shutdown.expect("close did not wait for a read in progress").expect("close");
4061        close.join().expect("close thread");
4062
4063        assert!(matches!(page.join().expect("page thread"), Err(Error::OpenedIndexClosed)));
4064        assert_eq!(opened.state.continuations.lock().expect("continuations").len(), 0);
4065    }
4066
4067    #[test]
4068    fn flat_continuation_retains_its_normalized_native_query() {
4069        let (_root, opened) = opened(Arc::new(TestControls::default()));
4070        opened
4071            .state
4072            .index
4073            .apply(&Observation::new(vec![
4074                Op::Upsert {
4075                    path: PathBuf::from("a.rs"),
4076                    kind: EntryKind::File,
4077                    attrs: crate::Attrs::default(),
4078                },
4079                Op::Upsert {
4080                    path: PathBuf::from("b.txt"),
4081                    kind: EntryKind::File,
4082                    attrs: crate::Attrs::default(),
4083                },
4084                Op::Upsert {
4085                    path: PathBuf::from("c.rs"),
4086                    kind: EntryKind::File,
4087                    attrs: crate::Attrs::default(),
4088                },
4089            ]))
4090            .expect("seed entries");
4091        let selection = crate::query::EntrySelection {
4092            query: crate::query::Selection {
4093                include: vec![crate::query::Pattern::parse("*.rs").expect("pattern")],
4094                ..crate::query::Selection::default()
4095            },
4096            ..crate::query::EntrySelection::default()
4097        };
4098
4099        let first = opened
4100            .read(crate::ReadRequest {
4101                projections: vec![crate::ReadProjection::Flat {
4102                    selection,
4103                    shape: crate::RowShape::Compact,
4104                    page: crate::PageRequest { limit: 1, max_work: 3 },
4105                }],
4106                ..crate::ReadRequest::default()
4107            })
4108            .expect("first page");
4109        let crate::ProjectionResult::Flat(first_page) = &first.results[0] else {
4110            panic!("flat page");
4111        };
4112        assert_eq!(first_page.rows[0].portable_path.as_str(), "a.rs");
4113
4114        let second = opened
4115            .read(crate::ReadRequest {
4116                projections: vec![crate::ReadProjection::Continue {
4117                    continuation: first_page.next.expect("continuation"),
4118                    page: crate::PageRequest { limit: 1, max_work: 2 },
4119                }],
4120                expected: Some(first.version),
4121            })
4122            .expect("continued page");
4123        let crate::ProjectionResult::Flat(second_page) = &second.results[0] else {
4124            panic!("continued flat page");
4125        };
4126        assert_eq!(second_page.rows[0].portable_path.as_str(), "c.rs");
4127        assert!(second_page.next.is_none());
4128        // The full first page stopped at `b.txt`, the first entry it had not examined, so
4129        // the resumed page is the one that pays to re-evaluate and skip it under the
4130        // retained `*.rs` selection before reaching `c.rs`.
4131        assert_eq!(second.work.rows_visited, 2);
4132        opened.close().expect("close");
4133    }
4134
4135    /// A consumed or evicted token refuses its `Continue` and the rest of the read answers;
4136    /// a foreign token or a version it no longer holds still fails the read (`fdu-l89e`).
4137    #[test]
4138    fn continuations_are_single_use_version_pinned_handle_local_and_bounded() {
4139        let (_root, opened) = opened(Arc::new(TestControls::default()));
4140        opened
4141            .state
4142            .index
4143            .apply(&Observation::new(vec![
4144                Op::Upsert {
4145                    path: PathBuf::from("a"),
4146                    kind: EntryKind::File,
4147                    attrs: crate::Attrs::default(),
4148                },
4149                Op::Upsert {
4150                    path: PathBuf::from("b"),
4151                    kind: EntryKind::File,
4152                    attrs: crate::Attrs::default(),
4153                },
4154            ]))
4155            .expect("seed entries");
4156        let page = crate::PageRequest { limit: 1, max_work: 2 };
4157        let new_token = || {
4158            let response = opened
4159                .read(crate::ReadRequest {
4160                    projections: vec![crate::ReadProjection::Flat {
4161                        selection: crate::query::EntrySelection::default(),
4162                        shape: crate::RowShape::Compact,
4163                        page,
4164                    }],
4165                    ..crate::ReadRequest::default()
4166                })
4167                .expect("first page");
4168            let crate::ProjectionResult::Flat(result) = &response.results[0] else {
4169                panic!("flat page");
4170            };
4171            result.next.expect("continuation")
4172        };
4173
4174        // A token this root issued and no longer holds costs only its own projection: the
4175        // lookup beside it still answers.
4176        let refuses_beside_a_lookup = |continuation| {
4177            let response = opened
4178                .read(crate::ReadRequest {
4179                    projections: vec![
4180                        crate::ReadProjection::Continue { continuation, page },
4181                        crate::ReadProjection::Lookup { path: PathBuf::from("a") },
4182                    ],
4183                    ..crate::ReadRequest::default()
4184                })
4185                .expect("a refused continuation does not fail the read");
4186            matches!(
4187                response.results.as_slice(),
4188                [
4189                    crate::ProjectionResult::Refused(
4190                        crate::ProjectionRefusal::ContinuationUnavailable
4191                    ),
4192                    crate::ProjectionResult::Lookup(crate::Knowledge::Present(_)),
4193                ]
4194            )
4195        };
4196
4197        let replay = new_token();
4198        opened
4199            .read(crate::ReadRequest {
4200                projections: vec![crate::ReadProjection::Continue { continuation: replay, page }],
4201                ..crate::ReadRequest::default()
4202            })
4203            .expect("first continuation use");
4204        assert!(refuses_beside_a_lookup(replay), "a consumed token refuses its page");
4205
4206        let stale = new_token();
4207        opened
4208            .state
4209            .index
4210            .apply(&Observation::new(vec![Op::Upsert {
4211                path: PathBuf::from("c"),
4212                kind: EntryKind::File,
4213                attrs: crate::Attrs::default(),
4214            }]))
4215            .expect("advance version");
4216        assert!(matches!(
4217            opened.read(crate::ReadRequest {
4218                projections: vec![crate::ReadProjection::Continue { continuation: stale, page }],
4219                ..crate::ReadRequest::default()
4220            }),
4221            Err(Error::ContinuationStale { .. })
4222        ));
4223
4224        let retryable = new_token();
4225        // Two rows on a budget of one: the budget runs out before the page fills, so there
4226        // is no position to name. A one-row page would fill on its first entry and return.
4227        let limited = opened
4228            .read(crate::ReadRequest {
4229                projections: vec![crate::ReadProjection::Continue {
4230                    continuation: retryable,
4231                    page: crate::PageRequest { limit: 2, max_work: 1 },
4232                }],
4233                ..crate::ReadRequest::default()
4234            })
4235            .expect("bounded continuation");
4236        assert!(matches!(limited.results[0], crate::ProjectionResult::Limit(_)));
4237        assert!(matches!(
4238            opened
4239                .read(crate::ReadRequest {
4240                    projections: vec![crate::ReadProjection::Continue {
4241                        continuation: retryable,
4242                        page: crate::PageRequest { limit: 1, max_work: 2 },
4243                    }],
4244                    ..crate::ReadRequest::default()
4245                })
4246                .expect("retry continuation")
4247                .results[0],
4248            crate::ProjectionResult::Flat(_)
4249        ));
4250
4251        let foreign = new_token();
4252        let (_other_root, other) = self::opened(Arc::new(TestControls::default()));
4253        assert!(matches!(
4254            other.read(crate::ReadRequest {
4255                projections: vec![crate::ReadProjection::Continue { continuation: foreign, page }],
4256                ..crate::ReadRequest::default()
4257            }),
4258            Err(Error::ContinuationUnavailable)
4259        ));
4260
4261        let oldest = new_token();
4262        for _ in 0..super::continuation::MAX_CONTINUATIONS {
4263            let _ = new_token();
4264        }
4265        assert!(refuses_beside_a_lookup(oldest), "an evicted token refuses its page");
4266        other.close().expect("close other");
4267        opened.close().expect("close");
4268    }
4269
4270    /// Level order, proved against the sequence pre-order would have produced.
4271    ///
4272    /// An order is only proved by a fixture whose answer differs between the plausible
4273    /// readings, and "parent-first" admits both. This tree is three levels deep and wide
4274    /// at the top, so the two disagree:
4275    ///
4276    /// ```text
4277    /// a/  a/a1/  a/a1/deep.txt  b/  b/b1/  z.txt
4278    /// ```
4279    ///
4280    /// Level order returns `a`, `b`, `z.txt`, then `a/a1`, `b/b1`, then `a/a1/deep.txt` —
4281    /// every level whole before descending. Pre-order would return `a`, `a/a1`,
4282    /// `a/a1/deep.txt`, `b`, `b/b1`, `z.txt`, burying `b` behind the whole of `a`'s
4283    /// subtree. A page bound cutting the pre-order sequence at three rows would hide the
4284    /// existence of `b` and `z.txt` entirely, which is what level order prevents.
4285    #[test]
4286    fn tree_pages_are_breadth_first_across_levels() {
4287        let (_root, opened) = opened(Arc::new(TestControls::default()));
4288        opened
4289            .state
4290            .index
4291            .apply(&Observation::new(vec![
4292                Op::Upsert {
4293                    path: PathBuf::from("a"),
4294                    kind: EntryKind::Dir,
4295                    attrs: crate::Attrs::default(),
4296                },
4297                Op::Upsert {
4298                    path: PathBuf::from("b"),
4299                    kind: EntryKind::Dir,
4300                    attrs: crate::Attrs::default(),
4301                },
4302                Op::Upsert {
4303                    path: PathBuf::from("z.txt"),
4304                    kind: EntryKind::File,
4305                    attrs: crate::Attrs { size: 1, ..crate::Attrs::default() },
4306                },
4307                Op::Upsert {
4308                    path: PathBuf::from("a/a1"),
4309                    kind: EntryKind::Dir,
4310                    attrs: crate::Attrs::default(),
4311                },
4312                Op::Upsert {
4313                    path: PathBuf::from("b/b1"),
4314                    kind: EntryKind::Dir,
4315                    attrs: crate::Attrs::default(),
4316                },
4317                Op::Upsert {
4318                    path: PathBuf::from("a/a1/deep.txt"),
4319                    kind: EntryKind::File,
4320                    attrs: crate::Attrs { size: 2, ..crate::Attrs::default() },
4321                },
4322            ]))
4323            .expect("seed tree");
4324
4325        let rows = |depth: crate::query::Bound| -> Vec<String> {
4326            let response = opened
4327                .read(crate::ReadRequest {
4328                    projections: vec![crate::ReadProjection::Tree {
4329                        path: PathBuf::new(),
4330                        depth,
4331                        include_ignored: true,
4332                        page: crate::PageRequest {
4333                            limit: crate::MAX_PAGE_ROWS,
4334                            max_work: crate::MAX_PAGE_WORK,
4335                        },
4336                    }],
4337                    ..crate::ReadRequest::default()
4338                })
4339                .expect("tree read");
4340            let crate::ProjectionResult::Tree(crate::Knowledge::Present(page)) =
4341                &response.results[0]
4342            else {
4343                panic!("tree page");
4344            };
4345            page.rows.iter().map(|row| row.portable_path.as_str().to_owned()).collect()
4346        };
4347
4348        // One level is this directory's own children: directories first, then files, each
4349        // partition in canonical byte order.
4350        assert_eq!(rows(crate::query::Bound::Limit(1)), vec!["a", "b", "z.txt"]);
4351
4352        // Two levels adds the next level whole, never a subtree at a time.
4353        assert_eq!(rows(crate::query::Bound::Limit(2)), vec!["a", "b", "z.txt", "a/a1", "b/b1"]);
4354
4355        // Unbounded reaches the leaf, still level by level. Pre-order would have placed
4356        // `a/a1` and `a/a1/deep.txt` before `b`.
4357        assert_eq!(
4358            rows(crate::query::Bound::All),
4359            vec!["a", "b", "z.txt", "a/a1", "b/b1", "a/a1/deep.txt"]
4360        );
4361
4362        opened.close().expect("close");
4363    }
4364
4365    /// Paging a multi-level tree one row at a time reassembles the same sequence.
4366    ///
4367    /// Resumption is where a level-order traversal can go wrong invisibly: the cursor
4368    /// holds one frame, and crossing a level boundary means re-deriving the position from
4369    /// the ancestor chain. A page bound that lands exactly on such a boundary is the case
4370    /// that would duplicate or drop a row, so this walks every boundary in the fixture.
4371    #[test]
4372    fn tree_pages_resume_across_level_boundaries() {
4373        let (_root, opened) = opened(Arc::new(TestControls::default()));
4374        opened
4375            .state
4376            .index
4377            .apply(&Observation::new(vec![
4378                Op::Upsert {
4379                    path: PathBuf::from("a"),
4380                    kind: EntryKind::Dir,
4381                    attrs: crate::Attrs::default(),
4382                },
4383                Op::Upsert {
4384                    path: PathBuf::from("b"),
4385                    kind: EntryKind::Dir,
4386                    attrs: crate::Attrs::default(),
4387                },
4388                Op::Upsert {
4389                    path: PathBuf::from("z.txt"),
4390                    kind: EntryKind::File,
4391                    attrs: crate::Attrs { size: 1, ..crate::Attrs::default() },
4392                },
4393                Op::Upsert {
4394                    path: PathBuf::from("a/a1"),
4395                    kind: EntryKind::Dir,
4396                    attrs: crate::Attrs::default(),
4397                },
4398                Op::Upsert {
4399                    path: PathBuf::from("b/b1"),
4400                    kind: EntryKind::Dir,
4401                    attrs: crate::Attrs::default(),
4402                },
4403                Op::Upsert {
4404                    path: PathBuf::from("a/a1/deep.txt"),
4405                    kind: EntryKind::File,
4406                    attrs: crate::Attrs { size: 2, ..crate::Attrs::default() },
4407                },
4408            ]))
4409            .expect("seed tree");
4410
4411        let whole = vec!["a", "b", "z.txt", "a/a1", "b/b1", "a/a1/deep.txt"];
4412        let page = crate::PageRequest { limit: 1, max_work: crate::MAX_PAGE_WORK };
4413
4414        let mut seen: Vec<String> = Vec::new();
4415        let response = opened
4416            .read(crate::ReadRequest {
4417                projections: vec![crate::ReadProjection::Tree {
4418                    path: PathBuf::new(),
4419                    depth: crate::query::Bound::All,
4420                    include_ignored: true,
4421                    page,
4422                }],
4423                ..crate::ReadRequest::default()
4424            })
4425            .expect("first page");
4426        let crate::ProjectionResult::Tree(crate::Knowledge::Present(first)) = &response.results[0]
4427        else {
4428            panic!("tree page");
4429        };
4430        seen.extend(first.rows.iter().map(|row| row.portable_path.as_str().to_owned()));
4431        let mut continuation = first.next;
4432
4433        while let Some(token) = continuation {
4434            let response = opened
4435                .read(crate::ReadRequest {
4436                    projections: vec![crate::ReadProjection::Continue {
4437                        continuation: token,
4438                        page,
4439                    }],
4440                    ..crate::ReadRequest::default()
4441                })
4442                .expect("resumed page");
4443            let crate::ProjectionResult::Tree(crate::Knowledge::Present(next)) =
4444                &response.results[0]
4445            else {
4446                panic!("tree page");
4447            };
4448            seen.extend(next.rows.iter().map(|row| row.portable_path.as_str().to_owned()));
4449            continuation = next.next;
4450            assert!(seen.len() <= whole.len(), "paging must terminate, saw {seen:?}");
4451        }
4452
4453        assert_eq!(seen, whole, "one row at a time reassembles the single-page order");
4454        opened.close().expect("close");
4455    }
4456
4457    /// A page the work budget stops must hand back a continuation.
4458    ///
4459    /// The row limit and the work budget are different stopping conditions, and only the
4460    /// row limit is reached inside `collect_children`, which knows the exact child it
4461    /// stopped at. The budget can also run out while *advancing* between parents, where
4462    /// no row has been reached to point at. Breaking there returns `next: None`, which a
4463    /// caller cannot tell apart from a traversal that finished — the tree simply comes
4464    /// back missing every level below the one that fit.
4465    ///
4466    /// Sweeping the budget rather than naming one keeps this from testing an arithmetic
4467    /// coincidence: every budget large enough to make progress must reassemble the whole
4468    /// tree, whichever of the two conditions happens to stop each page.
4469    #[test]
4470    fn a_tree_page_stopped_by_the_work_budget_is_resumable() {
4471        let (_root, opened) = opened(Arc::new(TestControls::default()));
4472        opened
4473            .state
4474            .index
4475            .apply(&Observation::new(vec![
4476                Op::Upsert {
4477                    path: PathBuf::from("a"),
4478                    kind: EntryKind::Dir,
4479                    attrs: crate::Attrs::default(),
4480                },
4481                Op::Upsert {
4482                    path: PathBuf::from("b"),
4483                    kind: EntryKind::Dir,
4484                    attrs: crate::Attrs::default(),
4485                },
4486                Op::Upsert {
4487                    path: PathBuf::from("z.txt"),
4488                    kind: EntryKind::File,
4489                    attrs: crate::Attrs { size: 1, ..crate::Attrs::default() },
4490                },
4491                Op::Upsert {
4492                    path: PathBuf::from("a/a1"),
4493                    kind: EntryKind::Dir,
4494                    attrs: crate::Attrs::default(),
4495                },
4496                Op::Upsert {
4497                    path: PathBuf::from("b/b1"),
4498                    kind: EntryKind::Dir,
4499                    attrs: crate::Attrs::default(),
4500                },
4501                Op::Upsert {
4502                    path: PathBuf::from("a/a1/deep.txt"),
4503                    kind: EntryKind::File,
4504                    attrs: crate::Attrs { size: 2, ..crate::Attrs::default() },
4505                },
4506            ]))
4507            .expect("seed tree");
4508
4509        let whole = ["a", "b", "z.txt", "a/a1", "b/b1", "a/a1/deep.txt"];
4510
4511        // Two is the smallest budget that can still afford a row after the path walk, so
4512        // it is the smallest at which paging is obliged to make progress at all.
4513        for max_work in 2..=14_u64 {
4514            let page = crate::PageRequest { limit: crate::MAX_PAGE_ROWS, max_work };
4515            let mut seen: Vec<String> = Vec::new();
4516            let mut continuation = None;
4517            let mut pages = 0;
4518
4519            loop {
4520                let projection = match continuation {
4521                    None => crate::ReadProjection::Tree {
4522                        path: PathBuf::new(),
4523                        depth: crate::query::Bound::All,
4524                        include_ignored: true,
4525                        page,
4526                    },
4527                    Some(token) => crate::ReadProjection::Continue { continuation: token, page },
4528                };
4529                let response = opened
4530                    .read(crate::ReadRequest {
4531                        projections: vec![projection],
4532                        ..crate::ReadRequest::default()
4533                    })
4534                    .expect("tree read");
4535                let current = match &response.results[0] {
4536                    crate::ProjectionResult::Tree(crate::Knowledge::Present(page)) => page,
4537                    // A limit here would mean a budget the tree cannot be read at, and
4538                    // re-asking cannot help: the search that overran would restart and
4539                    // overrun again. Every budget that can hold a row must finish.
4540                    other => panic!("unexpected result at budget {max_work}: {other:?}"),
4541                };
4542                seen.extend(current.rows.iter().map(|row| row.portable_path.as_str().to_owned()));
4543                continuation = current.next;
4544                pages += 1;
4545                assert!(
4546                    pages <= whole.len() * 4 + 16,
4547                    "budget {max_work} never finished paging, saw {seen:?}"
4548                );
4549                if continuation.is_none() {
4550                    break;
4551                }
4552            }
4553
4554            let mut sorted = seen.clone();
4555            sorted.sort();
4556            let mut expected: Vec<String> = whole.iter().map(|row| (*row).to_owned()).collect();
4557            expected.sort();
4558            assert_eq!(
4559                sorted, expected,
4560                "budget {max_work} finished with next: None while missing rows; saw {seen:?}"
4561            );
4562        }
4563
4564        opened.close().expect("close");
4565    }
4566
4567    /// A page must move, even when the path walk has already spent the budget.
4568    ///
4569    /// `spent` starts at the cost of walking to the requested directory, and only a walk
4570    /// strictly longer than the budget is refused outright. At exactly the budget the walk
4571    /// is allowed, and then the first child pushes `spent` over before any row is emitted:
4572    /// the page returns no rows and a cursor pointing at that same child, and resuming
4573    /// reproduces it exactly. The bound stops being "how much work per page" and becomes
4574    /// "no page ever finishes".
4575    ///
4576    /// A budget says where to stop, not whether to start. Every page therefore emits at
4577    /// least one row or ends the traversal, and this reads a nested directory so the path
4578    /// walk is expensive enough to collide with the budget at all.
4579    #[test]
4580    fn a_page_moves_even_when_the_path_walk_spends_the_budget() {
4581        let (_root, opened) = opened(Arc::new(TestControls::default()));
4582        opened
4583            .state
4584            .index
4585            .apply(&Observation::new(vec![
4586                Op::Upsert {
4587                    path: PathBuf::from("a"),
4588                    kind: EntryKind::Dir,
4589                    attrs: crate::Attrs::default(),
4590                },
4591                Op::Upsert {
4592                    path: PathBuf::from("a/b"),
4593                    kind: EntryKind::Dir,
4594                    attrs: crate::Attrs::default(),
4595                },
4596                Op::Upsert {
4597                    path: PathBuf::from("a/b/x.txt"),
4598                    kind: EntryKind::File,
4599                    attrs: crate::Attrs { size: 1, ..crate::Attrs::default() },
4600                },
4601                Op::Upsert {
4602                    path: PathBuf::from("a/b/y.txt"),
4603                    kind: EntryKind::File,
4604                    attrs: crate::Attrs { size: 2, ..crate::Attrs::default() },
4605                },
4606            ]))
4607            .expect("seed tree");
4608
4609        // Walking to `a/b` costs three, so three is the budget that is spent on arrival.
4610        // Sweeping upward from it keeps this a property rather than one arithmetic
4611        // coincidence: every budget the request is allowed to make must terminate.
4612        for max_work in 3..=12_u64 {
4613            let page = crate::PageRequest { limit: crate::MAX_PAGE_ROWS, max_work };
4614            let mut seen: Vec<String> = Vec::new();
4615            let mut continuation = None;
4616            let mut pages = 0;
4617
4618            loop {
4619                let projection = match continuation {
4620                    None => crate::ReadProjection::Tree {
4621                        path: PathBuf::from("a/b"),
4622                        depth: crate::query::Bound::All,
4623                        include_ignored: true,
4624                        page,
4625                    },
4626                    Some(token) => crate::ReadProjection::Continue { continuation: token, page },
4627                };
4628                let response = opened
4629                    .read(crate::ReadRequest {
4630                        projections: vec![projection],
4631                        ..crate::ReadRequest::default()
4632                    })
4633                    .expect("tree read");
4634                let crate::ProjectionResult::Tree(crate::Knowledge::Present(current)) =
4635                    &response.results[0]
4636                else {
4637                    panic!("tree page");
4638                };
4639                seen.extend(current.rows.iter().map(|row| row.portable_path.as_str().to_owned()));
4640                continuation = current.next;
4641                pages += 1;
4642                assert!(
4643                    pages <= 16,
4644                    "budget {max_work} never terminated; after {pages} pages saw {seen:?}"
4645                );
4646                if continuation.is_none() {
4647                    break;
4648                }
4649            }
4650
4651            assert_eq!(seen, vec!["a/b/x.txt", "a/b/y.txt"], "budget {max_work} lost rows");
4652        }
4653        opened.close().expect("close");
4654    }
4655
4656    /// Descending must not scan the level it is leaving.
4657    ///
4658    /// A level of leaf directories is the shape of every tree's last level, and searching
4659    /// it for a directory child asks every parent in order to conclude there is nothing
4660    /// below. Noticing the first directory while the level is emitted answers the same
4661    /// question for free, so the work a page reports has to stay proportional to the rows
4662    /// it returns rather than to the width of the level under it.
4663    #[test]
4664    fn descending_costs_nothing_on_a_level_of_leaves() {
4665        let (_root, opened) = opened(Arc::new(TestControls::default()));
4666        let mut ops = Vec::new();
4667        for index in 0..60 {
4668            ops.push(Op::Upsert {
4669                path: PathBuf::from(format!("d{index:03}")),
4670                kind: EntryKind::Dir,
4671                attrs: crate::Attrs::default(),
4672            });
4673            ops.push(Op::Upsert {
4674                path: PathBuf::from(format!("d{index:03}/leaf.txt")),
4675                kind: EntryKind::File,
4676                attrs: crate::Attrs { size: 1, ..crate::Attrs::default() },
4677            });
4678        }
4679        opened.state.index.apply(&Observation::new(ops)).expect("seed wide leaf level");
4680
4681        let response = opened
4682            .read(crate::ReadRequest {
4683                projections: vec![crate::ReadProjection::Tree {
4684                    path: PathBuf::new(),
4685                    depth: crate::query::Bound::All,
4686                    include_ignored: true,
4687                    page: crate::PageRequest {
4688                        limit: crate::MAX_PAGE_ROWS,
4689                        max_work: crate::MAX_PAGE_WORK,
4690                    },
4691                }],
4692                ..crate::ReadRequest::default()
4693            })
4694            .expect("tree read");
4695        let crate::ProjectionResult::Tree(crate::Knowledge::Present(page)) = &response.results[0]
4696        else {
4697            panic!("tree page");
4698        };
4699        assert_eq!(page.rows.len(), 120, "60 directories and their 60 files");
4700        assert!(page.next.is_none(), "one page holds the whole tree");
4701        // Three steps per directory and no more: emit the directory as a row at level
4702        // one, emit its file as a row at level two, and step past it to its sibling. The
4703        // slack covers the path walk and the two advances that end each level.
4704        //
4705        // Searching for the descent instead of remembering it adds a fourth step per
4706        // directory, because it asks every one of them for a directory child before
4707        // concluding there is no level below. That is what this bound rejects: measured,
4708        // it is 180 steps with the memo and 241 without, for the same 121 rows.
4709        let width = 60;
4710        assert!(
4711            response.work.rows_visited <= 3 * width + 10,
4712            "descent scanned the level it was leaving: {} steps for {} rows",
4713            response.work.rows_visited,
4714            response.work.rows_returned
4715        );
4716        opened.close().expect("close");
4717    }
4718
4719    /// Every row path of one unbounded tree page from the root.
4720    fn tree_rows(opened: &OpenedIndex, include_ignored: bool) -> Vec<String> {
4721        let response = opened
4722            .read(crate::ReadRequest {
4723                projections: vec![crate::ReadProjection::Tree {
4724                    path: PathBuf::new(),
4725                    depth: crate::query::Bound::All,
4726                    include_ignored,
4727                    page: crate::PageRequest {
4728                        limit: crate::MAX_PAGE_ROWS,
4729                        max_work: crate::MAX_PAGE_WORK,
4730                    },
4731                }],
4732                ..crate::ReadRequest::default()
4733            })
4734            .expect("tree read");
4735        let crate::ProjectionResult::Tree(crate::Knowledge::Present(page)) = &response.results[0]
4736        else {
4737            panic!("tree page");
4738        };
4739        page.rows.iter().map(|row| row.portable_path.as_str().to_owned()).collect()
4740    }
4741
4742    /// A tree read that excludes ignored entries reads each row's ignore bit without the
4743    /// observation check `Index::is_ignored` makes, on the invariant that an opened root
4744    /// always observes control state. The invariant is pinned here, in every build profile,
4745    /// and a read over a tree nothing ignores keeps every row.
4746    #[test]
4747    fn an_opened_tree_read_excluding_ignored_entries_relies_on_an_observing_root() {
4748        let (_root, opened) = opened(Arc::new(TestControls::default()));
4749        opened
4750            .state
4751            .index
4752            .apply(&Observation::new(vec![
4753                Op::Upsert {
4754                    path: PathBuf::from("src"),
4755                    kind: EntryKind::Dir,
4756                    attrs: crate::Attrs::default(),
4757                },
4758                Op::Upsert {
4759                    path: PathBuf::from("src/main.rs"),
4760                    kind: EntryKind::File,
4761                    attrs: crate::Attrs { size: 1, ..crate::Attrs::default() },
4762                },
4763            ]))
4764            .expect("seed tree");
4765
4766        let image = opened.state.index.snapshot().expect("snapshot");
4767        assert!(image.observes_controls());
4768        let excluded = tree_rows(&opened, false);
4769        assert_eq!(excluded, ["src", "src/main.rs"]);
4770        assert_eq!(excluded, tree_rows(&opened, true));
4771
4772        opened.close().expect("close");
4773    }
4774
4775    /// Excluding ignored entries prunes the subtree, not merely the row.
4776    ///
4777    /// Filtering the row and descending anyway is an equally reasonable reading of an
4778    /// unstated rule, and it is observably different: it would still return
4779    /// `vendor/keep.txt` while hiding the directory that explains where it came from.
4780    #[test]
4781    fn excluding_ignored_prunes_the_subtree() {
4782        let (_root, opened) = opened(Arc::new(TestControls::default()));
4783        opened
4784            .state
4785            .index
4786            .apply(&Observation::new(vec![
4787                Op::Upsert {
4788                    path: PathBuf::from("src"),
4789                    kind: EntryKind::Dir,
4790                    attrs: crate::Attrs::default(),
4791                },
4792                Op::Upsert {
4793                    path: PathBuf::from("src/main.rs"),
4794                    kind: EntryKind::File,
4795                    attrs: crate::Attrs { size: 1, ..crate::Attrs::default() },
4796                },
4797                Op::ControlUpsert {
4798                    path: PathBuf::from(".gitignore"),
4799                    source: b"vendor/\n".to_vec(),
4800                },
4801                Op::Upsert {
4802                    path: PathBuf::from("vendor"),
4803                    kind: EntryKind::Dir,
4804                    attrs: crate::Attrs::default(),
4805                },
4806                Op::Upsert {
4807                    path: PathBuf::from("vendor/keep.txt"),
4808                    kind: EntryKind::File,
4809                    attrs: crate::Attrs { size: 2, ..crate::Attrs::default() },
4810                },
4811            ]))
4812            .expect("seed tree");
4813
4814        let included = tree_rows(&opened, true);
4815        assert!(included.iter().any(|row| row == "vendor"));
4816        assert!(included.iter().any(|row| row == "vendor/keep.txt"));
4817
4818        let excluded = tree_rows(&opened, false);
4819        assert!(!excluded.iter().any(|row| row == "vendor"), "the row is gone");
4820        assert!(
4821            !excluded.iter().any(|row| row == "vendor/keep.txt"),
4822            "and so is everything beneath it, which is what pruning means"
4823        );
4824        assert!(excluded.iter().any(|row| row == "src/main.rs"), "unignored work is untouched");
4825
4826        opened.close().expect("close");
4827    }
4828
4829    /// The remembered descent must be the first *unpruned* directory, not the first one.
4830    ///
4831    /// Noticing the next level's first parent while emitting is only equivalent to
4832    /// searching for it if both apply the same pruning rule. `first_directory_child`
4833    /// skips ignored directories, so remembering one before the ignore check makes the
4834    /// two disagree and hands the traversal a parent the search would never have chosen.
4835    ///
4836    /// No row leaks when that happens — every child of a pruned directory is itself
4837    /// ignored, so the row filter catches them a second time. Which is the point: the
4838    /// mistake is invisible in the output and visible only in the work, and a fixture
4839    /// that checked rows alone would pass either way. Pruning means an excluded
4840    /// directory is never expanded; being saved by a second filter is not pruning.
4841    ///
4842    /// So the ignored directory sorts first and is given enough children that expanding
4843    /// it cannot hide in the noise.
4844    #[test]
4845    fn the_remembered_descent_skips_a_pruned_first_child() {
4846        let (_root, opened) = opened(Arc::new(TestControls::default()));
4847        opened
4848            .state
4849            .index
4850            .apply(&Observation::new(vec![
4851                Op::ControlUpsert {
4852                    path: PathBuf::from(".gitignore"),
4853                    source: b"aaa_vendor/\n".to_vec(),
4854                },
4855                Op::Upsert {
4856                    path: PathBuf::from("src"),
4857                    kind: EntryKind::Dir,
4858                    attrs: crate::Attrs::default(),
4859                },
4860                Op::Upsert {
4861                    path: PathBuf::from("src/aaa_vendor"),
4862                    kind: EntryKind::Dir,
4863                    attrs: crate::Attrs::default(),
4864                },
4865                Op::Upsert {
4866                    path: PathBuf::from("src/bbb_keep"),
4867                    kind: EntryKind::Dir,
4868                    attrs: crate::Attrs::default(),
4869                },
4870                Op::Upsert {
4871                    path: PathBuf::from("src/bbb_keep/kept.txt"),
4872                    kind: EntryKind::File,
4873                    attrs: crate::Attrs { size: 2, ..crate::Attrs::default() },
4874                },
4875            ]))
4876            .expect("seed tree");
4877        let buried: Vec<Op> = (0..40)
4878            .map(|index| Op::Upsert {
4879                path: PathBuf::from(format!("src/aaa_vendor/hidden{index:03}.txt")),
4880                kind: EntryKind::File,
4881                attrs: crate::Attrs { size: 1, ..crate::Attrs::default() },
4882            })
4883            .collect();
4884        opened.state.index.apply(&Observation::new(buried)).expect("seed the pruned subtree");
4885
4886        let response = opened
4887            .read(crate::ReadRequest {
4888                projections: vec![crate::ReadProjection::Tree {
4889                    path: PathBuf::new(),
4890                    depth: crate::query::Bound::All,
4891                    include_ignored: false,
4892                    page: crate::PageRequest {
4893                        limit: crate::MAX_PAGE_ROWS,
4894                        max_work: crate::MAX_PAGE_WORK,
4895                    },
4896                }],
4897                ..crate::ReadRequest::default()
4898            })
4899            .expect("tree read");
4900        let crate::ProjectionResult::Tree(crate::Knowledge::Present(page)) = &response.results[0]
4901        else {
4902            panic!("tree page");
4903        };
4904        let rows: Vec<String> =
4905            page.rows.iter().map(|row| row.portable_path.as_str().to_owned()).collect();
4906
4907        assert!(rows.iter().any(|row| row == "src/bbb_keep"), "the kept directory is listed");
4908        assert!(
4909            rows.iter().any(|row| row == "src/bbb_keep/kept.txt"),
4910            "and the descent reached the level below it"
4911        );
4912        assert!(
4913            !rows.iter().any(|row| row == "src/aaa_vendor"),
4914            "the pruned directory is not a row"
4915        );
4916        assert!(
4917            !rows.iter().any(|row| row.starts_with("src/aaa_vendor/")),
4918            "and nothing beneath it is listed"
4919        );
4920
4921        // The load-bearing assertion. Expanding the pruned directory charges a step for
4922        // each of its forty children before discarding every one of them, so the work
4923        // separates a remembered descent that prunes from one that does not, where the
4924        // rows above cannot.
4925        assert!(
4926            response.work.rows_visited < 40,
4927            "the pruned subtree was expanded: {} steps for {} rows",
4928            response.work.rows_visited,
4929            response.work.rows_returned
4930        );
4931
4932        opened.close().expect("close");
4933    }
4934
4935    /// A name whose bytes are not UTF-8 is escaped and listed, not omitted.
4936    ///
4937    /// This fixture used to prove the opposite. While a portable name was optional these
4938    /// two entries were retained, counted in roll-ups, and absent from every page, and
4939    /// the page reported an omission count with escaped examples so the loss was at least
4940    /// visible. It also meant a lookup below such a directory had to answer `unknown`
4941    /// rather than `absent`, because the name asked for might have been in the invisible
4942    /// set.
4943    ///
4944    /// The encoding is total now, so the same fixture must show the opposite: both rows
4945    /// appear, ordered pages and roll-ups agree on the population, and absence is
4946    /// answerable. `x\xff` becomes `x%FF`; the valid prefix survives as text and only the
4947    /// undecodable byte is escaped.
4948    #[cfg(unix)]
4949    #[test]
4950    fn non_utf8_names_are_escaped_into_pages_rather_than_omitted() {
4951        use std::os::unix::ffi::OsStringExt;
4952
4953        let (_root, opened) = opened(Arc::new(TestControls::default()));
4954        let invalid = PathBuf::from(OsString::from_vec(vec![b'x', 0xff]));
4955        // A literal `%` beside an escaped byte is the pair that proves injectivity: if
4956        // `%` were left alone, a file actually named `y%FE` and this one would collide.
4957        let literal_percent = PathBuf::from("y%FE");
4958        opened
4959            .state
4960            .index
4961            .apply(&Observation::new(vec![
4962                Op::Upsert {
4963                    path: invalid.clone(),
4964                    kind: EntryKind::File,
4965                    attrs: crate::Attrs { size: 9, ..crate::Attrs::default() },
4966                },
4967                Op::Upsert {
4968                    path: literal_percent.clone(),
4969                    kind: EntryKind::File,
4970                    attrs: crate::Attrs { size: 1, ..crate::Attrs::default() },
4971                },
4972            ]))
4973            .expect("seed non-utf8 and literal-percent entries");
4974
4975        let response = opened
4976            .read(crate::ReadRequest {
4977                projections: vec![
4978                    crate::ReadProjection::Lookup { path: PathBuf::from("missing") },
4979                    crate::ReadProjection::Tree {
4980                        path: PathBuf::new(),
4981                        depth: crate::query::Bound::Limit(1),
4982                        include_ignored: true,
4983                        page: crate::PageRequest {
4984                            limit: crate::MAX_PAGE_ROWS,
4985                            max_work: crate::MAX_PAGE_WORK,
4986                        },
4987                    },
4988                    crate::ReadProjection::Flat {
4989                        selection: crate::query::EntrySelection::default(),
4990                        shape: crate::RowShape::Compact,
4991                        page: crate::PageRequest {
4992                            limit: crate::MAX_PAGE_ROWS,
4993                            max_work: crate::MAX_PAGE_WORK,
4994                        },
4995                    },
4996                    crate::ReadProjection::RollUp { path: PathBuf::new() },
4997                ],
4998                ..crate::ReadRequest::default()
4999            })
5000            .expect("portable read");
5001
5002        // Absence is answerable: nothing can be hiding in an unlistable set any more.
5003        assert!(matches!(
5004            response.results[0],
5005            crate::ProjectionResult::Lookup(crate::Knowledge::Absent)
5006        ));
5007
5008        let crate::ProjectionResult::Tree(crate::Knowledge::Present(tree)) = &response.results[1]
5009        else {
5010            panic!("tree page");
5011        };
5012        assert!(tree.complete);
5013        let names: Vec<_> =
5014            tree.rows.iter().map(|row| row.portable_path.as_str().to_owned()).collect();
5015        assert_eq!(names, vec!["x%FF".to_owned(), "y%25FE".to_owned()]);
5016
5017        let crate::ProjectionResult::Flat(flat) = &response.results[2] else {
5018            panic!("flat page");
5019        };
5020        let flat_names: Vec<_> =
5021            flat.rows.iter().map(|row| row.portable_path.as_str().to_owned()).collect();
5022        assert_eq!(flat_names, names, "ordered pages agree on one population");
5023
5024        // The population the pages return is the population the roll-up counts.
5025        assert!(matches!(
5026            &response.results[3],
5027            crate::ProjectionResult::RollUp(crate::Knowledge::Present(rollup))
5028                if rollup.all.files == 2
5029                    && rollup.all.files
5030                        == u64::try_from(flat.rows.len()).expect("row count fits u64")
5031        ));
5032
5033        opened
5034            .state
5035            .index
5036            .apply(&Observation::new(vec![
5037                Op::Remove { path: invalid },
5038                Op::Remove { path: literal_percent },
5039            ]))
5040            .expect("remove escaped entries");
5041        let after = opened
5042            .read(crate::ReadRequest {
5043                projections: vec![crate::ReadProjection::Tree {
5044                    path: PathBuf::new(),
5045                    depth: crate::query::Bound::Limit(1),
5046                    include_ignored: true,
5047                    page: crate::PageRequest {
5048                        limit: crate::MAX_PAGE_ROWS,
5049                        max_work: crate::MAX_PAGE_WORK,
5050                    },
5051                }],
5052                ..crate::ReadRequest::default()
5053            })
5054            .expect("read after removal");
5055        let crate::ProjectionResult::Tree(crate::Knowledge::Present(tree)) = &after.results[0]
5056        else {
5057            panic!("tree page");
5058        };
5059        assert!(tree.rows.is_empty(), "removal is symmetric for escaped names too");
5060        opened.close().expect("close");
5061    }
5062
5063    #[test]
5064    fn read_bounds_are_validated_before_any_continuation_is_consumed() {
5065        let (_root, opened) = opened(Arc::new(TestControls::default()));
5066        opened
5067            .state
5068            .index
5069            .apply(&Observation::new(vec![
5070                Op::Upsert {
5071                    path: PathBuf::from("a"),
5072                    kind: EntryKind::File,
5073                    attrs: crate::Attrs::default(),
5074                },
5075                Op::Upsert {
5076                    path: PathBuf::from("b"),
5077                    kind: EntryKind::File,
5078                    attrs: crate::Attrs::default(),
5079                },
5080            ]))
5081            .expect("seed entries");
5082        let first = opened
5083            .read(crate::ReadRequest {
5084                projections: vec![crate::ReadProjection::Flat {
5085                    selection: crate::query::EntrySelection::default(),
5086                    shape: crate::RowShape::Compact,
5087                    page: crate::PageRequest { limit: 1, max_work: 2 },
5088                }],
5089                ..crate::ReadRequest::default()
5090            })
5091            .expect("first page");
5092        let crate::ProjectionResult::Flat(page) = &first.results[0] else {
5093            panic!("flat page");
5094        };
5095        let continuation = page.next.expect("continuation");
5096
5097        assert!(matches!(
5098            opened.read(crate::ReadRequest {
5099                projections: vec![
5100                    crate::ReadProjection::Continue {
5101                        continuation,
5102                        page: crate::PageRequest { limit: 1, max_work: 2 },
5103                    },
5104                    crate::ReadProjection::Tree {
5105                        path: PathBuf::new(),
5106                        depth: crate::query::Bound::Limit(1),
5107                        include_ignored: true,
5108                        page: crate::PageRequest { limit: 0, max_work: 1 },
5109                    },
5110                ],
5111                ..crate::ReadRequest::default()
5112            }),
5113            Err(Error::PageRowLimit { attempted: 0, .. })
5114        ));
5115        opened
5116            .read(crate::ReadRequest {
5117                projections: vec![crate::ReadProjection::Continue {
5118                    continuation,
5119                    page: crate::PageRequest { limit: 1, max_work: 2 },
5120                }],
5121                ..crate::ReadRequest::default()
5122            })
5123            .expect("validation preserved continuation");
5124
5125        assert!(matches!(
5126            opened.read(crate::ReadRequest {
5127                projections: vec![
5128                    crate::ReadProjection::Diagnostics;
5129                    crate::MAX_READ_PROJECTIONS + 1
5130                ],
5131                ..crate::ReadRequest::default()
5132            }),
5133            Err(Error::ReadProjectionLimit { .. })
5134        ));
5135        assert!(matches!(
5136            opened.read(crate::ReadRequest {
5137                projections: vec![crate::ReadProjection::Aggregate {
5138                    selection: crate::query::EntrySelection::default(),
5139                    count_cap: 0,
5140                    max_work: 1,
5141                }],
5142                ..crate::ReadRequest::default()
5143            }),
5144            Err(Error::CountCapLimit { attempted: 0, .. })
5145        ));
5146        // Zero levels is its own rejection, not a row-bound one. It once reported
5147        // `PageRowLimit { attempted: 0 }`, which named a bound the caller had not set and
5148        // sent them to inspect `page.limit` instead of `depth`.
5149        assert!(matches!(
5150            opened.read(crate::ReadRequest {
5151                projections: vec![crate::ReadProjection::Tree {
5152                    path: PathBuf::new(),
5153                    depth: crate::query::Bound::Limit(0),
5154                    include_ignored: true,
5155                    page: crate::PageRequest { limit: 1, max_work: 1 },
5156                }],
5157                ..crate::ReadRequest::default()
5158            }),
5159            Err(Error::TreeDepthZero)
5160        ));
5161        let bounded_path = opened
5162            .read(crate::ReadRequest {
5163                projections: vec![crate::ReadProjection::Tree {
5164                    path: PathBuf::from("missing/deep"),
5165                    depth: crate::query::Bound::Limit(1),
5166                    include_ignored: true,
5167                    page: crate::PageRequest { limit: 1, max_work: 1 },
5168                }],
5169                ..crate::ReadRequest::default()
5170            })
5171            .expect("bounded path traversal");
5172        assert!(matches!(
5173            bounded_path.results[0],
5174            crate::ProjectionResult::Limit(crate::QueryLimit {
5175                projection: crate::LimitedProjection::Tree,
5176                rows_visited: 1,
5177                ..
5178            })
5179        ));
5180        assert!(matches!(
5181            opened.read(crate::ReadRequest {
5182                projections: vec![crate::ReadProjection::Report(crate::ReportRequest {
5183                    query: crate::query::Query {
5184                        views: vec![crate::query::ViewSpec::Summary; crate::MAX_REPORT_VIEWS + 1],
5185                        ..crate::query::Query::default()
5186                    },
5187                    now: std::time::UNIX_EPOCH,
5188                    max_work: crate::MAX_PAGE_WORK,
5189                })],
5190                ..crate::ReadRequest::default()
5191            }),
5192            Err(Error::ReportViewLimit { .. })
5193        ));
5194        opened.close().expect("close");
5195    }
5196
5197    #[test]
5198    fn a_coherent_read_cannot_straddle_a_commit() {
5199        let (_root, opened) = opened(Arc::new(TestControls::default()));
5200        let stop = Arc::new(AtomicBool::new(false));
5201        let writer_stop = Arc::clone(&stop);
5202        let writer_index = opened.state.index.clone();
5203        let writer = std::thread::spawn(move || {
5204            for round in 0..400_u64 {
5205                if writer_stop.load(Ordering::Relaxed) {
5206                    break;
5207                }
5208                writer_index
5209                    .apply(&Observation::new(vec![Op::Upsert {
5210                        path: PathBuf::from(format!("file-{round}")),
5211                        kind: EntryKind::File,
5212                        attrs: crate::Attrs { size: 1, ..crate::Attrs::default() },
5213                    }]))
5214                    .expect("writer commit");
5215            }
5216        });
5217
5218        for _ in 0..500 {
5219            let response = opened
5220                .read(crate::ReadRequest {
5221                    projections: vec![
5222                        crate::ReadProjection::Tree {
5223                            path: PathBuf::new(),
5224                            depth: crate::query::Bound::Limit(1),
5225                            include_ignored: true,
5226                            page: crate::PageRequest {
5227                                limit: crate::MAX_PAGE_ROWS,
5228                                max_work: crate::MAX_PAGE_WORK,
5229                            },
5230                        },
5231                        crate::ReadProjection::RollUp { path: PathBuf::new() },
5232                    ],
5233                    ..crate::ReadRequest::default()
5234                })
5235                .expect("coherent read");
5236            let crate::ProjectionResult::Tree(crate::Knowledge::Present(tree)) =
5237                &response.results[0]
5238            else {
5239                panic!("tree page");
5240            };
5241            let crate::ProjectionResult::RollUp(crate::Knowledge::Present(rollup)) =
5242                &response.results[1]
5243            else {
5244                panic!("root roll-up");
5245            };
5246            assert!(tree.next.is_none());
5247            assert_eq!(
5248                tree.rows.iter().filter(|row| row.kind == EntryKind::File).count() as u64,
5249                rollup.all.files
5250            );
5251        }
5252        stop.store(true, Ordering::Relaxed);
5253        writer.join().expect("writer");
5254        opened.close().expect("close");
5255    }
5256
5257    #[test]
5258    fn aggregates_distinguish_maintained_exact_totals_from_capped_counts() {
5259        let (_root, opened) = opened(Arc::new(TestControls::default()));
5260        opened
5261            .state
5262            .index
5263            .apply(&Observation::new(
5264                ["a", "b", "c"]
5265                    .into_iter()
5266                    .map(|path| Op::Upsert {
5267                        path: PathBuf::from(path),
5268                        kind: EntryKind::File,
5269                        attrs: crate::Attrs::default(),
5270                    })
5271                    .collect(),
5272            ))
5273            .expect("seed entries");
5274
5275        let response = opened
5276            .read(crate::ReadRequest {
5277                projections: vec![
5278                    crate::ReadProjection::Aggregate {
5279                        selection: crate::query::EntrySelection::default(),
5280                        count_cap: 1,
5281                        max_work: 1,
5282                    },
5283                    crate::ReadProjection::Aggregate {
5284                        selection: crate::query::EntrySelection {
5285                            query: crate::query::Selection {
5286                                kinds: vec![EntryKind::File],
5287                                ..crate::query::Selection::default()
5288                            },
5289                            ..crate::query::EntrySelection::default()
5290                        },
5291                        count_cap: 2,
5292                        max_work: 3,
5293                    },
5294                ],
5295                ..crate::ReadRequest::default()
5296            })
5297            .expect("aggregate read");
5298
5299        assert!(matches!(
5300            response.results[0],
5301            crate::ProjectionResult::Aggregate(crate::CountResult::Exact(3))
5302        ));
5303        assert!(matches!(
5304            response.results[1],
5305            crate::ProjectionResult::Aggregate(crate::CountResult::AtLeast(2))
5306        ));
5307        opened.close().expect("close");
5308    }
5309
5310    #[test]
5311    fn report_projection_matches_the_existing_query_and_fails_closed_at_its_work_bound() {
5312        let (_root, opened) = opened(Arc::new(TestControls::default()));
5313        opened
5314            .state
5315            .index
5316            .apply(&Observation::new(vec![
5317                Op::Upsert {
5318                    path: PathBuf::from("a"),
5319                    kind: EntryKind::File,
5320                    attrs: crate::Attrs { size: 3, ..crate::Attrs::default() },
5321                },
5322                Op::Upsert {
5323                    path: PathBuf::from("b"),
5324                    kind: EntryKind::File,
5325                    attrs: crate::Attrs { size: 5, ..crate::Attrs::default() },
5326                },
5327            ]))
5328            .expect("seed entries");
5329        let query = crate::query::Query {
5330            views: vec![crate::query::ViewSpec::Summary],
5331            ..crate::query::Query::default()
5332        };
5333
5334        let response = opened
5335            .read(crate::ReadRequest {
5336                projections: vec![crate::ReadProjection::Report(crate::ReportRequest {
5337                    query: query.clone(),
5338                    now: std::time::UNIX_EPOCH,
5339                    max_work: 1,
5340                })],
5341                ..crate::ReadRequest::default()
5342            })
5343            .expect("maintained report");
5344        let crate::ProjectionResult::Report(report) = &response.results[0] else {
5345            panic!("report projection");
5346        };
5347        let crate::query::Section::Summary(summary) = &report.sections[0] else {
5348            panic!("summary section");
5349        };
5350        assert_eq!((summary.files, summary.bytes), (2, 8));
5351
5352        let limited = opened
5353            .read(crate::ReadRequest {
5354                projections: vec![crate::ReadProjection::Report(crate::ReportRequest {
5355                    query: crate::query::Query {
5356                        selection: crate::query::Selection {
5357                            kinds: vec![EntryKind::File],
5358                            ..crate::query::Selection::default()
5359                        },
5360                        views: vec![crate::query::ViewSpec::Summary],
5361                        ..crate::query::Query::default()
5362                    },
5363                    now: std::time::UNIX_EPOCH,
5364                    max_work: 1,
5365                })],
5366                ..crate::ReadRequest::default()
5367            })
5368            .expect("bounded report");
5369        assert!(matches!(
5370            limited.results[0],
5371            crate::ProjectionResult::Limit(crate::QueryLimit {
5372                projection: crate::LimitedProjection::Report,
5373                ..
5374            })
5375        ));
5376        opened.close().expect("close");
5377    }
5378
5379    #[test]
5380    fn first_refusal_stops_expansion_and_commits_the_budget_state_with_prior_facts() {
5381        let root = tempfile::tempdir().expect("temp root");
5382        std::fs::create_dir(root.path().join("nested")).expect("fixture directory");
5383        std::fs::write(root.path().join("one"), b"1").expect("fixture");
5384        std::fs::write(root.path().join("two"), b"2").expect("fixture");
5385        std::fs::write(root.path().join("nested/deep"), b"deep").expect("deep fixture");
5386        let opened = open_fixture(
5387            root.path(),
5388            OpenOptions {
5389                batch_size: 64,
5390                budget: DiscoveryBudget { max_files: Some(1) },
5391                ..OpenOptions::default()
5392            },
5393        )
5394        .expect("opened root");
5395
5396        let state = wait_until_settled(&opened);
5397        assert_eq!(state.phase, crate::LifecyclePhase::Stopped);
5398        assert_eq!(state.coverage, crate::Coverage::Partial(crate::CoverageReason::Budget));
5399        assert_eq!(state.progress.files_retained, 1);
5400        assert_eq!(opened.state.index.total().expect("total").files, 1);
5401        assert_eq!(opened.state.index.kind(Path::new("nested/deep")).expect("deep lookup"), None);
5402        assert_eq!(
5403            opened.state.index.directory_complete(Path::new("")).expect("root completeness"),
5404            Some(false)
5405        );
5406        let partial = opened.state.index.snapshot().expect("partial snapshot image");
5407        assert!(matches!(
5408            crate::snapshot::save(&partial, &root.path().join("partial.fdu")),
5409            Err(Error::Snapshot(_))
5410        ));
5411        assert!(matches!(
5412            opened.prioritize(&[PathBuf::from("nested")]),
5413            Err(Error::OpenedIndexStopped)
5414        ));
5415
5416        let terminal = opened
5417            .state
5418            .index
5419            .since(crate::Clock::ZERO)
5420            .expect("journal")
5421            .commits
5422            .into_iter()
5423            .find(|commit| {
5424                commit.state.iter().any(|transition| {
5425                    matches!(
5426                        transition,
5427                        crate::StateTransition::IndexState {
5428                            current: crate::IndexState {
5429                                coverage: crate::Coverage::Partial(crate::CoverageReason::Budget),
5430                                ..
5431                            },
5432                            ..
5433                        }
5434                    )
5435                })
5436            })
5437            .expect("budget commit");
5438        assert!(terminal.changes.iter().any(|change| matches!(
5439            change,
5440            crate::EffectiveChange::Inserted { kind: EntryKind::File, .. }
5441        )));
5442        opened.close().expect("close");
5443    }
5444
5445    #[test]
5446    fn refresh_receipt_counts_verified_no_op_work_without_a_fact_commit() {
5447        let root = tempfile::tempdir().expect("temp root");
5448        std::fs::write(root.path().join("stable.txt"), b"stable").expect("fixture");
5449        let opened = open_fixture(root.path(), OpenOptions::default()).expect("opened root");
5450        assert_eq!(wait_until_settled(&opened).phase, crate::LifecyclePhase::Ready);
5451
5452        let result = opened.refresh(&[PathBuf::from("stable.txt")]).expect("refresh");
5453
5454        assert_eq!(result.accepted, vec![PathBuf::from("stable.txt")]);
5455        assert!(result.rejected.is_empty());
5456        assert_eq!(result.work.observations, 1, "the verified observation is still work");
5457        assert_eq!(result.work.unchanged, 1, "the matching fact is reported as unchanged");
5458        assert_eq!(result.work.stale, 0);
5459        opened.close().expect("close");
5460    }
5461
5462    #[test]
5463    fn refresh_can_fill_remaining_file_budget_without_exceeding_it() {
5464        let root = tempfile::tempdir().expect("temp root");
5465        std::fs::write(root.path().join("one"), b"1").expect("fixture");
5466        let opened = open_fixture(
5467            root.path(),
5468            OpenOptions {
5469                budget: DiscoveryBudget { max_files: Some(2) },
5470                ..OpenOptions::default()
5471            },
5472        )
5473        .expect("opened root");
5474        assert_eq!(wait_until_settled(&opened).progress.files_retained, 1);
5475        std::fs::write(root.path().join("two"), b"2").expect("new file");
5476
5477        let result = opened.refresh(&[PathBuf::from("two")]).expect("refresh");
5478
5479        assert_eq!(result.accepted, vec![PathBuf::from("two")]);
5480        assert!(result.rejected.is_empty());
5481        assert_eq!(opened.state.index.total().expect("total").files, 2);
5482        assert_eq!(result.state.phase, crate::LifecyclePhase::Ready);
5483        assert_eq!(result.state.progress.files_retained, 2);
5484        opened.close().expect("close");
5485    }
5486
5487    #[test]
5488    fn refresh_classifies_paths_and_collapses_overlapping_walks() {
5489        let root = tempfile::tempdir().expect("temp root");
5490        std::fs::create_dir_all(root.path().join("visible/nested")).expect("fixture directories");
5491        std::fs::write(root.path().join("visible/nested/leaf"), b"leaf").expect("fixture");
5492        std::fs::create_dir(root.path().join(".hidden")).expect("hidden directory");
5493        std::fs::write(root.path().join(".hidden/leaf"), b"hidden").expect("hidden fixture");
5494        let opened = open_fixture(
5495            root.path(),
5496            OpenOptions {
5497                hidden: Some(Arc::new(crate::HiddenPolicy::prune_hidden::<[&str; 0], &str>([]))),
5498                ..OpenOptions::default()
5499            },
5500        )
5501        .expect("opened root");
5502        assert_eq!(wait_until_settled(&opened).phase, crate::LifecyclePhase::Ready);
5503
5504        let result = opened
5505            .refresh(&[
5506                PathBuf::from("visible/nested"),
5507                PathBuf::from("visible"),
5508                PathBuf::from("visible/nested"),
5509                PathBuf::from("../escape"),
5510                PathBuf::from(".hidden/leaf"),
5511            ])
5512            .expect("refresh");
5513
5514        assert_eq!(
5515            result.accepted,
5516            vec![PathBuf::from("visible"), PathBuf::from("visible/nested")]
5517        );
5518        assert_eq!(
5519            result.rejected,
5520            vec![
5521                crate::RejectedRefreshPath {
5522                    path: PathBuf::from("../escape"),
5523                    reason: crate::RefreshRejection::OutsideRoot,
5524                },
5525                crate::RejectedRefreshPath {
5526                    path: PathBuf::from(".hidden/leaf"),
5527                    reason: crate::RefreshRejection::NotAdmitted,
5528                },
5529            ]
5530        );
5531        assert_eq!(result.work.directories_read, 2, "the descendant was not walked twice");
5532        opened.close().expect("close");
5533    }
5534
5535    #[test]
5536    fn refresh_widens_through_a_replaced_ancestor_and_reports_exact_commits() {
5537        let root = tempfile::tempdir().expect("temp root");
5538        std::fs::create_dir_all(root.path().join("parent/child")).expect("fixture directories");
5539        std::fs::write(root.path().join("parent/child/leaf"), b"leaf").expect("fixture");
5540        let opened = open_fixture(root.path(), OpenOptions::default()).expect("opened root");
5541        assert_eq!(wait_until_settled(&opened).phase, crate::LifecyclePhase::Ready);
5542        std::fs::remove_dir_all(root.path().join("parent")).expect("remove old subtree");
5543        std::fs::write(root.path().join("parent"), b"replacement").expect("replacement file");
5544        let before = current_version(&opened);
5545
5546        let result =
5547            opened.refresh(&[PathBuf::from("parent/child/leaf")]).expect("refresh widened path");
5548        let poll = opened
5549            .changes(crate::ChangeRequest { after: before, timeout: std::time::Duration::ZERO })
5550            .expect("refresh commits");
5551
5552        assert_eq!(result.after, before);
5553        assert_eq!(result.version, poll.version);
5554        assert_eq!(result.accepted, vec![PathBuf::from("parent/child/leaf")]);
5555        assert_eq!(
5556            opened.state.index.kind(Path::new("parent")).expect("kind"),
5557            Some(EntryKind::File)
5558        );
5559        assert_eq!(
5560            opened.state.index.kind(Path::new("parent/child/leaf")).expect("removed child"),
5561            None
5562        );
5563        let crate::ChangeOutcome::Changes { commits, impact } = poll.outcome else {
5564            panic!("refresh must advance the journal");
5565        };
5566        assert!(!commits.is_empty());
5567        assert_eq!(impact, result.impact);
5568        assert!(commits.iter().all(|commit| {
5569            commit.clock.0 > result.after.sequence.0 && commit.clock.0 <= result.version.sequence.0
5570        }));
5571        opened.close().expect("close");
5572    }
5573
5574    #[cfg(unix)]
5575    #[test]
5576    fn refresh_rejects_symlink_shadowed_ancestry_without_aborting_other_paths() {
5577        use std::os::unix::fs::symlink;
5578
5579        let root = tempfile::tempdir().expect("temp root");
5580        let outside = tempfile::tempdir().expect("outside root");
5581        std::fs::create_dir_all(root.path().join("shadow/child")).expect("baseline ancestry");
5582        std::fs::write(root.path().join("good.txt"), b"before").expect("baseline file");
5583        let opened = open_fixture(root.path(), OpenOptions::default()).expect("opened root");
5584        assert_eq!(wait_until_settled(&opened).phase, crate::LifecyclePhase::Ready);
5585
5586        std::fs::remove_dir_all(root.path().join("shadow")).expect("remove ancestry");
5587        symlink(outside.path(), root.path().join("shadow")).expect("shadow with symlink");
5588        std::fs::write(root.path().join("good.txt"), b"after and larger").expect("mutate file");
5589
5590        let result = opened
5591            .refresh(&[PathBuf::from("shadow/child/leaf"), PathBuf::from("good.txt")])
5592            .expect("one unsafe path is a rejection, not a batch error");
5593
5594        assert_eq!(result.accepted, vec![PathBuf::from("good.txt")]);
5595        assert_eq!(result.rejected.len(), 1);
5596        assert_eq!(result.rejected[0].path, Path::new("shadow/child/leaf"));
5597        assert_eq!(result.rejected[0].reason, crate::RefreshRejection::UnsafeAncestry);
5598        assert_eq!(
5599            opened.state.index.attrs(Path::new("good.txt")).expect("attrs").expect("retained").size,
5600            16
5601        );
5602        opened.close().expect("close");
5603    }
5604
5605    #[test]
5606    fn refresh_refusal_is_atomic_with_the_shared_file_budget() {
5607        let root = tempfile::tempdir().expect("temp root");
5608        std::fs::write(root.path().join("one"), b"1").expect("fixture");
5609        let opened = open_fixture(
5610            root.path(),
5611            OpenOptions {
5612                budget: DiscoveryBudget { max_files: Some(2) },
5613                ..OpenOptions::default()
5614            },
5615        )
5616        .expect("opened root");
5617        assert_eq!(wait_until_settled(&opened).progress.files_retained, 1);
5618        std::fs::write(root.path().join("two"), b"2").expect("new file");
5619        std::fs::write(root.path().join("three"), b"3").expect("new file");
5620
5621        let result = opened
5622            .refresh(&[PathBuf::from("two"), PathBuf::from("three")])
5623            .expect("bounded refresh");
5624
5625        assert_eq!(result.accepted.len(), 2);
5626        assert!(result.rejected.is_empty());
5627        assert_eq!(result.work.observations, 2);
5628        assert_eq!(result.work.resource_refused, 1);
5629        assert_eq!(opened.state.index.total().expect("total").files, 2);
5630        assert_eq!(result.state.progress.files_retained, 2);
5631        assert_eq!(result.state.phase, crate::LifecyclePhase::Stopped);
5632        assert_eq!(result.state.coverage, crate::Coverage::Partial(crate::CoverageReason::Budget));
5633        assert_eq!(result.issues.len(), 1);
5634        assert_eq!(result.issues[0].kind, crate::IssueKind::ResourceBudget);
5635
5636        std::fs::write(root.path().join("four"), b"4").expect("later file");
5637        let stopped = opened.refresh(&[PathBuf::from("four")]).expect("stopped refresh receipt");
5638        assert!(stopped.accepted.is_empty());
5639        assert_eq!(
5640            stopped.rejected,
5641            vec![crate::RejectedRefreshPath {
5642                path: PathBuf::from("four"),
5643                reason: crate::RefreshRejection::ResourceBudget,
5644            }]
5645        );
5646        assert_eq!(stopped.work.entries_visited, 1, "the refusal reports its probe");
5647        assert_eq!(stopped.work.files_visited, 1);
5648        assert_eq!(stopped.work.bytes_visited, 1);
5649        assert_eq!(opened.state.index.total().expect("bounded total").files, 2);
5650        opened.close().expect("close");
5651    }
5652
5653    #[test]
5654    fn concurrent_discovery_and_refresh_share_one_atomic_file_budget() {
5655        let root = tempfile::tempdir().expect("temp root");
5656        std::fs::write(root.path().join("from-discovery"), b"discovery").expect("fixture");
5657        let controls = Arc::new(TestControls::default());
5658        controls.gate(TestPoint::BeforeDiscovery).arm();
5659        let opened = OpenedIndex::open_for_test(
5660            root.path(),
5661            OpenOptions {
5662                budget: DiscoveryBudget { max_files: Some(1) },
5663                ..OpenOptions::default()
5664            },
5665            Arc::clone(&controls),
5666        )
5667        .expect("opened root");
5668        controls.gate(TestPoint::BeforeDiscovery).wait_reached();
5669        std::fs::write(root.path().join("from-refresh"), b"refresh").expect("new file");
5670
5671        let refreshed =
5672            opened.refresh(&[PathBuf::from("from-refresh")]).expect("refresh during discovery");
5673        assert_eq!(refreshed.accepted, vec![PathBuf::from("from-refresh")]);
5674        assert_eq!(opened.state.index.total().expect("after refresh").files, 1);
5675        controls.gate(TestPoint::BeforeDiscovery).release();
5676
5677        let state = wait_until_settled(&opened);
5678        assert_eq!(opened.state.index.total().expect("bounded total").files, 1);
5679        assert_eq!(state.progress.files_retained, 1);
5680        assert_eq!(state.phase, crate::LifecyclePhase::Stopped);
5681        assert_eq!(state.coverage, crate::Coverage::Partial(crate::CoverageReason::Budget));
5682        opened.close().expect("close");
5683    }
5684
5685    /// A directory another producer removed while it waited in the frontier is stale work.
5686    ///
5687    /// Discovery queues `sub` from the root listing. While it waits -- minutes, on a wide
5688    /// breadth-first walk -- `sub` is deleted, a refresh commits the removal, and `sub` is
5689    /// recreated. Discovery then lists the new directory and commits it beneath a parent the
5690    /// index no longer holds. That rejection used to end discovery as `Failed`, so
5691    /// observation never started and `close()` reported a worker failure. Both commit shapes
5692    /// are covered: with a batch of one the first flush fails on the child's ancestry, and
5693    /// with the default batch the final commit fails naming the directory complete.
5694    #[test]
5695    fn a_refresh_racing_discovery_leaves_the_queued_directory_as_stale_work() {
5696        for batch_size in [1, OpenOptions::default().batch_size] {
5697            let controls = Arc::new(TestControls::default());
5698            controls.gate(TestPoint::AfterRootDirectory).arm();
5699            let root = tempfile::tempdir().expect("temp root");
5700            std::fs::create_dir(root.path().join("sub")).expect("sub");
5701            std::fs::write(root.path().join("sub/inner.txt"), b"x").expect("fixture");
5702            let opened = OpenedIndex::open_for_test(
5703                root.path(),
5704                OpenOptions { batch_size, ..OpenOptions::default() },
5705                Arc::clone(&controls),
5706            )
5707            .expect("open");
5708            controls.gate(TestPoint::AfterRootDirectory).wait_reached();
5709            assert_eq!(
5710                opened.state.index.kind(Path::new("sub")).expect("lookup"),
5711                Some(EntryKind::Dir),
5712                "the root listing queued `sub`"
5713            );
5714
5715            std::fs::remove_dir_all(root.path().join("sub")).expect("remove sub");
5716            let refreshed = opened.refresh(&[PathBuf::from("sub")]).expect("refresh");
5717            assert_eq!(refreshed.accepted, vec![PathBuf::from("sub")]);
5718            std::fs::create_dir(root.path().join("sub")).expect("recreate sub");
5719            std::fs::write(root.path().join("sub/again.txt"), b"y").expect("fixture");
5720            controls.gate(TestPoint::AfterRootDirectory).release();
5721
5722            let state = wait_until_settled(&opened);
5723            assert_eq!(state.phase, crate::LifecyclePhase::Ready, "batch size {batch_size}");
5724            assert_eq!(state.coverage, crate::Coverage::Complete, "batch size {batch_size}");
5725            assert_eq!(state.issues.retained, 0, "batch size {batch_size}");
5726            // The refresh's verified removal stands. The recreated directory belongs to the
5727            // next producer that verifies the path, not to a frontier entry older than it.
5728            assert_eq!(opened.state.index.kind(Path::new("sub")).expect("lookup"), None);
5729            opened.close().unwrap_or_else(|error| panic!("batch size {batch_size}: {error}"));
5730        }
5731    }
5732
5733    /// A directory removed or replaced between its parent's listing and its own is stale
5734    /// work, not an inaccessible boundary that outlives every later verification.
5735    #[test]
5736    fn a_directory_that_vanishes_during_discovery_is_not_inaccessible() {
5737        for replace_with_file in [false, true] {
5738            let controls = Arc::new(TestControls::default());
5739            controls.gate(TestPoint::AfterRootDirectory).arm();
5740            let root = tempfile::tempdir().expect("temp root");
5741            std::fs::create_dir(root.path().join("sub")).expect("sub");
5742            std::fs::write(root.path().join("sub/inner.txt"), b"x").expect("fixture");
5743            std::fs::write(root.path().join("keep.txt"), b"k").expect("fixture");
5744            let opened = OpenedIndex::open_for_test(
5745                root.path(),
5746                OpenOptions::default(),
5747                Arc::clone(&controls),
5748            )
5749            .expect("open");
5750            controls.gate(TestPoint::AfterRootDirectory).wait_reached();
5751            std::fs::remove_dir_all(root.path().join("sub")).expect("remove sub");
5752            if replace_with_file {
5753                std::fs::write(root.path().join("sub"), b"now a file").expect("replacement");
5754            }
5755            controls.gate(TestPoint::AfterRootDirectory).release();
5756
5757            let state = wait_until_settled(&opened);
5758            let case = if replace_with_file { "replaced by a file" } else { "removed" };
5759            assert_eq!(state.phase, crate::LifecyclePhase::Ready, "{case}");
5760            assert_eq!(state.coverage, crate::Coverage::Complete, "{case}");
5761            assert_eq!(state.freshness, crate::Freshness::Fresh, "{case}");
5762            assert_eq!(state.issues.retained, 0, "{case}");
5763            opened.close().unwrap_or_else(|error| panic!("{case}: {error}"));
5764        }
5765    }
5766
5767    /// A budget stop is terminal even when it lands in the middle of discovery.
5768    ///
5769    /// A refresh trips the shared budget while discovery still has a file-less directory
5770    /// queued. Nothing that directory holds is refused, so discovery used to run on to
5771    /// `Finish`, which set the phase `Ready` unconditionally: `prioritize` succeeded again and
5772    /// an observer could have reached `Watching` with budget-partial coverage.
5773    #[test]
5774    fn a_budget_stop_during_discovery_stays_terminal() {
5775        let controls = Arc::new(TestControls::default());
5776        controls.gate(TestPoint::AfterRootDirectory).arm();
5777        let root = tempfile::tempdir().expect("temp root");
5778        std::fs::write(root.path().join("a.txt"), b"a").expect("fixture");
5779        std::fs::create_dir(root.path().join("emptydir")).expect("fixture");
5780        let opened = OpenedIndex::open_for_test(
5781            root.path(),
5782            OpenOptions {
5783                budget: DiscoveryBudget { max_files: Some(1) },
5784                ..OpenOptions::default()
5785            },
5786            Arc::clone(&controls),
5787        )
5788        .expect("open");
5789        controls.gate(TestPoint::AfterRootDirectory).wait_reached();
5790        std::fs::write(root.path().join("b.txt"), b"b").expect("over budget");
5791        let refreshed = opened.refresh(&[PathBuf::from("b.txt")]).expect("refresh");
5792        assert_eq!(refreshed.work.resource_refused, 1);
5793        assert_eq!(refreshed.state.phase, crate::LifecyclePhase::Stopped);
5794        controls.gate(TestPoint::AfterRootDirectory).release();
5795
5796        wait_for_worker_exit(&opened, "discovery");
5797        let state = opened.state.index.state().expect("state");
5798        assert_eq!(state.phase, crate::LifecyclePhase::Stopped);
5799        assert_eq!(state.coverage, crate::Coverage::Partial(crate::CoverageReason::Budget));
5800        assert_eq!(
5801            opened.state.index.directory_complete(Path::new("emptydir")).expect("lookup"),
5802            Some(false),
5803            "a listing that arrives after the stop must not land"
5804        );
5805        assert!(matches!(
5806            opened.prioritize(&[PathBuf::from("emptydir")]),
5807            Err(Error::OpenedIndexStopped)
5808        ));
5809        opened.close().expect("close");
5810    }
5811
5812    /// A `.gitignore` in `a` over the default line limit, listed between `-early/` and
5813    /// `zzz.txt`, with one entry per batch so the control lands in a batch of its own.
5814    fn tree_with_a_guarded_control() -> tempfile::TempDir {
5815        let root = tempfile::tempdir().expect("temp root");
5816        let a = root.path().join("a");
5817        std::fs::create_dir_all(a.join("-early")).expect("fixture");
5818        std::fs::write(a.join("-early").join("leaf.txt"), b"l").expect("fixture");
5819        let mut line = b"*.txt\n".to_vec();
5820        line.extend(std::iter::repeat_n(b'x', crate::control::DEFAULT_CONTROL_LINE_LIMIT + 1));
5821        std::fs::write(a.join(crate::control::CONTROL_FILE_NAME), &line).expect("control");
5822        std::fs::write(a.join("zzz.txt"), b"z").expect("fixture");
5823        std::fs::create_dir(root.path().join("b")).expect("fixture");
5824        std::fs::write(root.path().join("b/kept.txt"), b"k").expect("fixture");
5825        root
5826    }
5827
5828    fn diagnostics(opened: &OpenedIndex) -> crate::ReadDiagnostics {
5829        let response = opened
5830            .read(crate::ReadRequest {
5831                projections: vec![crate::ReadProjection::Diagnostics],
5832                ..crate::ReadRequest::default()
5833            })
5834            .expect("read");
5835        let [crate::ProjectionResult::Diagnostics(diagnostics)] = response.results.as_slice()
5836        else {
5837            panic!("one diagnostics result: {:?}", response.results);
5838        };
5839        diagnostics.clone()
5840    }
5841
5842    /// A control over a limit is refused, and the listing that carried it commits.
5843    ///
5844    /// The directory's other entries land and it completes, so structural coverage stays
5845    /// complete and no issue is retained; the refusal is a control-coverage fact naming the
5846    /// file (fdu-1onj). The listing used to be refused with its control: the directory
5847    /// stayed incomplete, the rest of its entries were dropped, and coverage said
5848    /// `Inaccessible` about a file that was perfectly readable.
5849    #[test]
5850    fn a_control_over_a_bound_is_refused_while_its_listing_commits() {
5851        let controls = Arc::new(TestControls::default());
5852        controls.use_deterministic_discovery_order();
5853        let root = tree_with_a_guarded_control();
5854        let opened = OpenedIndex::open_for_test(
5855            root.path(),
5856            OpenOptions { batch_size: 1, ..OpenOptions::default() },
5857            Arc::clone(&controls),
5858        )
5859        .expect("open");
5860        let state = wait_until_settled(&opened);
5861        assert_eq!(state.phase, crate::LifecyclePhase::Ready);
5862        assert_eq!(state.coverage, crate::Coverage::Complete);
5863        assert_eq!(state.issues, crate::IssueSummary::default());
5864
5865        let index = &opened.state.index;
5866        assert_eq!(index.directory_complete(Path::new("a")).expect("lookup"), Some(true));
5867        assert_eq!(index.directory_complete(Path::new("a/-early")).expect("lookup"), Some(true));
5868        assert_eq!(index.kind(Path::new("a/zzz.txt")).expect("lookup"), Some(EntryKind::File));
5869        assert_eq!(index.kind(Path::new("b/kept.txt")).expect("lookup"), Some(EntryKind::File));
5870        assert_eq!(
5871            index
5872                .snapshot()
5873                .expect("snapshot")
5874                .is_ignored(Path::new("a/zzz.txt"))
5875                .expect("observed"),
5876            None,
5877            "the refused file could have governed this entry"
5878        );
5879        assert_eq!(
5880            diagnostics(&opened).controls,
5881            crate::control::ControlObservation {
5882                limits: crate::control::ControlLimits::default(),
5883                applied: 0,
5884                rules: 0,
5885                refused: 1,
5886                refusals: vec![crate::control::RefusedControl {
5887                    path: PathBuf::from("a/.gitignore"),
5888                    reason: crate::control::ControlRefusalReason::LineLimit,
5889                }],
5890            }
5891        );
5892        opened.close().expect("close");
5893    }
5894
5895    /// An opened root with no line limit applies the file the default limit refuses, and
5896    /// its scope says which limits it ran under.
5897    #[test]
5898    fn an_opened_root_without_a_line_limit_applies_what_the_default_limit_refuses() {
5899        let root = tree_with_a_guarded_control();
5900        let limits = crate::control::ControlLimits {
5901            line_limit: None,
5902            ..crate::control::ControlLimits::default()
5903        };
5904        let options = OpenOptions { control_limits: limits, ..OpenOptions::default() };
5905        let opened = open_fixture(root.path(), options).expect("open");
5906        wait_until_settled(&opened);
5907
5908        let diagnostics = diagnostics(&opened);
5909        assert_eq!((diagnostics.controls.limits, diagnostics.controls.applied), (limits, 1));
5910        assert_eq!(diagnostics.controls.refused, 0);
5911        assert_eq!(
5912            diagnostics.scope,
5913            ScanConfig { control_limits: limits, ..ScanConfig::default() }.scope()
5914        );
5915        let snapshot = opened.state.index.snapshot().expect("snapshot");
5916        assert_eq!(snapshot.is_ignored(Path::new("a/zzz.txt")).expect("observed"), Some(true));
5917        opened.close().expect("close");
5918    }
5919
5920    /// Deliver `hints` to a scripted observer, then block until a marker event sent after
5921    /// them has been applied, so every earlier event has been too.
5922    #[cfg(feature = "watch")]
5923    fn observe_then_settle(
5924        opened: &OpenedIndex,
5925        controls: &TestControls,
5926        root: &Path,
5927        hints: &str,
5928    ) {
5929        static MARKERS: AtomicUsize = AtomicUsize::new(0);
5930        let marker = format!("settled-{}.txt", MARKERS.fetch_add(1, Ordering::Relaxed));
5931        std::fs::write(root.join(&marker), b"marker").expect("marker");
5932        controls.send_observation_hints(&format!("{hints}create\t{marker}\n"));
5933        let deadline = std::time::Instant::now() + TEST_GATE_TIMEOUT;
5934        while opened.state.index.kind(Path::new(&marker)).expect("lookup") != Some(EntryKind::File)
5935        {
5936            assert!(std::time::Instant::now() < deadline, "{marker} was not applied");
5937            std::thread::yield_now();
5938        }
5939    }
5940
5941    /// A watched root over a control bound reaches `Watching`, and later events and
5942    /// refreshes over the refused directory degrade instead of failing the observer.
5943    ///
5944    /// The observation handoff's full reconciliation reads the same refused control, and
5945    /// used to fail with the error discovery had already walked past, ending the root
5946    /// `Failed` after an extra full walk. In steady state, the first event touching the
5947    /// control-bearing directory failed it the same way.
5948    #[cfg(feature = "watch")]
5949    #[test]
5950    fn a_watched_root_over_a_control_bound_keeps_watching_through_events() {
5951        let root = tree_with_a_guarded_control();
5952        let scripts = tempfile::tempdir().expect("script root");
5953        let script = scripts.path().join("events.script");
5954        std::fs::write(&script, b"").expect("script");
5955        let controls = Arc::new(TestControls::default());
5956        let opened = OpenedIndex::open_for_test(
5957            root.path(),
5958            scripted_options(&script),
5959            Arc::clone(&controls),
5960        )
5961        .expect("open");
5962        let state = wait_until_phase(&opened, crate::LifecyclePhase::Watching);
5963        assert_eq!(state.coverage, crate::Coverage::Complete);
5964        assert_eq!(diagnostics(&opened).controls.refused, 1);
5965
5966        // An edit that keeps the file over the guard, and a sibling created beside it.
5967        let control = root.path().join("a/.gitignore");
5968        let mut grown = std::fs::read(&control).expect("control");
5969        grown.extend_from_slice(b"\n*.md\n");
5970        std::fs::write(&control, grown).expect("grown control");
5971        std::fs::write(root.path().join("a/new.txt"), b"new").expect("fixture");
5972        let hints = "modify\ta/.gitignore\ncreate\ta/new.txt\n";
5973        observe_then_settle(&opened, &controls, root.path(), hints);
5974        let watching = crate::LifecyclePhase::Watching;
5975        assert_eq!(opened.state.index.state().expect("state").phase, watching);
5976        assert_eq!(
5977            opened.state.index.kind(Path::new("a/new.txt")).expect("lookup"),
5978            Some(EntryKind::File)
5979        );
5980        assert_eq!(diagnostics(&opened).controls.refused, 1);
5981
5982        // A refresh over the refused directory commits too.
5983        let refreshed = opened.refresh(&[PathBuf::from("a")]).expect("refresh over the refusal");
5984        assert_eq!(refreshed.state.phase, watching);
5985
5986        // Rules that fit lift the refusal and apply.
5987        std::fs::write(&control, b"*.txt\n").expect("fitting control");
5988        observe_then_settle(&opened, &controls, root.path(), "modify\ta/.gitignore\n");
5989        let coverage = diagnostics(&opened).controls;
5990        assert_eq!((coverage.applied, coverage.refused), (1, 0));
5991        let snapshot = opened.state.index.snapshot().expect("snapshot");
5992        assert_eq!(snapshot.is_ignored(Path::new("a/new.txt")).expect("observed"), Some(true));
5993        assert_eq!(opened.state.index.state().expect("state").phase, watching);
5994        opened.close().expect("close");
5995    }
5996
5997    /// A directory created after discovery is complete once a refresh has listed it.
5998    #[test]
5999    fn a_directory_created_after_discovery_is_complete_once_refreshed() {
6000        let root = tempfile::tempdir().expect("temp root");
6001        std::fs::write(root.path().join("before.txt"), b"b").expect("fixture");
6002        let opened = open_fixture(root.path(), OpenOptions::default()).expect("open");
6003        let settled = wait_until_settled(&opened);
6004        assert_eq!(settled.coverage, crate::Coverage::Complete);
6005
6006        std::fs::create_dir_all(root.path().join("later/deeper")).expect("fixture");
6007        std::fs::write(root.path().join("later/deeper/inner.txt"), b"i").expect("fixture");
6008        let receipt = opened.refresh(&[PathBuf::from("later")]).expect("refresh");
6009        assert!(receipt.issues.is_empty(), "{:?}", receipt.issues);
6010
6011        let index = &opened.state.index;
6012        assert_eq!(index.directory_complete(Path::new("later")).expect("lookup"), Some(true));
6013        assert_eq!(
6014            index.directory_complete(Path::new("later/deeper")).expect("lookup"),
6015            Some(true)
6016        );
6017        let response = opened
6018            .read(crate::ReadRequest {
6019                projections: vec![
6020                    crate::ReadProjection::Lookup { path: PathBuf::from("later/missing") },
6021                    crate::ReadProjection::Lookup { path: PathBuf::from("later/deeper/missing") },
6022                ],
6023                ..crate::ReadRequest::default()
6024            })
6025            .expect("read");
6026        assert_eq!(response.state.coverage, crate::Coverage::Complete);
6027        for result in &response.results {
6028            assert!(
6029                matches!(result, crate::ProjectionResult::Lookup(crate::Knowledge::Absent)),
6030                "{result:?}"
6031            );
6032        }
6033        let completed: Vec<_> = index
6034            .since(receipt.after.sequence)
6035            .expect("journal")
6036            .commits
6037            .iter()
6038            .flat_map(|commit| commit.state.iter())
6039            .filter_map(|transition| match transition {
6040                crate::StateTransition::DirectoryComplete { path } => Some(path.clone()),
6041                _ => None,
6042            })
6043            .collect();
6044        assert_eq!(completed, [PathBuf::from("later"), PathBuf::from("later/deeper")]);
6045        assert_eq!(
6046            receipt.state.progress.directories_complete,
6047            settled.progress.directories_complete + 2
6048        );
6049        opened.close().expect("close");
6050    }
6051
6052    #[test]
6053    fn refresh_rejects_an_unbounded_input_before_filesystem_work() {
6054        let root = tempfile::tempdir().expect("temp root");
6055        std::fs::write(root.path().join("same"), b"same").expect("fixture");
6056        let opened = open_fixture(root.path(), OpenOptions::default()).expect("opened root");
6057        assert_eq!(wait_until_settled(&opened).phase, crate::LifecyclePhase::Ready);
6058        let paths = vec![PathBuf::from("same"); MAX_REFRESH_PATHS + 2];
6059
6060        assert!(matches!(
6061            opened.refresh(&paths),
6062            Err(Error::RefreshPathLimit {
6063                attempted,
6064                limit: MAX_REFRESH_PATHS,
6065            }) if attempted == MAX_REFRESH_PATHS + 2
6066        ));
6067        opened.close().expect("close");
6068    }
6069
6070    #[test]
6071    fn refresh_rejects_stale_preparation_and_counts_the_lost_race() {
6072        let controls = Arc::new(TestControls::default());
6073        let (root, opened) = opened(Arc::clone(&controls));
6074        std::fs::write(root.path().join("race"), b"filesystem").expect("fixture");
6075        controls.gate(TestPoint::AfterRefreshVerification).arm();
6076        let refresher = opened.clone();
6077        let refresh = thread::spawn(move || refresher.refresh(&[PathBuf::from("race")]));
6078        controls.gate(TestPoint::AfterRefreshVerification).wait_reached();
6079        let concurrent = Observation::new(vec![
6080            Op::Upsert {
6081                path: PathBuf::from("race"),
6082                kind: EntryKind::File,
6083                attrs: crate::Attrs { size: 99, ..crate::Attrs::default() },
6084            },
6085            Op::Upsert {
6086                path: PathBuf::from("other"),
6087                kind: EntryKind::File,
6088                attrs: crate::Attrs { size: 5, ..crate::Attrs::default() },
6089            },
6090        ]);
6091        apply_and_notify(&opened, &concurrent);
6092        controls.gate(TestPoint::AfterRefreshVerification).release();
6093
6094        let result = refresh.join().expect("refresh thread").expect("refresh receipt");
6095        assert_eq!(result.work.observations, 1);
6096        assert_eq!(result.work.stale, 1);
6097        assert!(
6098            result.impact.dirty_paths.contains(&PathBuf::from("other")),
6099            "advancing to the receipt version must cover a concurrent producer"
6100        );
6101        assert_eq!(
6102            opened.state.index.attrs(Path::new("race")).expect("attrs").expect("retained").size,
6103            99
6104        );
6105        opened.close().expect("close");
6106    }
6107
6108    #[test]
6109    fn refresh_receipt_names_a_pass_retired_by_newer_verification() {
6110        let controls = Arc::new(TestControls::default());
6111        let (root, opened) = opened(Arc::clone(&controls));
6112        std::fs::write(root.path().join("stable.txt"), b"stable").expect("fixture");
6113        controls.gate(TestPoint::AfterRefreshVerification).arm();
6114        let refresher = opened.clone();
6115        let refresh = thread::spawn(move || refresher.refresh(&[PathBuf::from("")]));
6116        controls.gate(TestPoint::AfterRefreshVerification).wait_reached();
6117        // The held root pass keeps bounded evidence about newer passes: one entry per
6118        // entry present when it began. One more distinct newer scope than that retires
6119        // it, so its scope closes partial with no error of its own to report.
6120        let budget = opened.state.index.len().expect("entry count");
6121        for number in 0..=budget {
6122            let path = PathBuf::from(format!("missing-{number}"));
6123            let (epoch, _) = opened.state.index.begin_reconcile(&path).expect("begin newer");
6124            opened
6125                .state
6126                .index
6127                .finish_reconcile(
6128                    &path,
6129                    epoch,
6130                    true,
6131                    &[],
6132                    &[],
6133                    crate::index::ReconcileErrors {
6134                        errors: &[],
6135                        terminal: None,
6136                        disproves_old: true,
6137                    },
6138                )
6139                .expect("finish newer");
6140        }
6141        controls.gate(TestPoint::AfterRefreshVerification).release();
6142
6143        let result = refresh.join().expect("refresh thread").expect("refresh receipt");
6144
6145        assert_eq!(
6146            result.state.coverage,
6147            crate::Coverage::Partial(crate::CoverageReason::Inaccessible),
6148            "the retired pass publishes its scope partial"
6149        );
6150        assert!(
6151            result.issues.iter().any(|issue| issue.message.contains("retry")),
6152            "the receipt must name the retry the retired pass earned: {:?}",
6153            result.issues
6154        );
6155        opened.close().expect("close");
6156    }
6157
6158    /// One refresh over two subtrees, one of them unreadable: the readable subtree is
6159    /// verified on its own walk. One completion flag for the whole set marked it partial
6160    /// because its sibling could not be read.
6161    #[cfg(unix)]
6162    #[test]
6163    fn multi_path_refresh_closes_each_subtree_on_its_own_walk() {
6164        use std::os::unix::fs::PermissionsExt;
6165
6166        if !crate::test_support::require_permission_bits() {
6167            return;
6168        }
6169
6170        let root = tempfile::tempdir().expect("temp root");
6171        std::fs::create_dir(root.path().join("readable")).expect("readable directory");
6172        std::fs::write(root.path().join("readable/file"), b"ok").expect("readable fixture");
6173        let blocked = root.path().join("blocked");
6174        std::fs::create_dir(&blocked).expect("blocked directory");
6175        std::fs::write(blocked.join("secret"), b"secret").expect("blocked fixture");
6176        let opened = open_fixture(root.path(), OpenOptions::default()).expect("open");
6177        let settled = wait_until_settled(&opened);
6178        assert_eq!(settled.phase, crate::LifecyclePhase::Ready);
6179        std::fs::set_permissions(&blocked, std::fs::Permissions::from_mode(0o000))
6180            .expect("make directory unreadable");
6181
6182        let receipt = opened
6183            .refresh(&[PathBuf::from("readable"), PathBuf::from("blocked")])
6184            .expect("refresh");
6185        std::fs::set_permissions(&blocked, std::fs::Permissions::from_mode(0o700))
6186            .expect("restore directory permissions");
6187
6188        assert_eq!(receipt.issues.len(), 1, "{:?}", receipt.issues);
6189        let index = &opened.state.index;
6190        assert_eq!(
6191            index.freshness_at(Path::new("readable")).expect("freshness"),
6192            crate::Freshness::Fresh
6193        );
6194        assert_eq!(
6195            index.freshness_at(Path::new("blocked")).expect("freshness"),
6196            crate::Freshness::Partial
6197        );
6198        let since = index.since(receipt.after.sequence).expect("journal");
6199        assert!(
6200            since.commits.iter().flat_map(|commit| commit.state.iter()).any(|transition| {
6201                matches!(
6202                    transition,
6203                    crate::StateTransition::Verified { path } if path == Path::new("readable")
6204                )
6205            }),
6206            "the readable subtree was not verified"
6207        );
6208        opened.close().expect("close");
6209    }
6210
6211    /// Discovery records a child gone by its stat as it records one the listing never
6212    /// returned: absent, with no issue, under complete coverage.
6213    #[test]
6214    fn discovery_omits_a_child_deleted_between_listing_and_stat() {
6215        let root = tempfile::tempdir().expect("temp root");
6216        std::fs::write(root.path().join("kept"), b"kept").expect("kept fixture");
6217        std::fs::write(root.path().join("gone"), b"gone").expect("gone fixture");
6218        let hook = crate::scan::install_child_metadata_hook(root.path(), |path| {
6219            if path.file_name() == Some(std::ffi::OsStr::new("gone")) {
6220                std::fs::remove_file(path).expect("delete between listing and stat");
6221            }
6222            None
6223        });
6224        let opened = open_fixture(root.path(), OpenOptions::default()).expect("open");
6225        let settled = wait_until_settled(&opened);
6226        drop(hook);
6227
6228        assert_eq!(settled.coverage, crate::Coverage::Complete);
6229        assert_eq!(settled.issues.retained, 0);
6230        let index = &opened.state.index;
6231        assert!(index.kind(Path::new("gone")).expect("lookup").is_none());
6232        assert!(index.kind(Path::new("kept")).expect("lookup").is_some());
6233        opened.close().expect("close");
6234    }
6235
6236    /// One transient child error no longer withholds completeness from every directory
6237    /// the pass listed. Each is recorded on its own listing, as discovery decides, so a
6238    /// directory first listed by such a pass answers Absent below it instead of staying
6239    /// Unknown { Building } under a complete root.
6240    #[test]
6241    fn refresh_records_completeness_per_listed_directory_despite_a_child_error() {
6242        let root = tempfile::tempdir().expect("temp root");
6243        std::fs::create_dir(root.path().join("steady")).expect("steady directory");
6244        std::fs::write(root.path().join("steady/kept"), b"kept").expect("steady fixture");
6245        let opened = open_fixture(root.path(), OpenOptions::default()).expect("open");
6246        let settled = wait_until_settled(&opened);
6247        assert_eq!(settled.coverage, crate::Coverage::Complete);
6248        std::fs::create_dir(root.path().join("fresh")).expect("fresh directory");
6249        std::fs::write(root.path().join("fresh/new"), b"new").expect("fresh fixture");
6250
6251        let hook = crate::scan::install_child_metadata_hook(root.path(), |path| {
6252            (path.file_name() == Some(std::ffi::OsStr::new("kept"))).then(|| {
6253                std::io::Error::new(std::io::ErrorKind::PermissionDenied, "injected child error")
6254            })
6255        });
6256        let receipt = opened.refresh(&[PathBuf::new()]);
6257        drop(hook);
6258        let receipt = receipt.expect("refresh");
6259
6260        assert_eq!(receipt.issues.len(), 1, "{:?}", receipt.issues);
6261        let index = &opened.state.index;
6262        assert_eq!(
6263            index.freshness_at(Path::new("")).expect("freshness"),
6264            crate::Freshness::Partial
6265        );
6266        assert_eq!(index.directory_complete(Path::new("fresh")).expect("lookup"), Some(true));
6267        // Its enumerated child could not be verified, so the retained listing proof is withdrawn.
6268        assert_eq!(index.directory_complete(Path::new("steady")).expect("lookup"), Some(false));
6269        let lookup = opened
6270            .read(crate::ReadRequest {
6271                projections: vec![crate::ReadProjection::Lookup {
6272                    path: PathBuf::from("fresh/missing"),
6273                }],
6274                ..crate::ReadRequest::default()
6275            })
6276            .expect("lookup")
6277            .results
6278            .into_iter()
6279            .next()
6280            .expect("lookup result");
6281        assert!(
6282            matches!(lookup, crate::ProjectionResult::Lookup(crate::Knowledge::Absent)),
6283            "{lookup:?}"
6284        );
6285        opened.close().expect("close");
6286    }
6287
6288    /// A refresh on a Failed root keeps the issue that explains the failure. The failure
6289    /// is the state the root is in, so a clean walk below the issue's path disproves
6290    /// nothing; dropping it left a Failed root with no retained cause.
6291    #[test]
6292    fn refresh_on_a_failed_root_keeps_the_issue_that_explains_it() {
6293        let (root, opened) = opened(Arc::default());
6294        std::fs::create_dir(root.path().join("sub")).expect("fixture directory");
6295        opened
6296            .state
6297            .index
6298            .transition_discovery(DiscoveryTransition::Begin)
6299            .expect("begin discovery");
6300        let failure = crate::Issue::from_error_under(
6301            root.path(),
6302            &Error::io(root.path().join("sub"), std::io::Error::other("provider failed here")),
6303        );
6304        opened
6305            .state
6306            .index
6307            .transition_discovery(DiscoveryTransition::Failed(failure))
6308            .expect("fail discovery");
6309        let failed = opened.state.index.state().expect("state");
6310        assert_eq!(failed.phase, crate::LifecyclePhase::Failed);
6311        assert_eq!(failed.issues.retained, 1);
6312
6313        let receipt = opened.refresh(&[PathBuf::from("sub")]).expect("refresh on a failed root");
6314        assert_eq!(receipt.work.stale, 0);
6315
6316        let after = opened.state.index.state().expect("state");
6317        assert_eq!(after.phase, crate::LifecyclePhase::Failed);
6318        let issues = opened.state.index.issues().expect("issues");
6319        assert_eq!(issues.len(), 1, "{issues:?}");
6320        assert_eq!(issues[0].path.as_deref(), Some(Path::new("sub")));
6321        assert_eq!(after.issues.retained, 1);
6322        opened.close().expect("close");
6323    }
6324
6325    #[test]
6326    fn close_cancels_verified_refresh_before_its_conditional_commit() {
6327        let controls = Arc::new(TestControls::default());
6328        let (root, opened) = opened(Arc::clone(&controls));
6329        std::fs::write(root.path().join("late"), b"late").expect("fixture");
6330        controls.gate(TestPoint::AfterRefreshVerification).arm();
6331        let refresher = opened.clone();
6332        let refresh = thread::spawn(move || refresher.refresh(&[PathBuf::from("late")]));
6333        controls.gate(TestPoint::AfterRefreshVerification).wait_reached();
6334        let closer = opened.clone();
6335        let close = thread::spawn(move || closer.close());
6336        let deadline = std::time::Instant::now() + TEST_GATE_TIMEOUT;
6337        while !opened.state.cancellation.is_cancelled() {
6338            assert!(std::time::Instant::now() < deadline, "close did not cancel refresh");
6339            thread::yield_now();
6340        }
6341        controls.gate(TestPoint::AfterRefreshVerification).release();
6342
6343        assert!(matches!(refresh.join().expect("refresh thread"), Err(Error::OpenedIndexClosed)));
6344        close.join().expect("close thread").expect("joined close");
6345        assert_eq!(opened.state.index.kind(Path::new("late")).expect("lookup"), None);
6346        opened.close().expect("repeat close");
6347    }
6348
6349    #[test]
6350    fn refresh_tracks_hidden_control_creation_edit_and_deletion() {
6351        let root = tempfile::tempdir().expect("temp root");
6352        std::fs::write(root.path().join("debug.log"), b"log").expect("fixture");
6353        std::fs::write(root.path().join("keep.rs"), b"keep").expect("fixture");
6354        let opened = open_fixture(
6355            root.path(),
6356            OpenOptions {
6357                hidden: Some(Arc::new(crate::HiddenPolicy::prune_hidden::<[&str; 0], &str>([]))),
6358                ..OpenOptions::default()
6359            },
6360        )
6361        .expect("opened root");
6362        assert_eq!(wait_until_settled(&opened).phase, crate::LifecyclePhase::Ready);
6363
6364        std::fs::write(root.path().join(".gitignore"), b"*.log\n").expect("create control");
6365        let created = opened.refresh(&[PathBuf::from(".gitignore")]).expect("create refresh");
6366        assert_eq!(created.accepted, vec![PathBuf::from(".gitignore")]);
6367        let image = opened.state.index.snapshot().expect("snapshot");
6368        assert!(
6369            image
6370                .controls()
6371                .expect("control state observed")
6372                .source_is(Path::new(".gitignore"), b"*.log\n")
6373        );
6374        assert_eq!(
6375            image.is_ignored(Path::new("debug.log")).expect("control state observed"),
6376            Some(true)
6377        );
6378
6379        std::fs::write(root.path().join(".gitignore"), b"*.tmp\n").expect("edit control");
6380        opened.refresh(&[PathBuf::from(".gitignore")]).expect("edit refresh");
6381        let image = opened.state.index.snapshot().expect("snapshot");
6382        assert!(
6383            image
6384                .controls()
6385                .expect("control state observed")
6386                .source_is(Path::new(".gitignore"), b"*.tmp\n")
6387        );
6388        assert_eq!(
6389            image.is_ignored(Path::new("debug.log")).expect("control state observed"),
6390            Some(false)
6391        );
6392        let unchanged =
6393            opened.refresh(&[PathBuf::from(".gitignore")]).expect("unchanged control refresh");
6394        assert_eq!(unchanged.work.observations, 1);
6395
6396        std::fs::remove_file(root.path().join(".gitignore")).expect("delete control");
6397        opened.refresh(&[PathBuf::from(".gitignore")]).expect("delete refresh");
6398        let image = opened.state.index.snapshot().expect("snapshot");
6399        assert!(image.controls().expect("control state observed").is_empty());
6400        let partitions = image.partition_total().expect("control state observed");
6401        assert_eq!(partitions.all, partitions.unignored);
6402        opened.close().expect("close");
6403    }
6404
6405    #[cfg(feature = "watch")]
6406    #[test]
6407    #[cfg(unix)]
6408    fn inaccessible_baseline_enters_watching_with_partial_coverage() {
6409        use std::os::unix::fs::PermissionsExt;
6410
6411        if !crate::test_support::require_permission_bits() {
6412            return;
6413        }
6414
6415        let root = tempfile::tempdir().expect("temp root");
6416        let scripts = tempfile::tempdir().expect("script root");
6417        let inaccessible = root.path().join("inaccessible");
6418        std::fs::create_dir(&inaccessible).expect("inaccessible directory");
6419        std::fs::write(inaccessible.join("secret"), b"secret").expect("fixture");
6420        std::fs::set_permissions(&inaccessible, std::fs::Permissions::from_mode(0o000))
6421            .expect("make directory inaccessible");
6422        let script = scripts.path().join("events.script");
6423        std::fs::write(&script, b"").expect("script");
6424
6425        let opened = open_fixture(root.path(), scripted_options(&script)).expect("open");
6426        let deadline = std::time::Instant::now() + TEST_GATE_TIMEOUT;
6427        let state = loop {
6428            let state = opened.state.index.state().expect("read state");
6429            if matches!(
6430                state.phase,
6431                crate::LifecyclePhase::Watching | crate::LifecyclePhase::Failed
6432            ) {
6433                break state;
6434            }
6435            assert!(std::time::Instant::now() < deadline, "observation handoff did not settle");
6436            std::thread::yield_now();
6437        };
6438        std::fs::set_permissions(&inaccessible, std::fs::Permissions::from_mode(0o700))
6439            .expect("restore directory permissions");
6440
6441        assert_eq!(state.phase, crate::LifecyclePhase::Watching);
6442        assert_eq!(state.coverage, crate::Coverage::Partial(crate::CoverageReason::Inaccessible));
6443        assert_eq!(state.freshness, crate::Freshness::Partial);
6444        assert!(state.issues.retained > 0);
6445        opened.close().expect("close");
6446    }
6447
6448    /// A boundary discovery could not read, and the handoff then read cleanly, is gone.
6449    ///
6450    /// Discovery records `blocked` as inaccessible; it becomes readable before the
6451    /// observation handoff, whose full pass then lists it without an error. `Finish` only
6452    /// ever upgraded a building root and `Watching` never upgraded coverage at all, so the
6453    /// root stayed partial for the life of the session after a pass proving otherwise.
6454    #[cfg(all(unix, feature = "watch"))]
6455    #[test]
6456    fn watching_after_a_clean_handoff_rederives_complete_coverage() {
6457        use std::os::unix::fs::PermissionsExt;
6458
6459        if !crate::test_support::require_permission_bits() {
6460            return;
6461        }
6462        let root = tempfile::tempdir().expect("temp root");
6463        let scripts = tempfile::tempdir().expect("script root");
6464        let script = scripts.path().join("events.script");
6465        std::fs::write(&script, b"").expect("script");
6466        let blocked = root.path().join("blocked");
6467        std::fs::create_dir(&blocked).expect("blocked directory");
6468        std::fs::write(blocked.join("secret"), b"secret").expect("fixture");
6469        std::fs::set_permissions(&blocked, std::fs::Permissions::from_mode(0o000))
6470            .expect("make directory inaccessible");
6471        let controls = Arc::new(TestControls::default());
6472        controls.gate(TestPoint::BeforeObservationHandoff).arm();
6473        let opened = OpenedIndex::open_for_test(
6474            root.path(),
6475            scripted_options(&script),
6476            Arc::clone(&controls),
6477        )
6478        .expect("open");
6479        controls.gate(TestPoint::BeforeObservationHandoff).wait_reached();
6480        let discovered = opened.state.index.state().expect("state after discovery");
6481        std::fs::set_permissions(&blocked, std::fs::Permissions::from_mode(0o700))
6482            .expect("restore directory permissions");
6483        controls.gate(TestPoint::BeforeObservationHandoff).release();
6484
6485        let state = wait_until_phase(&opened, crate::LifecyclePhase::Watching);
6486        assert_eq!(
6487            discovered.coverage,
6488            crate::Coverage::Partial(crate::CoverageReason::Inaccessible)
6489        );
6490        assert_eq!(state.coverage, crate::Coverage::Complete);
6491        assert_eq!(state.freshness, crate::Freshness::Fresh);
6492        assert_eq!(
6493            opened.state.index.kind(Path::new("blocked/secret")).expect("lookup"),
6494            Some(EntryKind::File),
6495            "the handoff read the formerly inaccessible directory"
6496        );
6497        // Complete coverage has to hold one level down as well. The handoff listed `blocked`
6498        // in full, so a name it does not hold is absent; before a reconcile recorded
6499        // completeness this answered `Unknown { Building }` forever on a complete root.
6500        assert_eq!(
6501            opened.state.index.directory_complete(Path::new("blocked")).expect("lookup"),
6502            Some(true)
6503        );
6504        let response = opened
6505            .read(crate::ReadRequest {
6506                projections: vec![crate::ReadProjection::Lookup {
6507                    path: PathBuf::from("blocked/missing"),
6508                }],
6509                ..crate::ReadRequest::default()
6510            })
6511            .expect("read");
6512        assert_eq!(response.state.phase, crate::LifecyclePhase::Watching);
6513        assert_eq!(response.state.coverage, crate::Coverage::Complete);
6514        assert!(
6515            matches!(
6516                response.results[0],
6517                crate::ProjectionResult::Lookup(crate::Knowledge::Absent)
6518            ),
6519            "{:?}",
6520            response.results[0]
6521        );
6522        opened.close().expect("close");
6523    }
6524
6525    /// A directory the observer adds after discovery is complete once its relist finishes.
6526    ///
6527    /// Only discovery used to mark a directory complete, so one created while the root was
6528    /// watched stayed incomplete for the session, and a lookup below it on a complete root
6529    /// answered `Unknown { Building }`, which never resolves.
6530    #[cfg(feature = "watch")]
6531    #[test]
6532    fn a_directory_the_observer_adds_is_complete_after_its_relist() {
6533        let root = tempfile::tempdir().expect("temp root");
6534        let scripts = tempfile::tempdir().expect("script root");
6535        let script = scripts.path().join("events.script");
6536        std::fs::write(&script, b"").expect("script");
6537        let controls = Arc::new(TestControls::default());
6538        let opened = OpenedIndex::open_for_test(
6539            root.path(),
6540            scripted_options(&script),
6541            Arc::clone(&controls),
6542        )
6543        .expect("open scripted observer");
6544        wait_until_phase(&opened, crate::LifecyclePhase::Watching);
6545        let start = current_version(&opened);
6546        let watching = opened.read(crate::ReadRequest::default()).expect("read state").state;
6547        assert_eq!(watching.coverage, crate::Coverage::Complete);
6548
6549        std::fs::create_dir_all(root.path().join("later/deeper")).expect("fixture");
6550        std::fs::write(root.path().join("later/deeper/inner.txt"), b"i").expect("fixture");
6551        controls.send_observation_hints("create-dir\tlater\n");
6552        wait_for_observed_walk(&opened, &controls, root.path(), start, Path::new("later"), 1);
6553
6554        let index = &opened.state.index;
6555        assert_eq!(
6556            index.kind(Path::new("later/deeper/inner.txt")).expect("lookup"),
6557            Some(EntryKind::File)
6558        );
6559        assert_eq!(index.directory_complete(Path::new("later")).expect("lookup"), Some(true));
6560        assert_eq!(
6561            index.directory_complete(Path::new("later/deeper")).expect("lookup"),
6562            Some(true)
6563        );
6564        let response = opened
6565            .read(crate::ReadRequest {
6566                projections: vec![
6567                    crate::ReadProjection::Lookup { path: PathBuf::from("later/missing") },
6568                    crate::ReadProjection::Lookup { path: PathBuf::from("later/deeper/missing") },
6569                ],
6570                ..crate::ReadRequest::default()
6571            })
6572            .expect("read");
6573        assert_eq!(response.state.coverage, crate::Coverage::Complete);
6574        for result in &response.results {
6575            assert!(
6576                matches!(result, crate::ProjectionResult::Lookup(crate::Knowledge::Absent)),
6577                "{result:?}"
6578            );
6579        }
6580        let since = index.since(start.sequence).expect("journal");
6581        let completed: Vec<_> = since
6582            .commits
6583            .iter()
6584            .flat_map(|commit| commit.state.iter())
6585            .filter_map(|transition| match transition {
6586                crate::StateTransition::DirectoryComplete { path } => Some(path.clone()),
6587                _ => None,
6588            })
6589            .collect();
6590        assert_eq!(completed, [PathBuf::from("later"), PathBuf::from("later/deeper")]);
6591        assert_eq!(
6592            response.state.progress.directories_complete,
6593            watching.progress.directories_complete + 2,
6594            "the progress count agrees with the transitions"
6595        );
6596        opened.close().expect("close");
6597    }
6598
6599    /// A directory deleted during discovery leaves a watched root complete.
6600    #[cfg(feature = "watch")]
6601    #[test]
6602    fn a_directory_that_vanishes_during_discovery_leaves_a_watched_root_complete() {
6603        let controls = Arc::new(TestControls::default());
6604        controls.gate(TestPoint::AfterRootDirectory).arm();
6605        let root = tempfile::tempdir().expect("temp root");
6606        let scripts = tempfile::tempdir().expect("script root");
6607        let script = scripts.path().join("events.script");
6608        std::fs::write(&script, b"").expect("script");
6609        std::fs::create_dir(root.path().join("sub")).expect("sub");
6610        std::fs::write(root.path().join("sub/inner.txt"), b"x").expect("fixture");
6611        std::fs::write(root.path().join("keep.txt"), b"k").expect("fixture");
6612        let opened = OpenedIndex::open_for_test(
6613            root.path(),
6614            scripted_options(&script),
6615            Arc::clone(&controls),
6616        )
6617        .expect("open");
6618        controls.gate(TestPoint::AfterRootDirectory).wait_reached();
6619        std::fs::remove_dir_all(root.path().join("sub")).expect("remove sub during discovery");
6620        controls.gate(TestPoint::AfterRootDirectory).release();
6621
6622        let state = wait_until_phase(&opened, crate::LifecyclePhase::Watching);
6623        assert_eq!(state.coverage, crate::Coverage::Complete);
6624        assert_eq!(state.freshness, crate::Freshness::Fresh);
6625        assert_eq!(state.issues.retained, 0);
6626        assert_eq!(
6627            opened.state.index.kind(Path::new("sub")).expect("lookup"),
6628            None,
6629            "the handoff pass removed the vanished directory"
6630        );
6631        assert_eq!(opened.state.index.directory_complete(Path::new("")).expect("root"), Some(true));
6632        opened.close().expect("close");
6633    }
6634
6635    #[cfg(feature = "watch")]
6636    #[test]
6637    fn observation_is_captured_before_baseline_and_closes_the_handoff_gap() {
6638        let root = tempfile::tempdir().expect("temp root");
6639        let scripts = tempfile::tempdir().expect("script root");
6640        let path = root.path().join("during.txt");
6641        std::fs::write(&path, b"before").expect("fixture");
6642        let script = scripts.path().join("events.script");
6643        std::fs::write(&script, b"modify\tduring.txt\n").expect("script");
6644        let controls = Arc::new(TestControls::default());
6645        controls.gate(TestPoint::BeforeDiscovery).arm();
6646
6647        let opened = OpenedIndex::open_for_test(
6648            root.path(),
6649            scripted_options(&script),
6650            Arc::clone(&controls),
6651        )
6652        .expect("opened observed root");
6653        controls.gate(TestPoint::BeforeDiscovery).wait_reached();
6654        std::fs::write(&path, b"changed-during-baseline").expect("mutate during handoff");
6655        controls.gate(TestPoint::BeforeDiscovery).release();
6656
6657        let state = wait_until_phase(&opened, crate::LifecyclePhase::Watching);
6658        assert_eq!(state.freshness, crate::Freshness::Fresh);
6659        assert_eq!(state.coverage, crate::Coverage::Complete);
6660        assert_eq!(
6661            opened
6662                .state
6663                .index
6664                .attrs(Path::new("during.txt"))
6665                .expect("attrs")
6666                .expect("retained")
6667                .size,
6668            23
6669        );
6670        let since = opened.state.index.since(crate::Clock::ZERO).expect("journal");
6671        assert!(since.commits.iter().any(|commit| {
6672            commit.state.iter().any(|transition| {
6673                matches!(
6674                    transition,
6675                    crate::StateTransition::IndexState { current, .. }
6676                        if current.phase == crate::LifecyclePhase::Reconciling
6677                )
6678            })
6679        }));
6680        assert!(since.commits.iter().any(|commit| {
6681            commit.state.iter().any(|transition| {
6682                matches!(
6683                    transition,
6684                    crate::StateTransition::IndexState { current, .. }
6685                        if current.phase == crate::LifecyclePhase::Watching
6686                )
6687            })
6688        }));
6689        opened.close().expect("joined close");
6690    }
6691
6692    #[cfg(feature = "watch")]
6693    #[test]
6694    fn scripted_overflow_is_provider_recovery_not_a_consumer_reset() {
6695        let root = tempfile::tempdir().expect("temp root");
6696        let scripts = tempfile::tempdir().expect("script root");
6697        std::fs::create_dir(root.path().join("src")).expect("directory");
6698        std::fs::write(root.path().join("src/present"), b"present").expect("fixture");
6699        let script = scripts.path().join("events.script");
6700        std::fs::write(&script, b"rescan\tsrc\n").expect("script");
6701        let opened = open_fixture(root.path(), scripted_options(&script)).expect("open");
6702
6703        let state = wait_until_phase(&opened, crate::LifecyclePhase::Watching);
6704        assert_eq!(state.freshness, crate::Freshness::Fresh);
6705        assert!(state.issues.retained > 0);
6706        let since = opened.state.index.since(crate::Clock::ZERO).expect("journal");
6707        assert!(!since.truncated);
6708        assert!(since.commits.iter().any(|commit| commit.changes.iter().any(|change| {
6709            matches!(
6710                change,
6711                crate::EffectiveChange::Invalidated {
6712                    path,
6713                    reason: crate::InvalidateReason::WatchOverflow,
6714                } if path == Path::new("src")
6715            )
6716        })));
6717        assert!(
6718            opened
6719                .state
6720                .index
6721                .issues()
6722                .expect("issues")
6723                .iter()
6724                .any(|issue| issue.kind == crate::IssueKind::ObservationGap)
6725        );
6726        opened.close().expect("close");
6727    }
6728
6729    #[cfg(feature = "watch")]
6730    #[test]
6731    fn scripted_directory_creation_closes_the_registration_gap() {
6732        let root = tempfile::tempdir().expect("temp root");
6733        let scripts = tempfile::tempdir().expect("script root");
6734        std::fs::create_dir(root.path().join("newdir")).expect("directory");
6735        std::fs::write(root.path().join("newdir/child"), b"child").expect("fixture");
6736        let script = scripts.path().join("events.script");
6737        std::fs::write(&script, b"create-dir\tnewdir\n").expect("script");
6738        let opened = open_fixture(root.path(), scripted_options(&script)).expect("open");
6739
6740        wait_until_phase(&opened, crate::LifecyclePhase::Watching);
6741        let since = opened.state.index.since(crate::Clock::ZERO).expect("journal");
6742        assert!(since.commits.iter().any(|commit| commit.changes.iter().any(|change| {
6743            matches!(
6744                change,
6745                crate::EffectiveChange::Invalidated {
6746                    path,
6747                    reason: crate::InvalidateReason::WatchSetupRace,
6748                } if path == Path::new("newdir")
6749            )
6750        })));
6751        assert_eq!(
6752            opened.state.index.kind(Path::new("newdir/child")).expect("lookup"),
6753            Some(EntryKind::File)
6754        );
6755        opened.close().expect("close");
6756    }
6757
6758    #[cfg(feature = "watch")]
6759    #[test]
6760    fn scripted_observation_keeps_the_opened_index_live_after_handoff() {
6761        let root = tempfile::tempdir().expect("temp root");
6762        let scripts = tempfile::tempdir().expect("script root");
6763        let script = scripts.path().join("events.script");
6764        std::fs::write(&script, b"").expect("script");
6765        let controls = Arc::new(TestControls::default());
6766        let opened = OpenedIndex::open_for_test(
6767            root.path(),
6768            scripted_options(&script),
6769            Arc::clone(&controls),
6770        )
6771        .expect("open scripted observer");
6772        wait_until_phase(&opened, crate::LifecyclePhase::Watching);
6773        let before = current_version(&opened);
6774
6775        std::fs::write(root.path().join("live.txt"), b"live").expect("live mutation");
6776        controls.send_observation_hints("create\tlive.txt\n");
6777        let deadline = std::time::Instant::now() + TEST_GATE_TIMEOUT;
6778        loop {
6779            if opened.state.index.kind(Path::new("live.txt")).expect("lookup")
6780                == Some(EntryKind::File)
6781            {
6782                break;
6783            }
6784            assert!(std::time::Instant::now() < deadline, "scripted hint was not applied");
6785            std::thread::yield_now();
6786        }
6787        let since = opened.state.index.since(before.sequence).expect("journal");
6788        assert!(since.commits.iter().any(|commit| commit.changes.iter().any(|change| {
6789            matches!(
6790                change,
6791                crate::EffectiveChange::Inserted { path, .. }
6792                    if path == Path::new("live.txt")
6793            )
6794        })));
6795        opened.close().expect("close");
6796    }
6797
6798    #[cfg(feature = "watch")]
6799    #[test]
6800    fn live_observation_gap_recovers_before_reporting_freshness() {
6801        let root = tempfile::tempdir().expect("temp root");
6802        let scripts = tempfile::tempdir().expect("script root");
6803        let script = scripts.path().join("events.script");
6804        std::fs::write(&script, b"").expect("script");
6805        let controls = Arc::new(TestControls::default());
6806        let opened = OpenedIndex::open_for_test(
6807            root.path(),
6808            scripted_options(&script),
6809            Arc::clone(&controls),
6810        )
6811        .expect("open scripted observer");
6812        wait_until_phase(&opened, crate::LifecyclePhase::Watching);
6813        let before = current_version(&opened);
6814
6815        std::fs::write(root.path().join("recovered.txt"), b"recovered").expect("missed mutation");
6816        controls.send_observation_hints("rescan\t.\n");
6817        let deadline = std::time::Instant::now() + TEST_GATE_TIMEOUT;
6818        loop {
6819            let state = opened.state.index.state().expect("read state");
6820            let recovered = opened.state.index.kind(Path::new("recovered.txt")).expect("lookup")
6821                == Some(EntryKind::File);
6822            if recovered && state.freshness == crate::Freshness::Fresh {
6823                break;
6824            }
6825            assert!(std::time::Instant::now() < deadline, "gap recovery did not finish");
6826            std::thread::yield_now();
6827        }
6828
6829        let since = opened.state.index.since(before.sequence).expect("journal");
6830        assert!(since.commits.iter().any(|commit| commit.changes.iter().any(|change| {
6831            matches!(
6832                change,
6833                crate::EffectiveChange::Invalidated {
6834                    path,
6835                    reason: crate::InvalidateReason::WatchOverflow,
6836                } if path.as_os_str().is_empty()
6837            )
6838        })));
6839        assert!(since.commits.iter().any(|commit| commit.state.iter().any(|transition| {
6840            matches!(
6841                transition,
6842                crate::StateTransition::Freshness {
6843                    current: crate::Freshness::Reconciling | crate::Freshness::Stale,
6844                    ..
6845                }
6846            )
6847        })));
6848        assert!(since.commits.iter().any(|commit| commit.state.iter().any(|transition| {
6849            matches!(
6850                transition,
6851                crate::StateTransition::Freshness { current: crate::Freshness::Fresh, .. }
6852            )
6853        })));
6854        assert!(
6855            opened
6856                .state
6857                .index
6858                .issues()
6859                .expect("issues")
6860                .iter()
6861                .any(|issue| issue.kind == crate::IssueKind::ObservationGap)
6862        );
6863        opened.close().expect("close");
6864    }
6865
6866    #[cfg(feature = "watch")]
6867    fn reconciles_of(opened: &OpenedIndex, since: crate::EngineVersion, path: &Path) -> usize {
6868        opened
6869            .state
6870            .index
6871            .since(since.sequence)
6872            .expect("journal")
6873            .commits
6874            .iter()
6875            .flat_map(|commit| commit.state.iter())
6876            .filter(|transition| {
6877                matches!(
6878                    transition,
6879                    crate::StateTransition::Freshness { path: marked, current, .. }
6880                        if marked == path && *current == crate::Freshness::Reconciling
6881                )
6882            })
6883            .count()
6884    }
6885
6886    /// Wait until the observer has walked `path` `walks` times since `start` and published
6887    /// everything the last walk will.
6888    ///
6889    /// The observer applies one intent at a time and publishes a walk's outcome before it
6890    /// takes the next, so a marker event sent after the walk has begun lands only once that
6891    /// walk is over. A marker sent together with the event that caused the walk could
6892    /// coalesce into the same intent and land first.
6893    #[cfg(feature = "watch")]
6894    fn wait_for_observed_walk(
6895        opened: &OpenedIndex,
6896        controls: &TestControls,
6897        root: &Path,
6898        start: crate::EngineVersion,
6899        path: &Path,
6900        walks: usize,
6901    ) {
6902        let deadline = std::time::Instant::now() + TEST_GATE_TIMEOUT;
6903        while reconciles_of(opened, start, path) < walks {
6904            assert!(std::time::Instant::now() < deadline, "walk {walks} of {path:?} did not begin");
6905            std::thread::yield_now();
6906        }
6907        let marker = format!("marker-{}-{walks}.txt", path.display());
6908        std::fs::write(root.join(&marker), &marker).expect("marker");
6909        controls.send_observation_hints(&format!("create\t{marker}\n"));
6910        let deadline = std::time::Instant::now() + TEST_GATE_TIMEOUT;
6911        while opened.state.index.kind(Path::new(&marker)).expect("lookup") != Some(EntryKind::File)
6912        {
6913            assert!(std::time::Instant::now() < deadline, "{marker} was not applied");
6914            std::thread::yield_now();
6915        }
6916    }
6917
6918    /// A gap over an unreadable directory is walked once, and its cause is retained.
6919    ///
6920    /// The walk's permission error made the reconciliation incomplete, so the invalidation
6921    /// was restored, and the observer drains that queue after every event: each unrelated
6922    /// event re-walked the same unreadable subtree, forever -- a full-tree walk per event
6923    /// for a root escalation. The report was then discarded, so the resulting partial
6924    /// freshness had no issue to explain it.
6925    #[cfg(all(unix, feature = "watch"))]
6926    #[test]
6927    fn an_unreadable_gap_is_walked_once_and_explains_itself() {
6928        use std::os::unix::fs::PermissionsExt;
6929
6930        if !crate::test_support::require_permission_bits() {
6931            return;
6932        }
6933        let root = tempfile::tempdir().expect("temp root");
6934        let scripts = tempfile::tempdir().expect("script root");
6935        let script = scripts.path().join("events.script");
6936        std::fs::write(&script, b"").expect("script");
6937        let blocked = root.path().join("blocked");
6938        std::fs::create_dir(&blocked).expect("blocked");
6939        std::fs::write(blocked.join("secret"), b"s").expect("fixture");
6940        let controls = Arc::new(TestControls::default());
6941        let opened = OpenedIndex::open_for_test(
6942            root.path(),
6943            scripted_options(&script),
6944            Arc::clone(&controls),
6945        )
6946        .expect("open scripted observer");
6947        let watching = wait_until_phase(&opened, crate::LifecyclePhase::Watching);
6948        assert_eq!(watching.freshness, crate::Freshness::Fresh);
6949        let start = current_version(&opened);
6950
6951        std::fs::set_permissions(&blocked, std::fs::Permissions::from_mode(0o000))
6952            .expect("make directory inaccessible");
6953        controls.send_observation_hints("rescan\tblocked\n");
6954        let deadline = std::time::Instant::now() + TEST_GATE_TIMEOUT;
6955        while reconciles_of(&opened, start, Path::new("blocked")) < 1 {
6956            assert!(std::time::Instant::now() < deadline, "the gap was not reconciled");
6957            std::thread::yield_now();
6958        }
6959
6960        // Two later events elsewhere. Once the second has landed, the drain that followed
6961        // the first has finished, so any re-walk it did is already in the journal.
6962        for name in ["live.txt", "marker.txt"] {
6963            std::fs::write(root.path().join(name), name).expect("unrelated mutation");
6964            controls.send_observation_hints(&format!("create\t{name}\n"));
6965            let deadline = std::time::Instant::now() + TEST_GATE_TIMEOUT;
6966            while opened.state.index.kind(Path::new(name)).expect("lookup") != Some(EntryKind::File)
6967            {
6968                assert!(std::time::Instant::now() < deadline, "{name} was not applied");
6969                std::thread::yield_now();
6970            }
6971        }
6972
6973        let walks = reconciles_of(&opened, start, Path::new("blocked"));
6974        let state = opened.state.index.state().expect("state");
6975        let issues = opened.state.index.issues().expect("issues");
6976        std::fs::set_permissions(&blocked, std::fs::Permissions::from_mode(0o700))
6977            .expect("restore directory permissions");
6978        assert_eq!(walks, 1, "an unreadable subtree must not be re-walked per unrelated event");
6979        assert_eq!(state.phase, crate::LifecyclePhase::Watching);
6980        assert_eq!(state.freshness, crate::Freshness::Partial);
6981        assert!(
6982            issues.iter().any(|issue| issue.kind == crate::IssueKind::Permission
6983                && issue.path.as_deref().is_some_and(|path| path.ends_with("blocked"))),
6984            "{issues:?}"
6985        );
6986        opened.close().expect("close");
6987    }
6988
6989    /// Re-walking one unreadable boundary retains one issue per cause, and a clean re-walk
6990    /// drops the cause it disproved.
6991    ///
6992    /// Every provider gap over an unreadable directory retained another `ObservationGap`
6993    /// and another `Permission` issue, so a few dozen routine gaps filled the bounded list
6994    /// with copies and every later distinct issue was omitted with no text. The permission
6995    /// issue also named an absolute path where the gap named the root-relative one.
6996    #[cfg(all(unix, feature = "watch"))]
6997    #[test]
6998    fn repeated_unreadable_reconciles_retain_one_issue_per_boundary() {
6999        use std::os::unix::fs::PermissionsExt;
7000
7001        if !crate::test_support::require_permission_bits() {
7002            return;
7003        }
7004        let root = tempfile::tempdir().expect("temp root");
7005        let scripts = tempfile::tempdir().expect("script root");
7006        let script = scripts.path().join("events.script");
7007        std::fs::write(&script, b"").expect("script");
7008        let blocked = root.path().join("blocked");
7009        std::fs::create_dir(&blocked).expect("blocked");
7010        std::fs::write(blocked.join("secret"), b"s").expect("fixture");
7011        let controls = Arc::new(TestControls::default());
7012        let opened = OpenedIndex::open_for_test(
7013            root.path(),
7014            scripted_options(&script),
7015            Arc::clone(&controls),
7016        )
7017        .expect("open scripted observer");
7018        let watching = wait_until_phase(&opened, crate::LifecyclePhase::Watching);
7019        assert_eq!(watching.issues.retained, 0);
7020        let start = current_version(&opened);
7021        let walks = std::cell::Cell::new(0);
7022        let rescan_then_wait = || {
7023            controls.send_observation_hints("rescan\tblocked\n");
7024            walks.set(walks.get() + 1);
7025            wait_for_observed_walk(
7026                &opened,
7027                &controls,
7028                root.path(),
7029                start,
7030                Path::new("blocked"),
7031                walks.get(),
7032            );
7033        };
7034        let issues_of = |kind: crate::IssueKind| {
7035            opened
7036                .state
7037                .index
7038                .issues()
7039                .expect("issues")
7040                .into_iter()
7041                .filter(|issue| issue.kind == kind)
7042                .collect::<Vec<_>>()
7043        };
7044
7045        std::fs::set_permissions(&blocked, std::fs::Permissions::from_mode(0o000))
7046            .expect("make directory inaccessible");
7047        for _ in 0..5 {
7048            rescan_then_wait();
7049        }
7050        let state = opened.state.index.state().expect("state");
7051        let permission = issues_of(crate::IssueKind::Permission);
7052        let gaps = issues_of(crate::IssueKind::ObservationGap);
7053        std::fs::set_permissions(&blocked, std::fs::Permissions::from_mode(0o700))
7054            .expect("restore directory permissions");
7055        assert_eq!(state.phase, crate::LifecyclePhase::Watching);
7056        assert_eq!(permission.len(), 1, "{permission:?}");
7057        assert_eq!(permission[0].path.as_deref(), Some(Path::new("blocked")), "{permission:?}");
7058        assert_eq!(gaps.len(), 1, "{gaps:?}");
7059        assert_eq!(gaps[0].path.as_deref(), Some(Path::new("blocked")), "{gaps:?}");
7060        assert_eq!(state.issues, crate::IssueSummary { retained: 2, omitted: 0 });
7061
7062        // Readable again: the next re-walk completes, which disproves the permission issue.
7063        // The gap stays, because it records that observation lost precision there.
7064        rescan_then_wait();
7065        let state = opened.state.index.state().expect("state");
7066        assert_eq!(issues_of(crate::IssueKind::Permission), []);
7067        assert_eq!(issues_of(crate::IssueKind::ObservationGap).len(), 1);
7068        assert_eq!(state.issues, crate::IssueSummary { retained: 1, omitted: 0 });
7069        assert_eq!(state.freshness, crate::Freshness::Fresh);
7070        assert_eq!(
7071            opened.state.index.kind(Path::new("blocked/secret")).expect("lookup"),
7072            Some(EntryKind::File)
7073        );
7074        opened.close().expect("close");
7075    }
7076
7077    #[cfg(feature = "watch")]
7078    #[test]
7079    fn live_observation_shares_the_exact_opened_root_file_budget() {
7080        let root = tempfile::tempdir().expect("temp root");
7081        let scripts = tempfile::tempdir().expect("script root");
7082        std::fs::write(root.path().join("baseline.txt"), b"baseline").expect("fixture");
7083        let script = scripts.path().join("events.script");
7084        std::fs::write(&script, b"").expect("script");
7085        let controls = Arc::new(TestControls::default());
7086        let mut options = scripted_options(&script);
7087        options.budget.max_files = Some(1);
7088        let opened = OpenedIndex::open_for_test(root.path(), options, Arc::clone(&controls))
7089            .expect("open scripted observer");
7090        wait_until_phase(&opened, crate::LifecyclePhase::Watching);
7091
7092        std::fs::write(root.path().join("over-budget.txt"), b"refused")
7093            .expect("over-budget mutation");
7094        controls.send_observation_hints("create\tover-budget.txt\n");
7095        let state = wait_until_phase(&opened, crate::LifecyclePhase::Stopped);
7096
7097        assert_eq!(state.coverage, crate::Coverage::Partial(crate::CoverageReason::Budget));
7098        assert_eq!(opened.state.index.kind(Path::new("over-budget.txt")).expect("lookup"), None);
7099        assert_eq!(opened.state.index.total().expect("total").files, 1);
7100        assert!(
7101            opened
7102                .state
7103                .index
7104                .issues()
7105                .expect("issues")
7106                .iter()
7107                .any(|issue| issue.kind == crate::IssueKind::ResourceBudget)
7108        );
7109        opened.close().expect("close");
7110    }
7111
7112    #[cfg(feature = "watch")]
7113    #[test]
7114    fn close_after_observation_verification_prevents_publication() {
7115        let root = tempfile::tempdir().expect("temp root");
7116        let scripts = tempfile::tempdir().expect("script root");
7117        let script = scripts.path().join("events.script");
7118        std::fs::write(&script, b"").expect("script");
7119        let controls = Arc::new(TestControls::default());
7120        let opened = OpenedIndex::open_for_test(
7121            root.path(),
7122            scripted_options(&script),
7123            Arc::clone(&controls),
7124        )
7125        .expect("open scripted observer");
7126        wait_until_phase(&opened, crate::LifecyclePhase::Watching);
7127
7128        controls.gate(TestPoint::AfterObservationVerification).arm();
7129        std::fs::write(root.path().join("too-late.txt"), b"verified").expect("late mutation");
7130        controls.send_observation_hints("create\ttoo-late.txt\n");
7131        controls.gate(TestPoint::AfterObservationVerification).wait_reached();
7132
7133        let closer = opened.clone();
7134        let close = thread::spawn(move || closer.close());
7135        let deadline = std::time::Instant::now() + TEST_GATE_TIMEOUT;
7136        while !opened.state.cancellation.is_cancelled() {
7137            assert!(std::time::Instant::now() < deadline, "close did not cancel observation");
7138            thread::yield_now();
7139        }
7140        assert!(!close.is_finished(), "close returned before the commit boundary released");
7141        controls.gate(TestPoint::AfterObservationVerification).release();
7142        close.join().expect("close thread").expect("joined close");
7143
7144        assert_eq!(opened.state.index.kind(Path::new("too-late.txt")).expect("lookup"), None);
7145        opened.close().expect("repeat close");
7146    }
7147
7148    #[cfg(feature = "watch")]
7149    #[test]
7150    fn stopped_discovery_never_claims_to_be_watching() {
7151        let root = tempfile::tempdir().expect("temp root");
7152        let scripts = tempfile::tempdir().expect("script root");
7153        std::fs::write(root.path().join("one"), b"one").expect("fixture");
7154        std::fs::write(root.path().join("two"), b"two").expect("fixture");
7155        let script = scripts.path().join("events.script");
7156        std::fs::write(&script, b"").expect("script");
7157        let mut options = scripted_options(&script);
7158        options.budget.max_files = Some(1);
7159        let opened = open_fixture(root.path(), options).expect("open");
7160
7161        wait_until_phase(&opened, crate::LifecyclePhase::Stopped);
7162        let since = opened.state.index.since(crate::Clock::ZERO).expect("journal");
7163        assert!(!since.commits.iter().any(|commit| commit.state.iter().any(|transition| {
7164            matches!(
7165                transition,
7166                crate::StateTransition::IndexState { current, .. }
7167                    if current.phase == crate::LifecyclePhase::Watching
7168            )
7169        })));
7170        opened.close().expect("close");
7171    }
7172
7173    #[cfg(feature = "watch")]
7174    #[test]
7175    fn close_joins_an_observation_worker_blocked_at_a_named_boundary() {
7176        let root = tempfile::tempdir().expect("temp root");
7177        let scripts = tempfile::tempdir().expect("script root");
7178        let script = scripts.path().join("events.script");
7179        std::fs::write(&script, b"").expect("script");
7180        let controls = Arc::new(TestControls::default());
7181        controls.gate(TestPoint::BeforeObservationPoll).arm();
7182        let opened = OpenedIndex::open_for_test(
7183            root.path(),
7184            scripted_options(&script),
7185            Arc::clone(&controls),
7186        )
7187        .expect("open");
7188        controls.gate(TestPoint::BeforeObservationPoll).wait_reached();
7189
7190        let closer = opened.clone();
7191        let close = thread::spawn(move || closer.close());
7192        let deadline = std::time::Instant::now() + TEST_GATE_TIMEOUT;
7193        while !opened.state.cancellation.is_cancelled() {
7194            assert!(std::time::Instant::now() < deadline, "close did not cancel observation");
7195            thread::yield_now();
7196        }
7197        assert!(!close.is_finished(), "close returned before the owned worker was released");
7198        controls.gate(TestPoint::BeforeObservationPoll).release();
7199        close.join().expect("close thread").expect("joined close");
7200        opened.close().expect("repeat close");
7201    }
7202
7203    #[cfg(feature = "watch")]
7204    #[test]
7205    fn malformed_script_fails_before_discovery_starts() {
7206        let root = tempfile::tempdir().expect("temp root");
7207        let scripts = tempfile::tempdir().expect("script root");
7208        let script = scripts.path().join("events.script");
7209        std::fs::write(&script, b"teleport\tmissing\n").expect("script");
7210        let error = open_fixture(root.path(), scripted_options(&script))
7211            .expect_err("invalid observer configuration must fail open");
7212        assert!(matches!(error, Error::WatchScript(_)));
7213    }
7214
7215    #[cfg(feature = "watch")]
7216    #[test]
7217    fn observation_rejects_a_restricted_scope_before_open_returns() {
7218        let root = tempfile::tempdir().expect("temp root");
7219        let scripts = tempfile::tempdir().expect("script root");
7220        let script = scripts.path().join("events.script");
7221        std::fs::write(&script, b"").expect("script");
7222        let mut options = scripted_options(&script);
7223        options.one_filesystem = true;
7224
7225        let error = open_fixture(root.path(), options)
7226            .expect_err("unsupported observed scope must fail open");
7227        // A platform that cannot express this scope refuses it while planning,
7228        // before the observer's narrower live-scope contract is considered.
7229        #[cfg(not(unix))]
7230        assert!(matches!(
7231            error,
7232            Error::InvalidRequest(crate::query::RequestError::ScopeUnsupported {
7233                axis: crate::query::ScopeAxis::OneFilesystem,
7234                ..
7235            })
7236        ));
7237        #[cfg(unix)]
7238        assert!(matches!(error, Error::UnsupportedScanConfig(_)));
7239    }
7240
7241    #[cfg(feature = "watch")]
7242    #[test]
7243    fn handoff_retries_a_benign_refresh_conflict_before_watching() {
7244        let root = tempfile::tempdir().expect("temp root");
7245        let scripts = tempfile::tempdir().expect("script root");
7246        let path = root.path().join("shared.txt");
7247        std::fs::write(&path, b"before").expect("fixture");
7248        let script = scripts.path().join("events.script");
7249        std::fs::write(&script, b"").expect("script");
7250        let controls = Arc::new(TestControls::default());
7251        controls.gate(TestPoint::AfterObservationVerification).arm();
7252        let opened = OpenedIndex::open_for_test(
7253            root.path(),
7254            scripted_options(&script),
7255            Arc::clone(&controls),
7256        )
7257        .expect("open scripted observer");
7258        controls.gate(TestPoint::AfterObservationVerification).wait_reached();
7259
7260        std::fs::write(&path, b"updated-by-refresh").expect("concurrent mutation");
7261        let refreshed =
7262            opened.refresh(&[PathBuf::from("shared.txt")]).expect("overlapping refresh succeeds");
7263        assert_eq!(refreshed.work.stale, 0);
7264        controls.gate(TestPoint::AfterObservationVerification).release();
7265
7266        wait_until_phase(&opened, crate::LifecyclePhase::Watching);
7267        assert_eq!(
7268            opened
7269                .state
7270                .index
7271                .attrs(Path::new("shared.txt"))
7272                .expect("attrs")
7273                .expect("retained")
7274                .size,
7275            18
7276        );
7277        opened.close().expect("close");
7278    }
7279
7280    /// A refresh that commits the very facts the handoff is about to commit is not a
7281    /// conflict. The handoff's conditional upserts were refused as stale because their
7282    /// baselines had moved, so it walked the whole root again, and three such refreshes in
7283    /// a row failed the root over commits that would have applied as unchanged. A control
7284    /// file is the case with two ops on one baseline, the entry and its rules, and both
7285    /// must converge.
7286    #[cfg(feature = "watch")]
7287    #[test]
7288    fn handoff_settles_through_a_convergent_refresh_without_a_second_walk() {
7289        for name in ["shared.txt", ".gitignore"] {
7290            let root = tempfile::tempdir().expect("temp root");
7291            let scripts = tempfile::tempdir().expect("script root");
7292            let path = root.path().join(name);
7293            std::fs::write(&path, b"before").expect("fixture");
7294            let script = scripts.path().join("events.script");
7295            std::fs::write(&script, b"").expect("script");
7296            let controls = Arc::new(TestControls::default());
7297            controls.gate(TestPoint::BeforeObservationHandoff).arm();
7298            controls.gate(TestPoint::AfterObservationVerification).arm();
7299            let opened = OpenedIndex::open_for_test(
7300                root.path(),
7301                scripted_options(&script),
7302                Arc::clone(&controls),
7303            )
7304            .expect("open scripted observer");
7305            controls.gate(TestPoint::BeforeObservationHandoff).wait_reached();
7306            // Discovery retained six bytes; the handoff's walk is about to stat seven.
7307            std::fs::write(&path, b"changed").expect("mutation before the handoff walk");
7308            controls.gate(TestPoint::BeforeObservationHandoff).release();
7309            controls.gate(TestPoint::AfterObservationVerification).wait_reached();
7310
7311            // The refresh sees the same seven bytes and commits them first.
7312            let refreshed = opened.refresh(&[PathBuf::from(name)]).expect("refresh");
7313            assert_eq!(refreshed.work.stale, 0, "{name}");
7314            let ops = if name == crate::control::CONTROL_FILE_NAME { 2 } else { 1 };
7315            assert_eq!(refreshed.work.observations, ops, "{name}");
7316            controls.gate(TestPoint::AfterObservationVerification).release();
7317
7318            let state = wait_until_phase(&opened, crate::LifecyclePhase::Watching);
7319            assert_eq!(state.coverage, crate::Coverage::Complete, "{name}");
7320            assert_eq!(state.freshness, crate::Freshness::Fresh, "{name}");
7321            let since = opened.state.index.since(crate::Clock::ZERO).expect("journal");
7322            let transitions: Vec<&crate::StateTransition> =
7323                since.commits.iter().flat_map(|commit| commit.state.iter()).collect();
7324            // A refused pass leaves the root partial before the retry verifies it.
7325            assert!(
7326                !transitions.iter().any(|transition| matches!(
7327                    transition,
7328                    crate::StateTransition::Freshness { current: crate::Freshness::Partial, .. }
7329                )),
7330                "{name}: the handoff's first pass was refused: {transitions:?}"
7331            );
7332            assert_eq!(
7333                transitions
7334                    .iter()
7335                    .filter(|transition| matches!(
7336                        transition,
7337                        crate::StateTransition::Verified { path } if path.as_os_str().is_empty()
7338                    ))
7339                    .count(),
7340                1,
7341                "{name}"
7342            );
7343            let index = &opened.state.index;
7344            assert_eq!(
7345                index.attrs(Path::new(name)).expect("attrs").expect("retained").size,
7346                7,
7347                "{name}"
7348            );
7349            if name == crate::control::CONTROL_FILE_NAME {
7350                assert!(
7351                    index
7352                        .read_with(|index| {
7353                            index.controls().is_ok_and(|controls| {
7354                                controls.source_is(Path::new(name), b"changed")
7355                            })
7356                        })
7357                        .expect("controls"),
7358                    "the refreshed rules are retained"
7359                );
7360            }
7361            opened.close().expect("close");
7362        }
7363    }
7364
7365    #[cfg(feature = "watch")]
7366    #[test]
7367    fn budget_stop_wins_a_race_with_the_transition_to_watching() {
7368        let root = tempfile::tempdir().expect("temp root");
7369        let scripts = tempfile::tempdir().expect("script root");
7370        std::fs::write(root.path().join("baseline.txt"), b"baseline").expect("fixture");
7371        let script = scripts.path().join("events.script");
7372        std::fs::write(&script, b"").expect("script");
7373        let controls = Arc::new(TestControls::default());
7374        controls.gate(TestPoint::BeforeObservationWatching).arm();
7375        let mut options = scripted_options(&script);
7376        options.budget.max_files = Some(1);
7377        let opened =
7378            OpenedIndex::open_for_test(root.path(), options, Arc::clone(&controls)).expect("open");
7379        controls.gate(TestPoint::BeforeObservationWatching).wait_reached();
7380
7381        std::fs::write(root.path().join("over-budget.txt"), b"refused")
7382            .expect("over-budget mutation");
7383        let refreshed = opened
7384            .refresh(&[PathBuf::from("over-budget.txt")])
7385            .expect("resource refusal is a typed result");
7386        assert_eq!(refreshed.work.resource_refused, 1);
7387        assert_eq!(refreshed.state.phase, crate::LifecyclePhase::Stopped);
7388        controls.gate(TestPoint::BeforeObservationWatching).release();
7389
7390        let state = wait_until_phase(&opened, crate::LifecyclePhase::Stopped);
7391        assert_eq!(state.coverage, crate::Coverage::Partial(crate::CoverageReason::Budget));
7392        let since = opened.state.index.since(crate::Clock::ZERO).expect("journal");
7393        assert!(!since.commits.iter().any(|commit| commit.state.iter().any(|transition| {
7394            matches!(
7395                transition,
7396                crate::StateTransition::IndexState { current, .. }
7397                    if current.phase == crate::LifecyclePhase::Watching
7398            )
7399        })));
7400        opened.close().expect("close");
7401    }
7402}