use std::sync::Arc;
use async_trait::async_trait;
use diesel::{OptionalExtension, QueryDsl, RunQueryDsl};
use diesel::prelude::*;
use email_address::EmailAddress;
use fistinc_diesel_utils::async_pool;
use fistinc_diesel_utils::database_pool::DBConnectionPool;
use fistinc_errors::RepositoryResult;
use fistinc_paging::{Paging, QueryParams};
use phonenumber::PhoneNumber;
use uuid::Uuid;
use crate::domain::entity::User;
use crate::domain::UserRepository;
use crate::infrastructure::models::UserDiesel;
pub struct UserRepositoryImpl {
pool: Arc<DBConnectionPool>,
}
impl UserRepositoryImpl {
pub fn new(pool: Arc<DBConnectionPool>) -> Self {
Self {
pool
}
}
async fn total(&self) -> RepositoryResult<i64> {
use crate::infrastructure::schema::users::dsl::users;
async_pool::run(self.pool.clone(), move |conn| {
users.count().get_result::<i64>(&conn)
}).await
}
async fn fetch(&self, query: &dyn QueryParams) -> RepositoryResult<Vec<User>> {
use crate::infrastructure::schema::users::dsl::users;
let builder = users
.limit(query.limit())
.offset(query.offset());
let result: Vec<_> =
async_pool::run(self.pool.clone(), move |conn| {
builder.load::<UserDiesel>(&conn)
}).await?;
Ok(result.into_iter().map(|v| -> User { v.into() }).collect())
}
fn map_option_user(value: Option<UserDiesel>) -> Option<User> {
value.map(|v| -> User { v.into() })
}
}
#[async_trait]
impl UserRepository for UserRepositoryImpl {
async fn find_all(&self, params: &dyn QueryParams) -> RepositoryResult<Paging<User>> {
let total = self.total();
let users = self.fetch(params);
Ok(Paging {
total: total.await?,
items: users.await?,
})
}
async fn find(&self, user_id: Uuid) -> RepositoryResult<Option<User>> {
use crate::infrastructure::schema::users::dsl::{id, users};
let id_filter = user_id.to_string();
async_pool::run(self.pool.clone(), move |conn| {
users
.filter(id.eq(id_filter))
.first::<UserDiesel>(&conn)
.optional()
}).await
.map(Self::map_option_user)
}
async fn find_by_email(&self, user_email: &EmailAddress) -> RepositoryResult<Option<User>> {
use crate::infrastructure::schema::users::dsl::{email, users};
let email_filter = user_email.to_string();
async_pool::run(self.pool.clone(), move |conn| {
users
.filter(email.eq(email_filter))
.first::<UserDiesel>(&conn)
.optional()
}).await
.map(Self::map_option_user)
}
async fn find_by_phone(&self, user_phone: &PhoneNumber) -> RepositoryResult<Option<User>> {
use crate::infrastructure::schema::users::dsl::{phone, users};
let phone_filter = user_phone.to_string();
async_pool::run(self.pool.clone(), move |conn| {
users
.filter(phone.eq(phone_filter))
.first::<UserDiesel>(&conn)
.optional()
}).await
.map(Self::map_option_user)
}
async fn create(&self, user: &User) -> RepositoryResult<()> {
use crate::infrastructure::schema::users::dsl::users;
let user_diesel = UserDiesel::from(user.clone());
async_pool::run(self.pool.clone(), move |conn| {
diesel::insert_into(users)
.values(user_diesel)
.execute(&conn)
}).await?;
Ok(())
}
async fn update(&self, user: &User) -> RepositoryResult<()> {
use crate::infrastructure::schema::users::dsl::{id, users};
let id_filter = user.id.to_string();
let user_diesel = UserDiesel::from(user.clone());
async_pool::run(self.pool.clone(), move |conn| {
diesel::update(users)
.filter(id.eq(id_filter))
.set(user_diesel)
.execute(&conn)
}).await?;
Ok(())
}
async fn delete(&self, user_id: Uuid) -> RepositoryResult<()> {
use crate::infrastructure::schema::users::dsl::{id, users};
let id_filter = user_id.to_string();
async_pool::run(self.pool.clone(), move |conn| {
diesel::delete(users)
.filter(id.eq(id_filter))
.execute(&conn)
}).await?;
Ok(())
}
}