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.
30///
31/// Detection strategy (#808):
32/// 1. Check `last_compaction.json` marker file first (written atomically by
33///    the observe hook). This is only a few bytes, immune to being displaced
34///    by large events.
35/// 2. Fall back to scanning the tail of `context_radar.jsonl` (256 KB window)
36///    for older hook binaries that don't write the marker.
37pub 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    // Drop the persistent stub index too (#955): the conversation's context was
52    // summarised away, so neither a warm nor a cold stub may claim "you already
53    // have this". Writes the emptied index synchronously so a restart in the
54    // crash window can't resurrect a pre-compaction stub.
55    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
74// ---------------------------------------------------------------------------
75// Detection: marker file (primary, #808)
76// ---------------------------------------------------------------------------
77
78/// Read `last_compaction.json` — a few-byte file written atomically by the
79/// observe hook. Returns `Some(ts)` if the marker exists and its timestamp
80/// is newer than `since_ts`.
81fn 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
89// ---------------------------------------------------------------------------
90// Detection: radar tail scan (fallback for older hook binaries)
91// ---------------------------------------------------------------------------
92
93/// Scan only the tail of radar JSONL for a compaction event newer than `since_ts`.
94/// #808: window widened from 4 KB to 256 KB so large `agent_response` /
95/// `thinking` events (up to 50 000 chars each) cannot bury the compaction entry.
96fn 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; // #808: was 4096
104    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    // --- #808 new tests ---
208
209    /// The marker file alone (without a radar entry) triggers a reset,
210    /// and the same marker timestamp does not trigger it a second time.
211    #[test]
212    #[serial]
213    fn marker_alone_triggers_reset() {
214        let dir = TempDir::new().unwrap();
215        // No context_radar.jsonl — only the marker
216        let marker = dir.path().join("last_compaction.json");
217        std::fs::write(&marker, r#"{"ts":5000}"#).unwrap();
218        // Need a radar file to exist (sync_if_compacted checks marker OR radar)
219        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        // First call: marker is newer than last_seen (0) → reset
227        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        // Re-deliver, then call again: same marker ts → no reset
232        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    /// A compaction event buried under a large (10 KB+) event line is still
238    /// detected via the widened 256 KB fallback window. This reproduces the
239    /// original bug — the old 4 KB window would miss it.
240    #[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        // Write the compaction event
248        writeln!(f, r#"{{"ts":3000,"event_type":"compaction","tokens":0}}"#).unwrap();
249
250        // Bury it under a large agent_response event (10 KB content)
251        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        // Must find the compaction despite the large event after it
264        assert!(sync_if_compacted(&mut cache, dir.path()));
265        assert!(!cache.is_full_delivered("/tmp/buried.rs"));
266    }
267
268    /// Marker file takes priority over the radar scan: even if the radar
269    /// has no compaction event, the marker triggers a reset.
270    #[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    /// A corrupt or empty marker file is ignored — falls back to radar scan.
290    #[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        // Corrupt marker ignored → radar scan finds the compaction
305        assert!(sync_if_compacted(&mut cache, dir.path()));
306        assert!(!cache.is_full_delivered("/tmp/fallback.rs"));
307    }
308}