use std::collections::BTreeMap;
use std::path::PathBuf;
use std::time::{Duration, Instant};
use crate::engine_contract::{Commit, EffectiveChange, EntryKind, Error, Result};
use crate::index::IndexHandle;
use crate::query::{Basis, Delivery, Query, Report, Request, Selection, WatchDelivery, report};
use crate::scan::ScanConfig;
use crate::watch::{WatchConfig, Watcher};
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct Change {
pub path: PathBuf,
pub kind: ChangeKind,
pub entry_kind: Option<EntryKind>,
pub bytes: Option<u64>,
pub allocated: Option<u64>,
pub mtime_ns: Option<i64>,
pub ignored: Option<bool>,
pub clock: u64,
}
#[derive(Clone, Copy, PartialEq, Eq, Debug)]
pub enum ChangeKind {
Upsert,
Remove,
Invalidate,
}
#[derive(Clone, Debug, Default)]
pub struct Batch {
pub changes: Vec<Change>,
pub dirty: bool,
}
struct BatchFacts {
ignored: Option<BTreeMap<PathBuf, bool>>,
reclassified: BTreeMap<PathBuf, EntryFacts>,
}
impl BatchFacts {
fn is_ignored(&self, path: &std::path::Path) -> Option<bool> {
self.ignored.as_ref()?.get(path).copied()
}
}
#[derive(Clone, Copy)]
struct EntryFacts {
kind: EntryKind,
bytes: u64,
allocated: u64,
mtime_ns: i64,
}
#[derive(Debug)]
pub enum SaveOutcome {
Written,
Skipped,
Failed(Error),
}
struct Persistence {
pending: bool,
last_attempt: Instant,
}
impl Persistence {
fn persist_due(
&mut self,
now: Instant,
interval: Duration,
save: impl FnOnce() -> Result<bool>,
) -> SaveOutcome {
if !save_is_due(self.pending, now.saturating_duration_since(self.last_attempt), interval) {
return SaveOutcome::Skipped;
}
let outcome = match save() {
Ok(true) => SaveOutcome::Written,
Ok(false) => SaveOutcome::Skipped,
Err(error) => SaveOutcome::Failed(error),
};
self.pending = pending_after(&outcome);
self.last_attempt = now;
outcome
}
}
fn save_is_due(pending: bool, since_last_save: Duration, interval: Duration) -> bool {
pending && since_last_save >= interval
}
fn pending_after(outcome: &SaveOutcome) -> bool {
!matches!(outcome, SaveOutcome::Written)
}
pub struct Session {
index: IndexHandle,
watcher: Watcher,
scan: ScanConfig,
request: Request,
plan: crate::Plan,
persistence: Persistence,
startup_save_error: Option<Error>,
presented: Option<u128>,
}
struct RepaintDigest {
hash: u128,
section: u64,
}
impl RepaintDigest {
const OFFSET_BASIS: u128 = 0x6c62_272e_07bb_0142_62b8_2175_6295_c58d;
const PRIME: u128 = 0x0000_0000_0100_0000_0000_0000_0000_013b;
const fn new() -> Self {
Self { hash: Self::OFFSET_BASIS, section: 0 }
}
fn mix(&mut self, bytes: &[u8]) {
for byte in bytes {
self.hash ^= u128::from(*byte);
self.hash = self.hash.wrapping_mul(Self::PRIME);
}
}
fn end_section(&mut self) {
let length = std::mem::take(&mut self.section);
self.mix(&length.to_le_bytes());
}
fn finish(mut self) -> u128 {
self.end_section();
self.hash
}
}
impl std::io::Write for RepaintDigest {
fn write(&mut self, bytes: &[u8]) -> std::io::Result<usize> {
self.mix(bytes);
self.section = self.section.wrapping_add(u64::try_from(bytes.len()).unwrap_or(u64::MAX));
Ok(bytes.len())
}
fn flush(&mut self) -> std::io::Result<()> {
Ok(())
}
}
impl Session {
pub fn start(request: Request, delivery: Delivery) -> Result<Self> {
Self::start_observed(request, delivery, None)
}
pub fn start_with_progress(
request: Request,
delivery: Delivery,
progress: &crate::Progress,
) -> Result<Self> {
Self::start_observed(request, delivery, Some(progress))
}
fn start_observed(
request: Request,
mut delivery: Delivery,
progress: Option<&crate::Progress>,
) -> Result<Self> {
delivery.watch.get_or_insert_with(WatchDelivery::default);
let plan =
crate::plan(&request, &delivery, crate::Route::Watch).map_err(Error::InvalidRequest)?;
let (index, report, pending, _diagnostics) =
crate::execute(&plan, &request.basis, false, progress)?;
let startup_save_error = pending.join().err();
let index = std::sync::Arc::into_inner(index)
.expect("the joined writer released the only other reference");
let mut session = Self::new_observed(
IndexHandle::new(index),
request,
&delivery,
WatchConfig::default(),
progress,
)?;
session.persistence.pending |= startup_save_error.is_some() || !report.is_complete();
session.startup_save_error = startup_save_error;
Ok(session)
}
pub fn persist_due(&mut self, now: Instant) -> SaveOutcome {
if let Some(error) = self.startup_save_error.take() {
self.persistence.last_attempt = now;
return SaveOutcome::Failed(error);
}
let interval = self.plan.delivery().watch.expect("watch plan").interval;
self.persistence.persist_due(now, interval, || {
if !self.plan.persists() || self.plan.delivery().cache_path.is_none() {
return Ok(false);
}
let index = self.index.snapshot()?;
crate::persist_index(&index, &self.plan)
})
}
pub fn new(
index: IndexHandle,
request: Request,
delivery: &Delivery,
watch: WatchConfig,
) -> Result<Self> {
Self::new_observed(index, request, delivery, watch, None)
}
fn new_observed(
index: IndexHandle,
request: Request,
delivery: &Delivery,
watch: WatchConfig,
progress: Option<&crate::Progress>,
) -> Result<Self> {
let root = index.root_path()?;
let scan = request.basis.scope.scan_config(delivery);
let delivery =
Delivery { watch: Some(delivery.watch.unwrap_or_default()), ..delivery.clone() };
crate::plan(&request, &delivery, crate::Route::Watch).map_err(Error::InvalidRequest)?;
crate::validate_basis_root(&root, &request.basis)?;
scan.validate_for_scope(index.scope()?)?;
let held = Basis {
root: root.clone(),
scope: scan.clone().into(),
content: index.read_with(crate::Index::content_set)?,
};
request.validate_read(&held).map_err(Error::InvalidRequest)?;
let watcher = Watcher::new(&root, watch)?;
Self::finish_initial_handoff(index, request, &delivery, watcher, scan, progress)
}
fn finish_initial_handoff(
index: IndexHandle,
request: Request,
delivery: &Delivery,
watcher: Watcher,
scan: ScanConfig,
progress: Option<&crate::Progress>,
) -> Result<Self> {
if let Some(progress) = progress {
progress.begin_pass(crate::ProgressPhase::Revalidating);
}
let observed = ScanConfig { progress: progress.cloned(), ..scan.clone() };
let mut dirty = false;
let reconciliation = crate::scan::reconcile_handle(&index, &observed, &mut |commit| {
dirty |= !commit.changes.is_empty();
})?;
if !reconciliation.scan.is_complete() && !delivery.accept_partial {
return Err(Error::ObservationHandoffIncomplete);
}
dirty |= drain_initial_capture(&watcher, &index, &observed)?;
if !delivery.accept_partial
&& !index.read_with(|index| crate::query::TreeStatus::of(index, &request).complete)?
{
return Err(Error::ObservationHandoffIncomplete);
}
let plan =
crate::plan(&request, delivery, crate::Route::Watch).map_err(Error::InvalidRequest)?;
Ok(Self {
index,
watcher,
scan,
request,
plan,
persistence: Persistence { pending: dirty, last_attempt: Instant::now() },
startup_save_error: None,
presented: None,
})
}
pub fn request(&self) -> &Request {
&self.request
}
pub fn query(&self) -> &Query {
&self.request.query
}
pub fn report(&self, generated_at: std::time::SystemTime) -> Result<Report> {
let index = self.index.snapshot()?;
report(&index, &self.request, generated_at)
}
pub fn changed_report(
&mut self,
generated_at: std::time::SystemTime,
format: crate::report_format::Format,
options: crate::report_format::RenderOptions,
) -> Result<Option<Report>> {
use std::io::Write as _;
let mut report = self.report(generated_at)?;
let pinned = std::time::SystemTime::UNIX_EPOCH;
let stamped = std::mem::replace(&mut report.provenance.generated_at, pinned);
let mut identity = RepaintDigest::new();
let rendered =
crate::report_format::write_with_options(&report, format, options, &mut identity)
.and_then(|()| {
identity.end_section();
write!(
identity,
"{:?}\n{:?}\n{:?}",
report.status, report.provenance.source, report.provenance.freshness
)?;
identity.end_section();
let notes = crate::report_format::diagnostic_lines(&report).into_lines();
for line in notes.iter().chain(&crate::report_format::report_warnings(&report))
{
writeln!(identity, "{line}")?;
}
Ok(())
});
report.provenance.generated_at = stamped;
rendered.map_err(|error| Error::io(&self.request.basis.root, error))?;
let identity = identity.finish();
if self.presented == Some(identity) {
return Ok(None);
}
self.presented = Some(identity);
Ok(Some(report))
}
pub fn changed_report_with_progress(
&mut self,
generated_at: std::time::SystemTime,
format: crate::report_format::Format,
options: crate::report_format::RenderOptions,
progress: &crate::Progress,
) -> Result<Option<Report>> {
progress.enter(crate::ProgressPhase::Summarizing);
self.changed_report(generated_at, format, options)
}
pub fn index_snapshot(&self) -> Result<crate::Index> {
self.index.snapshot()
}
pub fn next_batch(&mut self, timeout: Duration) -> Result<Option<Batch>> {
let mut commits: Vec<Commit> = Vec::new();
let outcome =
self.watcher.apply_next(&self.index, &self.scan, timeout, &mut |commit: &Commit| {
commits.push(commit.clone());
});
self.persistence.pending |= commits.iter().any(|commit| !commit.changes.is_empty());
let Some(_report) = outcome? else {
return Ok(None);
};
let mut batch = Batch {
changes: Vec::new(),
dirty: commits.iter().any(|commit| !commit.changes.is_empty()),
};
let facts = self.batch_facts(&commits)?;
for commit in &commits {
for effective in &commit.changes {
if let Some(change) = self.change_for(effective, commit.clock.0, &facts) {
batch.changes.push(change);
}
}
}
Ok(Some(batch))
}
fn batch_facts(&self, commits: &[Commit]) -> Result<BatchFacts> {
let mut touched: Vec<&PathBuf> = Vec::new();
let mut removed: Vec<(&PathBuf, crate::EntryKind)> = Vec::new();
let mut reclassified: Vec<&PathBuf> = Vec::new();
for effective in commits.iter().flat_map(|commit| &commit.changes) {
match effective {
EffectiveChange::Inserted { path, .. } | EffectiveChange::Updated { path, .. } => {
touched.push(path);
}
EffectiveChange::Reclassified { path, .. } => {
reclassified.push(path);
}
EffectiveChange::Removed { path, kind, .. } => removed.push((path, *kind)),
EffectiveChange::Invalidated { .. }
| EffectiveChange::ControlUpdated { .. }
| EffectiveChange::ControlRefusalUpdated { .. } => {}
}
}
self.index.read_with(|index| {
let observed = index.observes_controls();
let entries = reclassified
.into_iter()
.filter_map(|path| {
let id = index.lookup(path)?;
let attrs = index.attrs_of(id)?;
Some((
path.clone(),
EntryFacts {
kind: index.kind_of(id)?,
bytes: attrs.size,
allocated: attrs.allocated,
mtime_ns: attrs.mtime_ns,
},
))
})
.collect();
BatchFacts {
ignored: observed.then(|| {
let mut ignored = touched
.into_iter()
.filter_map(|path| match index.is_ignored(path) {
Ok(Some(ignored)) => Some((path.clone(), ignored)),
Ok(None) | Err(_) => None,
})
.collect::<BTreeMap<_, _>>();
for (path, kind) in removed {
if index.control_classification_known(path) {
ignored.insert(
path.clone(),
index.control_table().is_ignored(path, kind.is_dir()),
);
}
}
ignored
}),
reclassified: entries,
}
})
}
fn change_for(
&self,
effective: &EffectiveChange,
clock: u64,
facts: &BatchFacts,
) -> Option<Change> {
match effective {
EffectiveChange::Inserted { path, kind, attrs } => {
let name = path.file_name()?.to_string_lossy().into_owned();
let candidate = crate::query::Candidate {
relative: path,
name: &name,
kind: *kind,
bytes: attrs.size,
allocated: attrs.allocated,
mtime_ns: attrs.mtime_ns,
ignored: facts.is_ignored(path).unwrap_or(false),
};
self.selection().admits(&candidate).then(|| Change {
path: path.clone(),
kind: ChangeKind::Upsert,
entry_kind: Some(*kind),
bytes: Some(attrs.size),
allocated: Some(attrs.allocated),
mtime_ns: Some(attrs.mtime_ns),
ignored: facts.is_ignored(path),
clock,
})
}
EffectiveChange::Updated { path, kind, previous: _, current } => {
let name = path.file_name()?.to_string_lossy().into_owned();
let ignored = facts.is_ignored(path).unwrap_or(false);
let candidate = crate::query::Candidate {
relative: path,
name: &name,
kind: *kind,
bytes: current.size,
allocated: current.allocated,
mtime_ns: current.mtime_ns,
ignored,
};
if self.selection().admits(&candidate) {
Some(Change {
path: path.clone(),
kind: ChangeKind::Upsert,
entry_kind: Some(*kind),
bytes: Some(current.size),
allocated: Some(current.allocated),
mtime_ns: Some(current.mtime_ns),
ignored: facts.is_ignored(path),
clock,
})
} else if self.admits_by_path(path, &name) {
Some(Change {
path: path.clone(),
kind: ChangeKind::Remove,
entry_kind: None,
bytes: None,
allocated: None,
mtime_ns: None,
ignored: None,
clock,
})
} else {
None
}
}
EffectiveChange::Removed { path, .. } => {
let name = path.file_name()?.to_string_lossy().into_owned();
self.admits_by_path(path, &name).then(|| Change {
path: path.clone(),
kind: ChangeKind::Remove,
entry_kind: None,
bytes: None,
allocated: None,
mtime_ns: None,
ignored: facts.is_ignored(path),
clock,
})
}
EffectiveChange::Invalidated { path, .. } => Some(Change {
path: path.clone(),
kind: ChangeKind::Invalidate,
entry_kind: None,
bytes: None,
allocated: None,
mtime_ns: None,
ignored: None,
clock,
}),
EffectiveChange::Reclassified { path, previous_ignored, current_ignored } => {
let name = path.file_name()?.to_string_lossy().into_owned();
let entry = facts.reclassified.get(path)?;
let admits = |ignored: bool| {
self.selection().admits(&crate::query::Candidate {
relative: path,
name: &name,
kind: entry.kind,
bytes: entry.bytes,
allocated: entry.allocated,
mtime_ns: entry.mtime_ns,
ignored,
})
};
match (admits(*previous_ignored), admits(*current_ignored)) {
(true, false) => Some(Change {
path: path.clone(),
kind: ChangeKind::Remove,
entry_kind: None,
bytes: None,
allocated: None,
mtime_ns: None,
ignored: Some(*current_ignored),
clock,
}),
(_, true) => Some(Change {
path: path.clone(),
kind: ChangeKind::Upsert,
entry_kind: Some(entry.kind),
bytes: Some(entry.bytes),
allocated: Some(entry.allocated),
mtime_ns: Some(entry.mtime_ns),
ignored: Some(*current_ignored),
clock,
}),
_ => None,
}
}
EffectiveChange::ControlUpdated { .. }
| EffectiveChange::ControlRefusalUpdated { .. } => None,
}
}
fn admits_by_path(&self, path: &std::path::Path, name: &str) -> bool {
let selection = self.selection();
if selection.exclude.iter().any(|pattern| pattern.matches(path, name)) {
return false;
}
selection.include.is_empty()
|| selection.include.iter().any(|pattern| pattern.matches(path, name))
}
fn selection(&self) -> &Selection {
&self.request.query.selection
}
}
fn drain_initial_capture(
watcher: &Watcher,
index: &IndexHandle,
scan: &ScanConfig,
) -> Result<bool> {
let mut dirty = false;
for _ in 0..2 {
watcher.flush_capture()?;
let mut drained = false;
for _ in 0..=watcher.capture_backlog_bound() {
if watcher
.apply_next(index, scan, Duration::ZERO, &mut |commit| {
dirty |= !commit.changes.is_empty();
})?
.is_none()
{
drained = true;
break;
}
}
if !drained {
return Err(Error::ObservationHandoffIncomplete);
}
}
Ok(dirty)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn a_retained_session_rejects_a_request_for_another_root_before_binding() {
let a = tempfile::tempdir().expect("root a");
let b = tempfile::tempdir().expect("root b");
let basis = Basis {
root: a.path().into(),
scope: crate::query::Scope::default(),
content: crate::content::AnalysisSet::NONE,
};
let delivery = Delivery::new(crate::CachePolicy::Off, None);
let (index, _) = crate::open(&basis, &delivery).expect("open a");
let handle = IndexHandle::new(index);
let before = handle.clock().expect("clock");
let request = Request::new(
Basis { root: b.path().into(), ..basis },
Query::default(),
std::time::SystemTime::now(),
);
assert!(matches!(
Session::new(handle.clone(), request, &delivery, WatchConfig::default()),
Err(Error::InvalidRequest(crate::query::RequestError::RootMismatch { .. }))
));
assert_eq!(handle.clock().expect("clock"), before);
}
#[test]
fn a_save_is_due_only_when_a_change_is_pending_and_the_throttle_has_elapsed() {
let interval = Duration::from_secs(1);
let cases = [
(true, Duration::from_secs(2), true, "pending and past the interval"),
(true, interval, true, "pending, exactly at the interval: inclusive"),
(true, Duration::from_millis(1), false, "pending but throttled"),
(false, Duration::from_secs(60), false, "nothing pending, however long it has been"),
(false, Duration::ZERO, false, "nothing pending and just saved"),
];
for (pending, since, want, case) in cases {
assert_eq!(save_is_due(pending, since, interval), want, "{case}");
}
}
#[test]
fn only_a_completed_write_clears_the_pending_change() {
assert!(!pending_after(&SaveOutcome::Written), "a completed write persists the change");
assert!(
pending_after(&SaveOutcome::Skipped),
"a skipped save wrote nothing, so the change is still owed to disk",
);
assert!(
pending_after(&SaveOutcome::Failed(Error::Snapshot("failed".into()))),
"a failed save must be retried, not forgotten"
);
}
#[test]
fn a_burst_then_a_quiet_tree_still_persists() {
let interval = Duration::from_secs(1);
let mut pending = true;
assert!(!save_is_due(pending, Duration::from_millis(50), interval));
assert!(pending, "the throttle must not consume the change");
assert!(save_is_due(pending, Duration::from_secs(3), interval));
pending = pending_after(&SaveOutcome::Skipped);
assert!(pending);
pending = pending_after(&SaveOutcome::Written);
assert!(!pending, "once written, the loop stops rewriting an unchanged index");
}
#[test]
fn skips_and_failures_retry_only_after_another_interval() {
let start = Instant::now();
let interval = Duration::from_secs(2);
let mut persistence = Persistence { pending: true, last_attempt: start };
assert!(matches!(
persistence.persist_due(start + interval / 2, interval, || panic!("throttled")),
SaveOutcome::Skipped
));
assert!(matches!(
persistence.persist_due(start + interval, interval, || Ok(false)),
SaveOutcome::Skipped
));
assert!(persistence.pending);
assert!(matches!(
persistence.persist_due(start + interval, interval, || panic!("skip was throttled")),
SaveOutcome::Skipped
));
assert!(matches!(
persistence.persist_due(start + interval * 2, interval, || {
Err(Error::Snapshot("disk unavailable".into()))
}),
SaveOutcome::Failed(_)
));
assert!(persistence.pending);
assert!(matches!(
persistence
.persist_due(start + interval * 2, interval, || panic!("failure was throttled")),
SaveOutcome::Skipped
));
assert!(matches!(
persistence.persist_due(start + interval * 3, interval, || Ok(true)),
SaveOutcome::Written
));
assert!(!persistence.pending);
assert!(matches!(
persistence.persist_due(start + interval * 4, interval, || panic!("already persisted")),
SaveOutcome::Skipped
));
}
#[test]
fn handoff_changes_are_persisted_after_the_tree_goes_quiet() {
let root = tempfile::tempdir().expect("root");
let cache = tempfile::tempdir().expect("cache");
let cache_path = cache.path().join("snapshot");
let scan = ScanConfig::default();
std::fs::write(root.path().join("before.txt"), b"before").expect("before");
let (index, _) = crate::scan::scan_into_index(root.path(), &scan).expect("scan");
let request = Request::new(
Basis {
root: root.path().to_path_buf(),
scope: scan.clone().into(),
content: crate::content::AnalysisSet::NONE,
},
Query::default(),
std::time::SystemTime::now(),
);
let interval = Duration::from_secs(2);
let delivery = Delivery {
stale_ok: false,
cache: crate::CachePolicy::Auto,
cache_path: Some(cache_path.clone()),
accept_partial: false,
watch: Some(WatchDelivery { interval }),
workers: crate::query::Workers::default(),
batch_size: ScanConfig::default().batch_size,
order: crate::scan::ScanOrder::default(),
};
let script = tempfile::NamedTempFile::new().expect("script");
let (watcher, _sender) =
Watcher::scripted(root.path(), WatchConfig::default(), script.path()).expect("watcher");
std::fs::write(root.path().join("during-handoff.txt"), b"handoff").expect("handoff change");
let mut session = Session::finish_initial_handoff(
IndexHandle::new(index),
request,
&delivery,
watcher,
scan,
None,
)
.expect("handoff");
let started = session.persistence.last_attempt;
assert!(session.persistence.pending, "handoff changes need persistence too");
assert!(matches!(session.persist_due(started), SaveOutcome::Skipped));
assert!(!cache_path.exists(), "throttle delays the write");
assert!(matches!(session.persist_due(started + interval), SaveOutcome::Written));
let restored = crate::snapshot::load(&cache_path).expect("load").expect("saved");
assert!(matches!(
restored.path_state(std::path::Path::new("during-handoff.txt")),
crate::PathState::Present { .. }
));
assert!(matches!(session.persist_due(started + interval * 2), SaveOutcome::Skipped));
}
#[test]
fn startup_save_failure_keeps_the_session_live_and_retries() {
let root = tempfile::tempdir().expect("root");
let cache = tempfile::tempdir().expect("cache");
let parent = cache.path().join("blocked");
#[cfg(unix)]
std::os::unix::fs::symlink(cache.path().join("missing"), &parent)
.expect("dangling parent blocks the write");
#[cfg(not(unix))]
std::fs::write(&parent, b"").expect("file parent blocks the write");
let restore = || {
#[cfg(unix)]
std::fs::create_dir(cache.path().join("missing")).expect("restore the parent");
#[cfg(not(unix))]
std::fs::remove_file(&parent).expect("restore the parent");
};
let cache_path = parent.join("snapshot");
std::fs::write(root.path().join("file.txt"), b"content").expect("file");
let interval = Duration::from_secs(2);
let request = Request::new(
Basis {
root: root.path().to_path_buf(),
scope: ScanConfig::default().into(),
content: crate::content::AnalysisSet::NONE,
},
Query::default(),
std::time::SystemTime::now(),
);
let delivery = Delivery {
stale_ok: false,
cache: crate::CachePolicy::On,
cache_path: Some(cache_path.clone()),
accept_partial: false,
watch: Some(WatchDelivery { interval }),
workers: crate::query::Workers::default(),
batch_size: ScanConfig::default().batch_size,
order: crate::scan::ScanOrder::default(),
};
let mut session = Session::start(request, delivery).expect("save failure is nonfatal");
assert!(session.report(std::time::SystemTime::now()).expect("live report").status.complete);
let now = Instant::now();
assert!(matches!(session.persist_due(now), SaveOutcome::Failed(_)));
restore();
assert!(matches!(session.persist_due(now), SaveOutcome::Skipped));
assert!(matches!(session.persist_due(now + interval), SaveOutcome::Written));
assert!(crate::snapshot::load(&cache_path).expect("read snapshot").is_some());
}
#[test]
fn an_update_that_leaves_attribute_selection_emits_remove() {
let root = tempfile::tempdir().expect("tempdir");
std::fs::write(root.path().join("file.txt"), b"12345678").expect("fixture");
let scan = ScanConfig::default();
let (index, report) = crate::scan::scan_into_index(root.path(), &scan).expect("scan");
assert!(report.is_complete());
let request = Request::new(
Basis {
root: root.path().to_path_buf(),
scope: scan.clone().into(),
content: crate::content::AnalysisSet::NONE,
},
Query {
selection: Selection { min_size: Some(4), ..Selection::default() },
..Query::default()
},
std::time::SystemTime::now(),
);
let delivery = Delivery {
stale_ok: false,
cache: crate::CachePolicy::Off,
cache_path: None,
accept_partial: false,
watch: Some(WatchDelivery { interval: Duration::from_millis(50) }),
workers: crate::query::Workers::default(),
batch_size: ScanConfig::default().batch_size,
order: crate::scan::ScanOrder::default(),
};
let session =
Session::new(IndexHandle::new(index), request, &delivery, WatchConfig::default())
.expect("session");
let path = PathBuf::from("file.txt");
let change = session
.change_for(
&EffectiveChange::Updated {
path: path.clone(),
kind: EntryKind::File,
previous: crate::Attrs { size: 8, allocated: 8, ..crate::Attrs::default() },
current: crate::Attrs { size: 1, allocated: 1, ..crate::Attrs::default() },
},
1,
&BatchFacts {
ignored: Some(BTreeMap::from([(path, false)])),
reclassified: BTreeMap::new(),
},
)
.expect("membership transition");
assert_eq!(change.kind, ChangeKind::Remove);
}
#[test]
fn initial_handoff_drains_a_sticky_overflow_after_a_full_intent_queue() {
let root = tempfile::tempdir().expect("root");
let script = tempfile::NamedTempFile::new().expect("script");
let scan = ScanConfig::default();
let (index, report) = crate::scan::scan_into_index(root.path(), &scan).expect("scan");
assert!(report.is_complete());
let handle = IndexHandle::new(index);
let config = WatchConfig {
settle: Duration::from_millis(1),
max_hold: Duration::from_millis(2),
event_capacity: 8,
batch_path_capacity: 1,
intent_capacity: 1,
..WatchConfig::default()
};
let (watcher, sender) =
Watcher::scripted(root.path(), config, script.path()).expect("scripted watcher");
std::fs::write(root.path().join("a.txt"), b"a").expect("a");
sender.send("create\ta.txt\n").expect("first event");
watcher.flush_capture().expect("first barrier fills the intent queue");
std::fs::write(root.path().join("b.txt"), b"b").expect("b");
sender.send("create\tb.txt\n").expect("second event");
watcher.flush_capture().expect("second barrier retains sticky overflow");
drain_initial_capture(&watcher, &handle, &scan).expect("bounded handoff");
assert!(
handle.snapshot().expect("snapshot").lookup(std::path::Path::new("a.txt")).is_some()
);
assert!(
handle.snapshot().expect("snapshot").lookup(std::path::Path::new("b.txt")).is_some()
);
assert!(
watcher
.apply_next(&handle, &scan, Duration::ZERO, &mut |_| {})
.expect("proof poll")
.is_none(),
"no queued or sticky pre-handoff work remains"
);
}
#[test]
fn initial_handoff_enforces_partial_acceptance_after_its_reconciliation() {
let root = tempfile::tempdir().expect("root");
std::fs::write(root.path().join("kept.txt"), b"kept").expect("fixture");
let scan = ScanConfig::default();
let (index, report) = crate::scan::scan_into_index(root.path(), &scan).expect("scan");
assert!(report.is_complete());
let request = Request::new(
Basis {
root: root.path().to_path_buf(),
scope: scan.clone().into(),
content: crate::content::AnalysisSet::NONE,
},
Query::default(),
std::time::SystemTime::now(),
);
let _fault = crate::scan::install_walk_hook(root.path(), |_| {
Some(std::io::Error::new(
std::io::ErrorKind::PermissionDenied,
"deterministic handoff refusal",
))
});
let delivery = Delivery {
stale_ok: false,
cache: crate::CachePolicy::Off,
cache_path: None,
accept_partial: false,
watch: Some(WatchDelivery { interval: Duration::from_millis(50) }),
workers: crate::query::Workers::default(),
batch_size: ScanConfig::default().batch_size,
order: crate::scan::ScanOrder::default(),
};
let Err(error) = Session::new(
IndexHandle::new(index.clone()),
request.clone(),
&delivery,
WatchConfig::default(),
) else {
panic!("a partial handoff is refused");
};
assert!(matches!(error, Error::ObservationHandoffIncomplete));
let accepted = Session::new(
IndexHandle::new(index),
request,
&Delivery { accept_partial: true, ..delivery },
WatchConfig::default(),
)
.expect("the caller explicitly accepts a partial handoff");
assert!(
!accepted.report(std::time::SystemTime::now()).expect("partial report").status.complete
);
}
#[test]
fn initial_handoff_rechecks_partial_acceptance_after_draining_capture() {
use std::sync::Arc;
use std::sync::atomic::{AtomicUsize, Ordering};
let root = tempfile::tempdir().expect("root");
std::fs::write(root.path().join("kept.txt"), b"kept").expect("fixture");
let scan = ScanConfig::default();
let (index, report) = crate::scan::scan_into_index(root.path(), &scan).expect("scan");
assert!(report.is_complete());
let request = Request::new(
Basis {
root: root.path().to_path_buf(),
scope: scan.clone().into(),
content: crate::content::AnalysisSet::NONE,
},
Query::default(),
std::time::SystemTime::now(),
);
let delivery = Delivery {
stale_ok: false,
cache: crate::CachePolicy::Off,
cache_path: None,
accept_partial: false,
watch: Some(WatchDelivery { interval: Duration::from_millis(50) }),
workers: crate::query::Workers::default(),
batch_size: ScanConfig::default().batch_size,
order: crate::scan::ScanOrder::default(),
};
let script = tempfile::NamedTempFile::new().expect("script");
let (watcher, sender) =
Watcher::scripted(root.path(), WatchConfig::default(), script.path()).expect("watcher");
sender.send("rescan\t.\n").expect("queue initial gap");
let attempts = Arc::new(AtomicUsize::new(0));
let hook_attempts = Arc::clone(&attempts);
let _fault = crate::scan::install_walk_hook(root.path(), move |_| {
(hook_attempts.fetch_add(1, Ordering::SeqCst) > 0).then(|| {
std::io::Error::new(
std::io::ErrorKind::PermissionDenied,
"deterministic drain-only refusal",
)
})
});
let Err(error) = Session::finish_initial_handoff(
IndexHandle::new(index),
request,
&delivery,
watcher,
scan,
None,
) else {
panic!("a partial state created while draining is refused");
};
assert!(matches!(error, Error::ObservationHandoffIncomplete));
assert!(attempts.load(Ordering::SeqCst) > 1, "the drain ran after startup reconciliation");
}
#[test]
fn the_handoff_pass_restarts_the_counts() {
let root = tempfile::tempdir().expect("root");
let mut bytes = 0;
for directory in 0..2 {
let dir = root.path().join(format!("d{directory}"));
std::fs::create_dir(&dir).expect("directory");
for file in 0..3 {
let size = directory * 3 + file + 1;
std::fs::write(dir.join(format!("f{file}.txt")), vec![b'.'; size]).expect("file");
bytes += size as u64;
}
}
let scan = ScanConfig::default();
let (index, report) = crate::scan::scan_into_index(root.path(), &scan).expect("scan");
assert!(report.is_complete());
let request = Request::new(
Basis {
root: root.path().to_path_buf(),
scope: scan.clone().into(),
content: crate::content::AnalysisSet::NONE,
},
Query { views: vec![crate::query::ViewSpec::Summary], ..Query::default() },
std::time::SystemTime::now(),
);
let delivery = Delivery {
stale_ok: false,
cache: crate::CachePolicy::Off,
cache_path: None,
accept_partial: false,
watch: Some(WatchDelivery { interval: Duration::from_millis(50) }),
workers: crate::query::Workers::default(),
batch_size: ScanConfig::default().batch_size,
order: crate::scan::ScanOrder::default(),
};
let script = tempfile::NamedTempFile::new().expect("script");
let (watcher, _sender) =
Watcher::scripted(root.path(), WatchConfig::default(), script.path()).expect("watcher");
let progress = crate::Progress::new();
progress.add_walked(100, 100, 100, 100);
progress.enter(crate::ProgressPhase::Saving);
let session = Session::finish_initial_handoff(
IndexHandle::new(index.clone()),
request.clone(),
&delivery,
watcher,
scan.clone(),
Some(&progress),
)
.expect("handoff");
let snapshot = progress.snapshot();
assert_eq!(snapshot.phase, crate::ProgressPhase::Revalidating);
assert_eq!(
(snapshot.directories, snapshot.files, snapshot.bytes),
(3, 6, bytes),
"one walk of the root and its two directories, the first pass not added in"
);
assert_eq!(snapshot.allocated, report.allocated_walked, "allocated restarts with them");
let (plain_watcher, _plain_sender) =
Watcher::scripted(root.path(), WatchConfig::default(), script.path()).expect("watcher");
let plain = Session::finish_initial_handoff(
IndexHandle::new(index),
request,
&delivery,
plain_watcher,
scan,
None,
)
.expect("plain handoff");
let generated_at = std::time::SystemTime::now();
let observed_report = session.report(generated_at).expect("observed report");
let mut plain_report = plain.report(generated_at).expect("plain report");
plain_report.provenance = observed_report.provenance.clone();
let json = |report: &Report| {
crate::report_format::render(report, crate::report_format::Format::Json, false)
.expect("render")
};
assert_eq!(json(&plain_report), json(&observed_report));
assert_eq!(progress.snapshot(), snapshot, "the second handoff and reads were not observed");
}
#[test]
fn a_record_claims_a_classification_only_for_an_entry_the_index_still_holds() {
let unobserved = BatchFacts { ignored: None, reclassified: BTreeMap::new() };
assert_eq!(unobserved.is_ignored(std::path::Path::new("any.txt")), None);
let observed = BatchFacts {
ignored: Some(BTreeMap::from([
(PathBuf::from("build/out.bin"), true),
(PathBuf::from("src/main.rs"), false),
])),
reclassified: BTreeMap::new(),
};
assert_eq!(observed.is_ignored(std::path::Path::new("build/out.bin")), Some(true));
assert_eq!(observed.is_ignored(std::path::Path::new("src/main.rs")), Some(false));
assert_eq!(
observed.is_ignored(std::path::Path::new("gone.tmp")),
None,
"an entry the batch removed is in neither partition, not in the unignored one"
);
}
#[test]
fn a_started_session_reports_its_second_pass_and_then_nothing() {
let root = tempfile::tempdir().expect("root");
let cache = tempfile::tempdir().expect("cache");
let mut bytes = 0;
for directory in 0..4 {
let dir = root.path().join(format!("d{directory}"));
std::fs::create_dir(&dir).expect("directory");
for file in 0..3 {
let size = directory * 3 + file + 1;
std::fs::write(dir.join(format!("f{file}.txt")), vec![b'.'; size]).expect("file");
bytes += size as u64;
}
}
let request = || {
Request::new(
Basis {
root: root.path().to_path_buf(),
scope: ScanConfig::default().into(),
content: crate::content::AnalysisSet::NONE,
},
Query::default(),
std::time::UNIX_EPOCH,
)
};
let delivery = Delivery {
stale_ok: false,
cache: crate::CachePolicy::Auto,
cache_path: Some(cache.path().join("snapshot")),
accept_partial: false,
watch: Some(WatchDelivery { interval: Duration::from_secs(2) }),
workers: crate::query::Workers::default(),
batch_size: 4,
order: crate::scan::ScanOrder::default(),
};
let progress = crate::Progress::new();
let session = Session::start_with_progress(request(), delivery.clone(), &progress)
.expect("observed start");
let after_start = progress.snapshot();
assert_eq!(
after_start.phase,
crate::ProgressPhase::Revalidating,
"the handoff revalidation follows the joined save"
);
assert!(after_start.directories >= 5, "the root and four children: {after_start:?}");
assert!(after_start.files >= 12, "{after_start:?}");
assert!(after_start.bytes >= bytes, "{after_start:?}");
assert_eq!(after_start.analysis, None);
let plain = Session::start(request(), delivery).expect("plain start");
let generated_at = std::time::SystemTime::now();
let observed_report = session.report(generated_at).expect("observed report");
let mut plain_report = plain.report(generated_at).expect("plain report");
let setup_race = format!("{:?}", crate::InvalidateReason::WatchSetupRace);
let errors = &observed_report.status.errors;
let setup_race_only = !errors.is_empty()
&& errors.iter().all(|issue| {
issue.kind == crate::IssueKind::ObservationGap
&& issue.message.ends_with(&setup_race)
});
assert!(
observed_report.status.complete || setup_race_only,
"the watched answer is complete or carries only setup-race gaps: {:?}",
observed_report.status
);
assert_eq!(observed_report.status.errors_omitted, 0, "{:?}", observed_report.status);
assert!(plain_report.status.complete, "{:?}", plain_report.status);
assert_eq!(
observed_report.provenance.freshness, plain_report.provenance.freshness,
"the watched answer is as fresh as the plain one"
);
if errors.is_empty() {
assert_eq!(
observed_report.status.coverage, plain_report.status.coverage,
"with no retained issue, the coverage is the plain answer's"
);
}
plain_report.provenance = observed_report.provenance.clone();
plain_report.status = observed_report.status.clone();
let json = |report: &Report| {
crate::report_format::render(report, crate::report_format::Format::Json, false)
.expect("render")
};
assert_eq!(json(&plain_report), json(&observed_report), "the same facts either way");
assert_eq!(progress.snapshot(), after_start, "the second start was not observed");
}
#[test]
fn the_first_answer_is_built_under_the_summarizing_phase() {
let root = tempfile::tempdir().expect("root");
std::fs::write(root.path().join("a.txt"), b"alpha").expect("file");
std::fs::create_dir(root.path().join("d")).expect("directory");
std::fs::write(root.path().join("d/b.txt"), b"beta").expect("nested file");
let (mut session, _sender) = scripted_session(root.path(), Query::default());
let format = crate::report_format::Format::Json;
let options = crate::report_format::RenderOptions { color: false, bar_size: 0 };
let generated_at = std::time::SystemTime::now();
let plain_answer = session.report(generated_at).expect("plain answer");
let progress = crate::Progress::new();
progress.enter(crate::ProgressPhase::Revalidating);
let observed = session
.changed_report_with_progress(generated_at, format, options, &progress)
.expect("observed first answer")
.expect("a session's first answer is always given");
assert_eq!(
progress.snapshot().phase,
crate::ProgressPhase::Summarizing,
"the build is reported as the answer's construction"
);
let json =
|report: &Report| crate::report_format::render(report, format, false).expect("render");
assert_eq!(json(&observed), json(&plain_answer), "the same answer either way");
assert!(
session.changed_report(generated_at, format, options).expect("repaint").is_none(),
"nothing changed since the observed build, so nothing repaints"
);
}
fn scripted_session(
root: &std::path::Path,
query: Query,
) -> (Session, crate::watch::ScriptedSender) {
scripted_session_under(root, query, ScanConfig::default())
}
fn scripted_session_under(
root: &std::path::Path,
query: Query,
scan: ScanConfig,
) -> (Session, crate::watch::ScriptedSender) {
let (index, report) = crate::scan::scan_into_index(root, &scan).expect("scan");
assert!(report.is_complete());
let request = Request::new(
Basis {
root: root.to_path_buf(),
scope: scan.clone().into(),
content: crate::content::AnalysisSet::NONE,
},
query,
std::time::SystemTime::now(),
);
let delivery = Delivery {
stale_ok: false,
cache: crate::CachePolicy::Off,
cache_path: None,
accept_partial: false,
watch: Some(WatchDelivery { interval: Duration::from_millis(50) }),
workers: crate::query::Workers::default(),
batch_size: ScanConfig::default().batch_size,
order: crate::scan::ScanOrder::default(),
};
let script = tempfile::NamedTempFile::new().expect("script");
let (watcher, sender) =
Watcher::scripted(root, WatchConfig::default(), script.path()).expect("watcher");
let session = Session::finish_initial_handoff(
IndexHandle::new(index),
request,
&delivery,
watcher,
scan,
None,
)
.expect("handoff");
(session, sender)
}
fn touch(path: &std::path::Path) {
let file = std::fs::File::options().write(true).open(path).expect("open for touch");
let modified = file.metadata().expect("metadata").modified().expect("mtime");
file.set_modified(modified + Duration::from_secs(5)).expect("set mtime");
}
#[test]
fn an_aggregate_repaint_is_skipped_when_a_reader_would_see_no_change() {
use crate::report_format::{Format, RenderOptions};
let root = tempfile::tempdir().expect("root");
let big = root.path().join("big.txt");
std::fs::write(&big, vec![b'x'; 4096]).expect("big");
std::fs::write(root.path().join("small.txt"), b"small").expect("small");
let query = Query {
views: vec![crate::query::ViewSpec::Summary],
selection: Selection {
size: crate::query::SizeMetric::Apparent,
min_size: Some(4096),
..Selection::default()
},
..Query::default()
};
let (mut session, sender) = scripted_session(root.path(), query);
let now = std::time::SystemTime::now;
let text = |session: &mut Session| {
session.changed_report(now(), Format::Text, RenderOptions::default()).expect("report")
};
let summary = |report: &Report| match report.sections.first() {
Some(crate::query::Section::Summary(row)) => (row.files, row.bytes),
other => panic!("expected a summary, got {other:?}"),
};
let first = text(&mut session).expect("the first answer is always given");
assert_eq!(summary(&first), (1, 4096));
assert!(session.next_batch(Duration::from_millis(50)).expect("idle").is_none());
assert!(text(&mut session).is_none(), "an idle tree repaints nothing");
touch(&big);
sender.send("modify\tbig.txt\n").expect("script a touch");
let batch = session.next_batch(Duration::from_secs(10)).expect("touch").expect("observed");
assert!(batch.dirty, "the index records the new modification time");
assert!(text(&mut session).is_none(), "a touch that changes no size repaints nothing");
std::fs::write(root.path().join("small.txt"), b"still small").expect("grow small");
sender.send("modify\tsmall.txt\n").expect("script a filtered change");
let batch = session.next_batch(Duration::from_secs(10)).expect("small").expect("observed");
assert!(batch.dirty);
assert!(text(&mut session).is_none(), "a filtered change repaints nothing");
std::fs::write(&big, vec![b'x'; 8192]).expect("grow big");
sender.send("modify\tbig.txt\n").expect("script a visible change");
let batch = session.next_batch(Duration::from_secs(10)).expect("big").expect("observed");
assert!(batch.dirty);
let repainted = text(&mut session).expect("a visible change repaints");
assert_eq!(summary(&repainted), (1, 8192));
assert!(text(&mut session).is_none(), "and only once");
assert_eq!(summary(&session.report(now()).expect("report")), (1, 8192));
assert!(text(&mut session).is_none());
}
#[test]
fn a_repaint_identity_is_what_the_format_renders() {
use crate::report_format::{Format, RenderOptions};
let root = tempfile::tempdir().expect("root");
let big = root.path().join("big.txt");
std::fs::write(&big, vec![b'x'; 4096]).expect("big");
let query = Query { views: vec![crate::query::ViewSpec::Summary], ..Query::default() };
let (mut session, sender) = scripted_session(root.path(), query);
let json = |session: &mut Session| {
session
.changed_report(
std::time::SystemTime::now(),
Format::Json,
RenderOptions::default(),
)
.expect("report")
};
assert!(json(&mut session).is_some());
assert!(json(&mut session).is_none(), "a second generation instant alone is no change");
touch(&big);
sender.send("modify\tbig.txt\n").expect("script a touch");
assert!(session.next_batch(Duration::from_secs(10)).expect("touch").is_some());
assert!(json(&mut session).is_some(), "JSON shows the newest modification time");
}
#[test]
fn a_repaint_digest_is_fnv1a_128_framed_by_section() {
use std::io::Write as _;
let mut raw = RepaintDigest::new();
raw.mix(b"a");
assert_eq!(raw.hash, 0xd228_cb69_6f1a_8caf_7891_2b70_4e4a_8964, "the published vector");
let digest = |sections: &[&[u8]]| {
let mut digest = RepaintDigest::new();
for (at, section) in sections.iter().enumerate() {
if at > 0 {
digest.end_section();
}
digest.write_all(section).expect("digest");
}
digest.finish()
};
assert_eq!(digest(&[b"ab", b"c"]), digest(&[b"ab", b"c"]));
assert_ne!(digest(&[b"ab", b"c"]), digest(&[b"a", b"bc"]));
assert_ne!(digest(&[b"abc", b""]), digest(&[b"", b"abc"]));
}
#[test]
fn a_diagnostic_that_appears_on_an_unchanged_tree_repaints() {
use crate::report_format::{Format, RenderOptions};
let root = tempfile::tempdir().expect("root");
std::fs::write(root.path().join("big.txt"), vec![b'x'; 4096]).expect("big");
let control = root.path().join(".gitignore");
std::fs::write(&control, b"a\nb\nc\nd\ne\n").expect("short lines");
let scan = ScanConfig {
control_limits: crate::control::ControlLimits {
line_limit: Some(4),
..crate::control::ControlLimits::default()
},
..ScanConfig::default()
};
let (mut session, sender) = scripted_session_under(root.path(), Query::default(), scan);
let text = |session: &mut Session| {
session
.changed_report(
std::time::SystemTime::now(),
Format::Text,
RenderOptions::default(),
)
.expect("report")
};
let first = text(&mut session).expect("the first answer is always given");
let notes = |report: &Report| crate::report_format::diagnostic_lines(report).into_lines();
assert!(
!notes(&first).iter().any(|line| line.contains("line limit")),
"{:?}",
notes(&first)
);
std::fs::write(&control, b"abcdefghi\n").expect("one long line");
sender.send("modify\t.gitignore\n").expect("script the edit");
let batch = session.next_batch(Duration::from_secs(10)).expect("edit").expect("observed");
assert!(batch.dirty);
let repainted = text(&mut session).expect("a new refusal note repaints");
assert_eq!(
format!("{:?}", repainted.status),
format!("{:?}", first.status),
"the status alone would not have repainted"
);
assert!(
notes(&repainted).iter().any(|line| line.contains("line limit")),
"{:?}",
notes(&repainted)
);
assert!(text(&mut session).is_none(), "and only once");
}
#[test]
fn an_invalidation_keeps_its_change_record_and_repaints_by_the_answer() {
use crate::report_format::{Format, RenderOptions};
let root = tempfile::tempdir().expect("root");
std::fs::write(root.path().join("a.txt"), b"aaaa").expect("a");
let query = Query { views: vec![crate::query::ViewSpec::Summary], ..Query::default() };
let (mut session, sender) = scripted_session(root.path(), query);
let text = |session: &mut Session| {
session
.changed_report(
std::time::SystemTime::now(),
Format::Text,
RenderOptions::default(),
)
.expect("report")
};
let first = text(&mut session).expect("first answer");
sender.send("rescan\t.\n").expect("script an invalidation");
let batch =
session.next_batch(Duration::from_secs(10)).expect("invalidation").expect("observed");
assert!(
batch.changes.iter().any(|change| change.kind == ChangeKind::Invalidate),
"the invalidation reaches the change stream: {:?}",
batch.changes
);
let after = session.report(std::time::SystemTime::now()).expect("report");
let repainted = text(&mut session);
assert_eq!(
repainted.is_some(),
format!("{:?}", after.status) != format!("{:?}", first.status),
"a repaint follows a status change and nothing else: {:?}",
after.status
);
sender.send("rescan\t.\n").expect("script another invalidation");
let batch =
session.next_batch(Duration::from_secs(10)).expect("invalidation").expect("observed");
assert!(batch.changes.iter().any(|change| change.kind == ChangeKind::Invalidate));
assert!(text(&mut session).is_none());
}
}