agentsight_capture/analyzers/
materializing.rs1use crate::analyzers::{Analyzer, AnalyzerError};
5use crate::model::ViewSink;
6use crate::runners::EventStream;
7use crate::view::SharedMaterializedView;
8use async_trait::async_trait;
9use futures::stream::StreamExt;
10
11pub struct MaterializingAnalyzer {
12 view: SharedMaterializedView,
13}
14
15impl MaterializingAnalyzer {
16 pub fn with_view(view: SharedMaterializedView) -> Self {
17 Self { view }
18 }
19
20 pub fn add_view_sink(self, sink: Box<dyn ViewSink>) -> Self {
21 if let Ok(mut view) = self.view.lock() {
22 view.add_sink(sink);
23 } else {
24 log::warn!("MaterializingAnalyzer: failed to acquire view lock while adding sink");
25 }
26 self
27 }
28}
29
30#[async_trait]
31impl Analyzer for MaterializingAnalyzer {
32 async fn process(&mut self, stream: EventStream) -> Result<EventStream, AnalyzerError> {
33 let view = self.view.clone();
34
35 let processed = stream.map(move |event| {
36 if let Ok(mut view) = view.lock() {
37 if let Err(error) = view.ingest_event(&event) {
38 log::warn!("MaterializingAnalyzer: failed to ingest event: {}", error);
39 }
40 } else {
41 log::warn!(
42 "MaterializingAnalyzer: failed to acquire view lock while ingesting event"
43 );
44 }
45 event
46 });
47
48 Ok(Box::pin(processed))
49 }
50}