From 49d42b57c7c727c329ff6c9ca3f738ce4279734e Mon Sep 17 00:00:00 2001 From: ScriptedAlchemy Date: Sun, 27 Sep 2026 12:45:52 +0000 Subject: [PATCH 1/6] fix(code-index): verify fresh waits with a cheap lock-free sweep --- .../code_index_scheduler/freshness_witness.rs | 507 ++++++++++++++++-- .../src/code_index_scheduler/reconcile.rs | 155 ++++-- .../registry/owner_signals.rs | 87 +-- .../code_index_scheduler/tests/reconcile.rs | 238 ++++++-- 4 files changed, 822 insertions(+), 165 deletions(-) diff --git a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/freshness_witness.rs b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/freshness_witness.rs index ba2b464709..b84b6c1b54 100644 --- a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/freshness_witness.rs +++ b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/freshness_witness.rs @@ -9,15 +9,23 @@ //! settled against the sealed generation's own per-file content digests //! (`SanitizedCodeSnapshotV1::files`), re-derived from the bytes on disk //! through the same bounded read + sanitize + digest path reconciliation uses. +//! +//! [`SourceSweepCacheV1`] makes re-settling cheap without weakening that +//! authority: a digest re-derived once is reused only while the file's full +//! stat identity, change time included, still holds. -use std::collections::{BTreeMap, BTreeSet}; +use std::collections::{BTreeMap, BTreeSet, HashMap}; +use std::fs::Metadata; use std::io::Read; +#[cfg(unix)] +use std::os::unix::fs::MetadataExt; use std::path::{Path, PathBuf}; use std::sync::Arc; use std::sync::atomic::{AtomicBool, Ordering}; -use std::time::UNIX_EPOCH; +use std::time::{Duration, SystemTime, UNIX_EPOCH}; use gix::bstr::BStr; +use gix::dir::walk::{Action, Delegate, ForDeletionMode}; use rayon::prelude::*; use sha2::{Digest, Sha256}; use tracedecay_code_index::production::CodeIndexIgnoredSourceAdmissionV1; @@ -60,30 +68,15 @@ pub fn worktree_stat_sweep( ) -> Result { let repository = tracedecay_runtime_core::git_open::open(project_root) .map_err(|error| CodeIndexSchedulerErrorV1::Git(error.to_string()))?; - let classification = classification::WorktreeChangeClassificationV1::classify(&repository) - .map_err(|error| CodeIndexSchedulerErrorV1::Git(error.to_string()))?; - let registry = StaticLanguageRegistry::new(); - let admitted_paths = ignored_source_admissions - .iter() - .map(|admission| admission.logical_path.as_str()) - .collect::>(); - let mut candidate_paths = classification.candidate_paths(); - candidate_paths.extend(admitted_paths.iter().map(|path| (*path).to_owned())); + let candidate_roster = source_candidates(&repository, ignored_source_admissions)?; // One sweep span plus an entries gauge: the stat walk is O(candidates) and // must never publish one profiler event per file. hotpath::gauge!("daemon.code_index.freshness.stat_signature.candidates") - .set(candidate_paths.len() as u64); + .set(candidate_roster.len() as u64); let mut buf = Vec::new(); let mut candidates = Vec::new(); - for logical_path in candidate_paths { - let absolute = project_root.join(&logical_path); - let Some(extension) = absolute.extension().and_then(|value| value.to_str()) else { - continue; - }; - let Some(descriptor) = registry.descriptor_for_extension(&extension.to_lowercase()) else { - continue; - }; - let Ok(metadata) = std::fs::metadata(&absolute) else { + for candidate in candidate_roster { + let Ok(metadata) = std::fs::metadata(project_root.join(&candidate.logical_path)) else { continue; }; if !metadata.is_file() { @@ -94,16 +87,12 @@ pub fn worktree_stat_sweep( .ok() .and_then(|time| time.duration_since(UNIX_EPOCH).ok()) .map_or(0u128, |elapsed| elapsed.as_nanos()); - buf.extend_from_slice(logical_path.as_bytes()); + buf.extend_from_slice(candidate.logical_path.as_bytes()); buf.push(0); buf.extend_from_slice(&metadata.len().to_le_bytes()); buf.extend_from_slice(&mtime_nanos.to_le_bytes()); buf.push(0xff); - candidates.push(StatCandidateV1 { - explicitly_admitted: admitted_paths.contains(logical_path.as_str()), - logical_path, - language: descriptor.language.clone(), - }); + candidates.push(candidate); } Ok(WorktreeStatSweepV1 { signature: encode_tagged_lowercase_hex("sha256:", &Sha256::digest(&buf)), @@ -111,6 +100,35 @@ pub fn worktree_stat_sweep( }) } +/// Every ordinary or explicitly admitted path with a registered language, +/// present or not: the roster a stat sweep checks. +fn source_candidates( + repository: &gix::Repository, + ignored_source_admissions: &[CodeIndexIgnoredSourceAdmissionV1], +) -> Result, CodeIndexSchedulerErrorV1> { + let classification = classification::WorktreeChangeClassificationV1::classify(repository) + .map_err(|error| CodeIndexSchedulerErrorV1::Git(error.to_string()))?; + let registry = StaticLanguageRegistry::new(); + let admitted_paths = ignored_source_admissions + .iter() + .map(|admission| admission.logical_path.as_str()) + .collect::>(); + let mut candidate_paths = classification.candidate_paths(); + candidate_paths.extend(admitted_paths.iter().map(|path| (*path).to_owned())); + Ok(candidate_paths + .into_iter() + .filter_map(|logical_path| { + let extension = Path::new(&logical_path).extension()?.to_str()?; + let descriptor = registry.descriptor_for_extension(&extension.to_lowercase())?; + Some(StatCandidateV1 { + explicitly_admitted: admitted_paths.contains(logical_path.as_str()), + language: descriptor.language.clone(), + logical_path, + }) + }) + .collect()) +} + /// The per-file content identities one sealed generation was reconciled /// against: `logical path → content digest` of the sanitized bytes for every /// present file in the generation's snapshot manifest, plus the snapshot's @@ -229,19 +247,422 @@ fn sanitized_digest( .map(|(bytes, _, _)| content_digest(&bytes)) } +/// What capture would seal for one candidate's bytes on disk. +#[derive(Clone, Debug, PartialEq, Eq)] +enum CandidateContentV1 { + Digest(ContentDigest), + /// The privacy boundary withholds this file from every generation. + Withheld, + Unreadable, +} + +impl CandidateContentV1 { + fn derive(project_root: &Path, candidate: &StatCandidateV1) -> Self { + match read_candidate(project_root, candidate) + .and_then(|raw| sanitized_digest(candidate, &raw)) + { + Ok(digest) => Self::Digest(digest), + Err(CodeIndexSchedulerErrorV1::Privacy(_)) => Self::Withheld, + Err(_) => Self::Unreadable, + } + } + + fn matches(&self, expected: Option<&ContentDigest>) -> bool { + match (self, expected) { + (Self::Digest(digest), Some(expected)) => digest == expected, + // A withheld file's absence from the manifest is the one + // consistent state. + (Self::Withheld, None) => true, + _ => false, + } + } +} + fn candidate_matches_manifest( project_root: &Path, candidate: &StatCandidateV1, manifest: &SourceContentManifestV1, ) -> bool { - let digest = - read_candidate(project_root, candidate).and_then(|raw| sanitized_digest(candidate, &raw)); - match (digest, manifest.files.get(&candidate.logical_path)) { - (Ok(digest), Some(expected)) => digest == *expected, - // The privacy boundary withholds this file from every generation, so - // its absence from the manifest is the one consistent state. - (Err(CodeIndexSchedulerErrorV1::Privacy(_)), None) => true, - _ => false, + CandidateContentV1::derive(project_root, candidate) + .matches(manifest.files.get(&candidate.logical_path)) +} + +/// Timestamps this close to the moment of a stat can still be shared by a +/// later write (coarse kernel clocks, two-second FAT times), so such a stat +/// cannot yet tell the file's current state from its next one. +const RACY_STAT_WINDOW: Duration = Duration::from_secs(2); + +/// One inode's stat identity. On Unix the change time is the kernel's own +/// record of every content or metadata write and cannot be set back +/// (`touch -d`, `cp --preserve`, `rsync -a` all advance it), so an equal +/// settled key proves the bytes behind it unchanged. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +struct StatKeyV1 { + device: u64, + inode: u64, + mode: u32, + size: u64, + modified_nanos: i128, + changed_nanos: i128, +} + +impl StatKeyV1 { + #[cfg(unix)] + fn of(metadata: &Metadata) -> Self { + Self { + device: metadata.dev(), + inode: metadata.ino(), + mode: metadata.mode(), + size: metadata.size(), + modified_nanos: i128::from(metadata.mtime()) * 1_000_000_000 + + i128::from(metadata.mtime_nsec()), + changed_nanos: i128::from(metadata.ctime()) * 1_000_000_000 + + i128::from(metadata.ctime_nsec()), + } + } + + // ponytail: without a change time a stat cannot prove unchanged bytes, so + // non-Unix keys never settle and every sweep re-derives every digest, as + // before this cache. Upgrade path: the Windows change time once std + // exposes it. + #[cfg(not(unix))] + fn of(metadata: &Metadata) -> Self { + Self { + device: 0, + inode: 0, + mode: 0, + size: metadata.len(), + modified_nanos: metadata + .modified() + .ok() + .and_then(|time| time.duration_since(UNIX_EPOCH).ok()) + .map_or(0, |elapsed| elapsed.as_nanos() as i128), + changed_nanos: 0, + } + } + + /// Whether this key, sampled at `sampled_at`, is old enough that any + /// later write must produce a different one. + fn settled(&self, sampled_at: SystemTime) -> bool { + cfg!(unix) + && sampled_at + .checked_sub(RACY_STAT_WINDOW) + .and_then(|horizon| horizon.duration_since(UNIX_EPOCH).ok()) + .is_some_and(|horizon| self.changed_nanos < horizon.as_nanos() as i128) + } +} + +/// The key of whatever is at `path` now, `None` when nothing is. +fn sample(path: &Path) -> std::io::Result> { + match std::fs::symlink_metadata(path) { + Ok(metadata) => Ok(Some(StatKeyV1::of(&metadata))), + Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(None), + Err(error) => Err(error), + } +} + +/// A candidate roster and the directory and ignore-rule evidence it was +/// enumerated from. Adding, removing or renaming an entry advances its +/// directory's change time, and ignore rules live in the recorded files, so +/// while every recorded key holds the roster is the one a fresh Git walk +/// would produce, and the probe skips that walk. +struct CachedCandidateRosterV1 { + git_metadata_signature: String, + admitted_paths: Vec, + evidence: Vec<(PathBuf, Option)>, + candidates: Arc>, +} + +impl CachedCandidateRosterV1 { + fn holds(&self, git_metadata_signature: &str, admitted_paths: &[String]) -> bool { + self.git_metadata_signature == git_metadata_signature + && self.admitted_paths == admitted_paths + && self + .evidence + .iter() + .all(|(path, key)| sample(path).is_ok_and(|now| now == *key)) + } +} + +/// Records every directory the Git walk descends into, keyed before the walk +/// reads it, plus the `.gitignore` each one holds. +struct DirectoryEvidenceV1 { + root: PathBuf, + sampled_at: SystemTime, + evidence: Vec<(PathBuf, Option)>, + settled: bool, +} + +impl DirectoryEvidenceV1 { + fn record(&mut self, path: PathBuf, directory: bool) { + let key = sample(&path); + let settled = match &key { + Ok(Some(key)) => key.settled(self.sampled_at), + Ok(None) => !directory, + Err(_) => false, + }; + match key { + Ok(key) if settled => self.evidence.push((path, key)), + _ => self.settled = false, + } + } + + fn record_directory(&mut self, directory: PathBuf) { + let ignore_file = directory.join(".gitignore"); + self.record(directory, true); + // An absent `.gitignore` needs no key: creating one advances the + // directory's change time. + if sample(&ignore_file).is_ok_and(|key| key.is_some()) { + self.record(ignore_file, false); + } + } +} + +impl Delegate for DirectoryEvidenceV1 { + fn emit( + &mut self, + _entry: gix::dir::EntryRef<'_>, + _collapsed_directory_status: Option, + ) -> Action { + if self.settled { + Action::Continue(()) + } else { + Action::Break(()) + } + } + + fn can_recurse( + &mut self, + entry: gix::dir::EntryRef<'_>, + for_deletion: Option, + worktree_root_is_repository: bool, + ) -> bool { + let recurse = entry.status.can_recurse( + entry.disk_kind, + entry.pathspec_match, + for_deletion, + worktree_root_is_repository, + ); + if recurse && self.settled { + let directory = self + .root + .join(gix::path::from_bstr(entry.rela_path.as_ref())); + self.record_directory(directory); + } + recurse + } +} + +/// Keys every directory and ignore-rule file the candidate roster depends on, +/// before the roster's own walk reads them. `None` when any of them changed +/// too recently to be keyed. +fn roster_evidence( + repository: &gix::Repository, + project_root: &Path, + sampled_at: SystemTime, +) -> Option)>> { + let mut recorder = DirectoryEvidenceV1 { + root: project_root.to_path_buf(), + sampled_at, + evidence: Vec::new(), + settled: true, + }; + let common_dir = repository.common_dir(); + recorder.record(common_dir.join("config"), false); + recorder.record(common_dir.join("info").join("exclude"), false); + let global_excludes = match repository + .config_snapshot() + .trusted_path("core.excludesFile") + { + Ok(Some(path)) => Some(path), + Ok(None) => gix::path::env::xdg_config("ignore", &mut |name| std::env::var_os(name)), + Err(_) => return None, + }; + if let Some(path) = global_excludes { + recorder.record(path, false); + } + recorder.record_directory(project_root.to_path_buf()); + let index = repository.index_or_empty().ok()?; + let options = repository + .dirwalk_options() + .ok()? + .emit_untracked(gix::dir::walk::EmissionMode::Matching); + let interrupt = AtomicBool::new(false); + repository + .dirwalk( + &index, + Vec::::new(), + &interrupt, + options, + &mut recorder, + ) + .ok()?; + recorder.settled.then_some(recorder.evidence) +} + +/// What one sweep of the witness checked, for the profiler and the proof. +#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)] +pub struct SourceSweepStatsV1 { + /// Whether the candidate roster came from a fresh Git walk. + pub walked: bool, + pub candidates: usize, + /// Files whose bytes were read because no settled key vouched for them. + pub hashed: usize, +} + +/// What earlier sweeps proved about the worktree: per-file content digests +/// that hold while each file's settled stat key does, and the last candidate +/// roster with the evidence that keeps it valid. A clean sweep then stats +/// files and directories and reads no bytes. It is a cache of the content +/// proof, not a second authority: every entry was derived from the bytes on +/// disk, and a disagreeing key falls back to re-deriving them. +#[derive(Default)] +pub struct SourceSweepCacheV1 { + contents: HashMap, + roster: Option, +} + +impl SourceSweepCacheV1 { + /// Whether the bytes on disk still carry exactly the manifest's content + /// identities, with the same verdict [`WorktreeStatSweepV1::content_matches`] + /// gives over a fresh sweep under the same roster. + #[hotpath::measure(label = "daemon.code_index.freshness.source_sweep")] + pub fn witness_matches( + &mut self, + project_root: &Path, + ignored_source_admissions: &[CodeIndexIgnoredSourceAdmissionV1], + git_metadata_signature: &str, + manifest: &SourceContentManifestV1, + shutting_down: &AtomicBool, + ) -> (bool, SourceSweepStatsV1) { + let mut stats = SourceSweepStatsV1::default(); + let sampled_at = SystemTime::now(); + let admitted_paths = ignored_source_admissions + .iter() + .map(|admission| admission.logical_path.clone()) + .collect::>(); + let candidates = match self.roster.as_ref() { + Some(roster) if roster.holds(git_metadata_signature, &admitted_paths) => { + Arc::clone(&roster.candidates) + } + _ => { + stats.walked = true; + let Ok(repository) = tracedecay_runtime_core::git_open::open(project_root) else { + return (false, stats); + }; + let evidence = roster_evidence(&repository, project_root, sampled_at); + let Ok(candidates) = source_candidates(&repository, ignored_source_admissions) + else { + return (false, stats); + }; + let candidates = Arc::new(candidates); + let live = candidates + .iter() + .map(|candidate| candidate.logical_path.as_str()) + .collect::>(); + self.contents.retain(|path, _| live.contains(path.as_str())); + self.roster = evidence.map(|evidence| CachedCandidateRosterV1 { + git_metadata_signature: git_metadata_signature.to_owned(), + admitted_paths, + evidence, + candidates: Arc::clone(&candidates), + }); + candidates + } + }; + let mut present = Vec::new(); + for candidate in candidates.iter().filter(|candidate| { + candidate.explicitly_admitted || !is_generated_path_segment(&candidate.logical_path) + }) { + let absolute = project_root.join(&candidate.logical_path); + let Ok(metadata) = std::fs::symlink_metadata(&absolute) else { + continue; + }; + // A link's own key says nothing about its target's bytes. + let key = if metadata.file_type().is_symlink() { + if !std::fs::metadata(&absolute).is_ok_and(|target| target.is_file()) { + continue; + } + None + } else if metadata.is_file() { + Some(StatKeyV1::of(&metadata)) + } else { + continue; + }; + present.push((candidate, key)); + } + stats.candidates = present.len(); + hotpath::gauge!("daemon.code_index.freshness.source_sweep.candidates") + .set(present.len() as u64); + let present_paths = present + .iter() + .map(|(candidate, _)| candidate.logical_path.as_str()) + .collect::>(); + if !manifest + .files + .keys() + .all(|logical_path| present_paths.contains(logical_path.as_str())) + { + return (false, stats); + } + let mut disputed = Vec::new(); + let mut unvouched = Vec::new(); + for (candidate, key) in present { + match self.contents.get(&candidate.logical_path) { + Some((cached, content)) if Some(*cached) == key => { + if !content.matches(manifest.files.get(&candidate.logical_path)) { + disputed.push(candidate); + } + } + _ => unvouched.push((candidate, key)), + } + } + stats.hashed = unvouched.len(); + hotpath::gauge!("daemon.code_index.freshness.source_sweep.hashed") + .set(unvouched.len() as u64); + let Ok(derived) = parallelism::install(|| { + unvouched + .par_iter() + .map(|(candidate, key)| { + parallelism::with_background_cpu_permit(|| { + if shutting_down.load(Ordering::Acquire) { + return (*candidate, *key, CandidateContentV1::Unreadable); + } + ( + *candidate, + *key, + CandidateContentV1::derive(project_root, candidate), + ) + }) + }) + .collect::>() + }) else { + return (false, stats); + }; + if shutting_down.load(Ordering::Acquire) { + return (false, stats); + } + for (candidate, key, content) in derived { + if !content.matches(manifest.files.get(&candidate.logical_path)) { + disputed.push(candidate); + } + // The key was sampled before the read, so a settled key vouches + // for these bytes: a write after the stat would change it. + match key { + Some(key) + if key.settled(sampled_at) && content != CandidateContentV1::Unreadable => + { + self.contents + .insert(candidate.logical_path.clone(), (key, content)); + } + _ => { + self.contents.remove(&candidate.logical_path); + } + } + } + let matches = disputed.is_empty() + || tracked_files_match_after_clean_filters(project_root, &disputed, manifest); + (matches, stats) } } @@ -304,22 +725,6 @@ impl ReconciledSourceWitnessV1 { content_manifest: SourceContentManifestV1::for_snapshot(snapshot), } } - - /// Metadata differs → not current, without reading a byte. Metadata equal - /// → current only when every candidate's content digest still matches - /// the sealed manifest. - pub fn matches_worktree( - &self, - project_root: &Path, - ignored_source_admissions: &[CodeIndexIgnoredSourceAdmissionV1], - shutting_down: &AtomicBool, - ) -> bool { - let Ok(sweep) = worktree_stat_sweep(project_root, ignored_source_admissions) else { - return false; - }; - sweep.signature == self.stat_signature - && sweep.content_matches(project_root, &self.content_manifest, shutting_down) - } } /// Durable binding between one sealed generation and the exact ordinary plus diff --git a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/reconcile.rs b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/reconcile.rs index 552809598d..c1481b7711 100644 --- a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/reconcile.rs +++ b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/reconcile.rs @@ -59,6 +59,7 @@ use crate::code_index::{ use super::freshness_witness::{ ReconciledSourceWitnessV1, RestoreFreshnessWitnessV1, SourceContentManifestV1, + SourceSweepCacheV1, }; use super::git_tree_capture::{CapturedFileOutcomeV1, CapturedFileRosterV1}; use super::publication_store::GenerationDecodeBudgetV1; @@ -502,6 +503,8 @@ pub(crate) struct SourceFreshnessFenceV1 { pub(super) state: Arc>, last_reconciled_at_micros: Arc, source_epoch: Arc, + /// Held only by source sweeps, never across a scheduler pass. + sweep_cache: Arc>, } #[derive(Clone)] @@ -511,6 +514,8 @@ pub(super) struct SourceFreshnessFenceStateV1 { /// The stat signature (negative cache) and sealed file digests (proof) /// the last completed reconcile established; `None` until one has. source_witness: Option, + /// The ignored-source roster that proof was established under. + source_roster: Vec, verified_against_source: bool, freshness_unknown: bool, reconciled_without_generation: bool, @@ -524,6 +529,7 @@ impl SourceFreshnessFenceV1 { git_metadata: identity::GitMetadataFingerprintV1::default(), last_reconciled_at: Instant::now(), source_witness: None, + source_roster: Vec::new(), verified_against_source: false, freshness_unknown: true, reconciled_without_generation: false, @@ -531,6 +537,7 @@ impl SourceFreshnessFenceV1 { })), last_reconciled_at_micros: Arc::new(AtomicI64::new(0)), source_epoch, + sweep_cache: Arc::default(), } } @@ -545,12 +552,14 @@ impl SourceFreshnessFenceV1 { &self, git_metadata: identity::GitMetadataFingerprintV1, source_witness: Option, + source_roster: &[CodeIndexIgnoredSourceAdmissionV1], reconciled_without_generation: bool, ) { let micros = now_micros().0; let mut state = self.state.lock().unwrap_or_else(PoisonError::into_inner); state.git_metadata = git_metadata; state.source_witness = source_witness; + state.source_roster = source_roster.to_vec(); state.freshness_unknown = false; state.last_reconciled_at = Instant::now(); state.verified_against_source = true; @@ -713,6 +722,86 @@ impl SourceFreshnessFenceV1 { self.last_reconciled_at_micros .store(micros, Ordering::Release); } + + /// Whether `freshness`'s source witness still describes the worktree: + /// every candidate's content digest, under the roster the proof was + /// established with, equals the sealed file manifest. + fn source_witness_matches( + &self, + freshness: &SourceFreshnessFenceStateV1, + git_metadata: &identity::GitMetadataFingerprintV1, + project_root: &Path, + shutting_down: &AtomicBool, + ) -> bool { + freshness.source_witness.as_ref().is_some_and(|witness| { + self.sweep_cache + .lock() + .unwrap_or_else(PoisonError::into_inner) + .witness_matches( + project_root, + &freshness.source_roster, + &git_metadata.stable_signature(), + &witness.content_manifest, + shutting_down, + ) + .0 + }) + } + + /// Sweep the current witness and report what the sweep had to read. + #[cfg(test)] + pub(super) fn source_sweep_for_test( + &self, + project_root: &Path, + shutting_down: &AtomicBool, + ) -> (bool, freshness_witness::SourceSweepStatsV1) { + let freshness = self.snapshot(); + let witness = freshness + .source_witness + .as_ref() + .expect("a reconciled fence carries a source witness"); + self.sweep_cache + .lock() + .unwrap_or_else(PoisonError::into_inner) + .witness_matches( + project_root, + &freshness.source_roster, + &identity::GitMetadataFingerprintV1::capture(project_root).stable_signature(), + &witness.content_manifest, + shutting_down, + ) + } + + /// The freshness ladder over this fence alone, so a caller verifying the + /// source never waits for a scheduler pass: unverified, then Git + /// metadata, then (unless the last proof is younger than + /// `trust_recent_reconcile`) the source witness. A matching witness + /// resets the probe clock. + pub(super) fn ladder_verdict( + &self, + project_root: &Path, + shutting_down: &AtomicBool, + trust_recent_reconcile: Option, + ) -> FreshnessProbeVerdictV1 { + let freshness = self.snapshot(); + if !freshness.verified_against_source { + return FreshnessProbeVerdictV1::Unverified; + } + let git_metadata = identity::GitMetadataFingerprintV1::capture(project_root); + if git_metadata.differs_from(&freshness.git_metadata) { + return FreshnessProbeVerdictV1::Moved; + } + if trust_recent_reconcile + .is_some_and(|threshold| freshness.last_reconciled_at.elapsed() < threshold) + { + return FreshnessProbeVerdictV1::Current; + } + if self.source_witness_matches(&freshness, &git_metadata, project_root, shutting_down) { + self.refresh_monotonic_clock(true); + return FreshnessProbeVerdictV1::Current; + } + FreshnessProbeVerdictV1::Moved + } } /// What the cheap Git/stat freshness ladder concluded about the retained @@ -2935,6 +3024,7 @@ impl CodeIndexWorktreeSchedulerV1 { self.freshness_fence.mark_reconciled( metadata, source_witness, + &self.ignored_source_admissions, reconciled_without_generation, ); } @@ -2944,8 +3034,12 @@ impl CodeIndexWorktreeSchedulerV1 { metadata: identity::GitMetadataFingerprintV1, source_witness: Option, ) { - self.freshness_fence - .mark_reconciled(metadata, source_witness, false); + self.freshness_fence.mark_reconciled( + metadata, + source_witness, + &self.ignored_source_admissions, + false, + ); } /// Record the restore-time freshness witness for the current active @@ -3062,16 +3156,14 @@ impl CodeIndexWorktreeSchedulerV1 { } /// Whether the last reconcile's source witness still describes the - /// worktree: unchanged stat metadata (the negative cache) and, only then, - /// every candidate's content digest equal to the sealed file manifest. + /// worktree: every candidate's content digest equals the sealed manifest. fn source_witness_matches_worktree(&self, freshness: &SourceFreshnessFenceStateV1) -> bool { - freshness.source_witness.as_ref().is_some_and(|witness| { - witness.matches_worktree( - &self.project_root, - &self.ignored_source_admissions, - &self.shutting_down, - ) - }) + self.freshness_fence.source_witness_matches( + freshness, + &identity::GitMetadataFingerprintV1::capture(&self.project_root), + &self.project_root, + &self.shutting_down, + ) } /// Mint the exact-source currency witness for one generation from the @@ -3235,35 +3327,11 @@ impl CodeIndexWorktreeSchedulerV1 { /// as movement here made a concurrent query escalate the targeted hint /// pass into an overflow rescan and relabel the arrival as its own. pub(super) fn freshness_probe_verdict(&mut self) -> FreshnessProbeVerdictV1 { - self.freshness_ladder_verdict(true) - } - - /// The ladder behind [`Self::freshness_probe_verdict`]. Without the - /// bounded-staleness shortcut it always sweeps the source witness, which - /// is what a caller waiting for the source as it is now asks for. - fn freshness_ladder_verdict( - &mut self, - trust_recent_reconcile: bool, - ) -> FreshnessProbeVerdictV1 { - let freshness = self.freshness_fence.snapshot(); - if !freshness.verified_against_source { - return FreshnessProbeVerdictV1::Unverified; - } - if identity::GitMetadataFingerprintV1::capture(&self.project_root) - .differs_from(&freshness.git_metadata) - { - return FreshnessProbeVerdictV1::Moved; - } - if trust_recent_reconcile - && freshness.last_reconciled_at.elapsed() < self.policy.staleness_threshold - { - return FreshnessProbeVerdictV1::Current; - } - if self.source_witness_matches_worktree(&freshness) { - self.freshness_fence.refresh_monotonic_clock(true); - return FreshnessProbeVerdictV1::Current; - } - FreshnessProbeVerdictV1::Moved + self.freshness_fence.ladder_verdict( + &self.project_root, + &self.shutting_down, + Some(self.policy.staleness_threshold), + ) } /// Decide whether the cheap Git/stat ladder requires an authoritative @@ -3306,13 +3374,6 @@ impl CodeIndexWorktreeSchedulerV1 { self.request_reconcile_for_verdict(verdict) } - /// [`Self::request_fresh_for_query_background`] against the source as it - /// is now: the source witness is swept even inside the staleness window. - pub fn request_fresh_now_background(&mut self) -> bool { - let verdict = self.freshness_ladder_verdict(false); - self.request_reconcile_for_verdict(verdict) - } - fn request_reconcile_for_verdict(&mut self, verdict: FreshnessProbeVerdictV1) -> bool { match verdict { FreshnessProbeVerdictV1::Current => false, diff --git a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry/owner_signals.rs b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry/owner_signals.rs index 09f116fe82..47b72c7f6d 100644 --- a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry/owner_signals.rs +++ b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry/owner_signals.rs @@ -1,6 +1,7 @@ //! Change signals for one project root, and the readiness wait built on them. use std::path::{Path, PathBuf}; +use std::sync::{Arc, PoisonError}; use std::time::Duration; use tracedecay_runtime_core::path_safety::canonical_existing_identity; @@ -16,7 +17,10 @@ use super::{ CodeIndexCadenceTriggerV1, CodeIndexOwnerActivityV1, CodeIndexSchedulerRegistryV1, unique_mounted_for_scope, }; -use crate::code_index_scheduler::{CodeIndexCadenceTelemetryV1, LatestCodeTextGenerationV1}; +use crate::code_index_scheduler::reconcile::FreshnessProbeVerdictV1; +use crate::code_index_scheduler::{ + CodeIndexCadenceTelemetryV1, CodeIndexWorktreeSchedulerV1, LatestCodeTextGenerationV1, +}; /// Why the source could not be proven current before waiting. #[derive(Clone, Copy, Debug, PartialEq, Eq)] @@ -155,9 +159,11 @@ impl CodeIndexOwnerSignalsV1 { impl CodeIndexSchedulerRegistryV1 { /// Sweep the source witness now and post a wake for any proven change, - /// so later freshness reads describe the source as of this call. An - /// unmounted root has nothing to sweep; its mount reconciles. Without - /// `sweep_source` only the publication park is checked. + /// so later freshness reads describe the source as of this call. The + /// sweep reads only the freshness fence, never the scheduler mutex, so a + /// pass in flight cannot delay it. An unmounted root has nothing to + /// sweep; its mount reconciles. Without `sweep_source` only the + /// publication park is checked. async fn request_fresh_now( &self, project_root: &Path, @@ -166,7 +172,7 @@ impl CodeIndexSchedulerRegistryV1 { let Ok(canonical) = canonical_existing_identity(project_root) else { return Ok(()); }; - let (scheduler, pending_wake, wake) = { + let (source_freshness, shutting_down, hints, epoch, pending_wake, wake) = { let mounted = self.mounted.lock().await; let Some(worktree) = mounted.get(&canonical) else { return Ok(()); @@ -178,22 +184,31 @@ impl CodeIndexSchedulerRegistryV1 { return Ok(()); } ( - std::sync::Arc::clone(&worktree.scheduler), - std::sync::Arc::clone(&worktree.pending_wake), - std::sync::Arc::clone(&worktree.wake), + worktree.source_freshness.clone(), + Arc::clone(&worktree.shutting_down), + Arc::clone(&worktree.hints), + Arc::clone(&worktree.epoch), + Arc::clone(&worktree.pending_wake), + Arc::clone(&worktree.wake), ) }; tokio::task::spawn_blocking(move || { - let mut scheduler = scheduler - .lock() - .unwrap_or_else(std::sync::PoisonError::into_inner); - if scheduler.request_fresh_now_background() { - Self::note_wake( - &pending_wake, - &wake, - CodeIndexCadenceTriggerV1::QueryAdmission, - ); + match source_freshness.ladder_verdict(&canonical, &shutting_down, None) { + FreshnessProbeVerdictV1::Current => return, + FreshnessProbeVerdictV1::Unverified => {} + FreshnessProbeVerdictV1::Moved => { + CodeIndexWorktreeSchedulerV1::record_background_reconcile_hint( + &mut hints.lock().unwrap_or_else(PoisonError::into_inner), + &epoch, + true, + ); + } } + Self::note_wake( + &pending_wake, + &wake, + CodeIndexCadenceTriggerV1::QueryAdmission, + ); }) .await .map_err(|_| CodeIndexFreshSweepRefusedV1::SweepFailed) @@ -232,18 +247,17 @@ impl CodeIndexSchedulerRegistryV1 { /// Wait until `project_root` reaches `target`, re-reading freshness only /// when the registry publishes a change, for at most `budget`. /// - /// A target the current reading already satisfies is reached at once: - /// that reading is the scheduler's last proof, the answer a plain status - /// read gives, so a save no hook reported is left to the backstop sweep - /// rather than swept inside the caller's budget. Otherwise the wait proves - /// freshness against the source as it is now: the bounded - /// Git/stat/content probe either refreshes the verified watermark or - /// posts the wake for a proven change, so a reading taken after it - /// cannot report an edit the scheduler has not yet seen as fresh. - /// `graph_ready` does not depend on freshness and skips that probe. An - /// unmounted root is waited through: a mount that lands inside the budget - /// reconciles the source as of that mount. Dropping the future abandons - /// the wait; a wake the probe posted is ordinary demand. + /// `fresh` means verified against the source as of the request: the wait + /// first sweeps the source witness, which either refreshes the verified + /// watermark or posts the wake for a proven change, so a reading taken + /// after it cannot report an edit no hook announced as fresh. The sweep + /// stats files against digests earlier sweeps proved and reads bytes + /// only where a stat moved, without the scheduler mutex. `ready` and + /// `graph_ready` accept a reading that already satisfies them, the answer + /// a plain status read gives; a pending `ready` still sweeps first. + /// An unmounted root is waited through: a mount that lands inside the + /// budget reconciles the source as of that mount. Dropping the future + /// abandons the wait; a wake the sweep posted is ordinary demand. pub async fn wait_for_readiness( &self, project_root: &Path, @@ -251,16 +265,19 @@ impl CodeIndexSchedulerRegistryV1 { budget: Duration, ) -> Result { let deadline = tokio::time::Instant::now() + budget; - if self - .dashboard_freshness_read(project_root) - .await? - .is_some_and(|freshness| freshness.readiness(target) == CodeIndexReadinessV1::Reached) + if target != CodeIndexReadinessTargetV1::Fresh + && self + .dashboard_freshness_read(project_root) + .await? + .is_some_and(|freshness| { + freshness.readiness(target) == CodeIndexReadinessV1::Reached + }) { return Ok(CodeIndexReadinessWaitReadV1::Reached); } let mut signals = CodeIndexOwnerSignalsV1::subscribe(self, project_root).await; - // The probe can take the scheduler mutex; the caller's budget bounds - // it, and an unproven source cannot be reported as reached. + // The caller's budget bounds the sweep, and an unproven source cannot + // be reported as reached. let sweep_source = target != CodeIndexReadinessTargetV1::GraphReady; match tokio::time::timeout_at(deadline, self.request_fresh_now(project_root, sweep_source)) .await diff --git a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/tests/reconcile.rs b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/tests/reconcile.rs index 4657f6b2ab..affbbefc07 100644 --- a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/tests/reconcile.rs +++ b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/tests/reconcile.rs @@ -58,6 +58,7 @@ use crate::{ classification::{WorktreeChangeClassV1, WorktreeChangeClassificationV1}, feedback_document_identity_from_generation, freshness_witness::RestoreFreshnessWitnessV1, + freshness_witness::SourceSweepStatsV1, registry::{ ColdMountOpenEventV1, ServingGenerationInstallationOutcomeV1, ServingGenerationRollbackOutcomeV1, dashboard_code_graph_serving, @@ -6742,12 +6743,10 @@ async fn readiness_wait_reaches_ready_exactly_when_the_held_graph_publishes() { registry.shutdown().await; } -/// A settled worktree already satisfies `fresh`, so a short wait reaches it -/// even while another holder owns the scheduler lock the source sweep would -/// need; the busy-read ladder still reports the served generation current. -#[tokio::test(flavor = "multi_thread", worker_threads = 4)] -async fn readiness_wait_reaches_a_target_the_current_reading_already_holds() { - let fixture = GitFixture::new(&[("src/lib.rs", "pub fn source() -> u32 { 1 }\n")]); +async fn settled_fresh_wait_fixture( + source: &str, +) -> (GitFixture, TempDir, CodeIndexSchedulerRegistryV1) { + let fixture = GitFixture::new(&[("src/lib.rs", source)]); let store = TempDir::new().expect("store root"); let registry = CodeIndexSchedulerRegistryV1::with_background_reconcile_permits(1, 1); registry @@ -6760,46 +6759,221 @@ async fn readiness_wait_reaches_a_target_the_current_reading_already_holds() { .expect("mount worktree"); wait_for_initial_generation(®istry, fixture.path()).await; settled_owner_with_idle_admission(®istry, fixture.path()).await; + (fixture, store, registry) +} + +async fn served_texts_containing( + registry: &CodeIndexSchedulerRegistryV1, + root: &Path, + needle: &str, +) -> Vec { + let mut texts = registry + .latest_complete_serving_for_test(root) + .await + .expect("served generation") + .lexical() + .iter() + .filter(|chunk| chunk.sanitized_text.as_str().contains(needle)) + .map(|chunk| chunk.sanitized_text.as_str().to_owned()) + .collect::>(); + texts.sort(); + texts +} + +/// `tracedecay_status wait_for { state: fresh }` on a current index verifies +/// the source and reaches it well inside a one-second budget. +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn fresh_wait_on_a_current_index_reaches_inside_one_second() { + let (fixture, _store, registry) = + settled_fresh_wait_fixture("pub fn source() -> u32 { 1 }\n").await; + let started = Instant::now(); + let outcome = registry + .wait_for_readiness( + fixture.path(), + tracedecay_contracts::code_index_freshness::CodeIndexReadinessTargetV1::Fresh, + Duration::from_secs(1), + ) + .await + .expect("freshness read"); + let elapsed = started.elapsed(); + assert!( + matches!( + outcome, + tracedecay_contracts::code_index_freshness::CodeIndexReadinessWaitReadV1::Reached + ), + "{outcome:?}" + ); + assert!(elapsed < Duration::from_secs(1), "{elapsed:?}"); + registry.shutdown().await; +} + +/// A save no hook reported is still part of "the source as of the request": +/// the `fresh` wait's own sweep finds it, and the wait returns only once the +/// index serves the saved bytes. +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn fresh_wait_catches_an_unreported_save_and_returns_after_its_reindex() { + let (fixture, _store, registry) = + settled_fresh_wait_fixture("pub fn source() -> u32 { 1 }\n").await; + let before = registry + .latest_generation_id(fixture.path()) + .await + .expect("initial generation"); + fixture.edit("src/lib.rs", "pub fn source() -> u32 { 2 }\n"); + + let outcome = registry + .wait_for_readiness( + fixture.path(), + tracedecay_contracts::code_index_freshness::CodeIndexReadinessTargetV1::Fresh, + SERVING_SEAT_FAILURE_CEILING, + ) + .await + .expect("freshness read"); + assert!( + matches!( + outcome, + tracedecay_contracts::code_index_freshness::CodeIndexReadinessWaitReadV1::Reached + ), + "{outcome:?}" + ); + assert_ne!( + registry.latest_generation_id(fixture.path()).await, + Some(before) + ); + assert_eq!( + served_texts_containing(®istry, fixture.path(), "fn source").await, + vec![ + "pub fn source() -> u32 { 2 }".to_owned(), + "pub fn source() -> u32 { 2 }\n".to_owned() + ] + ); + registry.shutdown().await; +} + +/// The `fresh` sweep never waits for the scheduler mutex. While another +/// holder owns it, a quiet tree is still verified and reached, and an +/// unreported save is still found: the wait times out on the old generation +/// refreshing instead of reaching it, and the save is indexed once the +/// holder lets go. +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn fresh_wait_verifies_the_source_while_a_pass_holds_the_scheduler() { + let (fixture, _store, registry) = + settled_fresh_wait_fixture("pub fn source() -> u32 { 1 }\n").await; let fresh = tracedecay_contracts::code_index_freshness::CodeIndexReadinessTargetV1::Fresh; - assert!(matches!( - registry - .wait_for_readiness(fixture.path(), fresh, SERVING_SEAT_FAILURE_CEILING) - .await - .expect("freshness read"), - tracedecay_contracts::code_index_freshness::CodeIndexReadinessWaitReadV1::Reached - )); + let before = registry + .latest_generation_id(fixture.path()) + .await + .expect("initial generation"); + let held = hold_scheduler_for_root(®istry, fixture.path()).await; - let handle = registry - .scheduler_handle(fixture.path()) + let quiet = registry + .wait_for_readiness(fixture.path(), fresh, Duration::from_secs(1)) .await - .expect("mounted scheduler"); - let (held_tx, held_rx) = std::sync::mpsc::channel(); - let (release_tx, release_rx) = std::sync::mpsc::channel::<()>(); - let lock_thread = std::thread::spawn(move || { - let _guard = handle - .lock() - .unwrap_or_else(std::sync::PoisonError::into_inner); - held_tx.send(()).expect("signal scheduler lock held"); - let _ = release_rx.recv(); - }); - held_rx.recv().expect("scheduler lock acquired"); + .expect("freshness read"); + fixture.edit("src/lib.rs", "pub fn source() -> u32 { 3 }\n"); + let edited = registry + .wait_for_readiness(fixture.path(), fresh, Duration::from_millis(500)) + .await + .expect("freshness read"); + held.release().await; + assert!( + matches!( + quiet, + tracedecay_contracts::code_index_freshness::CodeIndexReadinessWaitReadV1::Reached + ), + "{quiet:?}" + ); + let tracedecay_contracts::code_index_freshness::CodeIndexReadinessWaitReadV1::TimedOut { + last: Some(last), + } = edited + else { + panic!("an unreported save must not read as fresh: {edited:?}"); + }; + assert_eq!( + (last.staleness_state, last.latest_generation_id.as_deref()), + ( + Some(tracedecay_contracts::code_index_freshness::CodeIndexStalenessStateV1::Refreshing), + Some(before.as_str()) + ) + ); - let held = registry - .wait_for_readiness(fixture.path(), fresh, Duration::from_millis(200)) + let reindexed = registry + .wait_for_readiness(fixture.path(), fresh, SERVING_SEAT_FAILURE_CEILING) .await .expect("freshness read"); - release_tx.send(()).expect("release scheduler lock"); - lock_thread.join().expect("lock thread joins"); assert!( matches!( - held, + reindexed, tracedecay_contracts::code_index_freshness::CodeIndexReadinessWaitReadV1::Reached ), - "a current worktree must not time out behind the source sweep: {held:?}" + "{reindexed:?}" + ); + assert_eq!( + served_texts_containing(®istry, fixture.path(), "fn source").await, + vec![ + "pub fn source() -> u32 { 3 }".to_owned(), + "pub fn source() -> u32 { 3 }\n".to_owned() + ] ); registry.shutdown().await; } +/// Settled stats vouch for digests an earlier sweep derived, so re-sweeping a +/// quiet tree walks no directories and reads no source bytes, while a +/// same-length rewrite that restores its mtime still advances the change +/// time and is re-read and caught. +#[test] +fn source_sweep_rereads_only_files_whose_settled_stat_moved() { + let fixture = GitFixture::new(&[ + ("src/lib.rs", "pub fn alpha() -> u32 { 1 }\n"), + ("src/other.rs", "pub fn beta() -> u32 { 1 }\n"), + ]); + let store = TempDir::new().expect("store root"); + let mut scheduler = scheduler( + &fixture, + store.path().to_path_buf(), + Arc::new(SharedCodeIndexBytePoolV1::default()), + ); + let _ = published(scheduler.reconcile_now().expect("initial publish")); + let fence = scheduler.freshness_fence(); + let shutting_down = std::sync::atomic::AtomicBool::new(false); + std::thread::sleep(Duration::from_millis(2_100)); + + assert_eq!( + fence.source_sweep_for_test(fixture.path(), &shutting_down), + ( + true, + SourceSweepStatsV1 { + walked: true, + candidates: 2, + hashed: 2 + } + ) + ); + assert_eq!( + fence.source_sweep_for_test(fixture.path(), &shutting_down), + ( + true, + SourceSweepStatsV1 { + walked: false, + candidates: 2, + hashed: 0 + } + ) + ); + rewrite_preserving_stat(&fixture, "src/lib.rs", "pub fn alpha() -> u32 { 2 }\n"); + assert_eq!( + fence.source_sweep_for_test(fixture.path(), &shutting_down), + ( + false, + SourceSweepStatsV1 { + walked: false, + candidates: 2, + hashed: 1 + } + ) + ); +} + /// A pass can start and settle entirely between two reads of the running /// level. The owner-activity counts only grow, so a reader that looks after /// the pass still sees it, and the worker phase says the pass and its tail From 5076502f32f480297a8e7f8307963277bbfc43eb Mon Sep 17 00:00:00 2001 From: ScriptedAlchemy Date: Sun, 27 Sep 2026 13:20:40 +0000 Subject: [PATCH 2/6] fix(hooks): report a hook event delivered once the daemon admits it --- .../src/code_index_scheduler/reconcile.rs | 9 +++++ crates/tracedecay/src/daemon/core_hooks.rs | 33 +++++++++++++++---- 2 files changed, 36 insertions(+), 6 deletions(-) diff --git a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/reconcile.rs b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/reconcile.rs index c1481b7711..c320467721 100644 --- a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/reconcile.rs +++ b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/reconcile.rs @@ -3374,6 +3374,15 @@ impl CodeIndexWorktreeSchedulerV1 { self.request_reconcile_for_verdict(verdict) } + /// [`Self::request_fresh_for_query_background`] against the source as it + /// is now: the source witness is swept even inside the staleness window. + pub fn request_fresh_now_background(&mut self) -> bool { + let verdict = + self.freshness_fence + .ladder_verdict(&self.project_root, &self.shutting_down, None); + self.request_reconcile_for_verdict(verdict) + } + fn request_reconcile_for_verdict(&mut self, verdict: FreshnessProbeVerdictV1) -> bool { match verdict { FreshnessProbeVerdictV1::Current => false, diff --git a/crates/tracedecay/src/daemon/core_hooks.rs b/crates/tracedecay/src/daemon/core_hooks.rs index 098b54767c..81fdbba0c2 100644 --- a/crates/tracedecay/src/daemon/core_hooks.rs +++ b/crates/tracedecay/src/daemon/core_hooks.rs @@ -6,7 +6,7 @@ use std::path::Path; -use tokio::io::AsyncWriteExt; +use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader}; use tokio::time::{Duration, timeout}; use tracedecay_hooks::core_events::{DaemonHookEvent, HOOK_EVENT_METHOD, HookEventNotifyOutcomeV1}; @@ -66,8 +66,10 @@ async fn notify_hook_event_to_connection( return HookEventNotifyOutcomeV1::Malformed; }; // A stateless request, not a notification: its own daemon connection has - // no `initialize` session for a notification to ride. The result is not - // awaited, so delivery stays fire-and-forget. + // no `initialize` session for a notification to ride. The daemon answers + // once it has admitted the event, the edited paths into the code-index + // queue included, so `Delivered` means a status read issued afterwards + // already sees the save. let mut request = JsonRpcRequest { jsonrpc: "2.0".to_string(), id: Some(serde_json::Value::from(1)), @@ -81,7 +83,7 @@ async fn notify_hook_event_to_connection( let Ok(stream) = BrokerStream::connect(connection.endpoint()).await else { return HookEventNotifyOutcomeV1::Unavailable; }; - let (_reader, mut writer) = stream.into_owned_split(); + let (reader, mut writer) = stream.into_owned_split(); if write_daemon_preamble(&mut writer, &connection, &handshake) .await .is_err() @@ -94,10 +96,29 @@ async fn notify_hook_event_to_connection( if writer.write_all(b"\n").await.is_err() { return HookEventNotifyOutcomeV1::Unavailable; } - if writer.flush().await.is_err() || writer.shutdown().await.is_err() { + if writer.flush().await.is_err() { return HookEventNotifyOutcomeV1::Unavailable; } - HookEventNotifyOutcomeV1::Delivered + let mut reader = BufReader::new(reader); + let mut response = String::new(); + loop { + response.clear(); + match reader.read_line(&mut response).await { + Ok(0) | Err(_) => return HookEventNotifyOutcomeV1::Unavailable, + Ok(_) => {} + } + let Ok(value) = serde_json::from_str::(&response) else { + return HookEventNotifyOutcomeV1::Malformed; + }; + if value.get("id") != request.id.as_ref() { + continue; + } + return if value.get("result").is_some() { + HookEventNotifyOutcomeV1::Delivered + } else { + HookEventNotifyOutcomeV1::Malformed + }; + } } #[cfg(all(test, unix))] From 5df7efd181768f9d4ceb60e489f44ff71d3ed9ff Mon Sep 17 00:00:00 2001 From: ScriptedAlchemy Date: Sun, 27 Sep 2026 13:32:20 +0000 Subject: [PATCH 3/6] refactor(code-index): route every content proof through the sweep cache --- .../code_index_scheduler/freshness_witness.rs | 79 +----------- .../src/code_index_scheduler/reconcile.rs | 118 ++++++++++-------- 2 files changed, 74 insertions(+), 123 deletions(-) diff --git a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/freshness_witness.rs b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/freshness_witness.rs index b84b6c1b54..837d5f6908 100644 --- a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/freshness_witness.rs +++ b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/freshness_witness.rs @@ -51,11 +51,9 @@ struct StatCandidateV1 { } /// One stat sweep over every ordinary or explicitly admitted source -/// candidate: the cheap signature plus the candidate roster it hashed, so the -/// content proof compares exactly the files the signature covered. +/// candidate: a negative cache only, settled by [`SourceSweepCacheV1`]. pub struct WorktreeStatSweepV1 { pub signature: String, - candidates: Vec, } /// A cheap stat-level sweep over every ordinary or explicitly admitted source @@ -74,7 +72,6 @@ pub fn worktree_stat_sweep( hotpath::gauge!("daemon.code_index.freshness.stat_signature.candidates") .set(candidate_roster.len() as u64); let mut buf = Vec::new(); - let mut candidates = Vec::new(); for candidate in candidate_roster { let Ok(metadata) = std::fs::metadata(project_root.join(&candidate.logical_path)) else { continue; @@ -92,11 +89,9 @@ pub fn worktree_stat_sweep( buf.extend_from_slice(&metadata.len().to_le_bytes()); buf.extend_from_slice(&mtime_nanos.to_le_bytes()); buf.push(0xff); - candidates.push(candidate); } Ok(WorktreeStatSweepV1 { signature: encode_tagged_lowercase_hex("sha256:", &Sha256::digest(&buf)), - candidates, }) } @@ -162,63 +157,6 @@ impl SourceContentManifestV1 { } } -impl WorktreeStatSweepV1 { - /// Whether the bytes on disk still carry exactly the content identities - /// the generation sealed: every present manifest file is still a - /// candidate, and every candidate's re-derived canonical digest equals - /// its manifest entry (or the privacy boundary withholds it and the - /// manifest agrees). Read + sanitize + digest is per-file pure work over - /// independent paths, so it fans out across the indexing pool exactly as - /// capture does. - #[hotpath::measure(label = "daemon.code_index.freshness.content_verify")] - pub fn content_matches( - &self, - project_root: &Path, - manifest: &SourceContentManifestV1, - shutting_down: &AtomicBool, - ) -> bool { - let candidates = self - .candidates - .iter() - .filter(|candidate| { - candidate.explicitly_admitted || !is_generated_path_segment(&candidate.logical_path) - }) - .collect::>(); - hotpath::gauge!("daemon.code_index.freshness.content_verify.candidates") - .set(candidates.len() as u64); - let candidate_paths = candidates - .iter() - .map(|candidate| candidate.logical_path.as_str()) - .collect::>(); - if !manifest - .files - .keys() - .all(|logical_path| candidate_paths.contains(logical_path.as_str())) - { - return false; - } - let Ok(disputed) = parallelism::install(|| { - candidates - .par_iter() - .copied() - .filter(|candidate| { - parallelism::with_background_cpu_permit(|| { - shutting_down.load(Ordering::Acquire) - || !candidate_matches_manifest(project_root, candidate, manifest) - }) - }) - .collect::>() - }) else { - return false; - }; - if shutting_down.load(Ordering::Acquire) { - return false; - } - disputed.is_empty() - || tracked_files_match_after_clean_filters(project_root, &disputed, manifest) - } -} - fn read_candidate( project_root: &Path, candidate: &StatCandidateV1, @@ -278,15 +216,6 @@ impl CandidateContentV1 { } } -fn candidate_matches_manifest( - project_root: &Path, - candidate: &StatCandidateV1, - manifest: &SourceContentManifestV1, -) -> bool { - CandidateContentV1::derive(project_root, candidate) - .matches(manifest.files.get(&candidate.logical_path)) -} - /// Timestamps this close to the moment of a stat can still be shared by a /// later write (coarse kernel clocks, two-second FAT times), so such a stat /// cannot yet tell the file's current state from its next one. @@ -524,8 +453,10 @@ pub struct SourceSweepCacheV1 { impl SourceSweepCacheV1 { /// Whether the bytes on disk still carry exactly the manifest's content - /// identities, with the same verdict [`WorktreeStatSweepV1::content_matches`] - /// gives over a fresh sweep under the same roster. + /// identities: every present manifest file is still a candidate, and + /// every candidate's canonical digest equals its manifest entry (or the + /// privacy boundary withholds it and the manifest agrees). Digests no + /// settled key vouches for are re-derived across the indexing pool. #[hotpath::measure(label = "daemon.code_index.freshness.source_sweep")] pub fn witness_matches( &mut self, diff --git a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/reconcile.rs b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/reconcile.rs index c320467721..909f566c1a 100644 --- a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/reconcile.rs +++ b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/reconcile.rs @@ -734,20 +734,40 @@ impl SourceFreshnessFenceV1 { shutting_down: &AtomicBool, ) -> bool { freshness.source_witness.as_ref().is_some_and(|witness| { - self.sweep_cache - .lock() - .unwrap_or_else(PoisonError::into_inner) - .witness_matches( - project_root, - &freshness.source_roster, - &git_metadata.stable_signature(), - &witness.content_manifest, - shutting_down, - ) - .0 + self.sweep_sources( + project_root, + &freshness.source_roster, + git_metadata, + &witness.content_manifest, + shutting_down, + ) + .0 }) } + /// The one source sweep: every content proof, the freshness ladder's and + /// a reconcile's, goes through this cache, so digests one derives are + /// reused by the next. + fn sweep_sources( + &self, + project_root: &Path, + roster: &[CodeIndexIgnoredSourceAdmissionV1], + git_metadata: &identity::GitMetadataFingerprintV1, + manifest: &SourceContentManifestV1, + shutting_down: &AtomicBool, + ) -> (bool, freshness_witness::SourceSweepStatsV1) { + self.sweep_cache + .lock() + .unwrap_or_else(PoisonError::into_inner) + .witness_matches( + project_root, + roster, + &git_metadata.stable_signature(), + manifest, + shutting_down, + ) + } + /// Sweep the current witness and report what the sweep had to read. #[cfg(test)] pub(super) fn source_sweep_for_test( @@ -760,16 +780,13 @@ impl SourceFreshnessFenceV1 { .source_witness .as_ref() .expect("a reconciled fence carries a source witness"); - self.sweep_cache - .lock() - .unwrap_or_else(PoisonError::into_inner) - .witness_matches( - project_root, - &freshness.source_roster, - &identity::GitMetadataFingerprintV1::capture(project_root).stable_signature(), - &witness.content_manifest, - shutting_down, - ) + self.sweep_sources( + project_root, + &freshness.source_roster, + &identity::GitMetadataFingerprintV1::capture(project_root), + &witness.content_manifest, + shutting_down, + ) } /// The freshness ladder over this fence alone, so a caller verifying the @@ -1668,7 +1685,7 @@ impl CodeIndexWorktreeSchedulerV1 { else { return Ok(None); }; - if self.retained_frontier_stat_sweep(&pointer).is_none() { + if !self.retained_frontier_stat_sweep(&pointer) { return Ok(None); } let Some(generation) = self @@ -1730,7 +1747,7 @@ impl CodeIndexWorktreeSchedulerV1 { return Ok(None); } let source_manifest = SourceContentManifestV1::for_snapshot(generation.snapshot()); - if !sweep.content_matches(&self.project_root, &source_manifest, &self.shutting_down) { + if !self.sources_match_manifest(&source_manifest) { return Ok(None); } let snapshot_content_identity = generation.snapshot().content_identity.clone(); @@ -1796,13 +1813,9 @@ impl CodeIndexWorktreeSchedulerV1 { // the first branch and would keep swallowing the pass. let content_matches = decoded.as_ref().is_some_and(|generation| { self.retained_frontier_stat_sweep(&pointer) - .is_some_and(|sweep| { - sweep.content_matches( - &self.project_root, - &SourceContentManifestV1::for_snapshot(generation.snapshot()), - &self.shutting_down, - ) - }) + && self.sources_match_manifest(&SourceContentManifestV1::for_snapshot( + generation.snapshot(), + )) }); if !retained_empty_seat_settles_source(decoded.is_some(), content_matches) { return Ok(None); @@ -1836,24 +1849,22 @@ impl CodeIndexWorktreeSchedulerV1 { } /// Witness + git/stat fence that does not read sealed generation bytes. - /// `Some` is the negative cache only, the metadata the witness recorded - /// has not moved, and hands back the sweep so the caller can settle - /// currency against the retained generation's sealed file digests. - fn retained_frontier_stat_sweep( - &self, - pointer: &DurablePublicationPointerV1, - ) -> Option { - let witness = RestoreFreshnessWitnessV1::load(&self.store_root)?; + /// `true` is the negative cache only: the metadata the witness recorded + /// has not moved, and the caller still settles currency against the + /// retained generation's sealed file digests. + fn retained_frontier_stat_sweep(&self, pointer: &DurablePublicationPointerV1) -> bool { + let Some(witness) = RestoreFreshnessWitnessV1::load(&self.store_root) else { + return false; + }; if witness.generation_id != pointer.generation_id { - return None; + return false; } let metadata = identity::GitMetadataFingerprintV1::capture(&self.project_root); if witness.git_metadata_signature != metadata.stable_signature() { - return None; + return false; } self.worktree_stat_sweep() - .ok() - .filter(|sweep| witness.stat_signature == sweep.signature) + .is_ok_and(|sweep| witness.stat_signature == sweep.signature) } /// Verify an unchanged retained text generation without decoding the full @@ -2094,13 +2105,8 @@ impl CodeIndexWorktreeSchedulerV1 { // generation's sealed file digests; its matching stat signature is // the negative cache that lets a moved tree skip the byte comparison. let source_manifest = SourceContentManifestV1::for_snapshot(metadata.snapshot()); - let sealed_bytes_match = retained_is_reusable - && !has_hints - && sampled_sweep.content_matches( - &self.project_root, - &source_manifest, - &self.shutting_down, - ); + let sealed_bytes_match = + retained_is_reusable && !has_hints && self.sources_match_manifest(&source_manifest); let quiet_witness = sealed_bytes_match && witness.as_ref().is_some_and(|witness| { witness.git_metadata_signature == sampled_metadata.stable_signature() @@ -3155,6 +3161,20 @@ impl CodeIndexWorktreeSchedulerV1 { } } + /// Whether the bytes on disk carry exactly `manifest`'s digests under this + /// scheduler's ignored-source roster. + fn sources_match_manifest(&self, manifest: &SourceContentManifestV1) -> bool { + self.freshness_fence + .sweep_sources( + &self.project_root, + &self.ignored_source_admissions, + &identity::GitMetadataFingerprintV1::capture(&self.project_root), + manifest, + &self.shutting_down, + ) + .0 + } + /// Whether the last reconcile's source witness still describes the /// worktree: every candidate's content digest equals the sealed manifest. fn source_witness_matches_worktree(&self, freshness: &SourceFreshnessFenceStateV1) -> bool { From c2f4c9298dd0c493f346a071b2d47af2c3367502 Mon Sep 17 00:00:00 2001 From: ScriptedAlchemy Date: Sun, 27 Sep 2026 13:48:53 +0000 Subject: [PATCH 4/6] perf(code-index): stat sweep candidates across the indexing pool --- .../code_index_scheduler/freshness_witness.rs | 58 +++++++++++-------- 1 file changed, 33 insertions(+), 25 deletions(-) diff --git a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/freshness_witness.rs b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/freshness_witness.rs index 837d5f6908..996a9074ff 100644 --- a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/freshness_witness.rs +++ b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/freshness_witness.rs @@ -306,10 +306,12 @@ impl CachedCandidateRosterV1 { fn holds(&self, git_metadata_signature: &str, admitted_paths: &[String]) -> bool { self.git_metadata_signature == git_metadata_signature && self.admitted_paths == admitted_paths - && self - .evidence - .iter() - .all(|(path, key)| sample(path).is_ok_and(|now| now == *key)) + && parallelism::install(|| { + self.evidence + .par_iter() + .all(|(path, key)| sample(path).is_ok_and(|now| now == *key)) + }) + .unwrap_or(false) } } @@ -501,27 +503,33 @@ impl SourceSweepCacheV1 { candidates } }; - let mut present = Vec::new(); - for candidate in candidates.iter().filter(|candidate| { - candidate.explicitly_admitted || !is_generated_path_segment(&candidate.logical_path) - }) { - let absolute = project_root.join(&candidate.logical_path); - let Ok(metadata) = std::fs::symlink_metadata(&absolute) else { - continue; - }; - // A link's own key says nothing about its target's bytes. - let key = if metadata.file_type().is_symlink() { - if !std::fs::metadata(&absolute).is_ok_and(|target| target.is_file()) { - continue; - } - None - } else if metadata.is_file() { - Some(StatKeyV1::of(&metadata)) - } else { - continue; - }; - present.push((candidate, key)); - } + let Ok(present) = parallelism::install(|| { + candidates + .par_iter() + .filter(|candidate| { + candidate.explicitly_admitted + || !is_generated_path_segment(&candidate.logical_path) + }) + .filter_map(|candidate| { + let absolute = project_root.join(&candidate.logical_path); + let metadata = std::fs::symlink_metadata(&absolute).ok()?; + // A link's own key says nothing about its target's bytes. + let key = if metadata.file_type().is_symlink() { + if !std::fs::metadata(&absolute).is_ok_and(|target| target.is_file()) { + return None; + } + None + } else if metadata.is_file() { + Some(StatKeyV1::of(&metadata)) + } else { + return None; + }; + Some((candidate, key)) + }) + .collect::>() + }) else { + return (false, stats); + }; stats.candidates = present.len(); hotpath::gauge!("daemon.code_index.freshness.source_sweep.candidates") .set(present.len() as u64); From 54c37a29fb8378e598ebee1806ea2c4142f5e560 Mon Sep 17 00:00:00 2001 From: ScriptedAlchemy Date: Sun, 27 Sep 2026 14:32:30 +0000 Subject: [PATCH 5/6] fix(code-index): convert stat nanos without a wrapping cast --- .../src/code_index_scheduler/freshness_witness.rs | 6 ++++-- 1 file changed, 4 insertions(+), 2 deletions(-) diff --git a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/freshness_witness.rs b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/freshness_witness.rs index 996a9074ff..1a076582f5 100644 --- a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/freshness_witness.rs +++ b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/freshness_witness.rs @@ -265,7 +265,8 @@ impl StatKeyV1 { .modified() .ok() .and_then(|time| time.duration_since(UNIX_EPOCH).ok()) - .map_or(0, |elapsed| elapsed.as_nanos() as i128), + .and_then(|elapsed| i128::try_from(elapsed.as_nanos()).ok()) + .unwrap_or(0), changed_nanos: 0, } } @@ -277,7 +278,8 @@ impl StatKeyV1 { && sampled_at .checked_sub(RACY_STAT_WINDOW) .and_then(|horizon| horizon.duration_since(UNIX_EPOCH).ok()) - .is_some_and(|horizon| self.changed_nanos < horizon.as_nanos() as i128) + .and_then(|horizon| i128::try_from(horizon.as_nanos()).ok()) + .is_some_and(|horizon| self.changed_nanos < horizon) } } From 784aa6d78254da9c5e6204cda8dc3384f9edfce3 Mon Sep 17 00:00:00 2001 From: ScriptedAlchemy Date: Sun, 27 Sep 2026 14:41:28 +0000 Subject: [PATCH 6/6] perf(code-index): keep freshness stats off the indexing pool --- .../code_index_scheduler/freshness_witness.rs | 133 +++++++++++------- 1 file changed, 84 insertions(+), 49 deletions(-) diff --git a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/freshness_witness.rs b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/freshness_witness.rs index 1a076582f5..f8e4e02c6f 100644 --- a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/freshness_witness.rs +++ b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/freshness_witness.rs @@ -33,6 +33,7 @@ use tracedecay_domain::canonical_text::encode_tagged_lowercase_hex; use tracedecay_domain::{ ContentDigest, LanguageId, SanitizedCodeSnapshotV1, SnapshotFileDispositionV1, }; +use tracedecay_runtime_core::git_repository::GIT_STATUS_MODIFICATION_CHECK_THREADS; use super::{CodeIndexSchedulerErrorV1, classification, ignored_dependencies, privacy}; use crate::code_index::chunks::content_digest; @@ -221,6 +222,9 @@ impl CandidateContentV1 { /// cannot yet tell the file's current state from its next one. const RACY_STAT_WINDOW: Duration = Duration::from_secs(2); +/// At most this many digests are re-derived on the sweeping thread. +const INLINE_DIGEST_LIMIT: usize = 16; + /// One inode's stat identity. On Unix the change time is the kernel's own /// record of every content or metadata write and cannot be set back /// (`touch -d`, `cp --preserve`, `rsync -a` all advance it), so an equal @@ -283,6 +287,31 @@ impl StatKeyV1 { } } +/// Apply a stat-sized `probe` to every item on a few scoped threads, in order. +/// Stats never queue on the indexing pool, so a sweep is not delayed by the +/// builds that pool is running. +fn stat_concurrently(items: &[T], probe: impl Fn(&T) -> R + Sync) -> Vec { + let chunk = items + .len() + .div_ceil(GIT_STATUS_MODIFICATION_CHECK_THREADS) + .max(1); + std::thread::scope(|scope| { + let probe = &probe; + let workers = items + .chunks(chunk) + .map(|slice| scope.spawn(move || slice.iter().map(probe).collect::>())) + .collect::>(); + workers + .into_iter() + .flat_map(|worker| { + worker + .join() + .unwrap_or_else(|panic| std::panic::resume_unwind(panic)) + }) + .collect() + }) +} + /// The key of whatever is at `path` now, `None` when nothing is. fn sample(path: &Path) -> std::io::Result> { match std::fs::symlink_metadata(path) { @@ -308,12 +337,11 @@ impl CachedCandidateRosterV1 { fn holds(&self, git_metadata_signature: &str, admitted_paths: &[String]) -> bool { self.git_metadata_signature == git_metadata_signature && self.admitted_paths == admitted_paths - && parallelism::install(|| { - self.evidence - .par_iter() - .all(|(path, key)| sample(path).is_ok_and(|now| now == *key)) + && stat_concurrently(&self.evidence, |(path, key)| { + sample(path).is_ok_and(|now| now == *key) }) - .unwrap_or(false) + .into_iter() + .all(|holds| holds) } } @@ -505,33 +533,31 @@ impl SourceSweepCacheV1 { candidates } }; - let Ok(present) = parallelism::install(|| { - candidates - .par_iter() - .filter(|candidate| { - candidate.explicitly_admitted - || !is_generated_path_segment(&candidate.logical_path) - }) - .filter_map(|candidate| { - let absolute = project_root.join(&candidate.logical_path); - let metadata = std::fs::symlink_metadata(&absolute).ok()?; - // A link's own key says nothing about its target's bytes. - let key = if metadata.file_type().is_symlink() { - if !std::fs::metadata(&absolute).is_ok_and(|target| target.is_file()) { - return None; - } - None - } else if metadata.is_file() { - Some(StatKeyV1::of(&metadata)) - } else { - return None; - }; - Some((candidate, key)) - }) - .collect::>() - }) else { - return (false, stats); - }; + let eligible = candidates + .iter() + .filter(|candidate| { + candidate.explicitly_admitted || !is_generated_path_segment(&candidate.logical_path) + }) + .collect::>(); + let present = stat_concurrently(&eligible, |candidate| { + let absolute = project_root.join(&candidate.logical_path); + let metadata = std::fs::symlink_metadata(&absolute).ok()?; + // A link's own key says nothing about its target's bytes. + let key = if metadata.file_type().is_symlink() { + if !std::fs::metadata(&absolute).is_ok_and(|target| target.is_file()) { + return None; + } + None + } else if metadata.is_file() { + Some(StatKeyV1::of(&metadata)) + } else { + return None; + }; + Some((*candidate, key)) + }) + .into_iter() + .flatten() + .collect::>(); stats.candidates = present.len(); hotpath::gauge!("daemon.code_index.freshness.source_sweep.candidates") .set(present.len() as u64); @@ -561,29 +587,38 @@ impl SourceSweepCacheV1 { stats.hashed = unvouched.len(); hotpath::gauge!("daemon.code_index.freshness.source_sweep.hashed") .set(unvouched.len() as u64); - let Ok(derived) = parallelism::install(|| { + let derive = |candidate: &StatCandidateV1| { + if shutting_down.load(Ordering::Acquire) { + CandidateContentV1::Unreadable + } else { + CandidateContentV1::derive(project_root, candidate) + } + }; + // The few files an edit touches are read on this thread; only a bulk + // re-derivation (a cold cache) takes the indexing pool and its CPU + // permits. + let derived = if unvouched.len() <= INLINE_DIGEST_LIMIT { unvouched - .par_iter() - .map(|(candidate, key)| { - parallelism::with_background_cpu_permit(|| { - if shutting_down.load(Ordering::Acquire) { - return (*candidate, *key, CandidateContentV1::Unreadable); - } - ( - *candidate, - *key, - CandidateContentV1::derive(project_root, candidate), - ) - }) - }) + .iter() + .map(|(candidate, _)| derive(candidate)) .collect::>() - }) else { - return (false, stats); + } else { + let Ok(derived) = parallelism::install(|| { + unvouched + .par_iter() + .map(|(candidate, _)| { + parallelism::with_background_cpu_permit(|| derive(candidate)) + }) + .collect::>() + }) else { + return (false, stats); + }; + derived }; if shutting_down.load(Ordering::Acquire) { return (false, stats); } - for (candidate, key, content) in derived { + for ((candidate, key), content) in unvouched.into_iter().zip(derived) { if !content.matches(manifest.files.get(&candidate.logical_path)) { disputed.push(candidate); }