agentsight_capture_core/runners/
agent.rs1use super::{EventStream, Runner, RunnerError};
5use crate::analyzers::Analyzer;
6use async_trait::async_trait;
7use futures::stream::select_all;
8
9#[derive(Default)]
12pub struct AgentRunner {
13 runners: Vec<Box<dyn Runner>>,
14 analyzers: Vec<Box<dyn Analyzer>>,
15}
16
17impl AgentRunner {
18 pub fn new() -> Self {
19 Self::default()
20 }
21
22 pub fn add_runner(mut self, runner: Box<dyn Runner>) -> Self {
23 self.runners.push(runner);
24 self
25 }
26
27 pub fn add_global_analyzer(mut self, analyzer: Box<dyn Analyzer>) -> Self {
28 self.analyzers.push(analyzer);
29 self
30 }
31
32 pub fn runner_count(&self) -> usize {
33 self.runners.len()
34 }
35
36 pub fn analyzer_count(&self) -> usize {
37 self.analyzers.len()
38 }
39}
40
41#[async_trait]
42impl Runner for AgentRunner {
43 async fn run(&mut self) -> Result<EventStream, RunnerError> {
44 if self.runners.is_empty() {
45 return Err("No runners configured for AgentRunner".into());
46 }
47
48 let mut streams = Vec::new();
49 for runner in &mut self.runners {
50 streams.push(runner.run().await?);
51 }
52 let mut stream = Box::pin(select_all(streams)) as EventStream;
53 for analyzer in &mut self.analyzers {
54 stream = analyzer
55 .process(stream)
56 .await
57 .map_err(|error| format!("Global analyzer error: {error}"))?;
58 }
59 Ok(stream)
60 }
61
62 fn add_analyzer(mut self, analyzer: Box<dyn Analyzer>) -> Self {
63 self.analyzers.push(analyzer);
64 self
65 }
66}
67
68#[cfg(test)]
69mod tests {
70 use super::*;
71 use crate::analyzers::AnalyzerError;
72 use crate::event::Event;
73 use futures::{StreamExt, stream};
74
75 struct MemoryRunner(Vec<Event>);
76
77 #[async_trait]
78 impl Runner for MemoryRunner {
79 async fn run(&mut self) -> Result<EventStream, RunnerError> {
80 Ok(Box::pin(stream::iter(std::mem::take(&mut self.0))))
81 }
82
83 fn add_analyzer(self, _analyzer: Box<dyn Analyzer>) -> Self {
84 self
85 }
86 }
87
88 struct CountAnalyzer;
89
90 #[async_trait]
91 impl Analyzer for CountAnalyzer {
92 async fn process(&mut self, stream: EventStream) -> Result<EventStream, AnalyzerError> {
93 Ok(Box::pin(stream.map(|mut event| {
94 event.data["seen"] = serde_json::json!(true);
95 event
96 })))
97 }
98 }
99
100 fn event(pid: u32) -> Event {
101 Event::new_with_timestamp(1, "test".into(), pid, "agent".into(), serde_json::json!({}))
102 }
103
104 #[tokio::test]
105 async fn merges_runners_and_applies_stable_analyzer_boundary() {
106 let mut runner = AgentRunner::new()
107 .add_runner(Box::new(MemoryRunner(vec![event(1)])))
108 .add_runner(Box::new(MemoryRunner(vec![event(2)])))
109 .add_global_analyzer(Box::new(CountAnalyzer));
110
111 let events: Vec<_> = runner.run().await.unwrap().collect().await;
112 assert_eq!(events.len(), 2);
113 assert!(events.iter().all(|event| event.data["seen"] == true));
114 }
115
116 #[tokio::test]
117 async fn rejects_empty_runner_set_and_propagates_runner_errors() {
118 assert!(AgentRunner::new().run().await.is_err());
119
120 struct Failing;
121 #[async_trait]
122 impl Runner for Failing {
123 async fn run(&mut self) -> Result<EventStream, RunnerError> {
124 Err("runner failed".into())
125 }
126 fn add_analyzer(self, _analyzer: Box<dyn Analyzer>) -> Self { self }
127 }
128 assert!(AgentRunner::new().add_runner(Box::new(Failing)).run().await.is_err());
129 }
130}