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 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 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
442fn 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; 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 #[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}