use async_trait::async_trait; use sqlx::{Executor, Postgres, query, query_as}; use crate::{ core::{ models::{ unit::UnitId, user::{NewUser, User, UserId}, }, repositories::{RepositoryError, users_repository::UsersRepository}, }, services::database::SqlxDatabase, }; struct UserDB { pub id: i32, pub external_id: Option, pub firstname: String, pub name: String, pub email: String, pub oidc_sub: String, pub admin: bool, } impl UserDB { fn into_user(self, units: Vec) -> User { User { id: self.id, external_id: self.external_id, firstname: self.firstname, name: self.name, email: self.email, oidc_sub: self.oidc_sub, admin: self.admin, units: units.into_iter().map(|u| u.name).collect(), } } } struct UnitIdDB { pub name: UnitId, } impl SqlxDatabase { async fn user_with_units<'a, E>(user: UserDB, executor: E) -> Result where E: Executor<'a, Database = Postgres>, { let units = query_as!( UnitIdDB, r#"SELECT unit_name AS "name!" FROM units_users WHERE user_id = $1 ORDER BY unit_name"#, user.id ) .fetch_all(executor) .await?; Ok(user.into_user(units)) } } #[async_trait] impl UsersRepository for SqlxDatabase { async fn get_user(&self, id: UserId) -> Result { let mut tx = self.pool.begin().await?; let user = query_as!(UserDB, r#"SELECT * FROM users WHERE id = $1"#, id) .fetch_one(&mut *tx) .await?; let user = Self::user_with_units(user, &mut *tx).await?; tx.commit().await?; Ok(user) } async fn get_user_email(&self, email: String) -> Result { let mut tx = self.pool.begin().await?; let user = query_as!(UserDB, r#"SELECT * FROM users WHERE email = $1"#, email) .fetch_one(&mut *tx) .await?; let user = Self::user_with_units(user, &mut *tx).await?; tx.commit().await?; Ok(user) } async fn get_user_external_id(&self, external_id: String) -> Result { let mut tx = self.pool.begin().await?; let user = query_as!( UserDB, r#"SELECT * FROM users WHERE external_id = $1"#, external_id ) .fetch_one(&mut *tx) .await?; let user = Self::user_with_units(user, &mut *tx).await?; tx.commit().await?; Ok(user) } async fn get_user_oidc_sub(&self, oidc_sub: String) -> Result { let mut tx = self.pool.begin().await?; let user = query_as!( UserDB, r#"SELECT * FROM users WHERE oidc_sub = $1"#, oidc_sub ) .fetch_one(&mut *tx) .await?; let user = Self::user_with_units(user, &mut *tx).await?; tx.commit().await?; Ok(user) } async fn upsert_user(&self, user: NewUser) -> Result { let mut tx = self.pool.begin().await?; let user_db = query_as!( UserDB, r#"INSERT INTO users (external_id, firstname, "name", email, oidc_sub) VALUES ($1, $2, $3, $4, $5) ON CONFLICT (oidc_sub) DO UPDATE SET external_id = EXCLUDED.external_id, firstname = EXCLUDED.firstname, "name" = EXCLUDED.name, email = EXCLUDED.email RETURNING *"#, user.external_id, user.firstname, user.name, user.email, user.oidc_sub ) .fetch_one(&mut *tx) .await?; query!(r#"DELETE FROM units_users WHERE user_id = $1"#, user_db.id) .execute(&mut *tx) .await?; query!( r#"INSERT INTO units_users (user_id, unit_name) SELECT $1, UNNEST($2::text[])"#, user_db.id, &user.units ) .execute(&mut *tx) .await?; tx.commit().await?; Ok(User { id: user_db.id, external_id: user_db.external_id, firstname: user_db.firstname, name: user_db.name, email: user_db.email, oidc_sub: user_db.oidc_sub, units: user.units, admin: user_db.admin, }) } async fn set_user_admin(&self, id: UserId, admin: bool) -> Result<(), RepositoryError> { let result = query!(r#"UPDATE users SET admin = $2 WHERE id = $1"#, id, admin) .execute(&self.pool) .await?; if result.rows_affected() == 0 { return Err(RepositoryError::NotFound(format!("user {id}"))); } Ok(()) } }