use crate::event::Event;
use crate::event_bus::EventHandler;
use crate::event_store::EventStore;
use crate::read_model_store::ReadModelStore;
use async_trait::async_trait;
use std::sync::Arc;
use thiserror::Error;
#[derive(Debug, Error)]
pub enum ProjectionError {
#[error(
"Projection '{name}' failed to handle event {event_type} (event_id: {event_id}): {message}"
)]
HandleFailed {
name: String,
event_type: String,
event_id: String,
message: String,
},
#[error("Projection '{name}' failed to clear: {message}")]
ClearFailed { name: String, message: String },
#[error("Projection '{name}' failed to rebuild: {message}")]
RebuildFailed { name: String, message: String },
#[error("Failed to load events for rebuild: {0}")]
EventStoreError(String),
#[error("Read model store error in projection '{name}': {message}")]
ReadModelError { name: String, message: String },
#[error("Projection error: {message}")]
Other { message: String },
}
impl ProjectionError {
pub fn handle_failed(
name: impl Into<String>,
event_type: impl Into<String>,
event_id: impl Into<String>,
message: impl Into<String>,
) -> Self {
ProjectionError::HandleFailed {
name: name.into(),
event_type: event_type.into(),
event_id: event_id.into(),
message: message.into(),
}
}
pub fn clear_failed(name: impl Into<String>, message: impl Into<String>) -> Self {
ProjectionError::ClearFailed {
name: name.into(),
message: message.into(),
}
}
pub fn rebuild_failed(name: impl Into<String>, message: impl Into<String>) -> Self {
ProjectionError::RebuildFailed {
name: name.into(),
message: message.into(),
}
}
pub fn read_model_error(name: impl Into<String>, message: impl Into<String>) -> Self {
ProjectionError::ReadModelError {
name: name.into(),
message: message.into(),
}
}
pub fn other(message: impl Into<String>) -> Self {
ProjectionError::Other {
message: message.into(),
}
}
}
pub type ProjectionResult<T> = Result<T, ProjectionError>;
#[async_trait]
pub trait Projector: Send + Sync {
fn name(&self) -> &str;
fn handles(&self) -> Vec<String>;
async fn apply(&self, event: &Event, store: &dyn ReadModelStore) -> ProjectionResult<()>;
async fn init(&self, _store: &dyn ReadModelStore) -> ProjectionResult<()> {
Ok(())
}
}
#[async_trait]
pub trait Projection: Send + Sync {
fn name(&self) -> &str;
fn handles(&self) -> Vec<String>;
async fn handle(&self, event: &Event) -> ProjectionResult<()>;
async fn clear(&self) -> ProjectionResult<()>;
async fn rebuild(&self, events: Vec<Event>) -> ProjectionResult<()> {
self.clear().await?;
for event in events {
if self.handles().contains(&event.event_type) {
self.handle(&event).await?;
}
}
Ok(())
}
}
pub struct ProjectionUnit {
projector: Box<dyn Projector>,
store: Arc<dyn ReadModelStore>,
table: String,
}
impl ProjectionUnit {
pub fn new(
projector: Box<dyn Projector>,
store: Arc<dyn ReadModelStore>,
table: impl Into<String>,
) -> Self {
Self {
projector,
store,
table: table.into(),
}
}
}
#[async_trait]
impl Projection for ProjectionUnit {
fn name(&self) -> &str {
self.projector.name()
}
fn handles(&self) -> Vec<String> {
self.projector.handles()
}
async fn handle(&self, event: &Event) -> ProjectionResult<()> {
self.projector.apply(event, self.store.as_ref()).await
}
async fn clear(&self) -> ProjectionResult<()> {
self.store
.truncate(&self.table)
.await
.map_err(|e| ProjectionError::clear_failed(self.projector.name(), e.to_string()))
}
}
pub struct ProjectionEngine {
projections: Vec<Box<dyn Projection>>,
event_store: Box<dyn EventStore>,
}
impl ProjectionEngine {
pub fn new(event_store: Box<dyn EventStore>) -> Self {
Self {
projections: Vec::new(),
event_store,
}
}
pub fn register(&mut self, projection: Box<dyn Projection>) {
tracing::info!("Registering projection: {}", projection.name());
self.projections.push(projection);
}
pub fn register_projector(
&mut self,
projector: Box<dyn Projector>,
store: Arc<dyn ReadModelStore>,
table: impl Into<String>,
) {
let unit = ProjectionUnit::new(projector, store, table);
self.register(Box::new(unit));
}
pub async fn process(&self, event: &Event) -> ProjectionResult<()> {
for projection in &self.projections {
if projection.handles().contains(&event.event_type) {
tracing::debug!(
"Processing event {} ({}) in projection {}",
event.event_type,
event.event_id,
projection.name()
);
projection.handle(event).await.map_err(|e| {
ProjectionError::handle_failed(
projection.name(),
&event.event_type,
event.event_id.to_string(),
e.to_string(),
)
})?;
}
}
Ok(())
}
pub async fn process_batch(&self, events: Vec<Event>) -> ProjectionResult<()> {
for event in events {
self.process(&event).await?;
}
Ok(())
}
pub async fn rebuild_all(&self) -> ProjectionResult<()> {
tracing::info!("Rebuilding all projections");
let events = self
.event_store
.stream_all(0)
.await
.map_err(|e| ProjectionError::EventStoreError(e.to_string()))?;
tracing::info!("Loaded {} events for rebuild", events.len());
for projection in &self.projections {
tracing::info!("Rebuilding projection: {}", projection.name());
projection
.rebuild(events.clone())
.await
.map_err(|e| ProjectionError::rebuild_failed(projection.name(), e.to_string()))?;
tracing::info!("Rebuilt projection: {}", projection.name());
}
Ok(())
}
pub async fn rebuild_projection(&self, name: &str) -> ProjectionResult<()> {
tracing::info!("Rebuilding projection: {}", name);
let projection = self
.projections
.iter()
.find(|p| p.name() == name)
.ok_or_else(|| ProjectionError::other(format!("Projection not found: {}", name)))?;
let events = self
.event_store
.stream_all(0)
.await
.map_err(|e| ProjectionError::EventStoreError(e.to_string()))?;
projection
.rebuild(events)
.await
.map_err(|e| ProjectionError::rebuild_failed(name, e.to_string()))?;
tracing::info!("Rebuilt projection: {}", name);
Ok(())
}
pub fn projection_count(&self) -> usize {
self.projections.len()
}
pub fn projection_names(&self) -> Vec<String> {
self.projections
.iter()
.map(|p| p.name().to_string())
.collect()
}
pub fn all_handled_event_types(&self) -> Vec<String> {
let mut all: Vec<String> = self.projections.iter().flat_map(|p| p.handles()).collect();
all.sort();
all.dedup();
all
}
}
pub struct ProjectionEngineHandler {
engine: Arc<ProjectionEngine>,
handles: Vec<String>,
}
impl ProjectionEngineHandler {
pub fn new(engine: Arc<ProjectionEngine>) -> Self {
let handles = engine.all_handled_event_types();
Self { engine, handles }
}
}
#[async_trait]
impl EventHandler for ProjectionEngineHandler {
fn handles(&self) -> Vec<String> {
self.handles.clone()
}
async fn handle(&self, event: &Event) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
self.engine
.process(event)
.await
.map_err(|e| -> Box<dyn std::error::Error + Send + Sync> { Box::new(e) })
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::event_store::{EventStore, EventStoreResult, VersionCheck};
use crate::read_model_store::InMemoryReadModelStore;
use std::sync::{Arc, Mutex};
struct MockEventStore {
events: Arc<Mutex<Vec<Event>>>,
}
impl MockEventStore {
fn new() -> Self {
Self {
events: Arc::new(Mutex::new(Vec::new())),
}
}
fn add_event(&self, event: Event) {
self.events.lock().unwrap().push(event);
}
}
#[async_trait]
impl EventStore for MockEventStore {
async fn append(
&self,
_aggregate_id: &str,
_version_check: VersionCheck,
events: Vec<Event>,
) -> EventStoreResult<()> {
self.events.lock().unwrap().extend(events);
Ok(())
}
async fn load(&self, aggregate_id: &str) -> EventStoreResult<Vec<Event>> {
Ok(self
.events
.lock()
.unwrap()
.iter()
.filter(|e| e.aggregate_id == aggregate_id)
.cloned()
.collect())
}
async fn load_from(
&self,
aggregate_id: &str,
from_sequence: i64,
) -> EventStoreResult<Vec<Event>> {
Ok(self
.events
.lock()
.unwrap()
.iter()
.filter(|e| e.aggregate_id == aggregate_id && e.sequence >= from_sequence)
.cloned()
.collect())
}
async fn stream_all(&self, _from_position: i64) -> EventStoreResult<Vec<Event>> {
Ok(self.events.lock().unwrap().clone())
}
async fn get_version(&self, aggregate_id: &str) -> EventStoreResult<i64> {
Ok(self
.events
.lock()
.unwrap()
.iter()
.filter(|e| e.aggregate_id == aggregate_id)
.map(|e| e.sequence)
.max()
.unwrap_or(0))
}
}
struct MockProjector {
name: String,
handles_types: Vec<String>,
}
impl MockProjector {
fn new(name: &str, handles: Vec<String>) -> Self {
Self {
name: name.to_string(),
handles_types: handles,
}
}
}
#[async_trait]
impl Projector for MockProjector {
fn name(&self) -> &str {
&self.name
}
fn handles(&self) -> Vec<String> {
self.handles_types.clone()
}
async fn apply(&self, event: &Event, store: &dyn ReadModelStore) -> ProjectionResult<()> {
use crate::read_model_store::Upsert;
store
.upsert(Upsert::new(
"test_table",
event.event_id.to_string(),
serde_json::json!({
"id": event.event_id.to_string(),
"event_type": event.event_type,
"version": event.sequence,
}),
))
.await
.map_err(|e| {
ProjectionError::handle_failed(
&self.name,
&event.event_type,
event.event_id.to_string(),
e.to_string(),
)
})?;
Ok(())
}
}
fn make_projection(
name: &str,
handles: Vec<String>,
store: Arc<InMemoryReadModelStore>,
) -> Box<ProjectionUnit> {
Box::new(ProjectionUnit::new(
Box::new(MockProjector::new(name, handles)),
store,
"test_table",
))
}
#[tokio::test]
async fn test_projection_engine_new() {
let store = Box::new(MockEventStore::new());
let engine = ProjectionEngine::new(store);
assert_eq!(engine.projection_count(), 0);
}
#[tokio::test]
async fn test_register_projection() {
let store = Box::new(MockEventStore::new());
let mut engine = ProjectionEngine::new(store);
let rm_store = Arc::new(InMemoryReadModelStore::new());
let projection = make_projection("Test", vec!["TestEvent".to_string()], rm_store);
engine.register(projection);
assert_eq!(engine.projection_count(), 1);
assert_eq!(engine.projection_names(), vec!["Test"]);
}
#[tokio::test]
async fn test_process_event() {
let store = Box::new(MockEventStore::new());
let mut engine = ProjectionEngine::new(store);
let rm_store = Arc::new(InMemoryReadModelStore::new());
let projection = make_projection("Test", vec!["UserCreated".to_string()], rm_store.clone());
engine.register(projection);
let event = Event::new(
"User",
"user-1",
1,
"UserCreated",
serde_json::json!({"name": "Alice"}),
);
engine.process(&event).await.unwrap();
assert_eq!(rm_store.get_rows("test_table").len(), 1);
}
#[tokio::test]
async fn test_projection_filtering() {
let store = Box::new(MockEventStore::new());
let mut engine = ProjectionEngine::new(store);
let rm_store = Arc::new(InMemoryReadModelStore::new());
let projection = make_projection("Test", vec!["UserCreated".to_string()], rm_store.clone());
engine.register(projection);
let event1 = Event::new("User", "user-1", 1, "UserCreated", serde_json::json!({}));
engine.process(&event1).await.unwrap();
let event2 = Event::new("User", "user-1", 2, "UserDeleted", serde_json::json!({}));
engine.process(&event2).await.unwrap();
assert_eq!(rm_store.get_rows("test_table").len(), 1);
}
#[tokio::test]
async fn test_rebuild_all() {
let event_store = MockEventStore::new();
event_store.add_event(Event::new(
"User",
"user-1",
1,
"UserCreated",
serde_json::json!({}),
));
event_store.add_event(Event::new(
"User",
"user-2",
1,
"UserCreated",
serde_json::json!({}),
));
let mut engine = ProjectionEngine::new(Box::new(event_store));
let rm_store = Arc::new(InMemoryReadModelStore::new());
let projection = make_projection("Test", vec!["UserCreated".to_string()], rm_store.clone());
engine.register(projection);
engine.rebuild_all().await.unwrap();
assert_eq!(rm_store.get_rows("test_table").len(), 2);
}
#[tokio::test]
async fn test_multiple_projections() {
let store = Box::new(MockEventStore::new());
let mut engine = ProjectionEngine::new(store);
let rm_store1 = Arc::new(InMemoryReadModelStore::new());
let rm_store2 = Arc::new(InMemoryReadModelStore::new());
let proj1 = make_projection(
"Projection1",
vec!["UserCreated".to_string()],
rm_store1.clone(),
);
let proj2 = make_projection(
"Projection2",
vec!["UserCreated".to_string(), "UserDeleted".to_string()],
rm_store2.clone(),
);
engine.register(proj1);
engine.register(proj2);
let event = Event::new("User", "user-1", 1, "UserCreated", serde_json::json!({}));
engine.process(&event).await.unwrap();
assert_eq!(rm_store1.get_rows("test_table").len(), 1);
assert_eq!(rm_store2.get_rows("test_table").len(), 1);
}
#[tokio::test]
async fn test_process_batch() {
let store = Box::new(MockEventStore::new());
let mut engine = ProjectionEngine::new(store);
let rm_store = Arc::new(InMemoryReadModelStore::new());
let projection = make_projection("Test", vec!["UserCreated".to_string()], rm_store.clone());
engine.register(projection);
let events = vec![
Event::new("User", "user-1", 1, "UserCreated", serde_json::json!({})),
Event::new("User", "user-2", 1, "UserCreated", serde_json::json!({})),
Event::new("User", "user-3", 1, "UserCreated", serde_json::json!({})),
];
engine.process_batch(events).await.unwrap();
assert_eq!(rm_store.get_rows("test_table").len(), 3);
}
#[tokio::test]
async fn test_projector_init_default() {
let projector = MockProjector::new("Test", vec![]);
let store = InMemoryReadModelStore::new();
projector.init(&store).await.unwrap();
}
#[tokio::test]
async fn test_register_projector_convenience() {
let event_store = Box::new(MockEventStore::new());
let mut engine = ProjectionEngine::new(event_store);
let rm_store: Arc<dyn ReadModelStore> = Arc::new(InMemoryReadModelStore::new());
engine.register_projector(
Box::new(MockProjector::new("Convenient", vec!["X".to_string()])),
rm_store,
"my_table",
);
assert_eq!(engine.projection_count(), 1);
assert_eq!(engine.projection_names(), vec!["Convenient"]);
}
}