Skip to main content

lean_ctx/hook_handlers/
observe.rs

1//! Observe hook handler: records all IDE hook events for context awareness
2//! (event parsing, token estimation, model/transcript detection, radar log).
3//! Split out of `hook_handlers/mod.rs`; `use super::*` re-imports parent items.
4
5#[allow(clippy::wildcard_imports)]
6use super::*;
7
8// ---------------------------------------------------------------------------
9// Observe handler — records ALL hook events for context awareness
10// ---------------------------------------------------------------------------
11
12/// Unified observe handler for all IDE hook events.
13/// Reads JSON from stdin, normalizes to `ObserveEvent`, counts tokens,
14/// appends to `context_radar.jsonl`, and exits immediately.
15pub fn handle_observe() {
16    if is_disabled() {
17        return;
18    }
19    let Some(input) = read_stdin_with_timeout(HOOK_STDIN_TIMEOUT) else {
20        return;
21    };
22    // Dedicated rules-injection mode (#343): a Claude/Codex `SessionStart` hook
23    // injects the compact lean-ctx summary as `additionalContext` — the
24    // non-polluting stand-in for the (skipped) CLAUDE.md/AGENTS.md block. Both
25    // agents register `hook observe` on SessionStart, so this is the single
26    // emit point (the Codex-specific handler stays silent in dedicated mode).
27    emit_dedicated_session_context(&input);
28    let Some(event) = parse_observe_event(&input) else {
29        return;
30    };
31    append_radar_event(&event);
32
33    // Output-echo analysis (#501): measure how much of the agent's reply
34    // re-quotes content lean-ctx already delivered, and feed the adaptive
35    // mode policy with an automatic feedback event.
36    if event.event_type == "agent_response" {
37        if let Some(text) = event.content.as_deref() {
38            crate::core::output_echo::analyze_and_record(text);
39        }
40    }
41}
42
43fn emit_dedicated_session_context(input: &str) {
44    let Ok(v) = serde_json::from_str::<serde_json::Value>(input) else {
45        return;
46    };
47    if v.get("hook_event_name").and_then(|e| e.as_str()) != Some("SessionStart") {
48        return;
49    }
50    if !crate::core::config::Config::load().dedicated_session_context_active() {
51        return;
52    }
53    emit_session_start_additional_context(crate::rules_inject::dedicated_session_summary());
54}
55
56#[derive(serde::Serialize)]
57struct ObserveEvent {
58    ts: u64,
59    event_type: &'static str,
60    tokens: usize,
61    #[serde(skip_serializing_if = "Option::is_none")]
62    tool_name: Option<String>,
63    #[serde(skip_serializing_if = "Option::is_none")]
64    detail: Option<String>,
65    #[serde(skip_serializing_if = "Option::is_none")]
66    content: Option<String>,
67    #[serde(skip_serializing_if = "Option::is_none")]
68    model: Option<String>,
69    #[serde(skip_serializing_if = "Option::is_none")]
70    conversation_id: Option<String>,
71}
72
73const MAX_CONTENT_CHARS: usize = 50_000;
74
75fn parse_observe_event(input: &str) -> Option<ObserveEvent> {
76    let v: serde_json::Value = serde_json::from_str(input).ok()?;
77
78    let ts = std::time::SystemTime::now()
79        .duration_since(std::time::UNIX_EPOCH)
80        .unwrap_or_default()
81        .as_secs();
82
83    let model = v
84        .get("model")
85        .and_then(|m| m.as_str())
86        .filter(|m| !m.is_empty())
87        .map(String::from);
88    let conversation_id = v
89        .get("conversation_id")
90        .and_then(|c| c.as_str())
91        .filter(|c| !c.is_empty())
92        .map(String::from);
93
94    let transcript_path = v
95        .get("transcript_path")
96        .and_then(|t| t.as_str())
97        .filter(|t| !t.is_empty())
98        .map(String::from);
99
100    if let Some(ref m) = model {
101        persist_detected_model(m);
102    }
103    if let Some(ref tp) = transcript_path {
104        persist_transcript_path(tp, conversation_id.as_deref());
105    }
106
107    let mut event = detect_event_type(&v, ts)?;
108    event.model = model;
109    event.conversation_id = conversation_id;
110    Some(event)
111}
112
113fn detect_event_type(v: &serde_json::Value, ts: u64) -> Option<ObserveEvent> {
114    if let Some(result) = v
115        .get("result_json")
116        .or_else(|| v.get("result"))
117        .or_else(|| v.get("tool_response"))
118        .or_else(|| v.get("tool_output"))
119    {
120        let tool = v
121            .get("tool_name")
122            .and_then(|t| t.as_str())
123            .unwrap_or("unknown");
124        let tokens = estimate_tokens_json(result);
125        let content_str = match result {
126            serde_json::Value::String(s) => s.clone(),
127            other => other.to_string(),
128        };
129        return Some(ObserveEvent {
130            ts,
131            event_type: "mcp_call",
132            tokens,
133            tool_name: Some(tool.to_string()),
134            detail: v
135                .get("server_name")
136                .and_then(|s| s.as_str())
137                .map(String::from),
138            content: Some(cap_content(&content_str)),
139            model: None,
140            conversation_id: None,
141        });
142    }
143
144    if let Some(output) = v.get("output") {
145        let cmd = v
146            .get("command")
147            .and_then(|c| c.as_str())
148            .unwrap_or("")
149            .to_string();
150        let tokens = estimate_tokens_value(output);
151        let out_str = match output {
152            serde_json::Value::String(s) => s.clone(),
153            other => other.to_string(),
154        };
155        return Some(ObserveEvent {
156            ts,
157            event_type: "shell",
158            tokens,
159            tool_name: None,
160            detail: Some(truncate_str(&cmd, 80)),
161            content: Some(cap_content(&format!("$ {cmd}\n{out_str}"))),
162            model: None,
163            conversation_id: None,
164        });
165    }
166
167    if v.get("content").is_some() && v.get("file_path").is_some() {
168        let path = v
169            .get("file_path")
170            .and_then(|p| p.as_str())
171            .unwrap_or("")
172            .to_string();
173        let file_content = v.get("content").and_then(|c| c.as_str()).unwrap_or("");
174        let tokens = file_content.len() / 4;
175        return Some(ObserveEvent {
176            ts,
177            event_type: "file_read",
178            tokens,
179            tool_name: None,
180            detail: Some(truncate_str(&path, 120)),
181            content: Some(cap_content(file_content)),
182            model: None,
183            conversation_id: None,
184        });
185    }
186
187    if let Some(text) = v.get("text").and_then(|t| t.as_str()) {
188        let has_duration = v.get("duration_ms").is_some();
189        let event_type = if has_duration {
190            "thinking"
191        } else {
192            "agent_response"
193        };
194        let tokens = text.len() / 4;
195        return Some(ObserveEvent {
196            ts,
197            event_type,
198            tokens,
199            tool_name: None,
200            detail: None,
201            content: Some(cap_content(text)),
202            model: None,
203            conversation_id: None,
204        });
205    }
206
207    if let Some(prompt) = v.get("prompt").and_then(|p| p.as_str()) {
208        let tokens = prompt.len() / 4;
209        let mut full = prompt.to_string();
210        if let Some(attachments) = v.get("attachments").and_then(|a| a.as_array()) {
211            if !attachments.is_empty() {
212                full.push_str(&format!("\n\n[{} attachments]", attachments.len()));
213                for att in attachments {
214                    if let Some(name) = att.get("name").and_then(|n| n.as_str()) {
215                        full.push_str(&format!("\n  - {name}"));
216                    }
217                }
218            }
219        }
220        return Some(ObserveEvent {
221            ts,
222            event_type: "user_message",
223            tokens,
224            tool_name: None,
225            detail: v
226                .get("attachments")
227                .and_then(|a| a.as_array())
228                .map(|a| format!("{} attachments", a.len())),
229            content: Some(cap_content(&full)),
230            model: None,
231            conversation_id: None,
232        });
233    }
234
235    if v.get("tool_name").is_some() || v.get("tool_input").is_some() {
236        let tool = v
237            .get("tool_name")
238            .and_then(|t| t.as_str())
239            .unwrap_or("unknown")
240            .to_string();
241        let is_lctx = tool.starts_with("ctx_") || tool.starts_with("mcp__lean-ctx__");
242        let tokens = v.get("tool_input").map_or(0, estimate_tokens_json);
243        let input_str = v
244            .get("tool_input")
245            .map(std::string::ToString::to_string)
246            .unwrap_or_default();
247        return Some(ObserveEvent {
248            ts,
249            event_type: if is_lctx { "mcp_call" } else { "native_tool" },
250            tokens,
251            tool_name: Some(tool),
252            detail: None,
253            content: if input_str.is_empty() {
254                None
255            } else {
256                Some(cap_content(&input_str))
257            },
258            model: None,
259            conversation_id: None,
260        });
261    }
262
263    // Claude Code emits `hook_event_name: "PreCompact"` (code.claude.com/docs/
264    // en/hooks); the generic `event`/`compaction` shapes cover other hosts.
265    // This check must run BEFORE the `session_id` catch-all below: every
266    // Claude hook payload carries `session_id` as a common field, so the
267    // compaction branch was unreachable for Claude — compactions were never
268    // recorded, `sync_if_compacted` never reset delivery flags, and
269    // post-compaction re-reads kept answering with "[unchanged]" stubs that
270    // pointed at context the host had already evicted (GL #555). Agents then
271    // fell back to native Read to recover the content.
272    let is_compaction = v.get("compaction").is_some()
273        || v.get("messages_count").is_some()
274        || v.get("hook_event_name")
275            .and_then(|e| e.as_str())
276            .is_some_and(|e| e == "PreCompact")
277        || v.get("event")
278            .and_then(|e| e.as_str())
279            .is_some_and(|e| e == "compaction" || e == "compact");
280    if !is_compaction && v.get("session_id").is_some() {
281        return Some(ObserveEvent {
282            ts,
283            event_type: "session",
284            tokens: 0,
285            tool_name: None,
286            detail: v
287                .get("session_id")
288                .and_then(|s| s.as_str())
289                .map(String::from),
290            content: None,
291            model: None,
292            conversation_id: None,
293        });
294    }
295
296    if is_compaction {
297        return Some(ObserveEvent {
298            ts,
299            event_type: "compaction",
300            tokens: 0,
301            tool_name: None,
302            detail: None,
303            content: None,
304            model: None,
305            conversation_id: None,
306        });
307    }
308
309    None
310}
311
312fn estimate_tokens_json(v: &serde_json::Value) -> usize {
313    match v {
314        serde_json::Value::String(s) => s.len() / 4,
315        _ => v.to_string().len() / 4,
316    }
317}
318
319fn estimate_tokens_value(v: &serde_json::Value) -> usize {
320    match v {
321        serde_json::Value::String(s) => s.len() / 4,
322        _ => v.to_string().len() / 4,
323    }
324}
325
326fn persist_detected_model(model: &str) {
327    let m = model.to_lowercase();
328    let is_bg_model = m.contains("flash")
329        || m.contains("mini")
330        || m.contains("haiku")
331        || m.contains("fast")
332        || m.contains("nano")
333        || m.contains("small");
334    if is_bg_model {
335        return;
336    }
337
338    let Ok(data_dir) = crate::core::data_dir::lean_ctx_data_dir() else {
339        return;
340    };
341    let path = data_dir.join("detected_model.json");
342    let ts = std::time::SystemTime::now()
343        .duration_since(std::time::UNIX_EPOCH)
344        .unwrap_or_default()
345        .as_secs();
346    let window = model_context_window(model);
347    let payload = serde_json::json!({
348        "model": model,
349        "window_size": window,
350        "detected_at": ts,
351    });
352    if let Ok(json) = serde_json::to_string_pretty(&payload) {
353        let tmp = path.with_extension("tmp");
354        if std::fs::write(&tmp, &json).is_ok() {
355            let _ = std::fs::rename(&tmp, &path);
356        }
357    }
358}
359
360pub fn model_context_window(model: &str) -> usize {
361    crate::core::model_registry::context_window_for_model(model)
362}
363
364pub fn load_detected_model() -> Option<(String, usize)> {
365    let data_dir = crate::core::data_dir::lean_ctx_data_dir().ok()?;
366    let path = data_dir.join("detected_model.json");
367    let content = std::fs::read_to_string(&path).ok()?;
368    let v: serde_json::Value = serde_json::from_str(&content).ok()?;
369    let model = v.get("model")?.as_str()?.to_string();
370    let window = v.get("window_size")?.as_u64()? as usize;
371    let detected_at = v.get("detected_at")?.as_u64()?;
372    let now = std::time::SystemTime::now()
373        .duration_since(std::time::UNIX_EPOCH)
374        .unwrap_or_default()
375        .as_secs();
376    if now.saturating_sub(detected_at) > 7200 {
377        return None;
378    }
379    Some((model, window))
380}
381
382fn persist_transcript_path(path: &str, conversation_id: Option<&str>) {
383    let Ok(data_dir) = crate::core::data_dir::lean_ctx_data_dir() else {
384        return;
385    };
386    let meta_path = data_dir.join("active_transcript.json");
387    let ts = std::time::SystemTime::now()
388        .duration_since(std::time::UNIX_EPOCH)
389        .unwrap_or_default()
390        .as_secs();
391    let payload = serde_json::json!({
392        "transcript_path": path,
393        "conversation_id": conversation_id,
394        "updated_at": ts,
395    });
396    if let Ok(json) = serde_json::to_string_pretty(&payload) {
397        let tmp = meta_path.with_extension("tmp");
398        if std::fs::write(&tmp, &json).is_ok() {
399            let _ = std::fs::rename(&tmp, &meta_path);
400        }
401    }
402}
403
404pub fn load_active_transcript() -> Option<(String, Option<String>)> {
405    let data_dir = crate::core::data_dir::lean_ctx_data_dir().ok()?;
406    let path = data_dir.join("active_transcript.json");
407    let content = std::fs::read_to_string(&path).ok()?;
408    let v: serde_json::Value = serde_json::from_str(&content).ok()?;
409    let tp = v.get("transcript_path")?.as_str()?.to_string();
410    let conv = v
411        .get("conversation_id")
412        .and_then(|c| c.as_str())
413        .map(String::from);
414    let updated = v.get("updated_at")?.as_u64()?;
415    let now = std::time::SystemTime::now()
416        .duration_since(std::time::UNIX_EPOCH)
417        .unwrap_or_default()
418        .as_secs();
419    if now.saturating_sub(updated) > 7200 {
420        return None;
421    }
422    Some((tp, conv))
423}
424
425fn cap_content(s: &str) -> String {
426    if s.len() <= MAX_CONTENT_CHARS {
427        s.to_string()
428    } else {
429        let truncated = safe_truncate(s, MAX_CONTENT_CHARS);
430        format!("{}…\n\n[truncated: {} total chars]", truncated, s.len())
431    }
432}
433
434fn truncate_str(s: &str, max: usize) -> String {
435    if s.len() <= max {
436        s.to_string()
437    } else {
438        format!("{}...", safe_truncate(s, max))
439    }
440}
441
442/// Truncate a string at a char boundary <= max bytes. Never panics on multi-byte UTF-8.
443fn safe_truncate(s: &str, max: usize) -> &str {
444    if max >= s.len() {
445        return s;
446    }
447    let mut end = max;
448    while end > 0 && !s.is_char_boundary(end) {
449        end -= 1;
450    }
451    &s[..end]
452}
453
454fn append_radar_event(event: &ObserveEvent) {
455    let Ok(data_dir) = crate::core::data_dir::lean_ctx_data_dir() else {
456        return;
457    };
458    let radar_path = data_dir.join("context_radar.jsonl");
459
460    if event.event_type == "session" {
461        if let Ok(meta) = std::fs::metadata(&radar_path) {
462            const MAX_RADAR_SIZE: u64 = 10 * 1024 * 1024; // 10 MB
463            if meta.len() > MAX_RADAR_SIZE {
464                let prev = data_dir.join("context_radar.prev.jsonl");
465                let _ = std::fs::rename(&radar_path, &prev);
466            }
467        }
468    }
469
470    let Ok(line) = serde_json::to_string(event) else {
471        return;
472    };
473
474    use std::fs::OpenOptions;
475    use std::io::Write;
476    if let Ok(mut f) = OpenOptions::new()
477        .create(true)
478        .append(true)
479        .open(&radar_path)
480    {
481        let _ = writeln!(f, "{line}");
482    }
483}
484
485#[cfg(test)]
486mod tests {
487    use super::*;
488
489    #[test]
490    fn detect_event_type_tool_response_is_mcp_call() {
491        let v = serde_json::json!({
492            "tool_name": "ctx_read",
493            "tool_response": "file contents here"
494        });
495        let event = detect_event_type(&v, 1000).unwrap();
496        assert_eq!(event.event_type, "mcp_call");
497    }
498
499    #[test]
500    fn detect_event_type_tool_output_is_mcp_call() {
501        let v = serde_json::json!({
502            "tool_name": "ctx_search",
503            "tool_output": "search results"
504        });
505        let event = detect_event_type(&v, 1000).unwrap();
506        assert_eq!(event.event_type, "mcp_call");
507    }
508
509    #[test]
510    fn detect_event_type_ctx_prefix_is_mcp_call() {
511        let v = serde_json::json!({
512            "tool_name": "ctx_read",
513            "tool_input": {"path": "src/main.rs"}
514        });
515        let event = detect_event_type(&v, 1000).unwrap();
516        assert_eq!(event.event_type, "mcp_call");
517    }
518
519    #[test]
520    fn detect_event_type_mcp_prefix_is_mcp_call() {
521        let v = serde_json::json!({
522            "tool_name": "mcp__lean-ctx__ctx_read",
523            "tool_input": {"path": "src/main.rs"}
524        });
525        let event = detect_event_type(&v, 1000).unwrap();
526        assert_eq!(event.event_type, "mcp_call");
527    }
528
529    #[test]
530    fn detect_event_type_native_read_is_native_tool() {
531        let v = serde_json::json!({
532            "tool_name": "Read",
533            "tool_input": {"path": "src/main.rs"}
534        });
535        let event = detect_event_type(&v, 1000).unwrap();
536        assert_eq!(event.event_type, "native_tool");
537    }
538
539    #[test]
540    fn detect_event_type_result_json_is_mcp_call() {
541        let v = serde_json::json!({
542            "tool_name": "ctx_read",
543            "result_json": {"content": "..."}
544        });
545        let event = detect_event_type(&v, 1000).unwrap();
546        assert_eq!(event.event_type, "mcp_call");
547    }
548
549    /// Real Claude Code PreCompact payload (code.claude.com/docs/en/hooks):
550    /// carries `session_id` like every Claude hook, so the compaction check
551    /// must win over the generic session catch-all (GL #555).
552    #[test]
553    fn detect_event_type_claude_precompact_is_compaction() {
554        let v = serde_json::json!({
555            "session_id": "abc123",
556            "transcript_path": "/Users/u/.claude/projects/x/abc123.jsonl",
557            "cwd": "/Users/u/project",
558            "hook_event_name": "PreCompact",
559            "trigger": "auto",
560            "custom_instructions": ""
561        });
562        let event = detect_event_type(&v, 1000).unwrap();
563        assert_eq!(event.event_type, "compaction");
564    }
565
566    #[test]
567    fn detect_event_type_plain_session_event_still_session() {
568        let v = serde_json::json!({
569            "session_id": "abc123",
570            "hook_event_name": "SessionStart"
571        });
572        let event = detect_event_type(&v, 1000).unwrap();
573        assert_eq!(event.event_type, "session");
574    }
575}