lean_ctx/server/
compaction_sync.rs1use 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
9pub 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
27pub 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 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
69fn 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 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}