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 {
38 let last_seen = LAST_COMPACTION_TS.load(Ordering::Relaxed);
39
40 let latest_ts = check_compaction_marker(data_dir, last_seen)
41 .or_else(|| scan_radar_tail(data_dir, last_seen));
42
43 let Some(latest_compaction_ts) = latest_ts else {
44 return false;
45 };
46
47 LAST_COMPACTION_TS.store(latest_compaction_ts, Ordering::Relaxed);
48 crate::core::search_delta::reset();
49 let reset_count = cache.reset_delivery_flags();
50 crate::core::cache_telemetry::record_compaction(reset_count as u64);
51 crate::core::read_stub_index::reset_in_dir(data_dir);
56 if reset_count > 0 {
57 eprintln!(
58 "[lean-ctx] compaction detected — reset {reset_count} delivery flags for re-read"
59 );
60 }
61
62 std::thread::spawn(|| {
63 if let Some(session) = crate::core::session::SessionState::load_latest()
64 && let Some(ref root) = session.project_root
65 && (!session.findings.is_empty() || !session.decisions.is_empty())
66 {
67 crate::tools::startup::auto_consolidate_knowledge(root);
68 }
69 });
70
71 true
72}
73
74fn check_compaction_marker(data_dir: &Path, since_ts: u64) -> Option<u64> {
82 let marker_path = data_dir.join("last_compaction.json");
83 let content = std::fs::read_to_string(&marker_path).ok()?;
84 let parsed: serde_json::Value = serde_json::from_str(content.trim()).ok()?;
85 let ts = parsed.get("ts")?.as_u64()?;
86 if ts > since_ts { Some(ts) } else { None }
87}
88
89fn scan_radar_tail(data_dir: &Path, since_ts: u64) -> Option<u64> {
97 use std::io::{Read, Seek, SeekFrom};
98
99 let radar_path = data_dir.join("context_radar.jsonl");
100 let mut file = std::fs::File::open(&radar_path).ok()?;
101 let file_len = file.metadata().ok()?.len();
102
103 const TAIL_BYTES: u64 = 256 * 1024; let content = if file_len <= TAIL_BYTES {
105 let mut s = String::new();
106 file.read_to_string(&mut s).ok()?;
107 s
108 } else {
109 file.seek(SeekFrom::End(-(TAIL_BYTES as i64))).ok()?;
110 let mut buf = vec![0u8; TAIL_BYTES as usize];
111 let n = file.read(&mut buf).ok()?;
112 let s = String::from_utf8_lossy(&buf[..n]).into_owned();
113 if let Some(idx) = s.find('\n') {
114 s[idx + 1..].to_string()
115 } else {
116 s
117 }
118 };
119
120 for line in content.lines().rev() {
121 if line.is_empty() {
122 continue;
123 }
124 let event: RadarEvent = match serde_json::from_str(line) {
125 Ok(e) => e,
126 Err(_) => continue,
127 };
128 if event.ts <= since_ts {
129 break;
130 }
131 if event.event_type == "compaction" {
132 return Some(event.ts);
133 }
134 }
135 None
136}
137
138#[cfg(test)]
139mod tests {
140 use super::*;
141 use serial_test::serial;
142 use std::io::Write;
143 use tempfile::TempDir;
144
145 fn make_cache_with_delivered(paths: &[&str]) -> SessionCache {
146 let mut cache = SessionCache::default();
147 for p in paths {
148 cache.store(p, "hello world");
149 cache.mark_full_delivered(p);
150 }
151 cache
152 }
153
154 #[test]
155 #[serial]
156 fn no_reset_without_compaction_event() {
157 let dir = TempDir::new().unwrap();
158 let radar = dir.path().join("context_radar.jsonl");
159 let mut f = std::fs::File::create(&radar).unwrap();
160 writeln!(f, r#"{{"ts":1000,"event_type":"mcp_call","tokens":50}}"#).unwrap();
161 drop(f);
162
163 LAST_COMPACTION_TS.store(0, Ordering::Relaxed);
164 let mut cache = make_cache_with_delivered(&["/tmp/a.rs"]);
165 assert!(!sync_if_compacted(&mut cache, dir.path()));
166 assert!(cache.is_full_delivered("/tmp/a.rs"));
167 }
168
169 #[test]
170 #[serial]
171 fn resets_after_compaction() {
172 let dir = TempDir::new().unwrap();
173 let radar = dir.path().join("context_radar.jsonl");
174 let mut f = std::fs::File::create(&radar).unwrap();
175 writeln!(f, r#"{{"ts":1000,"event_type":"mcp_call","tokens":50}}"#).unwrap();
176 writeln!(f, r#"{{"ts":2000,"event_type":"compaction","tokens":0}}"#).unwrap();
177 drop(f);
178
179 LAST_COMPACTION_TS.store(0, Ordering::Relaxed);
180 let mut cache = make_cache_with_delivered(&["/tmp/a.rs", "/tmp/b.rs"]);
181
182 assert!(cache.is_full_delivered("/tmp/a.rs"));
183 assert!(sync_if_compacted(&mut cache, dir.path()));
184 assert!(!cache.is_full_delivered("/tmp/a.rs"));
185 assert!(!cache.is_full_delivered("/tmp/b.rs"));
186 }
187
188 #[test]
189 #[serial]
190 fn does_not_double_reset() {
191 let dir = TempDir::new().unwrap();
192 let radar = dir.path().join("context_radar.jsonl");
193 let mut f = std::fs::File::create(&radar).unwrap();
194 writeln!(f, r#"{{"ts":2000,"event_type":"compaction","tokens":0}}"#).unwrap();
195 drop(f);
196
197 LAST_COMPACTION_TS.store(0, Ordering::Relaxed);
198 let mut cache = make_cache_with_delivered(&["/tmp/a.rs"]);
199 assert!(sync_if_compacted(&mut cache, dir.path()));
200 assert!(!cache.is_full_delivered("/tmp/a.rs"));
201
202 cache.mark_full_delivered("/tmp/a.rs");
203 assert!(!sync_if_compacted(&mut cache, dir.path()));
204 assert!(cache.is_full_delivered("/tmp/a.rs"));
205 }
206
207 #[test]
212 #[serial]
213 fn marker_alone_triggers_reset() {
214 let dir = TempDir::new().unwrap();
215 let marker = dir.path().join("last_compaction.json");
217 std::fs::write(&marker, r#"{"ts":5000}"#).unwrap();
218 let radar = dir.path().join("context_radar.jsonl");
220 std::fs::write(&radar, "").unwrap();
221
222 LAST_COMPACTION_TS.store(0, Ordering::Relaxed);
223 let mut cache = make_cache_with_delivered(&["/tmp/x.rs"]);
224 assert!(cache.is_full_delivered("/tmp/x.rs"));
225
226 assert!(sync_if_compacted(&mut cache, dir.path()));
228 assert!(!cache.is_full_delivered("/tmp/x.rs"));
229 assert_eq!(LAST_COMPACTION_TS.load(Ordering::Relaxed), 5000);
230
231 cache.mark_full_delivered("/tmp/x.rs");
233 assert!(!sync_if_compacted(&mut cache, dir.path()));
234 assert!(cache.is_full_delivered("/tmp/x.rs"));
235 }
236
237 #[test]
241 #[serial]
242 fn compaction_buried_under_large_event_line_is_still_detected() {
243 let dir = TempDir::new().unwrap();
244 let radar = dir.path().join("context_radar.jsonl");
245 let mut f = std::fs::File::create(&radar).unwrap();
246
247 writeln!(f, r#"{{"ts":3000,"event_type":"compaction","tokens":0}}"#).unwrap();
249
250 let large_content = "x".repeat(10_000);
252 writeln!(
253 f,
254 r#"{{"ts":3001,"event_type":"agent_response","tokens":2500,"content":"{large_content}"}}"#
255 )
256 .unwrap();
257 drop(f);
258
259 LAST_COMPACTION_TS.store(0, Ordering::Relaxed);
260 let mut cache = make_cache_with_delivered(&["/tmp/buried.rs"]);
261 assert!(cache.is_full_delivered("/tmp/buried.rs"));
262
263 assert!(sync_if_compacted(&mut cache, dir.path()));
265 assert!(!cache.is_full_delivered("/tmp/buried.rs"));
266 }
267
268 #[test]
271 #[serial]
272 fn marker_takes_priority_over_empty_radar() {
273 let dir = TempDir::new().unwrap();
274 let radar = dir.path().join("context_radar.jsonl");
275 let mut f = std::fs::File::create(&radar).unwrap();
276 writeln!(f, r#"{{"ts":1000,"event_type":"mcp_call","tokens":50}}"#).unwrap();
277 drop(f);
278
279 let marker = dir.path().join("last_compaction.json");
280 std::fs::write(&marker, r#"{"ts":4000}"#).unwrap();
281
282 LAST_COMPACTION_TS.store(0, Ordering::Relaxed);
283 let mut cache = make_cache_with_delivered(&["/tmp/priority.rs"]);
284 assert!(sync_if_compacted(&mut cache, dir.path()));
285 assert!(!cache.is_full_delivered("/tmp/priority.rs"));
286 assert_eq!(LAST_COMPACTION_TS.load(Ordering::Relaxed), 4000);
287 }
288
289 #[test]
291 #[serial]
292 fn corrupt_marker_falls_back_to_radar() {
293 let dir = TempDir::new().unwrap();
294 let radar = dir.path().join("context_radar.jsonl");
295 let mut f = std::fs::File::create(&radar).unwrap();
296 writeln!(f, r#"{{"ts":6000,"event_type":"compaction","tokens":0}}"#).unwrap();
297 drop(f);
298
299 let marker = dir.path().join("last_compaction.json");
300 std::fs::write(&marker, "not valid json").unwrap();
301
302 LAST_COMPACTION_TS.store(0, Ordering::Relaxed);
303 let mut cache = make_cache_with_delivered(&["/tmp/fallback.rs"]);
304 assert!(sync_if_compacted(&mut cache, dir.path()));
306 assert!(!cache.is_full_delivered("/tmp/fallback.rs"));
307 }
308}