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;
}