1use std::collections::BTreeMap;
7
8use serde::{Deserialize, Serialize};
9
10use super::event::{EventType, SessionEvent};
11
12#[derive(Debug, Clone, Serialize, Deserialize)]
14pub struct FileAccess {
15 pub file_path: String,
16 pub agent_instance_id: String,
17 pub timestamp: String,
18 #[serde(skip_serializing_if = "Option::is_none")]
19 pub digest: Option<String>,
20 #[serde(default, skip_serializing_if = "Option::is_none")]
22 pub operation: Option<String>,
23 #[serde(default, skip_serializing_if = "Option::is_none")]
24 pub additions: Option<u32>,
25 #[serde(default, skip_serializing_if = "Option::is_none")]
26 pub deletions: Option<u32>,
27 #[serde(default, skip_serializing_if = "Option::is_none")]
44 pub source: Option<String>,
45}
46
47#[derive(Debug, Clone, Serialize, Deserialize)]
49pub struct PortAccess {
50 pub port: u16,
51 pub agent_instance_id: String,
52 pub timestamp: String,
53 #[serde(skip_serializing_if = "Option::is_none")]
54 pub protocol: Option<String>,
55}
56
57#[derive(Debug, Clone, Serialize, Deserialize)]
59pub struct NetworkConnection {
60 pub destination: String,
61 #[serde(skip_serializing_if = "Option::is_none")]
62 pub port: Option<u16>,
63 pub agent_instance_id: String,
64 pub timestamp: String,
65}
66
67#[derive(Debug, Clone, Serialize, Deserialize)]
69pub struct ProcessExecution {
70 pub process_name: String,
71 pub agent_instance_id: String,
72 pub started_at: String,
73 #[serde(skip_serializing_if = "Option::is_none")]
74 pub exit_code: Option<i32>,
75 #[serde(skip_serializing_if = "Option::is_none")]
76 pub duration_ms: Option<u64>,
77 #[serde(default, skip_serializing_if = "Option::is_none")]
79 pub command: Option<String>,
80 #[serde(default, skip_serializing_if = "Option::is_none")]
83 pub source: Option<String>,
84}
85
86#[derive(Debug, Clone, Serialize, Deserialize)]
88pub struct ToolInvocation {
89 pub tool_name: String,
90 pub agent_instance_id: String,
91 pub timestamp: String,
92 #[serde(skip_serializing_if = "Option::is_none")]
93 pub duration_ms: Option<u64>,
94}
95
96#[derive(Debug, Clone, Default, Serialize, Deserialize)]
98pub struct SideEffects {
99 pub files_read: Vec<FileAccess>,
100 pub files_written: Vec<FileAccess>,
101 pub ports_opened: Vec<PortAccess>,
102 pub network_connections: Vec<NetworkConnection>,
103 pub processes: Vec<ProcessExecution>,
104 pub tool_invocations: Vec<ToolInvocation>,
105}
106
107impl SideEffects {
108 pub fn from_events(events: &[SessionEvent]) -> Self {
110 let mut se = SideEffects::default();
111
112 let mut started_processes: BTreeMap<(String, String), usize> = BTreeMap::new();
115
116 for event in events {
117 match &event.event_type {
118 EventType::AgentReadFile { file_path, digest } => {
119 se.files_read.push(FileAccess {
120 file_path: file_path.clone(),
121 agent_instance_id: event.agent_instance_id.clone(),
122 timestamp: event.timestamp.clone(),
123 digest: digest.clone(),
124 operation: None,
125 additions: None,
126 deletions: None,
127 source: Some(source_from_meta(event, "hook")),
128 });
129 }
130
131 EventType::AgentWroteFile {
132 file_path,
133 digest,
134 operation,
135 additions,
136 deletions,
137 } => {
138 se.files_written.push(FileAccess {
139 file_path: file_path.clone(),
140 agent_instance_id: event.agent_instance_id.clone(),
141 timestamp: event.timestamp.clone(),
142 digest: digest.clone(),
143 operation: operation.clone(),
144 additions: *additions,
145 deletions: *deletions,
146 source: Some(source_from_meta(event, "hook")),
147 });
148 }
149
150 EventType::AgentOpenedPort { port, protocol } => {
151 se.ports_opened.push(PortAccess {
152 port: *port,
153 agent_instance_id: event.agent_instance_id.clone(),
154 timestamp: event.timestamp.clone(),
155 protocol: protocol.clone(),
156 });
157 }
158
159 EventType::AgentConnectedNetwork { destination, port } => {
160 se.network_connections.push(NetworkConnection {
161 destination: destination.clone(),
162 port: *port,
163 agent_instance_id: event.agent_instance_id.clone(),
164 timestamp: event.timestamp.clone(),
165 });
166 }
167
168 EventType::AgentStartedProcess {
169 process_name,
170 pid: _,
171 command,
172 } => {
173 let idx = se.processes.len();
174 se.processes.push(ProcessExecution {
175 process_name: process_name.clone(),
176 agent_instance_id: event.agent_instance_id.clone(),
177 started_at: event.timestamp.clone(),
178 exit_code: None,
179 duration_ms: None,
180 command: command.clone(),
181 source: Some(source_from_meta(event, "hook")),
182 });
183 started_processes
184 .insert((event.agent_instance_id.clone(), process_name.clone()), idx);
185 }
186
187 EventType::AgentCompletedProcess {
188 process_name,
189 exit_code,
190 duration_ms,
191 command,
192 } => {
193 let key = (event.agent_instance_id.clone(), process_name.clone());
194 if let Some(&idx) = started_processes.get(&key) {
195 if let Some(proc) = se.processes.get_mut(idx) {
196 proc.exit_code = *exit_code;
197 proc.duration_ms = *duration_ms;
198 if proc.command.is_none() {
199 proc.command = command.clone();
200 }
201 }
202 } else {
203 se.processes.push(ProcessExecution {
204 process_name: process_name.clone(),
205 agent_instance_id: event.agent_instance_id.clone(),
206 started_at: event.timestamp.clone(),
207 exit_code: *exit_code,
208 duration_ms: *duration_ms,
209 command: command.clone(),
210 source: Some(source_from_meta(event, "hook")),
211 });
212 }
213 }
214
215 EventType::AgentCalledTool {
216 tool_name,
217 duration_ms,
218 ..
219 } => {
220 se.tool_invocations.push(ToolInvocation {
221 tool_name: tool_name.clone(),
222 agent_instance_id: event.agent_instance_id.clone(),
223 timestamp: event.timestamp.clone(),
224 duration_ms: *duration_ms,
225 });
226
227 promote_mcp_called_tool(event, tool_name, &mut se);
244 }
245
246 _ => {}
247 }
248 }
249
250 se
251 }
252
253 pub fn summary(&self) -> SideEffectSummary {
255 SideEffectSummary {
256 files_read: self.files_read.len() as u32,
257 files_written: self.files_written.len() as u32,
258 ports_opened: self.ports_opened.len() as u32,
259 network_connections: self.network_connections.len() as u32,
260 processes: self.processes.len() as u32,
261 tool_invocations: self.tool_invocations.len() as u32,
262 }
263 }
264}
265
266#[derive(Debug, Clone, Serialize, Deserialize)]
268pub struct SideEffectSummary {
269 pub files_read: u32,
270 pub files_written: u32,
271 pub ports_opened: u32,
272 pub network_connections: u32,
273 pub processes: u32,
274 pub tool_invocations: u32,
275}
276
277fn meta_string(event: &SessionEvent, dotted_path: &str) -> Option<String> {
298 let mut cur = event.meta.as_ref()?;
299 for segment in dotted_path.split('.') {
300 cur = cur.get(segment)?;
301 }
302 cur.as_str().map(|s| s.to_string())
303}
304
305fn first_meta_string(event: &SessionEvent, paths: &[&str]) -> Option<String> {
307 for path in paths {
308 if let Some(v) = meta_string(event, path) {
309 if !v.is_empty() {
310 return Some(v);
311 }
312 }
313 }
314 None
315}
316
317fn source_from_meta(event: &SessionEvent, default: &str) -> String {
340 event
341 .meta
342 .as_ref()
343 .and_then(|m| m.get("source"))
344 .and_then(|v| v.as_str())
345 .filter(|s| !s.is_empty())
346 .map(|s| s.to_string())
347 .unwrap_or_else(|| default.to_string())
348}
349
350#[derive(Debug, Clone, Copy, PartialEq, Eq)]
355enum ToolCategory {
356 Read,
357 Write,
358 Process,
359 Unknown,
360}
361
362fn classify_tool(tool_name: &str) -> ToolCategory {
363 let t = tool_name.to_lowercase();
364 if t.contains("bash")
366 || t.contains("shell")
367 || t.contains("exec")
368 || t.contains("run_command")
369 || t.contains("ran_command")
370 {
371 return ToolCategory::Process;
372 }
373 if t.contains("write")
376 || t.contains("edit")
377 || t.contains("create_file")
378 || t.contains("modify")
379 || t.contains("patch")
380 || t.contains("save_file")
381 || t.contains("delete_file")
382 || t.contains("remove_file")
383 || t.contains("rename_file")
384 {
385 return ToolCategory::Write;
386 }
387 if t.contains("read")
389 || t.contains("view_file")
390 || t.contains("cat_file")
391 || t.contains("open_file")
392 || t.contains("get_file_contents")
393 {
394 return ToolCategory::Read;
395 }
396 ToolCategory::Unknown
397}
398
399fn promote_mcp_called_tool(event: &SessionEvent, tool_name: &str, se: &mut SideEffects) {
400 let category = classify_tool(tool_name);
401
402 let file_path = first_meta_string(
405 event,
406 &[
407 "tool_input.file_path",
408 "tool_input.path",
409 "tool_input.notebook_path",
410 "tool_input.target_file",
411 "file_path",
412 "path",
413 ],
414 );
415
416 let command = first_meta_string(
418 event,
419 &["tool_input.command", "command", "tool_input.cmd", "cmd"],
420 );
421
422 match (category, file_path, command) {
423 (ToolCategory::Read, Some(p), _) => {
424 se.files_read.push(FileAccess {
425 file_path: p,
426 agent_instance_id: event.agent_instance_id.clone(),
427 timestamp: event.timestamp.clone(),
428 digest: None,
429 operation: None,
430 additions: None,
431 deletions: None,
432 source: Some("mcp".into()),
433 });
434 }
435 (ToolCategory::Write, Some(p), _) => {
436 se.files_written.push(FileAccess {
437 file_path: p,
438 agent_instance_id: event.agent_instance_id.clone(),
439 timestamp: event.timestamp.clone(),
440 digest: None,
441 operation: None,
442 additions: None,
443 deletions: None,
444 source: Some("mcp".into()),
445 });
446 }
447 (ToolCategory::Process, _, Some(cmd)) => {
448 let short = cmd.chars().take(120).collect::<String>();
451 se.processes.push(ProcessExecution {
452 process_name: short,
453 agent_instance_id: event.agent_instance_id.clone(),
454 started_at: event.timestamp.clone(),
455 exit_code: None,
456 duration_ms: None,
457 command: Some(cmd),
458 source: Some("mcp".into()),
459 });
460 }
461 (ToolCategory::Unknown, Some(p), _) => {
462 se.files_written.push(FileAccess {
468 file_path: p,
469 agent_instance_id: event.agent_instance_id.clone(),
470 timestamp: event.timestamp.clone(),
471 digest: None,
472 operation: None,
473 additions: None,
474 deletions: None,
475 source: Some("mcp".into()),
476 });
477 }
478 _ => {
479 }
484 }
485}
486
487#[cfg(test)]
488mod tests {
489 use super::*;
490 use crate::session::event::*;
491
492 fn evt(event_type: EventType) -> SessionEvent {
493 SessionEvent {
494 session_id: "ssn_001".into(),
495 event_id: generate_event_id(),
496 timestamp: "2026-04-05T08:00:00Z".into(),
497 sequence_no: 0,
498 trace_id: "t".into(),
499 span_id: "s".into(),
500 parent_span_id: None,
501 agent_id: "agent://test".into(),
502 agent_instance_id: "ai_1".into(),
503 agent_name: "test".into(),
504 agent_role: None,
505 host_id: "h".into(),
506 tool_runtime_id: None,
507 event_type,
508 artifact_ref: None,
509 meta: None,
510 }
511 }
512
513 #[test]
514 fn aggregates_file_and_tool_events() {
515 let events = vec![
516 evt(EventType::AgentReadFile {
517 file_path: "src/main.rs".into(),
518 digest: None,
519 }),
520 evt(EventType::AgentWroteFile {
521 file_path: "src/lib.rs".into(),
522 digest: Some("sha256:abc".into()),
523 operation: Some("modified".into()),
524 additions: Some(10),
525 deletions: Some(3),
526 }),
527 evt(EventType::AgentCalledTool {
528 tool_name: "read_file".into(),
529 tool_input_digest: None,
530 tool_output_digest: None,
531 duration_ms: Some(10),
532 }),
533 evt(EventType::AgentCalledTool {
534 tool_name: "write_file".into(),
535 tool_input_digest: None,
536 tool_output_digest: None,
537 duration_ms: None,
538 }),
539 ];
540
541 let se = SideEffects::from_events(&events);
542 assert_eq!(se.files_read.len(), 1);
543 assert_eq!(se.files_written.len(), 1);
544 assert_eq!(se.tool_invocations.len(), 2);
545 let summary = se.summary();
546 assert_eq!(summary.tool_invocations, 2);
547 }
548
549 #[test]
550 fn matches_process_start_and_complete() {
551 let events = vec![
552 evt(EventType::AgentStartedProcess {
553 process_name: "npm test".into(),
554 pid: Some(1234),
555 command: Some("npm test --runInBand".into()),
556 }),
557 evt(EventType::AgentCompletedProcess {
558 process_name: "npm test".into(),
559 exit_code: Some(0),
560 duration_ms: Some(5000),
561 command: None,
562 }),
563 ];
564
565 let se = SideEffects::from_events(&events);
566 assert_eq!(se.processes.len(), 1);
567 assert_eq!(se.processes[0].exit_code, Some(0));
568 assert_eq!(se.processes[0].duration_ms, Some(5000));
569 }
570
571 fn called_tool_with_meta(tool_name: &str, meta: serde_json::Value) -> SessionEvent {
574 let mut e = evt(EventType::AgentCalledTool {
575 tool_name: tool_name.into(),
576 tool_input_digest: None,
577 tool_output_digest: None,
578 duration_ms: None,
579 });
580 e.meta = Some(meta);
581 e
582 }
583
584 #[test]
585 fn hook_file_events_carry_source_hook() {
586 let events = vec![
590 evt(EventType::AgentReadFile {
591 file_path: "src/a.rs".into(),
592 digest: None,
593 }),
594 evt(EventType::AgentWroteFile {
595 file_path: "src/b.rs".into(),
596 digest: None,
597 operation: None,
598 additions: None,
599 deletions: None,
600 }),
601 ];
602 let se = SideEffects::from_events(&events);
603 assert_eq!(se.files_read[0].source.as_deref(), Some("hook"));
604 assert_eq!(se.files_written[0].source.as_deref(), Some("hook"));
605 }
606
607 #[test]
608 fn mcp_called_tool_with_file_path_promotes_to_files_written() {
609 let events = vec![called_tool_with_meta(
615 "Edit",
616 serde_json::json!({
617 "source": "mcp-bridge",
618 "tool_input": { "file_path": "src/api/receipt.ts" },
619 }),
620 )];
621 let se = SideEffects::from_events(&events);
622 assert_eq!(
623 se.files_written.len(),
624 1,
625 "Edit with file_path must promote to files_written"
626 );
627 assert_eq!(se.files_written[0].file_path, "src/api/receipt.ts");
628 assert_eq!(se.files_written[0].source.as_deref(), Some("mcp"));
629 assert_eq!(se.tool_invocations.len(), 1);
632 }
633
634 #[test]
635 fn mcp_read_tool_promotes_to_files_read() {
636 let events = vec![called_tool_with_meta(
637 "Read",
638 serde_json::json!({ "tool_input": { "file_path": "package.json" } }),
639 )];
640 let se = SideEffects::from_events(&events);
641 assert_eq!(se.files_read.len(), 1);
642 assert_eq!(se.files_read[0].file_path, "package.json");
643 assert_eq!(se.files_read[0].source.as_deref(), Some("mcp"));
644 }
645
646 #[test]
647 fn mcp_bash_tool_promotes_to_processes() {
648 let events = vec![called_tool_with_meta(
649 "Bash",
650 serde_json::json!({ "tool_input": { "command": "bun test --run" } }),
651 )];
652 let se = SideEffects::from_events(&events);
653 assert_eq!(se.processes.len(), 1);
654 assert_eq!(se.processes[0].command.as_deref(), Some("bun test --run"));
655 assert_eq!(se.processes[0].source.as_deref(), Some("mcp"));
656 }
657
658 #[test]
659 fn mcp_unknown_tool_with_path_defaults_to_files_written() {
660 let events = vec![called_tool_with_meta(
665 "mcp__weird-vendor__do_thing",
666 serde_json::json!({ "tool_input": { "file_path": "config.toml" } }),
667 )];
668 let se = SideEffects::from_events(&events);
669 assert_eq!(se.files_written.len(), 1);
670 assert_eq!(se.files_written[0].file_path, "config.toml");
671 assert_eq!(se.files_written[0].source.as_deref(), Some("mcp"));
672 }
673
674 #[test]
675 fn mcp_called_tool_without_meta_does_not_promote() {
676 let events = vec![called_tool_with_meta(
681 "ls",
682 serde_json::json!({"source": "mcp-bridge"}),
683 )];
684 let se = SideEffects::from_events(&events);
685 assert_eq!(se.files_read.len(), 0);
686 assert_eq!(se.files_written.len(), 0);
687 assert_eq!(se.processes.len(), 0);
688 assert_eq!(se.tool_invocations.len(), 1);
689 }
690
691 #[test]
692 fn mcp_promotion_handles_alt_path_field_names() {
693 for path_field in &[
696 "tool_input.path",
697 "tool_input.target_file",
698 "file_path",
699 "path",
700 ] {
701 let mut meta_obj = serde_json::Map::new();
702 let parts: Vec<&str> = path_field.split('.').collect();
704 if parts.len() == 1 {
705 meta_obj.insert(parts[0].into(), serde_json::json!("x.txt"));
706 } else {
707 let inner = serde_json::json!({ parts[1]: "x.txt" });
708 meta_obj.insert(parts[0].into(), inner);
709 }
710 let events = vec![called_tool_with_meta(
711 "Edit",
712 serde_json::Value::Object(meta_obj),
713 )];
714 let se = SideEffects::from_events(&events);
715 assert_eq!(
716 se.files_written.len(),
717 1,
718 "expected promotion via {} but got nothing",
719 path_field,
720 );
721 }
722 }
723
724 fn evt_with_meta(et: EventType, meta: serde_json::Value) -> SessionEvent {
734 let mut e = evt(et);
735 e.meta = Some(meta);
736 e
737 }
738
739 #[test]
740 fn source_from_meta_preserves_session_event_cli() {
741 let events = vec![evt_with_meta(
745 EventType::AgentWroteFile {
746 file_path: "src/x.rs".into(),
747 digest: None,
748 operation: None,
749 additions: None,
750 deletions: None,
751 },
752 serde_json::json!({"source": "session-event-cli"}),
753 )];
754 let se = SideEffects::from_events(&events);
755 assert_eq!(
756 se.files_written[0].source.as_deref(),
757 Some("session-event-cli"),
758 "session-event-cli must be preserved verbatim, not downgraded to hook"
759 );
760 }
761
762 #[test]
763 fn source_from_meta_preserves_daemon_atime() {
764 let events = vec![evt_with_meta(
768 EventType::AgentReadFile {
769 file_path: "src/x.rs".into(),
770 digest: None,
771 },
772 serde_json::json!({"source": "daemon-atime"}),
773 )];
774 let se = SideEffects::from_events(&events);
775 assert_eq!(
776 se.files_read[0].source.as_deref(),
777 Some("daemon-atime"),
778 "daemon-atime must be preserved verbatim, not downgraded to hook"
779 );
780 }
781
782 #[test]
783 fn source_from_meta_preserves_arbitrary_unknown_label() {
784 let events = vec![evt_with_meta(
789 EventType::AgentWroteFile {
790 file_path: "x".into(),
791 digest: None,
792 operation: None,
793 additions: None,
794 deletions: None,
795 },
796 serde_json::json!({"source": "future-bridge-v2"}),
797 )];
798 let se = SideEffects::from_events(&events);
799 assert_eq!(
800 se.files_written[0].source.as_deref(),
801 Some("future-bridge-v2")
802 );
803 }
804
805 #[test]
806 fn source_from_meta_falls_back_when_meta_source_absent() {
807 let events = vec![evt(EventType::AgentWroteFile {
813 file_path: "x".into(),
814 digest: None,
815 operation: None,
816 additions: None,
817 deletions: None,
818 })];
819 let se = SideEffects::from_events(&events);
820 assert_eq!(se.files_written[0].source.as_deref(), Some("hook"));
821 }
822
823 #[test]
824 fn source_from_meta_treats_empty_string_as_absent() {
825 let events = vec![evt_with_meta(
828 EventType::AgentReadFile {
829 file_path: "x".into(),
830 digest: None,
831 },
832 serde_json::json!({"source": ""}),
833 )];
834 let se = SideEffects::from_events(&events);
835 assert_eq!(se.files_read[0].source.as_deref(), Some("hook"));
836 }
837}