post services integrations
This commit is contained in:
@@ -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<EventBus>` 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<EventBus>`.
|
||||
- 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<EventBus>`.
|
||||
- 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<EventBus>` 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.
|
||||
+1
-1
@@ -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());
|
||||
|
||||
@@ -37,17 +37,11 @@ impl CategoryRepository {
|
||||
|
||||
pub async fn update(&self, active: category::ActiveModel) -> AnyResult<category::Model> {
|
||||
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<category::Model> {
|
||||
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)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -32,13 +32,11 @@ impl ChannelRepository {
|
||||
|
||||
pub async fn update(&self, active: channel::ActiveModel) -> AnyResult<channel::Model> {
|
||||
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<channel::Model> {
|
||||
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)
|
||||
}
|
||||
|
||||
|
||||
@@ -48,38 +48,11 @@ impl MessageRepository {
|
||||
|
||||
pub async fn update(&self, active: message::ActiveModel) -> AnyResult<message::Model> {
|
||||
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<message::Model> {
|
||||
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<Scope> = Vec::new();
|
||||
scopes.push(Scope::uuid("channel", message.channel_id));
|
||||
// retrieve related channel and server
|
||||
let server_id: Option<uuid::Uuid> = channel::Entity::find_by_id(message.channel_id)
|
||||
.select_only()
|
||||
.column(channel::Column::ServerId)
|
||||
.into_tuple::<Option<uuid::Uuid>>()
|
||||
.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)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -24,8 +24,7 @@ mod user;
|
||||
|
||||
#[derive(Clone, Debug)]
|
||||
pub struct RepositoryContext {
|
||||
db: DatabaseConnection,
|
||||
events: Arc<EventBus>,
|
||||
pub db: DatabaseConnection,
|
||||
}
|
||||
|
||||
#[derive(Clone, Debug)]
|
||||
@@ -41,8 +40,8 @@ pub struct Repositories {
|
||||
}
|
||||
|
||||
impl Repositories {
|
||||
pub fn new(db: DatabaseConnection, events: Arc<EventBus>) -> Self {
|
||||
let context = Arc::new(RepositoryContext { db, events });
|
||||
pub fn new(db: DatabaseConnection) -> Self {
|
||||
let context = Arc::new(RepositoryContext { db });
|
||||
|
||||
Self {
|
||||
server: ServerRepository {
|
||||
|
||||
@@ -49,13 +49,11 @@ impl RoleRepository {
|
||||
|
||||
pub async fn create(&self, active: role::ActiveModel) -> AnyResult<role::Model> {
|
||||
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<role::Model> {
|
||||
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)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -31,7 +31,6 @@ impl ServerRepository {
|
||||
|
||||
pub async fn update(&self, active: server::ActiveModel) -> AnyResult<server::Model> {
|
||||
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)
|
||||
}
|
||||
|
||||
|
||||
@@ -40,7 +40,6 @@ impl ServerItemOrderRepository {
|
||||
})
|
||||
.await?;
|
||||
|
||||
self.context.events.emit("server_order_updated", server_id);
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
|
||||
@@ -53,13 +53,11 @@ impl UserRepository {
|
||||
|
||||
pub async fn update(&self, active: user::ActiveModel) -> AnyResult<user::Model> {
|
||||
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<user::Model> {
|
||||
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)
|
||||
}
|
||||
|
||||
|
||||
@@ -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<UpdateCategoryRequest>,
|
||||
) -> Result<Json<CategoryResponse>, 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<AppState>,
|
||||
Path(id): Path<Uuid>,
|
||||
) -> Result<StatusCode, HTTPError> {
|
||||
if state.repositories.category.delete(id).await? {
|
||||
if state.services.category.delete_category(id).await? {
|
||||
Ok(StatusCode::NO_CONTENT)
|
||||
} else {
|
||||
Err(HTTPError::NotFound)
|
||||
|
||||
@@ -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<AppState>,
|
||||
Path(id): Path<Uuid>,
|
||||
) -> Result<StatusCode, HTTPError> {
|
||||
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<SetChannelPermissionRequest>,
|
||||
) -> Result<Json<ChannelUserPermissionResponse>, 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<SetChannelPermissionRequest>,
|
||||
) -> Result<Json<ChannelRolePermissionResponse>, 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?;
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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<AppState>,
|
||||
Path(id): Path<Uuid>,
|
||||
) -> Result<StatusCode, HTTPError> {
|
||||
if state.repositories.role.delete(id).await? {
|
||||
if state.services.role.delete_role(id).await? {
|
||||
Ok(StatusCode::NO_CONTENT)
|
||||
} else {
|
||||
Err(HTTPError::NotFound)
|
||||
|
||||
@@ -83,8 +83,7 @@ pub async fn create(
|
||||
State(state): State<AppState>,
|
||||
Json(payload): Json<CreateServerRequest>,
|
||||
) -> Result<(StatusCode, Json<ServerResponse>), 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<UpdateServerRequest>,
|
||||
) -> Result<Json<ServerResponse>, 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<AppState>,
|
||||
Path(id): Path<Uuid>,
|
||||
) -> Result<StatusCode, HTTPError> {
|
||||
if state.repositories.server.delete(id).await? {
|
||||
if state.services.server.delete_server(id).await? {
|
||||
Ok(StatusCode::NO_CONTENT)
|
||||
} else {
|
||||
Err(HTTPError::NotFound)
|
||||
|
||||
@@ -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<AppState>,
|
||||
Path(id): Path<Uuid>,
|
||||
) -> Result<StatusCode, HTTPError> {
|
||||
if state.repositories.user.delete(id).await? {
|
||||
if state.services.user.delete_user(id).await? {
|
||||
Ok(StatusCode::NO_CONTENT)
|
||||
} else {
|
||||
Err(HTTPError::NotFound)
|
||||
|
||||
@@ -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<ServicesContext>,
|
||||
}
|
||||
|
||||
impl CategoryService {
|
||||
pub fn new(service_context: Arc<ServicesContext>) -> Self {
|
||||
Self { service_context }
|
||||
}
|
||||
|
||||
pub async fn create_category(
|
||||
&self,
|
||||
server_id: Uuid,
|
||||
name: String,
|
||||
) -> Result<category::Model, 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 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<category::Model, 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 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<bool, 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 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)
|
||||
}
|
||||
}
|
||||
@@ -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<ServicesContext>,
|
||||
}
|
||||
|
||||
impl ChannelService {
|
||||
pub fn new(service_context: Arc<ServicesContext>) -> Self {
|
||||
Self { service_context }
|
||||
}
|
||||
|
||||
pub async fn create_channel(
|
||||
&self,
|
||||
payload: CreateChannelRequest,
|
||||
) -> Result<channel::Model, anyhow::Error> {
|
||||
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<i64> = 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::<Option<i64>>()
|
||||
.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<channel::Model, 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 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<bool, anyhow::Error> {
|
||||
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(())
|
||||
}
|
||||
}
|
||||
@@ -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<ServicesContext>,
|
||||
}
|
||||
|
||||
impl MessageService {
|
||||
pub fn new(service_context: Arc<ServicesContext>) -> Self {
|
||||
Self { service_context }
|
||||
}
|
||||
|
||||
pub async fn create_message(
|
||||
&self,
|
||||
channel_id: Uuid,
|
||||
author_id: Uuid,
|
||||
content: String,
|
||||
) -> Result<message::Model, 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 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<Scope> = Vec::new();
|
||||
scopes.push(Scope::uuid("channel", msg.channel_id));
|
||||
|
||||
let server_id: Option<Uuid> = channel::Entity::find_by_id(msg.channel_id)
|
||||
.select_only()
|
||||
.column(channel::Column::ServerId)
|
||||
.into_tuple::<Option<Uuid>>()
|
||||
.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<message::Model, 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 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<bool, 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 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)
|
||||
}
|
||||
}
|
||||
@@ -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<PermissionSyncService>,
|
||||
pub server_order: Arc<ServerOrderService>,
|
||||
pub channel: Arc<ChannelService>,
|
||||
pub server: Arc<ServerService>,
|
||||
pub category: Arc<CategoryService>,
|
||||
pub message: Arc<MessageService>,
|
||||
pub user: Arc<UserService>,
|
||||
pub role: Arc<RoleService>,
|
||||
}
|
||||
|
||||
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
|
||||
|
||||
@@ -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<ServicesContext>,
|
||||
}
|
||||
|
||||
impl RoleService {
|
||||
pub fn new(service_context: Arc<ServicesContext>) -> Self {
|
||||
Self { service_context }
|
||||
}
|
||||
|
||||
pub async fn create_role(
|
||||
&self,
|
||||
active: role::ActiveModel,
|
||||
) -> Result<role::Model, 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 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<role::Model, 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 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<bool, 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 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)
|
||||
}
|
||||
}
|
||||
@@ -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<ServicesContext>,
|
||||
}
|
||||
|
||||
impl ServerService {
|
||||
pub fn new(service_context: Arc<ServicesContext>) -> Self {
|
||||
Self { service_context }
|
||||
}
|
||||
|
||||
pub async fn create_server(
|
||||
&self,
|
||||
name: String,
|
||||
is_default: bool,
|
||||
) -> Result<server::Model, 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 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<server::Model, 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 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<bool, 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 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<bool, 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::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::<Uuid>()
|
||||
.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<bool, 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 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)
|
||||
}
|
||||
}
|
||||
@@ -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<ServicesContext>,
|
||||
}
|
||||
|
||||
impl UserService {
|
||||
pub fn new(service_context: Arc<ServicesContext>) -> Self {
|
||||
Self { service_context }
|
||||
}
|
||||
|
||||
pub async fn create_user(
|
||||
&self,
|
||||
active: user::ActiveModel,
|
||||
) -> Result<user::Model, 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 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<user::Model, 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 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<bool, 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 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)
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user