Skip to main content

rac_engine/
freshness.rs

1//! Server-lifetime serving freshness (ADR-105) — port of
2//! `services/freshness.py` `FreshnessTracker` for the long-lived MCP server
3//! (INDEX-PLAN B6).
4//!
5//! Detection uses an event-driven clean accelerator where the platform can
6//! provide a synchronous barrier, otherwise the stat-manifest scan. Events
7//! never compute the changed set: any dirty or uncertain signal falls back to
8//! the authoritative scan. Whatever the rung, the served read-model is
9//! built from the tracker's incrementally maintained generations,
10//! byte-identical to a fresh whole-corpus walk at the current corpus state.
11
12use std::collections::{BTreeMap, BTreeSet, HashMap, HashSet};
13use std::path::PathBuf;
14
15use crate::delta_generation::{
16    DeltaDocuments, DeltaGeneration, GraphGeneration, IdentityGeneration, ScopeGeneration,
17    SearchGeneration, SummaryGeneration,
18};
19use crate::derived::{build_derived_index_from_items, DerivedIndex, SCHEMA_VERSION};
20use crate::derived_cache::{corpus_hash_from_complete_manifest, stat_scan};
21use crate::freshness_watch::EventWatch;
22use crate::index_store::{open_store, write_store, FileState, MmapIndexReader};
23use crate::relationships::CorpusItem;
24
25/// What the tracker currently serves: the memory-mapped base (delta-empty),
26/// or the re-derived snapshot bundle (the delta window).
27pub enum TrackerModel {
28    View(MmapIndexReader),
29    Snapshot(DerivedIndex),
30    /// P6 production generation. The document overlay is immutable for the
31    /// lifetime of this served model and is published only after every
32    /// incremental projection has been staged successfully.
33    Delta(Box<DeltaGeneration>),
34}
35
36pub struct FreshnessTracker {
37    cache_dir: PathBuf,
38    root_str: String,
39    threshold: Option<usize>,
40    watcher: EventWatch,
41
42    manifest: Vec<(String, FileState)>,
43    items: HashMap<String, CorpusItem>, // rel -> parsed snapshot entry
44    model: Option<TrackerModel>,
45    hash: Option<String>,
46    base_hash: Option<String>,
47    base_generation: u64,
48    /// Logical served-corpus generation. Unlike `base_generation`, this also
49    /// advances for mutation-window snapshots that have not compacted.
50    serving_generation: u64,
51    delta_paths: HashSet<String>,
52    /// Present for the production P6 delta lifecycle. The explicit snapshot
53    /// fallback leaves these absent.
54    delta_documents: Option<DeltaDocuments>,
55    delta_identity: Option<IdentityGeneration>,
56    delta_search: Option<SearchGeneration>,
57    delta_graph: Option<GraphGeneration>,
58    delta_scope: Option<ScopeGeneration>,
59    delta_summary: Option<SummaryGeneration>,
60    /// ADR-107 RSS finalization: after compaction the resident parsed
61    /// snapshot is shed and the mapped base is the whole answer; the next
62    /// change repopulates by a full re-parse on demand.
63    snapshot_shed: bool,
64    last_parse_workers: usize,
65    last_parse_files: usize,
66    last_detect_scanned: bool,
67}
68
69impl FreshnessTracker {
70    /// Production S1 freshness: immutable base plus cumulative delta.
71    pub fn new(cache_dir: PathBuf, root: &str, threshold: Option<usize>) -> Self {
72        Self::new_delta(cache_dir, root, threshold, true)
73    }
74
75    /// Explicit rollback path retained for the first S1 soak release.
76    pub fn new_snapshot(cache_dir: PathBuf, root: &str, threshold: Option<usize>) -> Self {
77        Self::new_with_watcher(cache_dir, root, threshold, true)
78    }
79
80    /// Force the authoritative stat rung. Used by fallback/parity tests and
81    /// remains the behavior on platforms without a synchronous watcher.
82    pub fn new_stat(cache_dir: PathBuf, root: &str, threshold: Option<usize>) -> Self {
83        Self::new_delta(cache_dir, root, threshold, false)
84    }
85
86    fn new_with_watcher(
87        cache_dir: PathBuf,
88        root: &str,
89        threshold: Option<usize>,
90        watcher_enabled: bool,
91    ) -> Self {
92        Self {
93            cache_dir,
94            root_str: root.to_string(),
95            threshold,
96            watcher: EventWatch::new(root, watcher_enabled),
97            manifest: Vec::new(),
98            items: HashMap::new(),
99            model: None,
100            hash: None,
101            base_hash: None,
102            base_generation: 0,
103            serving_generation: 0,
104            delta_paths: HashSet::new(),
105            delta_documents: None,
106            delta_identity: None,
107            delta_search: None,
108            delta_graph: None,
109            delta_scope: None,
110            delta_summary: None,
111            snapshot_shed: false,
112            last_parse_workers: 1,
113            last_parse_files: 0,
114            last_detect_scanned: false,
115        }
116    }
117
118    fn new_delta(
119        cache_dir: PathBuf,
120        root: &str,
121        threshold: Option<usize>,
122        watcher_enabled: bool,
123    ) -> Self {
124        let mut tracker = Self::new_with_watcher(cache_dir, root, threshold, watcher_enabled);
125        tracker.delta_documents = Some(DeltaDocuments::empty());
126        tracker.delta_identity = Some(IdentityGeneration::empty());
127        tracker.delta_search = Some(SearchGeneration::empty());
128        tracker.delta_graph = Some(GraphGeneration::empty());
129        tracker.delta_scope = Some(ScopeGeneration::empty());
130        tracker.delta_summary = Some(SummaryGeneration::empty());
131        tracker
132    }
133
134    // --- observable state (scorecards and pinning tests) ------------------
135
136    pub fn mode(&self) -> &'static str {
137        self.watcher.mode()
138    }
139
140    pub fn base_generation(&self) -> u64 {
141        self.base_generation
142    }
143
144    pub fn serving_generation(&self) -> u64 {
145        self.serving_generation
146    }
147
148    pub fn delta_size(&self) -> usize {
149        self.delta_documents
150            .as_ref()
151            .map_or_else(|| self.delta_paths.len(), DeltaDocuments::delta_len)
152    }
153
154    pub fn delta_enabled(&self) -> bool {
155        self.delta_documents.is_some()
156    }
157
158    pub fn delta_base_documents(&self) -> usize {
159        self.delta_documents
160            .as_ref()
161            .map_or(0, DeltaDocuments::base_len)
162    }
163
164    pub fn delta_upserts(&self) -> usize {
165        self.delta_documents
166            .as_ref()
167            .map_or(0, DeltaDocuments::upsert_len)
168    }
169
170    pub fn delta_tombstones(&self) -> usize {
171        self.delta_documents
172            .as_ref()
173            .map_or(0, DeltaDocuments::tombstone_len)
174    }
175
176    pub fn last_parse_files(&self) -> usize {
177        self.last_parse_files
178    }
179
180    pub fn corpus_hash(&self) -> Option<&str> {
181        self.hash.as_deref()
182    }
183
184    pub fn last_detect_scanned(&self) -> bool {
185        self.last_detect_scanned
186    }
187
188    // --- the serving surface ----------------------------------------------
189
190    /// The current read-model, freshened through the detection ladder. An
191    /// unchanged corpus returns the cached model with no re-derive.
192    pub fn read_model(&mut self, verify: bool) -> &TrackerModel {
193        let cold = self.model.is_none();
194        let detect_started = crate::timing::start();
195        let (changed, scanned) = self.detect(verify);
196        self.last_detect_scanned = scanned;
197        crate::timing::emit_since(
198            "tracker.detect",
199            detect_started,
200            &[
201                ("files", self.manifest.len() as u64),
202                ("changed", changed.len() as u64),
203                ("scanned", u64::from(scanned)),
204            ],
205        );
206        if changed.is_empty() && !cold {
207            return self.model.as_ref().expect("warm model");
208        }
209        if !cold {
210            let recompute_started = crate::timing::start();
211            if self.delta_documents.is_some() {
212                self.rebuild_delta(&changed);
213            } else {
214                self.apply(&changed);
215                self.rebuild_model();
216            }
217            self.maybe_compact();
218            crate::timing::emit_since(
219                "tracker.recompute",
220                recompute_started,
221                &[("changed", changed.len() as u64), ("cold", 0)],
222            );
223            return self.model.as_ref().expect("rebuilt model");
224        }
225        // Cold start: the whole corpus parsed from nothing; the three cold
226        // phases feed the DECIDED_TIMING scorecard (ADR-107).
227        let parse_start = std::time::Instant::now();
228        if self.delta_documents.is_some() {
229            self.rebuild_delta(&changed);
230        } else {
231            self.apply(&changed);
232        }
233        let derive_start = std::time::Instant::now();
234        if self.delta_documents.is_none() {
235            self.rebuild_model();
236        }
237        let write_start = std::time::Instant::now();
238        self.maybe_compact();
239        let end = std::time::Instant::now();
240        crate::timing::emit_since(
241            "tracker.recompute",
242            Some(parse_start),
243            &[("changed", changed.len() as u64), ("cold", 1)],
244        );
245        crate::parallel_build::emit_build_timing(&crate::parallel_build::BuildStats {
246            files: self.manifest.len(),
247            workers: self.last_parse_workers,
248            parse_ms: (derive_start - parse_start).as_secs_f64() * 1000.0,
249            derive_ms: (write_start - derive_start).as_secs_f64() * 1000.0,
250            write_ms: (end - write_start).as_secs_f64() * 1000.0,
251        });
252        self.model.as_ref().expect("cold model")
253    }
254
255    /// Freshen and return the logical corpus generation with its model. Server
256    /// lifetime derived views use the generation as their invalidation key.
257    pub fn read_model_with_generation(&mut self, verify: bool) -> (u64, &TrackerModel) {
258        self.read_model(verify);
259        (
260            self.serving_generation,
261            self.model.as_ref().expect("freshened model"),
262        )
263    }
264
265    // --- detection ----------------------------------------------------------
266
267    fn detect(&mut self, verify: bool) -> (std::collections::BTreeSet<String>, bool) {
268        let confirm_all = self.model.is_none() || verify;
269        if !confirm_all && self.watcher.is_clean() {
270            return (std::collections::BTreeSet::new(), false);
271        }
272
273        let mut all_changed = std::collections::BTreeSet::new();
274        // A stable bracket is the barrier: if an event arrives while scanning,
275        // scan again. Under continuous writes, leave the watcher unacknowledged
276        // after the bounded retries so the next call scans again.
277        for _ in 0..3 {
278            self.watcher.prepare_scan();
279            let before = self.watcher.checkpoint();
280            let (new_manifest, changed) =
281                stat_scan(&self.root_str, &self.manifest, confirm_all, true);
282            self.manifest = new_manifest;
283            all_changed.extend(changed);
284            let Some(before) = before else {
285                return (all_changed, true);
286            };
287            if self.watcher.acknowledge_if_stable(before) {
288                return (all_changed, true);
289            }
290        }
291        (all_changed, true)
292    }
293
294    // --- applying the changed set -------------------------------------------
295
296    fn apply(&mut self, changed: &std::collections::BTreeSet<String>) {
297        let current: HashSet<&str> = self.manifest.iter().map(|(rel, _)| rel.as_str()).collect();
298        if self.snapshot_shed {
299            self.reparse_full();
300            self.snapshot_shed = false;
301        } else {
302            let root = PathBuf::from(&self.root_str);
303            let present: Vec<PathBuf> = changed
304                .iter()
305                .filter(|rel| current.contains(rel.as_str()))
306                .map(|rel| root.join(rel))
307                .collect();
308            for rel in changed {
309                if !current.contains(rel.as_str()) {
310                    self.items.remove(rel); // removed
311                }
312            }
313            let (parsed, workers) = crate::parallel_build::parallel_parse_paths(&present);
314            self.last_parse_workers = workers;
315            self.last_parse_files = present.len();
316            for item in parsed {
317                let rel = rel_of(&self.root_str, &item.path);
318                self.items.insert(rel, item);
319            }
320        }
321        let current: HashSet<String> = self
322            .manifest
323            .iter()
324            .map(|(rel, _)| rel.clone())
325            .collect();
326        self.items.retain(|rel, _| current.contains(rel));
327        if !changed.is_empty() {
328            self.delta_paths.extend(changed.iter().cloned());
329        }
330        self.hash = Some(corpus_hash_from_complete_manifest(&self.manifest));
331    }
332
333    fn reparse_full(&mut self) {
334        let root = PathBuf::from(&self.root_str);
335        let paths: Vec<PathBuf> = self.manifest.iter().map(|(rel, _)| root.join(rel)).collect();
336        let (parsed, workers) = crate::parallel_build::parallel_parse_paths(&paths);
337        self.last_parse_workers = workers;
338        self.last_parse_files = paths.len();
339        self.items = parsed
340            .into_iter()
341            .map(|item| (rel_of(&self.root_str, &item.path), item))
342            .collect();
343    }
344
345    /// The snapshot in walk (sorted-path) order — the fresh-walk order.
346    fn ordered_items(&self) -> Vec<CorpusItem> {
347        crate::walk::find_markdown_files(&self.root_str, true)
348            .into_iter()
349            .filter_map(|entry| self.items.get(&entry.components.join("/")).cloned())
350            .collect()
351    }
352
353    fn rebuild_model(&mut self) {
354        let hash = self.hash.clone().expect("hash set by apply");
355        if Some(hash.as_str()) == self.base_hash.as_deref() && self.delta_paths.is_empty() {
356            if let Some(view) = open_store(&self.cache_dir, &hash, SCHEMA_VERSION) {
357                self.model = Some(TrackerModel::View(view));
358                return;
359            }
360        }
361        let derived =
362            build_derived_index_from_items(&self.root_str, &self.ordered_items(), true);
363        self.model = Some(TrackerModel::Snapshot(derived));
364        self.serving_generation += 1;
365    }
366
367    /// Build a complete candidate generation from staged overlays, then swap
368    /// every serving field only after all projections succeed.
369    fn rebuild_delta(&mut self, changed: &BTreeSet<String>) {
370        let current: HashSet<&str> = self.manifest.iter().map(|(rel, _)| rel.as_str()).collect();
371        let root = PathBuf::from(&self.root_str);
372        let present: Vec<PathBuf> = changed
373            .iter()
374            .filter(|rel| current.contains(rel.as_str()))
375            .map(|rel| root.join(rel))
376            .collect();
377        let (parsed, workers) = crate::parallel_build::parallel_parse_paths(&present);
378        self.last_parse_workers = workers;
379        self.last_parse_files = present.len();
380        let parsed: BTreeMap<String, CorpusItem> = parsed
381            .into_iter()
382            .map(|item| (rel_of(&self.root_str, &item.path), item))
383            .collect();
384
385        // A parser omission would make the staged generation incomplete.
386        // Reparse the current corpus from an empty base instead of publishing
387        // a partial overlay.
388        let (
389            candidate,
390            identity_candidate,
391            search_candidate,
392            graph_candidate,
393            scope_candidate,
394            summary_candidate,
395        ) =
396            if self.model.is_none() && parsed.len() == present.len() {
397                self.delta_candidate_from_parsed(parsed)
398            } else if parsed.len() == present.len() {
399                let identity = self
400                    .delta_identity
401                    .as_ref()
402                    .expect("delta identity")
403                    .stage(changed, &parsed);
404                let search = self
405                    .delta_search
406                    .as_ref()
407                    .expect("delta search")
408                    .stage(changed, &parsed);
409                let graph = self
410                    .delta_graph
411                    .as_ref()
412                    .expect("delta graph")
413                    .stage(changed, &parsed, &identity);
414                let scope = self
415                    .delta_scope
416                    .as_ref()
417                    .expect("delta scope")
418                    .stage(changed, &parsed);
419                let summary = self
420                    .delta_summary
421                    .as_ref()
422                    .expect("delta summary")
423                    .stage(changed, &parsed);
424                let documents = self
425                    .delta_documents
426                    .as_ref()
427                    .expect("delta documents")
428                    .stage(changed, parsed);
429                (documents, identity, search, graph, scope, summary)
430            } else {
431                self.full_delta_candidate()
432            };
433        let (
434            candidate,
435            identity_candidate,
436            search_candidate,
437            graph_candidate,
438            scope_candidate,
439            summary_candidate,
440        ) =
441            if candidate.live_len() == self.manifest.len() {
442                (
443                    candidate,
444                    identity_candidate,
445                    search_candidate,
446                    graph_candidate,
447                    scope_candidate,
448                    summary_candidate,
449                )
450            } else {
451                self.full_delta_candidate()
452            };
453        let hash = corpus_hash_from_complete_manifest(&self.manifest);
454        let serving_generation = self.serving_generation + 1;
455        let generation = DeltaGeneration {
456            base_generation: self.base_generation,
457            serving_generation,
458            changed_paths: candidate.changed_paths(),
459            identity: identity_candidate.clone(),
460            search: search_candidate.clone(),
461            graph: graph_candidate.clone(),
462            scope: scope_candidate.clone(),
463            summary: summary_candidate.clone(),
464        };
465
466        self.delta_documents = Some(candidate);
467        self.delta_identity = Some(identity_candidate);
468        self.delta_search = Some(search_candidate);
469        self.delta_graph = Some(graph_candidate);
470        self.delta_scope = Some(scope_candidate);
471        self.delta_summary = Some(summary_candidate);
472        self.hash = Some(hash);
473        self.serving_generation = serving_generation;
474        self.model = Some(TrackerModel::Delta(Box::new(generation)));
475    }
476
477    fn full_delta_candidate(
478        &mut self,
479    ) -> (
480        DeltaDocuments,
481        IdentityGeneration,
482        SearchGeneration,
483        GraphGeneration,
484        ScopeGeneration,
485        SummaryGeneration,
486    ) {
487        let root = PathBuf::from(&self.root_str);
488        let paths: Vec<PathBuf> = self
489            .manifest
490            .iter()
491            .map(|(rel, _)| root.join(rel))
492            .collect();
493        let (parsed, workers) = crate::parallel_build::parallel_parse_paths(&paths);
494        self.last_parse_workers = workers;
495        self.last_parse_files = paths.len();
496        let parsed: BTreeMap<String, CorpusItem> = parsed
497            .into_iter()
498            .map(|item| (rel_of(&self.root_str, &item.path), item))
499            .collect();
500        self.delta_candidate_from_parsed(parsed)
501    }
502
503    fn delta_candidate_from_parsed(
504        &self,
505        parsed: BTreeMap<String, CorpusItem>,
506    ) -> (
507        DeltaDocuments,
508        IdentityGeneration,
509        SearchGeneration,
510        GraphGeneration,
511        ScopeGeneration,
512        SummaryGeneration,
513    ) {
514        let identity = IdentityGeneration::from_items(
515            parsed.iter().map(|(path, item)| (path.as_str(), item)),
516        );
517        let search = SearchGeneration::from_items(
518            parsed.iter().map(|(path, item)| (path.as_str(), item)),
519        );
520        let graph = GraphGeneration::from_items(
521            parsed.iter().map(|(path, item)| (path.as_str(), item)),
522            &identity,
523        );
524        let scope = ScopeGeneration::from_items(
525            parsed.iter().map(|(path, item)| (path.as_str(), item)),
526        );
527        let summary = SummaryGeneration::from_items(
528            parsed.iter().map(|(path, item)| (path.as_str(), item)),
529        );
530        let changed = self.manifest.iter().map(|(rel, _)| rel.clone()).collect();
531        (
532            DeltaDocuments::empty().stage(&changed, parsed),
533            identity,
534            search,
535            graph,
536            scope,
537            summary,
538        )
539    }
540
541    // --- compaction -----------------------------------------------------------
542
543    fn threshold_for(&self, base_count: usize) -> usize {
544        self.threshold.unwrap_or_else(|| 10_000.max(base_count / 100))
545    }
546
547    fn maybe_compact(&mut self) {
548        if self.base_hash.is_none() {
549            self.compact(); // cold: establish the first base
550            return;
551        }
552        if self.delta_size() >= self.threshold_for(self.manifest.len()) {
553            self.compact();
554        }
555    }
556
557    fn compact(&mut self) {
558        let hash = self.hash.clone().expect("hash set");
559        let derived_owned;
560        let derived = match &self.model {
561            Some(TrackerModel::Snapshot(derived)) => derived,
562            Some(TrackerModel::Delta(generation)) => {
563                derived_owned = generation.materialize_derived(&self.root_str, true);
564                &derived_owned
565            }
566            _ => {
567                derived_owned =
568                    build_derived_index_from_items(&self.root_str, &self.ordered_items(), true);
569                &derived_owned
570            }
571        };
572        if !write_store(&self.cache_dir, &hash, SCHEMA_VERSION, derived) {
573            return; // unwritable cache dir: keep serving the snapshot (ADR-080)
574        }
575        crate::derived_cache::write_marker_public(&self.cache_dir, &hash);
576        let Some(view) = open_store(&self.cache_dir, &hash, SCHEMA_VERSION) else {
577            return;
578        };
579        self.model = Some(TrackerModel::View(view));
580        self.base_hash = Some(hash);
581        self.base_generation += 1;
582        self.delta_paths.clear();
583        if let Some(documents) = self.delta_documents.as_mut() {
584            let ordered_paths: Vec<&str> = self.manifest.iter().map(|(rel, _)| rel.as_str()).collect();
585            documents.promote(ordered_paths);
586            self.delta_identity
587                .as_mut()
588                .expect("delta identity")
589                .promote();
590            self.delta_search
591                .as_mut()
592                .expect("delta search")
593                .promote();
594            self.delta_graph
595                .as_mut()
596                .expect("delta graph")
597                .promote();
598            self.delta_scope
599                .as_mut()
600                .expect("delta scope")
601                .promote();
602            self.delta_summary
603                .as_mut()
604                .expect("delta summary")
605                .promote();
606            // P6 removes snapshot shedding for its parsed document base so
607            // the first post-compaction edit remains change-bound.
608            self.snapshot_shed = false;
609        } else {
610            // ADR-107 RSS finalization for the established default path.
611            self.items = HashMap::new();
612            self.snapshot_shed = true;
613        }
614    }
615}
616
617fn rel_of(root_str: &str, display_path: &str) -> String {
618    let root = crate::walk::normalize_root(root_str);
619    display_path
620        .strip_prefix(&format!("{root}/"))
621        .unwrap_or(display_path)
622        .to_string()
623}