diff --git a/.junie/plans/refactor-write-services-and-read-repositories.md b/.junie/plans/refactor-write-services-and-read-repositories.md new file mode 100644 index 0000000..5dffac4 --- /dev/null +++ b/.junie/plans/refactor-write-services-and-read-repositories.md @@ -0,0 +1,110 @@ +--- +sessionId: session-260729-174324-a4dc +--- + +# Requirements + +### Overview & Goals +Currently, write operations (mutations) and `EventBus` emissions take place inside `Repositories`. This design poses two main challenges: +1. **Lack of Transactional Atomicity**: Multi-table writes (e.g., creating a channel while simultaneously inserting its order record into `server_item_order`) cannot share an atomic database transaction. +2. **Premature / Ghost Events**: Events are emitted inside Repositories before confirming whether higher-level multi-step operations or surrounding DB transactions committed successfully. + +To solve this, we are refactoring to a lightweight **Command/Query separation**: +- **Services (Commands / Write Operations)**: Perform write operations, manage SeaORM database transactions (`db.begin()`), and emit `EventBus` events strictly **after** transaction commits. +- **Repositories (Queries / Read Operations)**: Focus on complex reads, queries, and filters. Event bus emissions are completely removed from Repositories. + +### Scope +- **In Scope**: + - Removing `events: Arc` from `RepositoryContext` and `Repositories::new`. + - Stripping all `.events.emit(...)` calls from `CategoryRepository`, `ChannelRepository`, `MessageRepository`, `RoleRepository`, `ServerRepository`, `ServerItemOrderRepository`, and `UserRepository`. + - Creating/expanding dedicated domain services under `src/services/` (`ChannelService`, `ServerService`, `CategoryService`, `MessageService`, `UserService`, `RoleService`) with SeaORM transaction support and post-commit event emissions. + - Registering all domain services in `src/services/mod.rs` and `Services`. + - Updating Axum HTTP route handlers in `src/routes/` to delegate write operations (POST, PUT, DELETE) to `state.services` while keeping read operations (GET) on `state.repositories`. +- **Out of Scope**: + - Changing API DTO contracts or client-facing response schemas. + - Modifying underlying SeaORM database entities or table schemas. + +### Functional Requirements +- **FR1**: Repositories must be strictly read/query-focused and contain zero event emissions or `EventBus` references. +- **FR2**: Write mutations (create, update, delete) and permission management must be executed inside domain Services. +- **FR3**: Multi-step writes (such as creating a server/channel and updating `server_item_order`) must execute within an atomic SeaORM transaction (`db.begin().await?`). +- **FR4**: `EventBus` events must only be emitted after the database transaction successfully commits. +- **FR5**: Axum route handlers must invoke service methods for all state-mutating requests (POST, PUT, DELETE) and repository methods for read requests (GET). + +# Technical Design + +### Current Implementation +- `RepositoryContext` in `src/repositories/mod.rs` holds both `db: DatabaseConnection` and `events: Arc`. +- Repositories (`ChannelRepository`, `ServerRepository`, `CategoryRepository`, `MessageRepository`, `UserRepository`, `RoleRepository`, `ServerItemOrderRepository`) execute `active.insert()`, `active.update()`, and delete operations directly and emit events immediately inside repository methods. +- Axum route handlers in `src/routes/*/handlers.rs` call repository write methods directly. + +### Key Decisions +1. **Command / Query Responsibility Segregation**: + - Repositories handle data access, queries, filters, and read models. + - Services handle business logic, transactional boundaries (`db.begin()`), and event dispatching. +2. **Post-Commit Event Emission**: + - Events are only emitted after `txn.commit().await?` succeeds, preventing ghost/premature events on transaction rollback. + +### Proposed Changes & Affected Files +1. **`src/repositories/mod.rs` & Repository Modules**: + - Modify `RepositoryContext` to remove `events: Arc`. + - Remove `.events.emit(...)` calls from: + - `src/repositories/category.rs` + - `src/repositories/channel.rs` + - `src/repositories/message.rs` + - `src/repositories/role.rs` + - `src/repositories/server.rs` + - `src/repositories/server_item_order.rs` + - `src/repositories/user.rs` +2. **`src/services/` Modules**: + - Expand `src/services/` with new service files: + - `channel.rs` (`ChannelService`) + - `server.rs` (`ServerService`) + - `category.rs` (`CategoryService`) + - `message.rs` (`MessageService`) + - `user.rs` (`UserService`) + - `role.rs` (`RoleService`) +3. **`src/services/mod.rs`**: + - Update `Services` struct and `Services::new` to initialize and expose all domain services. +4. **`src/routes/` Handlers**: + - Update write handlers across `src/routes/{channel,server,category,message,user,role}/handlers.rs` to invoke `state.services.*`. + +### Architecture Diagram +```mermaid +graph LR + HTTP[Axum Handlers] -->|Write mutations| Services[Domain Services] + HTTP[Read requests] -->|Query/Filter| Repositories[Read Repositories] + Services -->|Transaction & DB mutations| DB[(SeaORM Database)] + Services -->|Post-commit emit| Events[EventBus] + Repositories -->|Read query| DB +``` + +# Delivery Steps + +### ✓ Step 1: Clean up Repositories and Remove Write Event Emissions +Clean up Repositories and RepositoryContext +- Remove `events: Arc` from `RepositoryContext` in `src/repositories/mod.rs` and update `Repositories::new`. +- Strip write event emissions (`self.context.events.emit(...)`) from all repositories (`CategoryRepository`, `ChannelRepository`, `MessageRepository`, `RoleRepository`, `ServerRepository`, `ServerItemOrderRepository`, `UserRepository`). +- Ensure repository write methods operate purely on database connections/ActiveModels without triggering event bus emissions. + +### ✓ Step 2: Create Domain Services with Transactional Atomicity and Post-Commit Events +Create Domain Services for Write Operations +- Create domain services in `src/services/` for channels, servers, categories, messages, users, and roles (e.g., `ChannelService`, `ServerService`, `CategoryService`, `MessageService`, `UserService`, `RoleService`). +- Implement SeaORM transaction management (`db.begin().await?`) in service write methods. +- Ensure `EventBus` events are emitted strictly **after** transaction commits. +- Handle multi-table transactional writes such as creating a channel while inserting into `server_item_order`. + +### ✓ Step 3: Register Services in Services Container +Register Services in Services Container and Update State/Context +- Update `src/services/mod.rs` to include and initialize the new services (`channel`, `server`, `category`, `message`, `user`, `role`) within the `Services` struct and `ServicesContext`. +- Expose the updated `Services` container via `AppState` / `ServicesContext`. + +### ✓ Step 4: Refactor Axum HTTP Route Handlers to Use Services +Refactor Axum HTTP Route Handlers +- Update write endpoints (POST, PUT, DELETE) across `src/routes/` to invoke service methods on `state.services` instead of repositories directly. +- Keep read endpoints (GET) using repositories for complex queries, filters, and tree generation. + +### ✓ Step 5: Verification and Testing +Verification and Testing +- Run cargo build/check to ensure clean compilation across all modules. +- Verify integration checks: channel creation populates `server_item_order`, events are only triggered upon successful DB transaction commit. \ No newline at end of file diff --git a/src/core/mod.rs b/src/core/mod.rs index c31da1e..6f42a52 100644 --- a/src/core/mod.rs +++ b/src/core/mod.rs @@ -32,7 +32,7 @@ impl App { let event_bus = Arc::new(EventBus::with_capacity(1024)); // Initialize shared repositories - let repositories = Arc::new(Repositories::new(db.clone(), event_bus.clone())); + let repositories = Arc::new(Repositories::new(db.clone())); // Initialize gateway manager let gateway = Arc::new(GatewayManager::default()); diff --git a/src/repositories/category.rs b/src/repositories/category.rs index 2bfbb27..661de40 100644 --- a/src/repositories/category.rs +++ b/src/repositories/category.rs @@ -37,17 +37,11 @@ impl CategoryRepository { pub async fn update(&self, active: category::ActiveModel) -> AnyResult { let category = active.update(&self.context.db).await?; - self.context - .events - .emit("category_updated", category.clone()); Ok(category) } pub async fn create(&self, active: category::ActiveModel) -> AnyResult { let category = active.insert(&self.context.db).await?; - self.context - .events - .emit("category_created", category.clone()); Ok(category) } @@ -55,7 +49,6 @@ impl CategoryRepository { let res = category::Entity::delete_by_id(id) .exec(&self.context.db) .await?; - self.context.events.emit("category_deleted", id); Ok(res.rows_affected > 0) } } diff --git a/src/repositories/channel.rs b/src/repositories/channel.rs index b63630a..07f8919 100644 --- a/src/repositories/channel.rs +++ b/src/repositories/channel.rs @@ -32,13 +32,11 @@ impl ChannelRepository { pub async fn update(&self, active: channel::ActiveModel) -> AnyResult { let channel = active.update(&self.context.db).await?; - self.context.events.emit("channel_updated", channel.clone()); Ok(channel) } pub async fn create(&self, active: channel::ActiveModel) -> AnyResult { let channel = active.insert(&self.context.db).await?; - self.context.events.emit("channel_created", channel.clone()); Ok(channel) } @@ -46,7 +44,6 @@ impl ChannelRepository { let res = channel::Entity::delete_by_id(id) .exec(&self.context.db) .await?; - self.context.events.emit("channel_deleted", id); Ok(res.rows_affected > 0) } diff --git a/src/repositories/message.rs b/src/repositories/message.rs index aacf878..1a0ef14 100644 --- a/src/repositories/message.rs +++ b/src/repositories/message.rs @@ -48,38 +48,11 @@ impl MessageRepository { pub async fn update(&self, active: message::ActiveModel) -> AnyResult { let message = active.update(&self.context.db).await?; - self.context.events.emit("message_updated", message.clone()); Ok(message) } pub async fn create(&self, active: message::ActiveModel) -> AnyResult { let message = active.insert(&self.context.db).await?; - - // self.context.events.emit("message_created", message.clone()); - - // todo : test - // Ici l'évènement est déclencher sur les topic suivant : - // message_created - // channel:_channel_uuid_:message_created - // si server : server:_server_uuid_:message_created - // scoped event - let mut scopes: Vec = Vec::new(); - scopes.push(Scope::uuid("channel", message.channel_id)); - // retrieve related channel and server - let server_id: Option = channel::Entity::find_by_id(message.channel_id) - .select_only() - .column(channel::Column::ServerId) - .into_tuple::>() - .one(&self.context.db) - .await? - .flatten(); - - if let Some(server_id) = server_id { - scopes.push(Scope::uuid("server", server_id)); - } - self.context - .events - .emit_scoped("message_created", scopes, message.clone()); Ok(message) } @@ -88,9 +61,6 @@ impl MessageRepository { .exec(&self.context.db) .await?; let deleted = result.rows_affected > 0; - if deleted { - self.context.events.emit("message_deleted", id); - } Ok(deleted) } } diff --git a/src/repositories/mod.rs b/src/repositories/mod.rs index ac599b0..f329348 100644 --- a/src/repositories/mod.rs +++ b/src/repositories/mod.rs @@ -24,8 +24,7 @@ mod user; #[derive(Clone, Debug)] pub struct RepositoryContext { - db: DatabaseConnection, - events: Arc, + pub db: DatabaseConnection, } #[derive(Clone, Debug)] @@ -41,8 +40,8 @@ pub struct Repositories { } impl Repositories { - pub fn new(db: DatabaseConnection, events: Arc) -> Self { - let context = Arc::new(RepositoryContext { db, events }); + pub fn new(db: DatabaseConnection) -> Self { + let context = Arc::new(RepositoryContext { db }); Self { server: ServerRepository { diff --git a/src/repositories/role.rs b/src/repositories/role.rs index eb3efbe..5f2c866 100644 --- a/src/repositories/role.rs +++ b/src/repositories/role.rs @@ -49,13 +49,11 @@ impl RoleRepository { pub async fn create(&self, active: role::ActiveModel) -> AnyResult { let group = active.insert(&self.context.db).await?; - self.context.events.emit("group_created", group.clone()); Ok(group) } pub async fn update(&self, active: role::ActiveModel) -> AnyResult { let group = active.update(&self.context.db).await?; - self.context.events.emit("group_updated", group.clone()); Ok(group) } @@ -63,7 +61,6 @@ impl RoleRepository { let res = role::Entity::delete_by_id(id) .exec(&self.context.db) .await?; - self.context.events.emit("group_deleted", id); Ok(res.rows_affected > 0) } } diff --git a/src/repositories/server.rs b/src/repositories/server.rs index 2963ad4..b7fbb9d 100644 --- a/src/repositories/server.rs +++ b/src/repositories/server.rs @@ -31,7 +31,6 @@ impl ServerRepository { pub async fn update(&self, active: server::ActiveModel) -> AnyResult { let server = active.update(&self.context.db).await?; - self.context.events.emit("server_updated", server.clone()); Ok(server) } @@ -47,7 +46,6 @@ impl ServerRepository { }; default_group.insert(&self.context.db).await?; - self.context.events.emit("server_created", server.clone()); Ok(server) } @@ -83,9 +81,6 @@ impl ServerRepository { .await? .ok_or_else(|| anyhow::anyhow!("Rôle par défaut introuvable"))?; - self.context - .events - .emit("server_user_created", (server_id, user_id)); Ok(true) } @@ -93,7 +88,6 @@ impl ServerRepository { let res = server::Entity::delete_by_id(id) .exec(&self.context.db) .await?; - self.context.events.emit("server_deleted", id); Ok(res.rows_affected > 0) } diff --git a/src/repositories/server_item_order.rs b/src/repositories/server_item_order.rs index 884b144..7751a48 100644 --- a/src/repositories/server_item_order.rs +++ b/src/repositories/server_item_order.rs @@ -40,7 +40,6 @@ impl ServerItemOrderRepository { }) .await?; - self.context.events.emit("server_order_updated", server_id); Ok(()) } } diff --git a/src/repositories/user.rs b/src/repositories/user.rs index 20ed70d..a714f57 100644 --- a/src/repositories/user.rs +++ b/src/repositories/user.rs @@ -53,13 +53,11 @@ impl UserRepository { pub async fn update(&self, active: user::ActiveModel) -> AnyResult { let user = active.update(&self.context.db).await?; - self.context.events.emit("user_updated", user.clone()); Ok(user) } pub async fn create(&self, active: user::ActiveModel) -> AnyResult { let user = active.insert(&self.context.db).await?; - self.context.events.emit("user_created", user.clone()); Ok(user) } @@ -80,9 +78,8 @@ impl UserRepository { active.password = Set(password); - let user = self.update(active).await?; + let _user = self.update(active).await?; - self.context.events.emit("user_changed", user); Ok(()) } @@ -91,9 +88,6 @@ impl UserRepository { .exec(&self.context.db) .await?; let deleted = result.rows_affected > 0; - if deleted { - self.context.events.emit("user_deleted", id); - } Ok(deleted) } diff --git a/src/routes/category/handlers.rs b/src/routes/category/handlers.rs index 5dadd2a..6a7f954 100644 --- a/src/routes/category/handlers.rs +++ b/src/routes/category/handlers.rs @@ -94,8 +94,7 @@ pub async fn create( .await? .ok_or(HTTPError::BadRequest("Server not found".to_string()))?; - let active_model = mapper::create_request_to_am(payload); - let category = state.repositories.category.create(active_model).await?; + let category = state.services.category.create_category(payload.server_id, payload.name).await?; Ok(( StatusCode::CREATED, Json(mapper::category_model_to_category_response(category)), @@ -127,15 +126,14 @@ pub async fn update( Json(payload): Json, ) -> Result, HTTPError> { // Vérifier l'existence - let category = state + let _category = state .repositories .category .get_by_id(id) .await? .ok_or(HTTPError::NotFound)?; - let active_model = mapper::update_request_to_am(category.id, category.server_id, payload); - let category = state.repositories.category.update(active_model).await?; + let category = state.services.category.update_category(id, payload.name).await?; Ok(Json(mapper::category_model_to_category_response(category))) } @@ -162,7 +160,7 @@ pub async fn delete( State(state): State, Path(id): Path, ) -> Result { - if state.repositories.category.delete(id).await? { + if state.services.category.delete_category(id).await? { Ok(StatusCode::NO_CONTENT) } else { Err(HTTPError::NotFound) diff --git a/src/routes/channel/handlers.rs b/src/routes/channel/handlers.rs index af81e3e..d13b146 100644 --- a/src/routes/channel/handlers.rs +++ b/src/routes/channel/handlers.rs @@ -109,8 +109,7 @@ pub async fn create( .ok_or(HTTPError::BadRequest("Category not found".to_string()))?; } - let active_model = mapper::create_request_to_am(payload); - let channel = state.repositories.channel.create(active_model).await?; + let channel = state.services.channel.create_channel(payload).await?; Ok(( StatusCode::CREATED, Json(mapper::channel_model_to_channel_response(channel)), @@ -170,8 +169,7 @@ pub async fn update( .ok_or(HTTPError::BadRequest("Category not found".to_string()))?; } - let active_model = mapper::update_request_to_am(id, payload); - let channel = state.repositories.channel.update(active_model).await?; + let channel = state.services.channel.update_channel(id, payload).await?; Ok(Json(mapper::channel_model_to_channel_response(channel))) } @@ -198,7 +196,7 @@ pub async fn delete( State(state): State, Path(id): Path, ) -> Result { - if state.repositories.channel.delete(id).await? { + if state.services.channel.delete_channel(id).await? { Ok(StatusCode::NO_CONTENT) } else { Err(HTTPError::NotFound) @@ -257,7 +255,7 @@ pub async fn set_user_permission( Json(payload): Json, ) -> Result, HTTPError> { state - .repositories + .services .channel .set_user_permission(channel_id, user_id, payload.permissions) .await?; @@ -304,7 +302,7 @@ pub async fn remove_user_permission( } state - .repositories + .services .channel .remove_user_permission(channel_id, user_id) .await?; @@ -364,7 +362,7 @@ pub async fn set_role_permission( Json(payload): Json, ) -> Result, HTTPError> { state - .repositories + .services .channel .set_role_permission(channel_id, role_id, payload.permissions) .await?; @@ -411,7 +409,7 @@ pub async fn remove_role_permission( } state - .repositories + .services .channel .remove_role_permission(channel_id, role_id) .await?; diff --git a/src/routes/message/handlers.rs b/src/routes/message/handlers.rs index e5ecc81..9402ce2 100644 --- a/src/routes/message/handlers.rs +++ b/src/routes/message/handlers.rs @@ -105,8 +105,7 @@ pub async fn create( ))?; } - let active_model = mapper::create_request_to_am(user.id, payload); - let message = state.repositories.message.create(active_model).await?; + let message = state.services.message.create_message(payload.channel_id, user.id, payload.content).await?; Ok(( StatusCode::CREATED, Json(mapper::message_model_to_message_response(message)), @@ -151,8 +150,7 @@ pub async fn update( return Err(HTTPError::Forbidden); } - let active_model = mapper::update_request_to_am(message, payload); - let message = state.repositories.message.update(active_model).await?; + let message = state.services.message.update_message(id, payload.content).await?; Ok(Json(mapper::message_model_to_message_response(message))) } @@ -188,12 +186,11 @@ pub async fn delete( .await? .ok_or(HTTPError::NotFound)?; - // Autoriser si auteur ou superuser if message.user_id != user.id && !user.is_superuser { return Err(HTTPError::Forbidden); } - if state.repositories.message.delete(id).await? { + if state.services.message.delete_message(id).await? { Ok(StatusCode::NO_CONTENT) } else { Err(HTTPError::NotFound) diff --git a/src/routes/role/handlers.rs b/src/routes/role/handlers.rs index c454120..4aec081 100644 --- a/src/routes/role/handlers.rs +++ b/src/routes/role/handlers.rs @@ -87,7 +87,7 @@ pub async fn create( .ok_or(HTTPError::BadRequest("Server not found".to_string()))?; let active_model = mapper::create_request_to_am(payload); - let group = state.repositories.role.create(active_model).await?; + let group = state.services.role.create_role(active_model).await?; Ok(( StatusCode::CREATED, Json(mapper::group_model_to_group_response(group)), @@ -127,7 +127,7 @@ pub async fn update( .ok_or(HTTPError::NotFound)?; let active_model = mapper::update_request_to_am(group.id, group.server_id, payload); - let group = state.repositories.role.update(active_model).await?; + let group = state.services.role.update_role(active_model).await?; Ok(Json(mapper::group_model_to_group_response(group))) } @@ -154,7 +154,7 @@ pub async fn delete( State(state): State, Path(id): Path, ) -> Result { - if state.repositories.role.delete(id).await? { + if state.services.role.delete_role(id).await? { Ok(StatusCode::NO_CONTENT) } else { Err(HTTPError::NotFound) diff --git a/src/routes/server/handlers.rs b/src/routes/server/handlers.rs index 4016dff..08ed590 100644 --- a/src/routes/server/handlers.rs +++ b/src/routes/server/handlers.rs @@ -83,8 +83,7 @@ pub async fn create( State(state): State, Json(payload): Json, ) -> Result<(StatusCode, Json), HTTPError> { - let active_model = mapper::create_request_to_am(payload); - let server = state.repositories.server.create(active_model).await?; + let server = state.services.server.create_server(payload.name, payload.is_default).await?; Ok(( StatusCode::CREATED, Json(mapper::server_model_to_server_response(server)), @@ -116,15 +115,14 @@ pub async fn update( Json(payload): Json, ) -> Result, HTTPError> { // Vérifier l'existence - let server = state + let _server = state .repositories .server .get_by_id(id) .await? .ok_or(HTTPError::NotFound)?; - let active_model = mapper::update_request_to_am(server.id, payload); - let server = state.repositories.server.update(active_model).await?; + let server = state.services.server.update_server(id, payload.name, payload.is_default).await?; Ok(Json(mapper::server_model_to_server_response(server))) } @@ -151,7 +149,7 @@ pub async fn delete( State(state): State, Path(id): Path, ) -> Result { - if state.repositories.server.delete(id).await? { + if state.services.server.delete_server(id).await? { Ok(StatusCode::NO_CONTENT) } else { Err(HTTPError::NotFound) diff --git a/src/routes/user/handlers.rs b/src/routes/user/handlers.rs index 5eb23ae..f6af161 100644 --- a/src/routes/user/handlers.rs +++ b/src/routes/user/handlers.rs @@ -107,7 +107,7 @@ pub async fn create( let active_model = mapper::create_request_to_am(payload) .map_err(|e| HTTPError::InternalServerError(e.to_string()))?; - let user = state.repositories.user.create(active_model).await?; + let user = state.services.user.create_user(active_model).await?; Ok(( StatusCode::CREATED, Json(mapper::user_model_to_user_response(user)), @@ -160,7 +160,7 @@ pub async fn update( } let active_model = mapper::update_request_to_am(user.id, payload); - let user = state.repositories.user.update(active_model).await?; + let user = state.services.user.update_user(active_model).await?; Ok(Json(mapper::user_model_to_user_response(user))) } @@ -189,7 +189,7 @@ pub async fn delete( State(state): State, Path(id): Path, ) -> Result { - if state.repositories.user.delete(id).await? { + if state.services.user.delete_user(id).await? { Ok(StatusCode::NO_CONTENT) } else { Err(HTTPError::NotFound) diff --git a/src/services/category.rs b/src/services/category.rs new file mode 100644 index 0000000..af89a81 --- /dev/null +++ b/src/services/category.rs @@ -0,0 +1,88 @@ +use crate::services::ServicesContext; +use crate::models::category; +use sea_orm::{ActiveModelTrait, ColumnTrait, EntityTrait, QueryFilter, TransactionTrait, Set}; +use std::sync::Arc; +use uuid::Uuid; + +#[derive(Debug, Clone)] +pub struct CategoryService { + service_context: Arc, +} + +impl CategoryService { + pub fn new(service_context: Arc) -> Self { + Self { service_context } + } + + pub async fn create_category( + &self, + server_id: Uuid, + name: String, + ) -> Result { + let db = &self.service_context.repositories.server.context.db; + let event_bus = &self.service_context.event_bus; + + let txn = db.begin().await?; + + let active = category::ActiveModel { + server_id: Set(server_id), + name: Set(name), + ..Default::default() + }; + let cat = active.insert(&txn).await?; + + txn.commit().await?; + + event_bus.emit("category_created", cat.clone()); + + Ok(cat) + } + + pub async fn update_category( + &self, + id: Uuid, + name: String, + ) -> Result { + let db = &self.service_context.repositories.server.context.db; + let event_bus = &self.service_context.event_bus; + + let txn = db.begin().await?; + + let existing = category::Entity::find_by_id(id) + .one(&txn) + .await? + .ok_or_else(|| anyhow::anyhow!("Category not found"))?; + + let mut active: category::ActiveModel = existing.into(); + active.name = Set(name); + + let cat = active.update(&txn).await?; + + txn.commit().await?; + + event_bus.emit("category_updated", cat.clone()); + + Ok(cat) + } + + pub async fn delete_category(&self, id: Uuid) -> Result { + let db = &self.service_context.repositories.server.context.db; + let event_bus = &self.service_context.event_bus; + + let txn = db.begin().await?; + + let res = category::Entity::delete_by_id(id) + .exec(&txn) + .await?; + + let deleted = res.rows_affected > 0; + + txn.commit().await?; + + if deleted { + event_bus.emit("category_deleted", id); + } + + Ok(deleted) + } +} diff --git a/src/services/channel.rs b/src/services/channel.rs new file mode 100644 index 0000000..af91ffe --- /dev/null +++ b/src/services/channel.rs @@ -0,0 +1,263 @@ +use crate::repositories::Repositories; +use crate::services::ServicesContext; +use crate::domain::dto::channel::{CreateChannelRequest, UpdateChannelRequest}; +use crate::models::channel; +use event_bus::Scope; +use sea_orm::{ActiveModelTrait, ColumnTrait, EntityTrait, QueryFilter, QuerySelect, QueryOrder, TransactionTrait, Set, PaginatorTrait}; +use std::sync::Arc; +use uuid::Uuid; + +#[derive(Debug, Clone)] +pub struct ChannelService { + service_context: Arc, +} + +impl ChannelService { + pub fn new(service_context: Arc) -> Self { + Self { service_context } + } + + pub async fn create_channel( + &self, + payload: CreateChannelRequest, + ) -> Result { + let db = &self.service_context.repositories.server.context.db; + let event_bus = &self.service_context.event_bus; + + // Start transaction + let txn = db.begin().await?; + + // 1. Insert channel within transaction + let active_model = channel::ActiveModel { + server_id: Set(payload.server_id), + category_id: Set(payload.category_id), + channel_type: Set(payload.channel_type), + name: Set(payload.name), + ..Default::default() + }; + + let channel = active_model.insert(&txn).await?; + + // 2. If server_id is present, insert into server_item_order within transaction + if let Some(server_id) = payload.server_id { + // Get max order key or determine order key + let max_order: Option = crate::models::server_item_order::Entity::find() + .filter(crate::models::server_item_order::Column::ServerId.eq(server_id)) + .select_only() + .column_as(crate::models::server_item_order::Column::OrderKey.max(), "max_key") + .into_tuple::>() + .one(&txn) + .await? + .flatten(); + + let next_order = max_order.unwrap_or(0) + 1; + + let order_item = crate::models::server_item_order::ActiveModel { + server_id: Set(server_id), + resource_id: Set(channel.id), + resource_type: Set(crate::models::server_item_order::OrderedResourceType::Channel), + parent_category_id: Set(payload.category_id), + order_key: Set(next_order), + ..Default::default() + }; + order_item.insert(&txn).await?; + } + + // Commit transaction + txn.commit().await?; + + // Post-commit event emission + event_bus.emit("channel_created", channel.clone()); + + Ok(channel) + } + + pub async fn update_channel( + &self, + id: Uuid, + payload: UpdateChannelRequest, + ) -> Result { + let db = &self.service_context.repositories.server.context.db; + let event_bus = &self.service_context.event_bus; + + let txn = db.begin().await?; + + let existing = channel::Entity::find_by_id(id) + .one(&txn) + .await? + .ok_or_else(|| anyhow::anyhow!("Channel not found"))?; + + let mut active: channel::ActiveModel = existing.into(); + active.server_id = Set(payload.server_id); + active.category_id = Set(payload.category_id); + active.channel_type = Set(payload.channel_type); + active.name = Set(payload.name); + + let channel = active.update(&txn).await?; + + // Update server_item_order parent_category_id or server_id if needed + if let Some(server_id) = channel.server_id { + crate::models::server_item_order::Entity::update_many() + .set(crate::models::server_item_order::ActiveModel { + parent_category_id: Set(channel.category_id), + ..Default::default() + }) + .filter(crate::models::server_item_order::Column::ResourceId.eq(channel.id)) + .exec(&txn) + .await?; + } + + txn.commit().await?; + + event_bus.emit("channel_updated", channel.clone()); + + Ok(channel) + } + + pub async fn delete_channel(&self, id: Uuid) -> Result { + let db = &self.service_context.repositories.server.context.db; + let event_bus = &self.service_context.event_bus; + + let txn = db.begin().await?; + + // Delete associated order records + crate::models::server_item_order::Entity::delete_many() + .filter(crate::models::server_item_order::Column::ResourceId.eq(id)) + .exec(&txn) + .await?; + + let res = channel::Entity::delete_by_id(id) + .exec(&txn) + .await?; + + let deleted = res.rows_affected > 0; + + txn.commit().await?; + + if deleted { + event_bus.emit("channel_deleted", id); + } + + Ok(deleted) + } + + pub async fn set_user_permission( + &self, + channel_id: Uuid, + user_id: Uuid, + permissions: u64, + ) -> Result<(), anyhow::Error> { + let db = &self.service_context.repositories.server.context.db; + let event_bus = &self.service_context.event_bus; + + let txn = db.begin().await?; + + let permission = crate::models::channel_user_permission::ActiveModel { + channel_id: Set(channel_id), + user_id: Set(user_id), + permissions: Set(permissions as i64), + ..Default::default() + }; + + crate::models::channel_user_permission::Entity::insert(permission) + .on_conflict( + sea_orm::sea_query::OnConflict::columns([ + crate::models::channel_user_permission::Column::ChannelId, + crate::models::channel_user_permission::Column::UserId, + ]) + .update_columns([crate::models::channel_user_permission::Column::Permissions]) + .to_owned(), + ) + .exec(&txn) + .await?; + + txn.commit().await?; + + event_bus.emit("channel_user_permission_created", (channel_id, user_id, permissions)); + + Ok(()) + } + + pub async fn remove_user_permission( + &self, + channel_id: Uuid, + user_id: Uuid, + ) -> Result<(), anyhow::Error> { + let db = &self.service_context.repositories.server.context.db; + let event_bus = &self.service_context.event_bus; + + let txn = db.begin().await?; + + crate::models::channel_user_permission::Entity::delete_many() + .filter(crate::models::channel_user_permission::Column::ChannelId.eq(channel_id)) + .filter(crate::models::channel_user_permission::Column::UserId.eq(user_id)) + .exec(&txn) + .await?; + + txn.commit().await?; + + event_bus.emit("channel_user_permission_deleted", (channel_id, user_id)); + + Ok(()) + } + + pub async fn set_role_permission( + &self, + channel_id: Uuid, + role_id: Uuid, + permissions: u64, + ) -> Result<(), anyhow::Error> { + let db = &self.service_context.repositories.server.context.db; + let event_bus = &self.service_context.event_bus; + + let txn = db.begin().await?; + + let permission = crate::models::channel_role_permission::ActiveModel { + channel_id: Set(channel_id), + role_id: Set(role_id), + permissions: Set(permissions as i64), + ..Default::default() + }; + + crate::models::channel_role_permission::Entity::insert(permission) + .on_conflict( + sea_orm::sea_query::OnConflict::columns([ + crate::models::channel_role_permission::Column::ChannelId, + crate::models::channel_role_permission::Column::RoleId, + ]) + .update_columns([crate::models::channel_role_permission::Column::Permissions]) + .to_owned(), + ) + .exec(&txn) + .await?; + + txn.commit().await?; + + event_bus.emit("channel_role_permission_created", (channel_id, role_id, permissions)); + + Ok(()) + } + + pub async fn remove_role_permission( + &self, + channel_id: Uuid, + role_id: Uuid, + ) -> Result<(), anyhow::Error> { + let db = &self.service_context.repositories.server.context.db; + let event_bus = &self.service_context.event_bus; + + let txn = db.begin().await?; + + crate::models::channel_role_permission::Entity::delete_many() + .filter(crate::models::channel_role_permission::Column::ChannelId.eq(channel_id)) + .filter(crate::models::channel_role_permission::Column::RoleId.eq(role_id)) + .exec(&txn) + .await?; + + txn.commit().await?; + + event_bus.emit("channel_role_permission_deleted", (channel_id, role_id)); + + Ok(()) + } +} diff --git a/src/services/message.rs b/src/services/message.rs new file mode 100644 index 0000000..a88c811 --- /dev/null +++ b/src/services/message.rs @@ -0,0 +1,106 @@ +use crate::services::ServicesContext; +use crate::models::{channel, message}; +use event_bus::Scope; +use sea_orm::{ActiveModelTrait, ColumnTrait, EntityTrait, QueryFilter, QuerySelect, TransactionTrait, Set}; +use std::sync::Arc; +use uuid::Uuid; + +#[derive(Debug, Clone)] +pub struct MessageService { + service_context: Arc, +} + +impl MessageService { + pub fn new(service_context: Arc) -> Self { + Self { service_context } + } + + pub async fn create_message( + &self, + channel_id: Uuid, + author_id: Uuid, + content: String, + ) -> Result { + let db = &self.service_context.repositories.server.context.db; + let event_bus = &self.service_context.event_bus; + + let txn = db.begin().await?; + + let active = message::ActiveModel { + channel_id: Set(channel_id), + user_id: Set(author_id), + content: Set(content), + ..Default::default() + }; + let msg = active.insert(&txn).await?; + + txn.commit().await?; + + let mut scopes: Vec = Vec::new(); + scopes.push(Scope::uuid("channel", msg.channel_id)); + + let server_id: Option = channel::Entity::find_by_id(msg.channel_id) + .select_only() + .column(channel::Column::ServerId) + .into_tuple::>() + .one(db) + .await? + .flatten(); + + if let Some(server_id) = server_id { + scopes.push(Scope::uuid("server", server_id)); + } + + event_bus.emit_scoped("message_created", scopes, msg.clone()); + + Ok(msg) + } + + pub async fn update_message( + &self, + id: Uuid, + content: String, + ) -> Result { + let db = &self.service_context.repositories.server.context.db; + let event_bus = &self.service_context.event_bus; + + let txn = db.begin().await?; + + let existing = message::Entity::find_by_id(id) + .one(&txn) + .await? + .ok_or_else(|| anyhow::anyhow!("Message not found"))?; + + let mut active: message::ActiveModel = existing.into(); + active.content = Set(content); + + let msg = active.update(&txn).await?; + + txn.commit().await?; + + event_bus.emit("message_updated", msg.clone()); + + Ok(msg) + } + + pub async fn delete_message(&self, id: Uuid) -> Result { + let db = &self.service_context.repositories.server.context.db; + let event_bus = &self.service_context.event_bus; + + let txn = db.begin().await?; + + let res = message::Entity::delete_by_id(id) + .exec(&txn) + .await?; + + let deleted = res.rows_affected > 0; + + txn.commit().await?; + + if deleted { + event_bus.emit("message_deleted", id); + } + + Ok(deleted) + } +} diff --git a/src/services/mod.rs b/src/services/mod.rs index c3f7d61..7584e85 100644 --- a/src/services/mod.rs +++ b/src/services/mod.rs @@ -1,11 +1,23 @@ use crate::repositories::Repositories; use crate::services::permission_sync::PermissionSyncService; use crate::services::server_order::ServerOrderService; +use crate::services::channel::ChannelService; +use crate::services::server::ServerService; +use crate::services::category::CategoryService; +use crate::services::message::MessageService; +use crate::services::user::UserService; +use crate::services::role::RoleService; use event_bus::EventBus; use std::sync::{Arc, OnceLock}; pub mod permission_sync; mod server_order; +pub mod channel; +pub mod server; +pub mod category; +pub mod message; +pub mod user; +pub mod role; #[derive(Debug, Clone)] pub struct ServicesContext { @@ -18,6 +30,12 @@ pub struct ServicesContext { pub struct Services { pub permission_sync: Arc, pub server_order: Arc, + pub channel: Arc, + pub server: Arc, + pub category: Arc, + pub message: Arc, + pub user: Arc, + pub role: Arc, } impl Services { @@ -29,10 +47,22 @@ impl Services { }); let permission_sync = Arc::new(PermissionSyncService::new(service_context.clone())); let server_order = Arc::new(ServerOrderService::new(service_context.clone())); + let channel = Arc::new(ChannelService::new(service_context.clone())); + let server = Arc::new(ServerService::new(service_context.clone())); + let category = Arc::new(CategoryService::new(service_context.clone())); + let message = Arc::new(MessageService::new(service_context.clone())); + let user = Arc::new(UserService::new(service_context.clone())); + let role = Arc::new(RoleService::new(service_context.clone())); let services = Self { permission_sync, server_order, + channel, + server, + category, + message, + user, + role, }; let _ = service_context.services.set(services.clone()); services diff --git a/src/services/role.rs b/src/services/role.rs new file mode 100644 index 0000000..53fdfde --- /dev/null +++ b/src/services/role.rs @@ -0,0 +1,73 @@ +use crate::services::ServicesContext; +use crate::models::role; +use sea_orm::{ActiveModelTrait, ColumnTrait, EntityTrait, QueryFilter, TransactionTrait, Set}; +use std::sync::Arc; +use uuid::Uuid; + +#[derive(Debug, Clone)] +pub struct RoleService { + service_context: Arc, +} + +impl RoleService { + pub fn new(service_context: Arc) -> Self { + Self { service_context } + } + + pub async fn create_role( + &self, + active: role::ActiveModel, + ) -> Result { + let db = &self.service_context.repositories.server.context.db; + let event_bus = &self.service_context.event_bus; + + let txn = db.begin().await?; + + let group = active.insert(&txn).await?; + + txn.commit().await?; + + event_bus.emit("group_created", group.clone()); + + Ok(group) + } + + pub async fn update_role( + &self, + active: role::ActiveModel, + ) -> Result { + let db = &self.service_context.repositories.server.context.db; + let event_bus = &self.service_context.event_bus; + + let txn = db.begin().await?; + + let group = active.update(&txn).await?; + + txn.commit().await?; + + event_bus.emit("group_updated", group.clone()); + + Ok(group) + } + + pub async fn delete_role(&self, id: Uuid) -> Result { + let db = &self.service_context.repositories.server.context.db; + let event_bus = &self.service_context.event_bus; + + let txn = db.begin().await?; + + let res = role::Entity::delete_by_id(id) + .exec(&txn) + .await?; + + let deleted = res.rows_affected > 0; + + txn.commit().await?; + + if deleted { + event_bus.emit("group_deleted", id); + } + + Ok(deleted) + } +} diff --git a/src/services/server.rs b/src/services/server.rs new file mode 100644 index 0000000..b66ed32 --- /dev/null +++ b/src/services/server.rs @@ -0,0 +1,153 @@ +use crate::repositories::Repositories; +use crate::services::ServicesContext; +use crate::models::{role, server}; +use sea_orm::{ActiveModelTrait, ColumnTrait, EntityTrait, QueryFilter, QuerySelect, QueryOrder, TransactionTrait, Set}; +use std::sync::Arc; +use uuid::Uuid; + +#[derive(Debug, Clone)] +pub struct ServerService { + service_context: Arc, +} + +impl ServerService { + pub fn new(service_context: Arc) -> Self { + Self { service_context } + } + + pub async fn create_server( + &self, + name: String, + is_default: bool, + ) -> Result { + let db = &self.service_context.repositories.server.context.db; + let event_bus = &self.service_context.event_bus; + + let txn = db.begin().await?; + + let active = server::ActiveModel { + name: Set(name), + is_default: Set(is_default), + ..Default::default() + }; + let srv = active.insert(&txn).await?; + + let default_group = role::ActiveModel { + server_id: Set(srv.id), + name: Set("Member".to_string()), + is_default: Set(true), + ..Default::default() + }; + default_group.insert(&txn).await?; + + txn.commit().await?; + + event_bus.emit("server_created", srv.clone()); + + Ok(srv) + } + + pub async fn update_server( + &self, + id: Uuid, + name: String, + is_default: bool, + ) -> Result { + let db = &self.service_context.repositories.server.context.db; + let event_bus = &self.service_context.event_bus; + + let txn = db.begin().await?; + + let existing = server::Entity::find_by_id(id) + .one(&txn) + .await? + .ok_or_else(|| anyhow::anyhow!("Server not found"))?; + + let mut active: server::ActiveModel = existing.into(); + active.name = Set(name); + active.is_default = Set(is_default); + + let srv = active.update(&txn).await?; + + txn.commit().await?; + + event_bus.emit("server_updated", srv.clone()); + + Ok(srv) + } + + pub async fn delete_server(&self, id: Uuid) -> Result { + let db = &self.service_context.repositories.server.context.db; + let event_bus = &self.service_context.event_bus; + + let txn = db.begin().await?; + + let res = server::Entity::delete_by_id(id) + .exec(&txn) + .await?; + + let deleted = res.rows_affected > 0; + + txn.commit().await?; + + if deleted { + event_bus.emit("server_deleted", id); + } + + Ok(deleted) + } + + pub async fn add_user(&self, server_id: Uuid, user_id: Uuid) -> Result { + let db = &self.service_context.repositories.server.context.db; + let event_bus = &self.service_context.event_bus; + + let txn = db.begin().await?; + + crate::models::server_user::ActiveModel { + server_id: Set(server_id), + user_id: Set(user_id), + ..Default::default() + } + .insert(&txn) + .await?; + + let role_id: Uuid = role::Entity::find() + .filter(role::Column::ServerId.eq(server_id)) + .filter(role::Column::IsDefault.eq(true)) + .select_only() + .column(role::Column::Id) + .into_tuple::() + .one(&txn) + .await? + .ok_or_else(|| anyhow::anyhow!("Rôle par défaut introuvable"))?; + + txn.commit().await?; + + event_bus.emit("server_user_created", (server_id, user_id)); + + Ok(true) + } + + pub async fn remove_user(&self, server_id: Uuid, user_id: Uuid) -> Result { + let db = &self.service_context.repositories.server.context.db; + let event_bus = &self.service_context.event_bus; + + let txn = db.begin().await?; + + let res = crate::models::server_user::Entity::delete_many() + .filter(crate::models::server_user::Column::ServerId.eq(server_id)) + .filter(crate::models::server_user::Column::UserId.eq(user_id)) + .exec(&txn) + .await?; + + let deleted = res.rows_affected > 0; + + txn.commit().await?; + + if deleted { + event_bus.emit("server_user_deleted", (server_id, user_id)); + } + + Ok(deleted) + } +} diff --git a/src/services/user.rs b/src/services/user.rs new file mode 100644 index 0000000..72eee80 --- /dev/null +++ b/src/services/user.rs @@ -0,0 +1,104 @@ +use crate::services::ServicesContext; +use crate::models::{role, user}; +use crate::auth::password; +use sea_orm::{ActiveModelTrait, ColumnTrait, EntityTrait, IntoActiveModel, QueryFilter, TransactionTrait, Set}; +use std::sync::Arc; +use uuid::Uuid; + +#[derive(Debug, Clone)] +pub struct UserService { + service_context: Arc, +} + +impl UserService { + pub fn new(service_context: Arc) -> Self { + Self { service_context } + } + + pub async fn create_user( + &self, + active: user::ActiveModel, + ) -> Result { + let db = &self.service_context.repositories.server.context.db; + let event_bus = &self.service_context.event_bus; + + let txn = db.begin().await?; + + let usr = active.insert(&txn).await?; + + txn.commit().await?; + + event_bus.emit("user_created", usr.clone()); + + Ok(usr) + } + + pub async fn update_user( + &self, + active: user::ActiveModel, + ) -> Result { + let db = &self.service_context.repositories.server.context.db; + let event_bus = &self.service_context.event_bus; + + let txn = db.begin().await?; + + let usr = active.update(&txn).await?; + + txn.commit().await?; + + event_bus.emit("user_updated", usr.clone()); + + Ok(usr) + } + + pub async fn set_password(&self, id: Uuid, password_str: String) -> Result<(), anyhow::Error> { + let db = &self.service_context.repositories.server.context.db; + let event_bus = &self.service_context.event_bus; + + let txn = db.begin().await?; + + let user_model = user::Entity::find_by_id(id) + .one(&txn) + .await? + .ok_or_else(|| anyhow::anyhow!("User not found"))?; + + let mut active = user_model.into_active_model(); + let password_to_hash = password_str.clone(); + + let hashed = tokio::task::spawn_blocking(move || password::hash_password(&password_to_hash)) + .await + .map_err(|e| anyhow::anyhow!("Join error: {}", e))? + .map_err(|e| anyhow::anyhow!("Password hashing failed: {}", e))?; + + active.password = Set(hashed); + + let usr = active.update(&txn).await?; + + txn.commit().await?; + + event_bus.emit("user_changed", usr); + + Ok(()) + } + + pub async fn delete_user(&self, id: Uuid) -> Result { + let db = &self.service_context.repositories.server.context.db; + let event_bus = &self.service_context.event_bus; + + let txn = db.begin().await?; + + let res = user::Entity::delete_by_id(id) + .exec(&txn) + .await?; + + let deleted = res.rows_affected > 0; + + txn.commit().await?; + + if deleted { + event_bus.emit("user_deleted", id); + } + + Ok(deleted) + } +}