Skip to main content

agentsight_capture_core/runners/
agent.rs

1// SPDX-License-Identifier: MIT
2// Copyright (c) 2026 eunomia-bpf org.
3
4use super::{EventStream, Runner, RunnerError};
5use crate::analyzers::Analyzer;
6use async_trait::async_trait;
7use futures::stream::select_all;
8
9/// AgentRunner composes multiple runners into a single unified stream
10/// with optional global analyzers applied to the merged stream.
11#[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}