use std::collections::HashMap;
use std::sync::Arc;
use tokio::sync::{mpsc, RwLock};
use dashmap::DashMap;
use anyhow::{Result, Context};
use tracing::{info, warn, error, instrument};
use futures::stream::{self, StreamExt};
use metrics::{counter, histogram};
use kotoba_ocel::OcelEvent;
use kotoba_storage::KeyValueStore;
use crate::materializer::Materializer;
pub struct EventProcessor<T: KeyValueStore> {
materializer: Arc<Materializer<T>>,
projections: Arc<DashMap<String, ProjectionHandler>>,
event_queue: mpsc::Sender<Vec<OcelEvent>>,
batch_size: usize,
stats: Arc<RwLock<ProcessorStats>>,
}
struct ProjectionHandler {
name: String,
event_types: Vec<String>,
processor: Box<dyn Fn(EventEnvelope) -> Result<()> + Send + Sync>,
}
#[derive(Debug, Clone)]
pub struct ProcessorStats {
pub events_received: u64,
pub events_processed: u64,
pub processing_errors: u64,
pub avg_processing_time_ms: f64,
pub queue_size: usize,
}
impl<T: KeyValueStore + 'static> EventProcessor<T> {
pub fn new(materializer: Arc<Materializer<T>>, batch_size: usize) -> Self {
let (tx, rx) = mpsc::channel(1000);
let materializer_clone = materializer.clone();
let processor = Self {
materializer,
projections: Arc::new(DashMap::new()),
event_queue: tx,
batch_size,
stats: Arc::new(RwLock::new(ProcessorStats::default())),
};
tokio::spawn(Self::process_events_task(processor.stats.clone(), rx, materializer_clone));
processor
}
pub async fn start(&self) -> Result<()> {
info!("Starting Event Processor");
Ok(())
}
pub async fn stop(&self) -> Result<()> {
info!("Stopping Event Processor");
drop(self.event_queue.clone());
Ok(())
}
#[instrument(skip(self))]
pub async fn register_projection(&self, name: &str) -> Result<()> {
info!("Registering projection: {}", name);
let handler = ProjectionHandler {
name: name.to_string(),
event_types: vec![
"node.created".to_string(),
"node.updated".to_string(),
"node.deleted".to_string(),
"edge.created".to_string(),
"edge.updated".to_string(),
"edge.deleted".to_string(),
],
processor: Box::new(move |event: EventEnvelope| {
Ok(())
}),
};
self.projections.insert(name.to_string(), handler);
info!("Projection registered: {}", name);
Ok(())
}
#[instrument(skip(self))]
pub async fn unregister_projection(&self, name: &str) -> Result<()> {
info!("Unregistering projection: {}", name);
self.projections.remove(name);
info!("Projection unregistered: {}", name);
Ok(())
}
#[instrument(skip(self, events))]
pub async fn process_batch(&self, events: Vec<OcelEvent>) -> Result<()> {
{
let mut stats = self.stats.write().await;
stats.events_received += events.len() as u64;
stats.queue_size += events.len();
}
if let Err(e) = self.event_queue.send(events).await {
error!("Failed to send events to processing queue: {}", e);
return Err(anyhow::anyhow!("Failed to queue events: {}", e));
}
Ok(())
}
pub async fn get_statistics(&self) -> ProcessorStats {
self.stats.read().await.clone()
}
async fn process_events_task(
stats: Arc<RwLock<ProcessorStats>>,
mut rx: mpsc::Receiver<Vec<OcelEvent>>,
materializer: Arc<Materializer<T>>,
) {
info!("Starting OCEL event processing task");
while let Some(events) = rx.recv().await {
let start_time = std::time::Instant::now();
let batch_size = events.len();
let result = Self::process_ocel_event_batch(&materializer, events).await;
let processing_time = start_time.elapsed();
let mut stats = stats.write().await;
stats.events_processed += batch_size as u64;
stats.queue_size = stats.queue_size.saturating_sub(batch_size);
if let Err(e) = result {
stats.processing_errors += 1;
error!("Error processing OCEL event batch: {}", e);
} else {
let total_time_ms = processing_time.as_millis() as f64;
let avg_time = (stats.avg_processing_time_ms + total_time_ms / batch_size as f64) / 2.0;
stats.avg_processing_time_ms = avg_time;
if batch_size > 0 {
}
}
}
info!("OCEL event processing task stopped");
}
async fn process_ocel_event_batch(
materializer: &Arc<Materializer<T>>,
events: Vec<OcelEvent>,
) -> Result<()> {
for event in events {
Self::route_ocel_event(materializer, event).await?;
}
Ok(())
}
async fn route_ocel_event(materializer: &Arc<Materializer<T>>, event: OcelEvent) -> Result<()> {
materializer.process_ocel_event(event).await?;
Ok(())
}
}
impl Default for ProcessorStats {
fn default() -> Self {
Self {
events_received: 0,
events_processed: 0,
processing_errors: 0,
avg_processing_time_ms: 0.0,
queue_size: 0,
}
}
}
pub type EventEnvelope = serde_json::Value;
#[cfg(test)]
mod tests {
use super::*;
use crate::materializer::Materializer;
#[tokio::test]
async fn test_event_processor_creation() {
let materializer = Arc::new(Materializer::default());
let processor = EventProcessor::new(materializer, 10);
let stats = processor.get_statistics().await;
assert_eq!(stats.events_received, 0);
assert_eq!(stats.events_processed, 0);
}
#[tokio::test]
async fn test_projection_registration() {
let materializer = Arc::new(Materializer::default());
let processor = EventProcessor::new(materializer, 10);
let result = processor.register_projection("test_projection").await;
assert!(result.is_ok(), "Projection should be registered successfully");
let result = processor.unregister_projection("test_projection").await;
assert!(result.is_ok(), "Projection should be unregistered successfully");
}
#[tokio::test]
async fn test_event_processing() {
let materializer = Arc::new(Materializer::default());
let processor = EventProcessor::new(materializer, 10);
let event = serde_json::json!({
"event_type": "node.created",
"aggregate_id": "test-123",
"data": {
"name": "Test Node"
}
});
let result = processor.process_batch(vec![event]).await;
assert!(result.is_ok(), "Event batch should be processed successfully");
let stats = processor.get_statistics().await;
assert_eq!(stats.events_received, 1);
}
}