use std::sync::Arc;
use anyhow::Result;
use tracing::{info, warn};
use dashmap::DashMap;
use chrono::Utc;
use kotoba_ocel::OcelEvent;
use kotoba_storage::KeyValueStore;
use crate::EventEnvelope;
pub struct Materializer<T: KeyValueStore> {
storage: Arc<T>,
prefix: String,
active_tasks: Arc<DashMap<String, MaterializationTask>>,
}
#[derive(Debug, Clone)]
struct MaterializationTask {
id: String,
projection_name: String,
status: TaskStatus,
start_time: chrono::DateTime<chrono::Utc>,
}
#[derive(Debug, Clone)]
enum TaskStatus {
Running,
Completed,
Failed(String),
}
#[derive(Debug, Clone)]
pub struct MaterializationResult {
pub projection_name: String,
pub events_processed: u64,
pub processing_time_ms: u64,
pub errors: Vec<String>,
}
impl<T: KeyValueStore + 'static> Materializer<T> {
pub fn new(
storage: Arc<T>,
prefix: String,
) -> Self {
Self {
storage,
prefix,
active_tasks: Arc::new(DashMap::new()),
}
}
pub async fn process_ocel_event(&self, ocel_event: OcelEvent) -> Result<()> {
info!("Processing OCEL event: {} ({})", ocel_event.id, ocel_event.activity);
let event_key = format!("{}:event:{}", self.prefix, ocel_event.id);
let event_data = serde_json::json!({
"id": ocel_event.id,
"activity": ocel_event.activity,
"timestamp": ocel_event.timestamp,
"omap": ocel_event.omap,
"vmap": ocel_event.vmap
});
let serialized_data = serde_json::to_vec(&event_data)?;
self.storage.put(event_key.as_bytes(), &serialized_data).await?;
info!("Successfully processed OCEL event: {}", ocel_event.id);
Ok(())
}
pub async fn process_event(&self, event_type: String, event: EventEnvelope) -> Result<()> {
warn!("Legacy event processing is deprecated. Use process_ocel_event instead.");
Ok(())
}
}