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