framework-cqrs-lib 0.5.4

handle state-machine with data persist in journal and store mongo for restfull actix api
Documentation
use crate::cqrs::core::daos::DAO;
use crate::cqrs::core::data::Entity;
use crate::cqrs::core::repositories::can_fetch_all::CanFetchAll;
use crate::cqrs::core::repositories::entities::{ReadOnlyEntityRepo, RepositoryEntity, WriteOnlyEntityRepo};
use crate::cqrs::core::repositories::query::Query;
use crate::cqrs::core::repositories::CanFetchMany;
use crate::cqrs::infra::daos::dbos::EntityDBO;
use crate::cqrs::models::errors::ResultErr;
use async_trait::async_trait;
use futures::lock::Mutex;
use std::sync::Arc;

pub struct MongoEntityRepository<DBO> {
    pub dao: Arc<Mutex<dyn DAO<EntityDBO<DBO, String>, String>>>,
}

#[async_trait]
impl<
    DATA: CanTransform<DBO> + Clone + Sync + Send,
    DBO: CanTransform<DATA> + Clone + Sync + Send
> ReadOnlyEntityRepo<DATA, String> for MongoEntityRepository<DBO> {
    async fn fetch_one(&self, id: &String) -> ResultErr<Option<Entity<DATA, String>>> {
        self.dao
            .lock().await
            .fetch_one(id).await
            .map(|maybedata| maybedata.map(|dbo| {
                Entity {
                    entity_id: dbo.entity_id,
                    data: dbo.data.transform_into_other(),
                    version: dbo.version,
                }
            }))
    }
}

#[async_trait]
impl<
    DATA: CanTransform<DBO> + Clone + Sync + Send,
    DBO: CanTransform<DATA> + Clone + Sync + Send
> CanFetchAll<Entity<DATA, String>> for MongoEntityRepository<DBO> {
    async fn fetch_all(&self, query: Query) -> ResultErr<Vec<Entity<DATA, String>>> {
        self.dao
            .lock().await
            .fetch_all(query)
            .await
            .map(|items| {
                items
                    .into_iter()
                    .map(|dbo| Entity {
                        entity_id: dbo.entity_id,
                        data: dbo.data.transform_into_other(),
                        version: dbo.version,
                    })
                    .collect()
            })
    }
}

#[async_trait]
impl<
    DATA: CanTransform<DBO> + Clone + Sync + Send,
    DBO: CanTransform<DATA> + Clone + Sync + Send
> CanFetchMany<Entity<DATA, String>> for MongoEntityRepository<DBO> {}

#[async_trait]
impl<
    DATA: CanTransform<DBO> + Clone + Sync + Send,
    DBO: CanTransform<DATA> + Clone + Sync + Send
> WriteOnlyEntityRepo<DATA, String> for MongoEntityRepository<DBO> {
    async fn insert(&self, entity: &Entity<DATA, String>) -> ResultErr<String> {
        let entity_dbo: EntityDBO<DBO, String> = EntityDBO {
            id_mongo: None,
            version: entity.version.clone(),
            entity_id: entity.entity_id.clone(),
            data: entity.data.transform_into_other(),
        };

        let sanitize_version: EntityDBO<DBO, String> = EntityDBO {
            version: Some(0),
            ..entity_dbo.clone()
        };

        self.dao
            .lock().await
            .insert(&sanitize_version, &entity_dbo.entity_id).await
    }

    async fn update(&self, id: &String, entity: &Entity<DATA, String>) -> ResultErr<String> {
        let entity_dbo: EntityDBO<DBO, String> = EntityDBO {
            id_mongo: None,
            version: entity.version.clone(),
            entity_id: entity.entity_id.clone(),
            data: entity.data.transform_into_other(),
        };
        let sanitize_version: EntityDBO<DBO, String> = EntityDBO {
            version: entity_dbo.version.map(|current| current + 1),
            ..entity_dbo
        };

        self.dao
            .lock().await
            .update(id, &sanitize_version).await
    }
}

#[async_trait]
impl<
    DATA: CanTransform<DBO> + Clone + Sync + Send,
    DBO: CanTransform<DATA> + Clone + Sync + Send
> RepositoryEntity<DATA, String> for MongoEntityRepository<DBO> {}

pub trait CanTransform<B> {
    fn transform_into_other(&self) -> B;
}