Skip to main content

lean_ctx/server/
compaction_sync.rs

1use std::path::Path;
2use std::sync::atomic::{AtomicU64, Ordering};
3
4use crate::core::cache::SessionCache;
5use crate::core::context_radar::RadarEvent;
6
7pub static LAST_COMPACTION_TS: AtomicU64 = AtomicU64::new(0);
8
9/// Effective cache policy: "aggressive" (default), "safe", or "off".
10pub fn effective_cache_policy() -> &'static str {
11    static POLICY: std::sync::OnceLock<String> = std::sync::OnceLock::new();
12    POLICY.get_or_init(|| {
13        if let Ok(v) = std::env::var("LEAN_CTX_CACHE_POLICY") {
14            let v = v.trim().to_lowercase();
15            if matches!(v.as_str(), "aggressive" | "safe" | "off") {
16                return v;
17            }
18        }
19        let cfg = crate::core::config::Config::load();
20        cfg.cache_policy
21            .as_deref()
22            .unwrap_or("aggressive")
23            .to_lowercase()
24    })
25}
26
27/// Check if a host compaction event occurred since our last check.
28/// If so, reset all `full_content_delivered` flags so the next read
29/// delivers full content instead of a stub.
30pub fn sync_if_compacted(cache: &mut SessionCache, data_dir: &Path) -> bool {
31    let last_seen = LAST_COMPACTION_TS.load(Ordering::Relaxed);
32    let radar_path = data_dir.join("context_radar.jsonl");
33
34    if !radar_path.exists() {
35        return false;
36    }
37
38    let Some(latest_compaction_ts) = find_latest_compaction(&radar_path, last_seen) else {
39        return false;
40    };
41
42    LAST_COMPACTION_TS.store(latest_compaction_ts, Ordering::Relaxed);
43    crate::core::search_delta::reset();
44    let reset_count = cache.reset_delivery_flags();
45    crate::core::cache_telemetry::record_compaction(reset_count as u64);
46    // Drop the persistent stub index too (#955): the conversation's context was
47    // summarised away, so neither a warm nor a cold stub may claim "you already
48    // have this". Writes the emptied index synchronously so a restart in the
49    // crash window can't resurrect a pre-compaction stub.
50    crate::core::read_stub_index::reset_in_dir(data_dir);
51    if reset_count > 0 {
52        eprintln!(
53            "[lean-ctx] compaction detected — reset {reset_count} delivery flags for re-read"
54        );
55    }
56
57    std::thread::spawn(|| {
58        if let Some(session) = crate::core::session::SessionState::load_latest()
59            && let Some(ref root) = session.project_root
60            && (!session.findings.is_empty() || !session.decisions.is_empty())
61        {
62            crate::tools::startup::auto_consolidate_knowledge(root);
63        }
64    });
65
66    true
67}
68
69/// Scan only the tail of radar JSONL for a compaction event newer than `since_ts`.
70/// Reads at most 4KB from the end to avoid unbounded I/O on large radar files.
71fn find_latest_compaction(radar_path: &Path, since_ts: u64) -> Option<u64> {
72    use std::io::{Read, Seek, SeekFrom};
73
74    let mut file = std::fs::File::open(radar_path).ok()?;
75    let file_len = file.metadata().ok()?.len();
76
77    const TAIL_BYTES: u64 = 4096;
78    let content = if file_len <= TAIL_BYTES {
79        let mut s = String::new();
80        file.read_to_string(&mut s).ok()?;
81        s
82    } else {
83        file.seek(SeekFrom::End(-(TAIL_BYTES as i64))).ok()?;
84        let mut buf = vec![0u8; TAIL_BYTES as usize];
85        let n = file.read(&mut buf).ok()?;
86        let s = String::from_utf8_lossy(&buf[..n]).into_owned();
87        // Skip first partial line (we seeked into the middle of it)
88        if let Some(idx) = s.find('\n') {
89            s[idx + 1..].to_string()
90        } else {
91            s
92        }
93    };
94
95    for line in content.lines().rev() {
96        if line.is_empty() {
97            continue;
98        }
99        let event: RadarEvent = match serde_json::from_str(line) {
100            Ok(e) => e,
101            Err(_) => continue,
102        };
103        if event.ts <= since_ts {
104            break;
105        }
106        if event.event_type == "compaction" {
107            return Some(event.ts);
108        }
109    }
110    None
111}
112
113#[cfg(test)]
114mod tests {
115    use super::*;
116    use serial_test::serial;
117    use std::io::Write;
118    use tempfile::TempDir;
119
120    fn make_cache_with_delivered(paths: &[&str]) -> SessionCache {
121        let mut cache = SessionCache::default();
122        for p in paths {
123            cache.store(p, "hello world");
124            cache.mark_full_delivered(p);
125        }
126        cache
127    }
128
129    #[test]
130    #[serial]
131    fn no_reset_without_compaction_event() {
132        let dir = TempDir::new().unwrap();
133        let radar = dir.path().join("context_radar.jsonl");
134        let mut f = std::fs::File::create(&radar).unwrap();
135        writeln!(f, r#"{{"ts":1000,"event_type":"mcp_call","tokens":50}}"#).unwrap();
136        drop(f);
137
138        LAST_COMPACTION_TS.store(0, Ordering::Relaxed);
139        let mut cache = make_cache_with_delivered(&["/tmp/a.rs"]);
140        assert!(!sync_if_compacted(&mut cache, dir.path()));
141        assert!(cache.is_full_delivered("/tmp/a.rs"));
142    }
143
144    #[test]
145    #[serial]
146    fn resets_after_compaction() {
147        let dir = TempDir::new().unwrap();
148        let radar = dir.path().join("context_radar.jsonl");
149        let mut f = std::fs::File::create(&radar).unwrap();
150        writeln!(f, r#"{{"ts":1000,"event_type":"mcp_call","tokens":50}}"#).unwrap();
151        writeln!(f, r#"{{"ts":2000,"event_type":"compaction","tokens":0}}"#).unwrap();
152        drop(f);
153
154        LAST_COMPACTION_TS.store(0, Ordering::Relaxed);
155        let mut cache = make_cache_with_delivered(&["/tmp/a.rs", "/tmp/b.rs"]);
156
157        assert!(cache.is_full_delivered("/tmp/a.rs"));
158        assert!(sync_if_compacted(&mut cache, dir.path()));
159        assert!(!cache.is_full_delivered("/tmp/a.rs"));
160        assert!(!cache.is_full_delivered("/tmp/b.rs"));
161    }
162
163    #[test]
164    #[serial]
165    fn does_not_double_reset() {
166        let dir = TempDir::new().unwrap();
167        let radar = dir.path().join("context_radar.jsonl");
168        let mut f = std::fs::File::create(&radar).unwrap();
169        writeln!(f, r#"{{"ts":2000,"event_type":"compaction","tokens":0}}"#).unwrap();
170        drop(f);
171
172        LAST_COMPACTION_TS.store(0, Ordering::Relaxed);
173        let mut cache = make_cache_with_delivered(&["/tmp/a.rs"]);
174        assert!(sync_if_compacted(&mut cache, dir.path()));
175        assert!(!cache.is_full_delivered("/tmp/a.rs"));
176
177        cache.mark_full_delivered("/tmp/a.rs");
178        assert!(!sync_if_compacted(&mut cache, dir.path()));
179        assert!(cache.is_full_delivered("/tmp/a.rs"));
180    }
181}