1#[allow(clippy::wildcard_imports)]
6use super::*;
7
8pub 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 emit_dedicated_session_context(&input);
28
29 super::edit_health::maybe_emit(&input);
33
34 let Some(event) = parse_observe_event(&input) else {
35 return;
36 };
37 if event.event_type == "compaction"
41 && let Ok(v) = serde_json::from_str::<serde_json::Value>(&input)
42 && let Some(session_id) = v.get("session_id").and_then(|s| s.as_str())
43 {
44 super::read_dedup::purge_session(session_id);
45 }
46 append_radar_event(&event);
47
48 if event.event_type == "agent_response"
52 && let Some(text) = event.content.as_deref()
53 {
54 crate::core::output_echo::analyze_and_record(text);
55 }
56}
57
58fn emit_dedicated_session_context(input: &str) {
59 let Ok(v) = serde_json::from_str::<serde_json::Value>(input) else {
60 return;
61 };
62 if !session_start_honours_additional_context(&v) {
63 return;
64 }
65 let cfg = crate::core::config::Config::load();
66 if !cfg.dedicated_session_context_active() {
67 return;
68 }
69 let summary = crate::core::rules_canonical::render(
73 cfg.shadow_mode,
74 crate::core::rules_canonical::Wrapper::Bare,
75 crate::core::config::CompressionLevel::Off,
76 );
77 emit_session_start_additional_context(&summary);
78}
79
80fn session_start_honours_additional_context(v: &serde_json::Value) -> bool {
89 let is_session_start =
90 v.get("hook_event_name").and_then(|e| e.as_str()) == Some("SessionStart");
91 let is_cursor = v
92 .get("conversation_id")
93 .and_then(|c| c.as_str())
94 .is_some_and(|c| !c.is_empty());
95 is_session_start && !is_cursor
96}
97
98#[derive(serde::Serialize)]
99struct ObserveEvent {
100 ts: u64,
101 event_type: &'static str,
102 tokens: usize,
103 #[serde(skip_serializing_if = "Option::is_none")]
104 tool_name: Option<String>,
105 #[serde(skip_serializing_if = "Option::is_none")]
106 detail: Option<String>,
107 #[serde(skip_serializing_if = "Option::is_none")]
108 content: Option<String>,
109 #[serde(skip_serializing_if = "Option::is_none")]
110 model: Option<String>,
111 #[serde(skip_serializing_if = "Option::is_none")]
112 conversation_id: Option<String>,
113}
114
115const MAX_CONTENT_CHARS: usize = 50_000;
116
117fn parse_observe_event(input: &str) -> Option<ObserveEvent> {
118 let v: serde_json::Value = serde_json::from_str(input).ok()?;
119
120 let ts = std::time::SystemTime::now()
121 .duration_since(std::time::UNIX_EPOCH)
122 .unwrap_or_default()
123 .as_secs();
124
125 let model = v
126 .get("model")
127 .and_then(|m| m.as_str())
128 .filter(|m| !m.is_empty())
129 .map(String::from);
130 let conversation_id = v
131 .get("conversation_id")
132 .and_then(|c| c.as_str())
133 .filter(|c| !c.is_empty())
134 .map(String::from);
135
136 let transcript_path = v
137 .get("transcript_path")
138 .and_then(|t| t.as_str())
139 .filter(|t| !t.is_empty())
140 .map(String::from);
141
142 if let Some(ref m) = model {
143 persist_detected_model(m);
144 }
145 if let Some(ref tp) = transcript_path {
146 persist_transcript_path(tp, conversation_id.as_deref());
147 }
148
149 let mut event = detect_event_type(&v, ts)?;
150 event.model = model;
151 event.conversation_id = conversation_id;
152 Some(event)
153}
154
155fn detect_event_type(v: &serde_json::Value, ts: u64) -> Option<ObserveEvent> {
156 if let Some(result) = v.get("toolResult") {
161 let tool = super::payload::resolve_tool_name(v).unwrap_or_else(|| "unknown".to_string());
162 let args = super::payload::resolve_tool_args(v);
163 let command = args
164 .as_ref()
165 .and_then(|a| a.get("command"))
166 .and_then(|c| c.as_str());
167 let result_text = result
168 .get("textResultForLlm")
169 .and_then(|t| t.as_str())
170 .map_or_else(|| result.to_string(), String::from);
171 let tokens = result_text.len() / 4;
172 let is_lctx = tool.starts_with("ctx_") || tool.starts_with("mcp__lean-ctx__");
173 let event_type = if is_lctx {
174 "mcp_call"
175 } else if command.is_some() {
176 "shell"
177 } else {
178 "native_tool"
179 };
180 let content = match command {
181 Some(cmd) => format!("$ {cmd}\n{result_text}"),
182 None => result_text,
183 };
184 return Some(ObserveEvent {
185 ts,
186 event_type,
187 tokens,
188 tool_name: Some(tool),
189 detail: command.map(|c| truncate_str(c, 80)),
190 content: Some(cap_content(&content)),
191 model: None,
192 conversation_id: None,
193 });
194 }
195
196 if let Some(result) = v
197 .get("result_json")
198 .or_else(|| v.get("result"))
199 .or_else(|| v.get("tool_response"))
200 .or_else(|| v.get("tool_output"))
201 {
202 let tool = v
203 .get("tool_name")
204 .and_then(|t| t.as_str())
205 .unwrap_or("unknown");
206 let tokens = estimate_tokens_json(result);
207 let content_str = match result {
208 serde_json::Value::String(s) => s.clone(),
209 other => other.to_string(),
210 };
211 return Some(ObserveEvent {
212 ts,
213 event_type: "mcp_call",
214 tokens,
215 tool_name: Some(tool.to_string()),
216 detail: v
217 .get("server_name")
218 .and_then(|s| s.as_str())
219 .map(String::from),
220 content: Some(cap_content(&content_str)),
221 model: None,
222 conversation_id: None,
223 });
224 }
225
226 if let Some(output) = v.get("output") {
227 let cmd = v
228 .get("command")
229 .and_then(|c| c.as_str())
230 .unwrap_or("")
231 .to_string();
232 let tokens = estimate_tokens_value(output);
233 let out_str = match output {
234 serde_json::Value::String(s) => s.clone(),
235 other => other.to_string(),
236 };
237 return Some(ObserveEvent {
238 ts,
239 event_type: "shell",
240 tokens,
241 tool_name: None,
242 detail: Some(truncate_str(&cmd, 80)),
243 content: Some(cap_content(&format!("$ {cmd}\n{out_str}"))),
244 model: None,
245 conversation_id: None,
246 });
247 }
248
249 if v.get("content").is_some() && v.get("file_path").is_some() {
250 let path = v
251 .get("file_path")
252 .and_then(|p| p.as_str())
253 .unwrap_or("")
254 .to_string();
255 let file_content = v.get("content").and_then(|c| c.as_str()).unwrap_or("");
256 let tokens = file_content.len() / 4;
257 return Some(ObserveEvent {
258 ts,
259 event_type: "file_read",
260 tokens,
261 tool_name: None,
262 detail: Some(truncate_str(&path, 120)),
263 content: Some(cap_content(file_content)),
264 model: None,
265 conversation_id: None,
266 });
267 }
268
269 if let Some(text) = v.get("text").and_then(|t| t.as_str()) {
270 let has_duration = v.get("duration_ms").is_some();
271 let event_type = if has_duration {
272 "thinking"
273 } else {
274 "agent_response"
275 };
276 let tokens = text.len() / 4;
277 return Some(ObserveEvent {
278 ts,
279 event_type,
280 tokens,
281 tool_name: None,
282 detail: None,
283 content: Some(cap_content(text)),
284 model: None,
285 conversation_id: None,
286 });
287 }
288
289 if let Some(prompt) = v.get("prompt").and_then(|p| p.as_str()) {
290 let tokens = prompt.len() / 4;
291 let mut full = prompt.to_string();
292 if let Some(attachments) = v.get("attachments").and_then(|a| a.as_array())
293 && !attachments.is_empty()
294 {
295 full.push_str(&format!("\n\n[{} attachments]", attachments.len()));
296 for att in attachments {
297 if let Some(name) = att.get("name").and_then(|n| n.as_str()) {
298 full.push_str(&format!("\n - {name}"));
299 }
300 }
301 }
302 return Some(ObserveEvent {
303 ts,
304 event_type: "user_message",
305 tokens,
306 tool_name: None,
307 detail: v
308 .get("attachments")
309 .and_then(|a| a.as_array())
310 .map(|a| format!("{} attachments", a.len())),
311 content: Some(cap_content(&full)),
312 model: None,
313 conversation_id: None,
314 });
315 }
316
317 if v.get("tool_name").is_some() || v.get("tool_input").is_some() {
318 let tool = v
319 .get("tool_name")
320 .and_then(|t| t.as_str())
321 .unwrap_or("unknown")
322 .to_string();
323 let is_lctx = tool.starts_with("ctx_") || tool.starts_with("mcp__lean-ctx__");
324 let tokens = v.get("tool_input").map_or(0, estimate_tokens_json);
325 let input_str = v
326 .get("tool_input")
327 .map(std::string::ToString::to_string)
328 .unwrap_or_default();
329 return Some(ObserveEvent {
330 ts,
331 event_type: if is_lctx { "mcp_call" } else { "native_tool" },
332 tokens,
333 tool_name: Some(tool),
334 detail: None,
335 content: if input_str.is_empty() {
336 None
337 } else {
338 Some(cap_content(&input_str))
339 },
340 model: None,
341 conversation_id: None,
342 });
343 }
344
345 let is_compaction = v.get("compaction").is_some()
355 || v.get("messages_count").is_some()
356 || v.get("hook_event_name")
357 .and_then(|e| e.as_str())
358 .is_some_and(|e| e == "PreCompact")
359 || v.get("event")
360 .and_then(|e| e.as_str())
361 .is_some_and(|e| e == "compaction" || e == "compact");
362 if !is_compaction && v.get("session_id").is_some() {
363 return Some(ObserveEvent {
364 ts,
365 event_type: "session",
366 tokens: 0,
367 tool_name: None,
368 detail: v
369 .get("session_id")
370 .and_then(|s| s.as_str())
371 .map(String::from),
372 content: None,
373 model: None,
374 conversation_id: None,
375 });
376 }
377
378 if is_compaction {
379 return Some(ObserveEvent {
380 ts,
381 event_type: "compaction",
382 tokens: 0,
383 tool_name: None,
384 detail: None,
385 content: None,
386 model: None,
387 conversation_id: None,
388 });
389 }
390
391 None
392}
393
394fn estimate_tokens_json(v: &serde_json::Value) -> usize {
395 match v {
396 serde_json::Value::String(s) => s.len() / 4,
397 _ => v.to_string().len() / 4,
398 }
399}
400
401fn estimate_tokens_value(v: &serde_json::Value) -> usize {
402 match v {
403 serde_json::Value::String(s) => s.len() / 4,
404 _ => v.to_string().len() / 4,
405 }
406}
407
408fn persist_detected_model(model: &str) {
409 let m = model.to_lowercase();
410 let is_bg_model = m.contains("flash")
411 || m.contains("mini")
412 || m.contains("haiku")
413 || m.contains("fast")
414 || m.contains("nano")
415 || m.contains("small");
416 if is_bg_model {
417 return;
418 }
419
420 let Ok(data_dir) = crate::core::data_dir::lean_ctx_data_dir() else {
421 return;
422 };
423 let path = data_dir.join("detected_model.json");
424 let ts = std::time::SystemTime::now()
425 .duration_since(std::time::UNIX_EPOCH)
426 .unwrap_or_default()
427 .as_secs();
428 let window = model_context_window(model);
429 let payload = serde_json::json!({
430 "model": model,
431 "window_size": window,
432 "detected_at": ts,
433 });
434 if let Ok(json) = serde_json::to_string_pretty(&payload) {
435 let tmp = path.with_extension("tmp");
436 if std::fs::write(&tmp, &json).is_ok() {
437 let _ = std::fs::rename(&tmp, &path);
438 }
439 }
440}
441
442pub fn model_context_window(model: &str) -> usize {
443 crate::core::model_registry::context_window_for_model(model)
444}
445
446pub fn load_detected_model() -> Option<(String, usize)> {
447 let data_dir = crate::core::data_dir::lean_ctx_data_dir().ok()?;
448 let path = data_dir.join("detected_model.json");
449 let content = std::fs::read_to_string(&path).ok()?;
450 let v: serde_json::Value = serde_json::from_str(&content).ok()?;
451 let model = v.get("model")?.as_str()?.to_string();
452 let window = v.get("window_size")?.as_u64()? as usize;
453 let detected_at = v.get("detected_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(detected_at) > 7200 {
459 return None;
460 }
461 Some((model, window))
462}
463
464fn persist_transcript_path(path: &str, conversation_id: Option<&str>) {
465 let Ok(data_dir) = crate::core::data_dir::lean_ctx_data_dir() else {
466 return;
467 };
468 let meta_path = data_dir.join("active_transcript.json");
469 let ts = std::time::SystemTime::now()
470 .duration_since(std::time::UNIX_EPOCH)
471 .unwrap_or_default()
472 .as_secs();
473 let payload = serde_json::json!({
474 "transcript_path": path,
475 "conversation_id": conversation_id,
476 "updated_at": ts,
477 });
478 if let Ok(json) = serde_json::to_string_pretty(&payload) {
479 let tmp = meta_path.with_extension("tmp");
480 if std::fs::write(&tmp, &json).is_ok() {
481 let _ = std::fs::rename(&tmp, &meta_path);
482 }
483 }
484}
485
486pub fn load_active_transcript() -> Option<(String, Option<String>)> {
487 let data_dir = crate::core::data_dir::lean_ctx_data_dir().ok()?;
488 let path = data_dir.join("active_transcript.json");
489 let content = std::fs::read_to_string(&path).ok()?;
490 let v: serde_json::Value = serde_json::from_str(&content).ok()?;
491 let tp = v.get("transcript_path")?.as_str()?.to_string();
492 let conv = v
493 .get("conversation_id")
494 .and_then(|c| c.as_str())
495 .map(String::from);
496 let updated = v.get("updated_at")?.as_u64()?;
497 let now = std::time::SystemTime::now()
498 .duration_since(std::time::UNIX_EPOCH)
499 .unwrap_or_default()
500 .as_secs();
501 if now.saturating_sub(updated) > 7200 {
502 return None;
503 }
504 Some((tp, conv))
505}
506
507fn cap_content(s: &str) -> String {
508 if s.len() <= MAX_CONTENT_CHARS {
509 s.to_string()
510 } else {
511 let truncated = safe_truncate(s, MAX_CONTENT_CHARS);
512 format!("{}…\n\n[truncated: {} total chars]", truncated, s.len())
513 }
514}
515
516fn truncate_str(s: &str, max: usize) -> String {
517 if s.len() <= max {
518 s.to_string()
519 } else {
520 format!("{}...", safe_truncate(s, max))
521 }
522}
523
524fn safe_truncate(s: &str, max: usize) -> &str {
526 if max >= s.len() {
527 return s;
528 }
529 let mut end = max;
530 while end > 0 && !s.is_char_boundary(end) {
531 end -= 1;
532 }
533 &s[..end]
534}
535
536fn append_radar_event(event: &ObserveEvent) {
537 let Ok(data_dir) = crate::core::data_dir::lean_ctx_data_dir() else {
538 return;
539 };
540 let radar_path = data_dir.join("context_radar.jsonl");
541
542 if event.event_type == "session"
543 && let Ok(meta) = std::fs::metadata(&radar_path)
544 {
545 const MAX_RADAR_SIZE: u64 = 10 * 1024 * 1024; if meta.len() > MAX_RADAR_SIZE {
547 let prev = data_dir.join("context_radar.prev.jsonl");
548 let _ = std::fs::rename(&radar_path, &prev);
549 }
550 }
551
552 let Ok(line) = serde_json::to_string(event) else {
553 return;
554 };
555
556 use std::fs::OpenOptions;
557 use std::io::Write;
558 if let Ok(mut f) = OpenOptions::new()
559 .create(true)
560 .append(true)
561 .open(&radar_path)
562 {
563 let _ = writeln!(f, "{line}");
564 }
565}
566
567#[must_use]
575pub fn radar_event_count() -> usize {
576 let Ok(data_dir) = crate::core::data_dir::lean_ctx_data_dir() else {
577 return 0;
578 };
579 let Ok(file) = std::fs::File::open(data_dir.join("context_radar.jsonl")) else {
580 return 0;
581 };
582 use std::io::{BufRead, BufReader};
583 BufReader::new(file)
584 .lines()
585 .map_while(Result::ok)
586 .filter(|l| !l.trim().is_empty())
587 .count()
588}
589
590#[cfg(test)]
591mod tests {
592 use super::*;
593
594 #[test]
595 fn detect_event_type_tool_response_is_mcp_call() {
596 let v = serde_json::json!({
597 "tool_name": "ctx_read",
598 "tool_response": "file contents here"
599 });
600 let event = detect_event_type(&v, 1000).unwrap();
601 assert_eq!(event.event_type, "mcp_call");
602 }
603
604 #[test]
605 fn detect_event_type_tool_output_is_mcp_call() {
606 let v = serde_json::json!({
607 "tool_name": "ctx_search",
608 "tool_output": "search results"
609 });
610 let event = detect_event_type(&v, 1000).unwrap();
611 assert_eq!(event.event_type, "mcp_call");
612 }
613
614 #[test]
615 fn detect_event_type_ctx_prefix_is_mcp_call() {
616 let v = serde_json::json!({
617 "tool_name": "ctx_read",
618 "tool_input": {"path": "src/main.rs"}
619 });
620 let event = detect_event_type(&v, 1000).unwrap();
621 assert_eq!(event.event_type, "mcp_call");
622 }
623
624 #[test]
625 fn detect_event_type_mcp_prefix_is_mcp_call() {
626 let v = serde_json::json!({
627 "tool_name": "mcp__lean-ctx__ctx_read",
628 "tool_input": {"path": "src/main.rs"}
629 });
630 let event = detect_event_type(&v, 1000).unwrap();
631 assert_eq!(event.event_type, "mcp_call");
632 }
633
634 #[test]
635 fn detect_event_type_native_read_is_native_tool() {
636 let v = serde_json::json!({
637 "tool_name": "Read",
638 "tool_input": {"path": "src/main.rs"}
639 });
640 let event = detect_event_type(&v, 1000).unwrap();
641 assert_eq!(event.event_type, "native_tool");
642 }
643
644 #[test]
645 fn detect_event_type_copilot_bash_posttooluse_is_shell() {
646 let v = serde_json::json!({
649 "toolName": "bash",
650 "toolArgs": "{\"command\":\"npm test\"}",
651 "toolResult": {
652 "resultType": "success",
653 "textResultForLlm": "All tests passed (15/15)"
654 }
655 });
656 let event = detect_event_type(&v, 1000).unwrap();
657 assert_eq!(event.event_type, "shell");
658 assert_eq!(event.tool_name.as_deref(), Some("bash"));
659 assert_eq!(event.detail.as_deref(), Some("npm test"));
660 assert!(event.content.unwrap().contains("All tests passed"));
661 }
662
663 #[test]
664 fn detect_event_type_copilot_ctx_tool_is_mcp_call() {
665 let v = serde_json::json!({
666 "toolName": "ctx_read",
667 "toolArgs": "{\"path\":\"src/main.rs\"}",
668 "toolResult": { "textResultForLlm": "file contents" }
669 });
670 let event = detect_event_type(&v, 1000).unwrap();
671 assert_eq!(event.event_type, "mcp_call");
672 assert_eq!(event.tool_name.as_deref(), Some("ctx_read"));
673 }
674
675 #[test]
676 fn detect_event_type_result_json_is_mcp_call() {
677 let v = serde_json::json!({
678 "tool_name": "ctx_read",
679 "result_json": {"content": "..."}
680 });
681 let event = detect_event_type(&v, 1000).unwrap();
682 assert_eq!(event.event_type, "mcp_call");
683 }
684
685 #[test]
689 fn detect_event_type_claude_precompact_is_compaction() {
690 let v = serde_json::json!({
691 "session_id": "abc123",
692 "transcript_path": "/Users/u/.claude/projects/x/abc123.jsonl",
693 "cwd": "/Users/u/project",
694 "hook_event_name": "PreCompact",
695 "trigger": "auto",
696 "custom_instructions": ""
697 });
698 let event = detect_event_type(&v, 1000).unwrap();
699 assert_eq!(event.event_type, "compaction");
700 }
701
702 #[test]
703 fn detect_event_type_plain_session_event_still_session() {
704 let v = serde_json::json!({
705 "session_id": "abc123",
706 "hook_event_name": "SessionStart"
707 });
708 let event = detect_event_type(&v, 1000).unwrap();
709 assert_eq!(event.event_type, "session");
710 }
711
712 #[test]
713 fn session_start_honoured_for_claude_payload() {
714 let v = serde_json::json!({
716 "hook_event_name": "SessionStart",
717 "session_id": "abc123",
718 "source": "startup"
719 });
720 assert!(session_start_honours_additional_context(&v));
721 }
722
723 #[test]
724 fn session_start_skipped_for_cursor_payload() {
725 let v = serde_json::json!({
728 "hook_event_name": "SessionStart",
729 "conversation_id": "0e1f4ed8-d858-4557-9fc5-6cbf5298eb8b",
730 "model": "claude-opus"
731 });
732 assert!(!session_start_honours_additional_context(&v));
733 }
734
735 #[test]
736 fn session_start_skipped_for_non_session_event() {
737 let v = serde_json::json!({ "hook_event_name": "PreToolUse", "session_id": "x" });
738 assert!(!session_start_honours_additional_context(&v));
739 }
740}