framework-cqrs-lib 0.4.0

handle state-machine with data persist in journal and store mongo for restfull actix api
Documentation
use std::sync::Arc;

use uuid::Uuid;

use crate::cqrs::core::context::Context;
use crate::cqrs::core::data::{Entity, EntityEvent};
use crate::cqrs::core::event_sourcing::CommandHandler;
use crate::cqrs::core::reducer::Reducer;
use crate::cqrs::core::repositories::entities::RepositoryEntity;
use crate::cqrs::core::repositories::events::RepositoryEvents;
use crate::cqrs::models::errors::{Error, ResultErr};

pub struct Engine<STATE: Clone, COMMAND, EVENT> {
    pub handlers: Vec<CommandHandler<STATE, COMMAND, EVENT>>,
    pub reducer: Reducer<EVENT, STATE>,
    pub store: Arc<dyn RepositoryEntity<STATE, String>>,
    pub journal: Arc<dyn RepositoryEvents<EVENT, String>>,
}

impl<STATE, COMMAND, EVENT> Engine<STATE, COMMAND, EVENT>
where
    STATE: Clone,
    EVENT: Clone,
{
    pub async fn compute(&self, command: COMMAND, entity_id: String, name: String, context: &Context) -> ResultErr<(EntityEvent<EVENT, String>, Entity<STATE, String>)> {
        let command_handler_found = self.handlers
            .iter().find(|handler| {
            match handler {
                CommandHandler::Create(created) => created.name() == name,
                CommandHandler::Update(updated) => updated.name() == name
            }
        })
            .ok_or(Error::Simple("pas de handler pour cette commande".to_string()))?;

        let maybe_entity = self.store.fetch_one(&entity_id).await?;
        let maybe_state = maybe_entity.clone().map(|entity| entity.data);

        let event = match command_handler_found {
            CommandHandler::Create(x) => x.on_command(entity_id.clone(), command, context).await,
            CommandHandler::Update(x) => {
                let state = maybe_state.clone().ok_or(Error::Simple("resource not found".to_string()))?;

                x.on_command(entity_id.clone(), state, command, context).await
            }
        }?;

        let new_state = (self.reducer.compute_new_state)(maybe_state, event.clone())
            .ok_or(Error::Simple("transition etat impossible".to_string()))?;
        let version = maybe_entity.clone()
            .map(|x| x.version.unwrap_or(0));
        let new_entity = Entity {
            entity_id: entity_id.clone(),
            data: new_state.clone(),
            version,
        };

        if maybe_entity.is_none() {
            self.store.insert(&new_entity).await?;
        } else {
            self.store.update(&entity_id, &new_entity).await?;
        }

        let event_entity = EntityEvent {
            entity_id: entity_id.clone(),
            event_id: Self::generate_id(),
            data: event.clone(),
        };
        self.journal.insert(&event_entity).await?;
        Ok((event_entity, new_entity))
    }

    fn generate_id() -> String {
        Uuid::new_v4().to_string()
    }
}