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 let Some(event) = parse_observe_event(&input) else {
29 return;
30 };
31 append_radar_event(&event);
32
33 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 }
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 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 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
481fn 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; 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 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 #[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}