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/CodeBuddy `SessionStart` hook
23    // injects the compact lean-ctx summary as `additionalContext` — the
24    // non-polluting stand-in for the (skipped) CLAUDE.md/CODEBUDDY.md/AGENTS.md block. All
25    // three 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        && let Some(text) = event.content.as_deref()
38    {
39        crate::core::output_echo::analyze_and_record(text);
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    // Session start additional context removed — the MCP instructions
52    // already carry the compact rules block.
53}
54
55#[derive(serde::Serialize)]
56struct ObserveEvent {
57    ts: u64,
58    event_type: &'static str,
59    tokens: usize,
60    #[serde(skip_serializing_if = "Option::is_none")]
61    tool_name: Option<String>,
62    #[serde(skip_serializing_if = "Option::is_none")]
63    detail: Option<String>,
64    #[serde(skip_serializing_if = "Option::is_none")]
65    content: Option<String>,
66    #[serde(skip_serializing_if = "Option::is_none")]
67    model: Option<String>,
68    #[serde(skip_serializing_if = "Option::is_none")]
69    conversation_id: Option<String>,
70}
71
72const MAX_CONTENT_CHARS: usize = 50_000;
73
74fn parse_observe_event(input: &str) -> Option<ObserveEvent> {
75    let v: serde_json::Value = serde_json::from_str(input).ok()?;
76
77    let ts = std::time::SystemTime::now()
78        .duration_since(std::time::UNIX_EPOCH)
79        .unwrap_or_default()
80        .as_secs();
81
82    let model = v
83        .get("model")
84        .and_then(|m| m.as_str())
85        .filter(|m| !m.is_empty())
86        .map(String::from);
87    let conversation_id = v
88        .get("conversation_id")
89        .and_then(|c| c.as_str())
90        .filter(|c| !c.is_empty())
91        .map(String::from);
92
93    let transcript_path = v
94        .get("transcript_path")
95        .and_then(|t| t.as_str())
96        .filter(|t| !t.is_empty())
97        .map(String::from);
98
99    if let Some(ref m) = model {
100        persist_detected_model(m);
101    }
102    if let Some(ref tp) = transcript_path {
103        persist_transcript_path(tp, conversation_id.as_deref());
104    }
105
106    let mut event = detect_event_type(&v, ts)?;
107    event.model = model;
108    event.conversation_id = conversation_id;
109    Some(event)
110}
111
112fn detect_event_type(v: &serde_json::Value, ts: u64) -> Option<ObserveEvent> {
113    // GitHub Copilot CLI postToolUse: camelCase `toolName` + `toolArgs`
114    // (JSON-encoded string) + `toolResult`. None of the snake_case branches
115    // below match this shape, so without a dedicated arm Copilot telemetry
116    // (heatmap, token savings, radar) is silently dropped (#551).
117    if let Some(result) = v.get("toolResult") {
118        let tool = super::payload::resolve_tool_name(v).unwrap_or_else(|| "unknown".to_string());
119        let args = super::payload::resolve_tool_args(v);
120        let command = args
121            .as_ref()
122            .and_then(|a| a.get("command"))
123            .and_then(|c| c.as_str());
124        let result_text = result
125            .get("textResultForLlm")
126            .and_then(|t| t.as_str())
127            .map_or_else(|| result.to_string(), String::from);
128        let tokens = result_text.len() / 4;
129        let is_lctx = tool.starts_with("ctx_") || tool.starts_with("mcp__lean-ctx__");
130        let event_type = if is_lctx {
131            "mcp_call"
132        } else if command.is_some() {
133            "shell"
134        } else {
135            "native_tool"
136        };
137        let content = match command {
138            Some(cmd) => format!("$ {cmd}\n{result_text}"),
139            None => result_text,
140        };
141        return Some(ObserveEvent {
142            ts,
143            event_type,
144            tokens,
145            tool_name: Some(tool),
146            detail: command.map(|c| truncate_str(c, 80)),
147            content: Some(cap_content(&content)),
148            model: None,
149            conversation_id: None,
150        });
151    }
152
153    if let Some(result) = v
154        .get("result_json")
155        .or_else(|| v.get("result"))
156        .or_else(|| v.get("tool_response"))
157        .or_else(|| v.get("tool_output"))
158    {
159        let tool = v
160            .get("tool_name")
161            .and_then(|t| t.as_str())
162            .unwrap_or("unknown");
163        let tokens = estimate_tokens_json(result);
164        let content_str = match result {
165            serde_json::Value::String(s) => s.clone(),
166            other => other.to_string(),
167        };
168        return Some(ObserveEvent {
169            ts,
170            event_type: "mcp_call",
171            tokens,
172            tool_name: Some(tool.to_string()),
173            detail: v
174                .get("server_name")
175                .and_then(|s| s.as_str())
176                .map(String::from),
177            content: Some(cap_content(&content_str)),
178            model: None,
179            conversation_id: None,
180        });
181    }
182
183    if let Some(output) = v.get("output") {
184        let cmd = v
185            .get("command")
186            .and_then(|c| c.as_str())
187            .unwrap_or("")
188            .to_string();
189        let tokens = estimate_tokens_value(output);
190        let out_str = match output {
191            serde_json::Value::String(s) => s.clone(),
192            other => other.to_string(),
193        };
194        return Some(ObserveEvent {
195            ts,
196            event_type: "shell",
197            tokens,
198            tool_name: None,
199            detail: Some(truncate_str(&cmd, 80)),
200            content: Some(cap_content(&format!("$ {cmd}\n{out_str}"))),
201            model: None,
202            conversation_id: None,
203        });
204    }
205
206    if v.get("content").is_some() && v.get("file_path").is_some() {
207        let path = v
208            .get("file_path")
209            .and_then(|p| p.as_str())
210            .unwrap_or("")
211            .to_string();
212        let file_content = v.get("content").and_then(|c| c.as_str()).unwrap_or("");
213        let tokens = file_content.len() / 4;
214        return Some(ObserveEvent {
215            ts,
216            event_type: "file_read",
217            tokens,
218            tool_name: None,
219            detail: Some(truncate_str(&path, 120)),
220            content: Some(cap_content(file_content)),
221            model: None,
222            conversation_id: None,
223        });
224    }
225
226    if let Some(text) = v.get("text").and_then(|t| t.as_str()) {
227        let has_duration = v.get("duration_ms").is_some();
228        let event_type = if has_duration {
229            "thinking"
230        } else {
231            "agent_response"
232        };
233        let tokens = text.len() / 4;
234        return Some(ObserveEvent {
235            ts,
236            event_type,
237            tokens,
238            tool_name: None,
239            detail: None,
240            content: Some(cap_content(text)),
241            model: None,
242            conversation_id: None,
243        });
244    }
245
246    if let Some(prompt) = v.get("prompt").and_then(|p| p.as_str()) {
247        let tokens = prompt.len() / 4;
248        let mut full = prompt.to_string();
249        if let Some(attachments) = v.get("attachments").and_then(|a| a.as_array())
250            && !attachments.is_empty()
251        {
252            full.push_str(&format!("\n\n[{} attachments]", attachments.len()));
253            for att in attachments {
254                if let Some(name) = att.get("name").and_then(|n| n.as_str()) {
255                    full.push_str(&format!("\n  - {name}"));
256                }
257            }
258        }
259        return Some(ObserveEvent {
260            ts,
261            event_type: "user_message",
262            tokens,
263            tool_name: None,
264            detail: v
265                .get("attachments")
266                .and_then(|a| a.as_array())
267                .map(|a| format!("{} attachments", a.len())),
268            content: Some(cap_content(&full)),
269            model: None,
270            conversation_id: None,
271        });
272    }
273
274    if v.get("tool_name").is_some() || v.get("tool_input").is_some() {
275        let tool = v
276            .get("tool_name")
277            .and_then(|t| t.as_str())
278            .unwrap_or("unknown")
279            .to_string();
280        let is_lctx = tool.starts_with("ctx_") || tool.starts_with("mcp__lean-ctx__");
281        let tokens = v.get("tool_input").map_or(0, estimate_tokens_json);
282        let input_str = v
283            .get("tool_input")
284            .map(std::string::ToString::to_string)
285            .unwrap_or_default();
286        return Some(ObserveEvent {
287            ts,
288            event_type: if is_lctx { "mcp_call" } else { "native_tool" },
289            tokens,
290            tool_name: Some(tool),
291            detail: None,
292            content: if input_str.is_empty() {
293                None
294            } else {
295                Some(cap_content(&input_str))
296            },
297            model: None,
298            conversation_id: None,
299        });
300    }
301
302    // Claude Code emits `hook_event_name: "PreCompact"` (code.claude.com/docs/
303    // en/hooks); the generic `event`/`compaction` shapes cover other hosts.
304    // This check must run BEFORE the `session_id` catch-all below: every
305    // Claude hook payload carries `session_id` as a common field, so the
306    // compaction branch was unreachable for Claude — compactions were never
307    // recorded, `sync_if_compacted` never reset delivery flags, and
308    // post-compaction re-reads kept answering with "[unchanged]" stubs that
309    // pointed at context the host had already evicted (GL #555). Agents then
310    // fell back to native Read to recover the content.
311    let is_compaction = v.get("compaction").is_some()
312        || v.get("messages_count").is_some()
313        || v.get("hook_event_name")
314            .and_then(|e| e.as_str())
315            .is_some_and(|e| e == "PreCompact")
316        || v.get("event")
317            .and_then(|e| e.as_str())
318            .is_some_and(|e| e == "compaction" || e == "compact");
319    if !is_compaction && v.get("session_id").is_some() {
320        return Some(ObserveEvent {
321            ts,
322            event_type: "session",
323            tokens: 0,
324            tool_name: None,
325            detail: v
326                .get("session_id")
327                .and_then(|s| s.as_str())
328                .map(String::from),
329            content: None,
330            model: None,
331            conversation_id: None,
332        });
333    }
334
335    if is_compaction {
336        return Some(ObserveEvent {
337            ts,
338            event_type: "compaction",
339            tokens: 0,
340            tool_name: None,
341            detail: None,
342            content: None,
343            model: None,
344            conversation_id: None,
345        });
346    }
347
348    None
349}
350
351fn estimate_tokens_json(v: &serde_json::Value) -> usize {
352    match v {
353        serde_json::Value::String(s) => s.len() / 4,
354        _ => v.to_string().len() / 4,
355    }
356}
357
358fn estimate_tokens_value(v: &serde_json::Value) -> usize {
359    match v {
360        serde_json::Value::String(s) => s.len() / 4,
361        _ => v.to_string().len() / 4,
362    }
363}
364
365fn persist_detected_model(model: &str) {
366    let m = model.to_lowercase();
367    let is_bg_model = m.contains("flash")
368        || m.contains("mini")
369        || m.contains("haiku")
370        || m.contains("fast")
371        || m.contains("nano")
372        || m.contains("small");
373    if is_bg_model {
374        return;
375    }
376
377    let Ok(data_dir) = crate::core::data_dir::lean_ctx_data_dir() else {
378        return;
379    };
380    let path = data_dir.join("detected_model.json");
381    let ts = std::time::SystemTime::now()
382        .duration_since(std::time::UNIX_EPOCH)
383        .unwrap_or_default()
384        .as_secs();
385    let window = model_context_window(model);
386    let payload = serde_json::json!({
387        "model": model,
388        "window_size": window,
389        "detected_at": ts,
390    });
391    if let Ok(json) = serde_json::to_string_pretty(&payload) {
392        let tmp = path.with_extension("tmp");
393        if std::fs::write(&tmp, &json).is_ok() {
394            let _ = std::fs::rename(&tmp, &path);
395        }
396    }
397}
398
399pub fn model_context_window(model: &str) -> usize {
400    crate::core::model_registry::context_window_for_model(model)
401}
402
403pub fn load_detected_model() -> Option<(String, usize)> {
404    let data_dir = crate::core::data_dir::lean_ctx_data_dir().ok()?;
405    let path = data_dir.join("detected_model.json");
406    let content = std::fs::read_to_string(&path).ok()?;
407    let v: serde_json::Value = serde_json::from_str(&content).ok()?;
408    let model = v.get("model")?.as_str()?.to_string();
409    let window = v.get("window_size")?.as_u64()? as usize;
410    let detected_at = v.get("detected_at")?.as_u64()?;
411    let now = std::time::SystemTime::now()
412        .duration_since(std::time::UNIX_EPOCH)
413        .unwrap_or_default()
414        .as_secs();
415    if now.saturating_sub(detected_at) > 7200 {
416        return None;
417    }
418    Some((model, window))
419}
420
421fn persist_transcript_path(path: &str, conversation_id: Option<&str>) {
422    let Ok(data_dir) = crate::core::data_dir::lean_ctx_data_dir() else {
423        return;
424    };
425    let meta_path = data_dir.join("active_transcript.json");
426    let ts = std::time::SystemTime::now()
427        .duration_since(std::time::UNIX_EPOCH)
428        .unwrap_or_default()
429        .as_secs();
430    let payload = serde_json::json!({
431        "transcript_path": path,
432        "conversation_id": conversation_id,
433        "updated_at": ts,
434    });
435    if let Ok(json) = serde_json::to_string_pretty(&payload) {
436        let tmp = meta_path.with_extension("tmp");
437        if std::fs::write(&tmp, &json).is_ok() {
438            let _ = std::fs::rename(&tmp, &meta_path);
439        }
440    }
441}
442
443pub fn load_active_transcript() -> Option<(String, Option<String>)> {
444    let data_dir = crate::core::data_dir::lean_ctx_data_dir().ok()?;
445    let path = data_dir.join("active_transcript.json");
446    let content = std::fs::read_to_string(&path).ok()?;
447    let v: serde_json::Value = serde_json::from_str(&content).ok()?;
448    let tp = v.get("transcript_path")?.as_str()?.to_string();
449    let conv = v
450        .get("conversation_id")
451        .and_then(|c| c.as_str())
452        .map(String::from);
453    let updated = v.get("updated_at")?.as_u64()?;
454    let now = std::time::SystemTime::now()
455        .duration_since(std::time::UNIX_EPOCH)
456        .unwrap_or_default()
457        .as_secs();
458    if now.saturating_sub(updated) > 7200 {
459        return None;
460    }
461    Some((tp, conv))
462}
463
464fn cap_content(s: &str) -> String {
465    if s.len() <= MAX_CONTENT_CHARS {
466        s.to_string()
467    } else {
468        let truncated = safe_truncate(s, MAX_CONTENT_CHARS);
469        format!("{}…\n\n[truncated: {} total chars]", truncated, s.len())
470    }
471}
472
473fn truncate_str(s: &str, max: usize) -> String {
474    if s.len() <= max {
475        s.to_string()
476    } else {
477        format!("{}...", safe_truncate(s, max))
478    }
479}
480
481/// Truncate a string at a char boundary <= max bytes. Never panics on multi-byte UTF-8.
482fn safe_truncate(s: &str, max: usize) -> &str {
483    if max >= s.len() {
484        return s;
485    }
486    let mut end = max;
487    while end > 0 && !s.is_char_boundary(end) {
488        end -= 1;
489    }
490    &s[..end]
491}
492
493fn append_radar_event(event: &ObserveEvent) {
494    let Ok(data_dir) = crate::core::data_dir::lean_ctx_data_dir() else {
495        return;
496    };
497    let radar_path = data_dir.join("context_radar.jsonl");
498
499    if event.event_type == "session"
500        && let Ok(meta) = std::fs::metadata(&radar_path)
501    {
502        const MAX_RADAR_SIZE: u64 = 10 * 1024 * 1024; // 10 MB
503        if meta.len() > MAX_RADAR_SIZE {
504            let prev = data_dir.join("context_radar.prev.jsonl");
505            let _ = std::fs::rename(&radar_path, &prev);
506        }
507    }
508
509    let Ok(line) = serde_json::to_string(event) else {
510        return;
511    };
512
513    use std::fs::OpenOptions;
514    use std::io::Write;
515    if let Ok(mut f) = OpenOptions::new()
516        .create(true)
517        .append(true)
518        .open(&radar_path)
519    {
520        let _ = writeln!(f, "{line}");
521    }
522}
523
524#[cfg(test)]
525mod tests {
526    use super::*;
527
528    #[test]
529    fn detect_event_type_tool_response_is_mcp_call() {
530        let v = serde_json::json!({
531            "tool_name": "ctx_read",
532            "tool_response": "file contents here"
533        });
534        let event = detect_event_type(&v, 1000).unwrap();
535        assert_eq!(event.event_type, "mcp_call");
536    }
537
538    #[test]
539    fn detect_event_type_tool_output_is_mcp_call() {
540        let v = serde_json::json!({
541            "tool_name": "ctx_search",
542            "tool_output": "search results"
543        });
544        let event = detect_event_type(&v, 1000).unwrap();
545        assert_eq!(event.event_type, "mcp_call");
546    }
547
548    #[test]
549    fn detect_event_type_ctx_prefix_is_mcp_call() {
550        let v = serde_json::json!({
551            "tool_name": "ctx_read",
552            "tool_input": {"path": "src/main.rs"}
553        });
554        let event = detect_event_type(&v, 1000).unwrap();
555        assert_eq!(event.event_type, "mcp_call");
556    }
557
558    #[test]
559    fn detect_event_type_mcp_prefix_is_mcp_call() {
560        let v = serde_json::json!({
561            "tool_name": "mcp__lean-ctx__ctx_read",
562            "tool_input": {"path": "src/main.rs"}
563        });
564        let event = detect_event_type(&v, 1000).unwrap();
565        assert_eq!(event.event_type, "mcp_call");
566    }
567
568    #[test]
569    fn detect_event_type_native_read_is_native_tool() {
570        let v = serde_json::json!({
571            "tool_name": "Read",
572            "tool_input": {"path": "src/main.rs"}
573        });
574        let event = detect_event_type(&v, 1000).unwrap();
575        assert_eq!(event.event_type, "native_tool");
576    }
577
578    #[test]
579    fn detect_event_type_copilot_bash_posttooluse_is_shell() {
580        // #551: Copilot CLI postToolUse — camelCase `toolName` + JSON-string
581        // `toolArgs` + `toolResult`. Was dropped before the fix; now recorded.
582        let v = serde_json::json!({
583            "toolName": "bash",
584            "toolArgs": "{\"command\":\"npm test\"}",
585            "toolResult": {
586                "resultType": "success",
587                "textResultForLlm": "All tests passed (15/15)"
588            }
589        });
590        let event = detect_event_type(&v, 1000).unwrap();
591        assert_eq!(event.event_type, "shell");
592        assert_eq!(event.tool_name.as_deref(), Some("bash"));
593        assert_eq!(event.detail.as_deref(), Some("npm test"));
594        assert!(event.content.unwrap().contains("All tests passed"));
595    }
596
597    #[test]
598    fn detect_event_type_copilot_ctx_tool_is_mcp_call() {
599        let v = serde_json::json!({
600            "toolName": "ctx_read",
601            "toolArgs": "{\"path\":\"src/main.rs\"}",
602            "toolResult": { "textResultForLlm": "file contents" }
603        });
604        let event = detect_event_type(&v, 1000).unwrap();
605        assert_eq!(event.event_type, "mcp_call");
606        assert_eq!(event.tool_name.as_deref(), Some("ctx_read"));
607    }
608
609    #[test]
610    fn detect_event_type_result_json_is_mcp_call() {
611        let v = serde_json::json!({
612            "tool_name": "ctx_read",
613            "result_json": {"content": "..."}
614        });
615        let event = detect_event_type(&v, 1000).unwrap();
616        assert_eq!(event.event_type, "mcp_call");
617    }
618
619    /// Real Claude Code PreCompact payload (code.claude.com/docs/en/hooks):
620    /// carries `session_id` like every Claude hook, so the compaction check
621    /// must win over the generic session catch-all (GL #555).
622    #[test]
623    fn detect_event_type_claude_precompact_is_compaction() {
624        let v = serde_json::json!({
625            "session_id": "abc123",
626            "transcript_path": "/Users/u/.claude/projects/x/abc123.jsonl",
627            "cwd": "/Users/u/project",
628            "hook_event_name": "PreCompact",
629            "trigger": "auto",
630            "custom_instructions": ""
631        });
632        let event = detect_event_type(&v, 1000).unwrap();
633        assert_eq!(event.event_type, "compaction");
634    }
635
636    #[test]
637    fn detect_event_type_plain_session_event_still_session() {
638        let v = serde_json::json!({
639            "session_id": "abc123",
640            "hook_event_name": "SessionStart"
641        });
642        let event = detect_event_type(&v, 1000).unwrap();
643        assert_eq!(event.event_type, "session");
644    }
645}