diff --git a/.junie/plans/gateway-event-bus-triggers.md b/.junie/plans/gateway-event-bus-triggers.md new file mode 100644 index 0000000..c66b134 --- /dev/null +++ b/.junie/plans/gateway-event-bus-triggers.md @@ -0,0 +1,179 @@ +--- +sessionId: session-260704-155358-1bpw +--- + +# Requirements + +### Overview & Goals +Ajouter le support des événements temps-réel dans le gateway WebSocket (`src/routes/gateway`) en s'abonnant aux événements émis par les repositories via l'`EventBus`, puis en les diffusant à tous les clients connectés sous forme de messages JSON structurés. + +### Scope +**In Scope** +- Abonnement aux événements de 4 namespaces : `Server`, `Channel`, `Category`, `Message` (create/update/delete pour chacun) +- Sérialisation en JSON : `{"namespace": "Message", "action": "add", "content": {...DTO}}` +- Diffusion à tous les clients WebSocket connectés via le `GatewayManager` +- Correspondance actions : `created` → `add`, `updated` → `update`, `deleted` → `remove` + +**Out of Scope** +- Filtrage par permissions +- Ciblage par serveur/channel (broadcast global pour l'instant) +- Autres namespaces (User, Group, Attachment…) + +### Format JSON +```json +{"namespace": "Message", "action": "add", "content": {"id": "...", "channel_id": "...", ...}} +{"namespace": "Server", "action": "remove", "content": ""} +``` +Les actions `add`/`update` transportent le DTO complet du modèle ; l'action `remove` transporte uniquement l'UUID. + +# Technical Design + +### Current Implementation +- `src/routes/gateway/mod.rs` : `GatewayManager` (RwLock sur `HashMap>`) + `GatewayClient` avec un `mpsc::UnboundedSender` +- `src/routes/gateway/handlers.rs` : `handle_socket` crée le client, spawn deux tâches (send/recv), mais **n'écoute pas encore l'EventBus** +- `src/core/state.rs` : `AppState` expose `gateway: Arc` et `event_bus: Arc` +- **Repositories** : chaque opération émet sur l'EventBus avec des topics comme `message_created`, `channel_deleted`, `server_updated`, etc. — les payloads sont les `Model` SeaORM (ou `Uuid` pour les deletes) +- **EventBus** : API `on_async(topic, |payload| async { ... })` avec `T: Any + Send + Sync + Clone` + +### Key Decisions + Décision | Choix retenu | Raison | +---|---|---| + Où s'abonner à l'EventBus | Dans `GatewayClient` lui-même (méthode `subscribe_to_events`) | Le client est le « consumer » responsable de ses abonnements, cohérent avec le modèle Django Channels | + Mapping Model → DTO dans l'event | Utiliser les mappers existants (`server_model_to_server_response`, etc.) | Cohérence avec l'API REST | + Type du payload `deleted` | `Uuid` brut | C'est ce qu'émettent les repositories | + Envoi au client | Via `self.sender` (mpsc) directement depuis les callbacks EventBus | Chaque client gère son propre flux | + Cycle de vie des abonnements | Les `SubscriptionHandle` sont stockés dans `GatewayClient` (champ `_event_handles`) | Drop automatique quand le client est supprimé du `GatewayManager` | + +### Proposed Changes + +#### 1. `src/routes/gateway/events.rs` — Nouveau fichier +Contient uniquement le struct `GatewayEvent` sérialisable en JSON : +```rust +#[derive(Serialize)] +pub struct GatewayEvent { + pub namespace: &'static str, + pub action: &'static str, + pub content: T, +} +``` + +#### 2. `src/routes/gateway/mod.rs` — Enrichir `GatewayClient` +`GatewayClient` reçoit : +- Un champ `event_bus: Arc` (copié depuis `AppState` à la création) +- Un champ `_event_handles: Vec` (privé, initialisé vide) +- Une méthode `subscribe_to_events(&mut self)` qui s'abonne aux 12 topics via `self.event_bus.on_async`, clone `self.sender` dans chaque callback, et stocke les handles dans `self._event_handles` + +```rust +pub struct GatewayClient { + pub user_id: UserId, + pub conn_id: ConnId, + pub sender: mpsc::UnboundedSender, + pub event_bus: Arc, + _event_handles: Vec, +} + +impl GatewayClient { + pub fn subscribe_to_events(&mut self) { + let sender = self.sender.clone(); + let handle = self.event_bus.on_async("message_created", move |payload: message::Model| { + let sender = sender.clone(); + async move { + let event = GatewayEvent { namespace: "Message", action: "add", content: map_message(payload) }; + let _ = sender.send(Message::Text(serde_json::to_string(&event).unwrap())); + } + }); + self._event_handles.push(handle); + // ... répété pour les 11 autres topics + } +} +``` + +#### 3. `src/routes/gateway/handlers.rs` — Appel de `subscribe_to_events` dans `handle_socket` +```rust +async fn handle_socket(socket: WebSocket, state: AppState, user: User) { + // ... (existant) + let mut client = GatewayClient::new(user, tx, Arc::clone(&state.event_bus)); + client.subscribe_to_events(); // le client s'abonne lui-même + state.gateway.add_client(client.clone()); + + // ... tasks send/recv (inchangées) + + // À la déconnexion, client est drop-pé → _event_handles supprimés → désinscription auto + state.gateway.remove_client(&client); +} +``` + +### Data Models / Contracts +```rust +// src/routes/gateway/events.rs +#[derive(Serialize)] +pub struct GatewayEvent { + pub namespace: &'static str, + pub action: &'static str, // "add" | "update" | "remove" + pub content: T, +} +``` + +Topics écoutés et types associés : + Topic EventBus | Namespace | Action | Type payload | +---|---|---|---| + `message_created` | `Message` | `add` | `message::Model` | + `message_updated` | `Message` | `update` | `message::Model` | + `message_deleted` | `Message` | `remove` | `Uuid` | + `channel_created` | `Channel` | `add` | `channel::Model` | + `channel_updated` | `Channel` | `update` | `channel::Model` | + `channel_deleted` | `Channel` | `remove` | `Uuid` | + `category_created` | `Category` | `add` | `category::Model` | + `category_updated` | `Category` | `update` | `category::Model` | + `category_deleted` | `Category` | `remove` | `Uuid` | + `server_created` | `Server` | `add` | `server::Model` | + `server_updated` | `Server` | `update` | `server::Model` | + `server_deleted` | `Server` | `remove` | `Uuid` | + +### Architecture Diagram +```mermaid +graph LR + Repo[Repository] -- emit topic --> EB[EventBus] + GC[GatewayClient] -- possède Arc --> EB + GC -- subscribe_to_events --> EB + EB -- on_async callback --> GC + GC -- sender.send JSON --> WS[WebSocket Client] + HS[handle_socket] -- crée GatewayClient avec event_bus --> GC + GC -- drop _event_handles à déconnexion --> EB +``` + +### File Structure +``` +src/routes/gateway/ +├── mod.rs ← modifié : GatewayClient + champs event_bus/_event_handles + méthode subscribe_to_events() +├── handlers.rs ← modifié : GatewayClient::new reçoit event_bus, appel client.subscribe_to_events() +├── routes.rs ← inchangé +└── events.rs ← nouveau : GatewayEvent +src/main.rs ← inchangé +``` + +# Delivery Steps + +### ✓ Step 1: Créer events.rs et enrichir GatewayClient avec event_bus et subscribe_to_events() +Le struct `GatewayClient` possède son propre accès à l'EventBus et sa méthode d'abonnement. + +- Créer `src/routes/gateway/events.rs` avec uniquement `GatewayEvent` (namespace, action, content) +- Déclarer `pub mod events;` dans `src/routes/gateway/mod.rs` +- Ajouter à `GatewayClient` : `event_bus: Arc` et `_event_handles: Vec` +- Mettre à jour `GatewayClient::new` pour accepter `event_bus: Arc` en paramètre +- Implémenter `pub fn subscribe_to_events(&mut self)` sur `GatewayClient` : + - S'abonner via `self.event_bus.on_async` aux 12 topics (message/channel/category/server × created/updated/deleted) + - Dans chaque callback, cloner `self.sender` et envoyer le JSON `GatewayEvent` + - Pour `_created` / `_updated` : mapper `Model` → DTO via les mappers existants + - Pour `_deleted` : content = `Uuid` brut + - Stocker les handles dans `self._event_handles` + +### ✓ Step 2: Brancher subscribe_to_events() dans handle_socket +Chaque client WebSocket s'abonne à l'EventBus dès sa création et se désinscrit automatiquement à la déconnexion. + +- Modifier `src/routes/gateway/handlers.rs` +- Passer `Arc::clone(&state.event_bus)` à `GatewayClient::new` +- Appeler `client.subscribe_to_events()` avant `state.gateway.add_client()` +- Les abonnements vivent le temps de la connexion via `_event_handles` dans le struct +- Vérifier la compilation (`cargo check`) +- Tester : connecter un client WebSocket, créer/modifier/supprimer une entité via l'API REST, vérifier la réception du JSON `{namespace, action, content}` \ No newline at end of file diff --git a/src/routes/gateway/events.rs b/src/routes/gateway/events.rs new file mode 100644 index 0000000..42dba8f --- /dev/null +++ b/src/routes/gateway/events.rs @@ -0,0 +1,8 @@ +use serde::Serialize; + +#[derive(Serialize)] +pub struct GatewayEvent { + pub namespace: &'static str, + pub action: &'static str, + pub content: T, +} diff --git a/src/routes/gateway/handlers.rs b/src/routes/gateway/handlers.rs index aef0c51..e338f07 100644 --- a/src/routes/gateway/handlers.rs +++ b/src/routes/gateway/handlers.rs @@ -40,9 +40,9 @@ pub async fn ws_handler( async fn handle_socket(socket: WebSocket, state: AppState, user: User) { let (mut sender, mut receiver) = socket.split(); let (tx, mut rx) = mpsc::unbounded_channel::(); - let event_bus = state.event_bus.clone(); - let client = GatewayClient::new(user, tx); + let mut client = GatewayClient::new(user, tx, state.event_bus.clone()); + client.subscribe_to_events(); client.on_connect().await; state.gateway.add_client(client.clone()); diff --git a/src/routes/gateway/mod.rs b/src/routes/gateway/mod.rs index 1be3a30..8071d99 100644 --- a/src/routes/gateway/mod.rs +++ b/src/routes/gateway/mod.rs @@ -1,10 +1,23 @@ +use event_bus::EventBus; +use crate::models::category; +use crate::models::channel; +use crate::models::message; +use crate::models::server; use crate::models::user::Model as User; +use crate::routes::category::mapper::category_model_to_category_response; +use crate::routes::channel::mapper::channel_model_to_channel_response; +use crate::routes::message::mapper::message_model_to_message_response; +use crate::routes::server::mapper::server_model_to_server_response; use axum::extract::ws::Message; +use events::GatewayEvent; use parking_lot::RwLock; use std::collections::HashMap; +use std::sync::Arc; use tokio::sync::mpsc; +use tokio::task::JoinHandle; use uuid::Uuid; +pub mod events; pub mod handlers; pub mod routes; @@ -18,7 +31,9 @@ pub struct GatewayManager { pub struct GatewayClient { user: User, connection_id: Uuid, - sender: mpsc::UnboundedSender, + pub sender: mpsc::UnboundedSender, + pub event_bus: Arc, + _event_handles: Vec>>, } impl GatewayManager { @@ -40,15 +55,243 @@ impl GatewayManager { } impl GatewayClient { - pub fn new(user: User, sender: mpsc::UnboundedSender) -> Self { + pub fn new(user: User, sender: mpsc::UnboundedSender, event_bus: Arc) -> Self { let connection_id = Uuid::new_v4(); Self { user, connection_id, sender, + event_bus, + _event_handles: Vec::new(), } } + pub fn subscribe_to_events(&mut self) { + let mut handles: Vec>> = Vec::new(); + + // ── Message ────────────────────────────────────────────────────────── + let sender = self.sender.clone(); + handles.push(Arc::new(self.event_bus.on_async::( + "message_created", + move |payload| { + let sender = sender.clone(); + async move { + let event = GatewayEvent { + namespace: "Message", + action: "add", + content: message_model_to_message_response(payload), + }; + if let Ok(json) = serde_json::to_string(&event) { + let _ = sender.send(Message::Text(json.into())); + } + } + }, + ))); + + let sender = self.sender.clone(); + handles.push(Arc::new(self.event_bus.on_async::( + "message_updated", + move |payload| { + let sender = sender.clone(); + async move { + let event = GatewayEvent { + namespace: "Message", + action: "update", + content: message_model_to_message_response(payload), + }; + if let Ok(json) = serde_json::to_string(&event) { + let _ = sender.send(Message::Text(json.into())); + } + } + }, + ))); + + let sender = self.sender.clone(); + handles.push(Arc::new(self.event_bus.on_async::( + "message_deleted", + move |payload| { + let sender = sender.clone(); + async move { + let event = GatewayEvent { + namespace: "Message", + action: "remove", + content: payload, + }; + if let Ok(json) = serde_json::to_string(&event) { + let _ = sender.send(Message::Text(json.into())); + } + } + }, + ))); + + // ── Channel ─────────────────────────────────────────────────────────── + let sender = self.sender.clone(); + handles.push(Arc::new(self.event_bus.on_async::( + "channel_created", + move |payload| { + let sender = sender.clone(); + async move { + let event = GatewayEvent { + namespace: "Channel", + action: "add", + content: channel_model_to_channel_response(payload), + }; + if let Ok(json) = serde_json::to_string(&event) { + let _ = sender.send(Message::Text(json.into())); + } + } + }, + ))); + + let sender = self.sender.clone(); + handles.push(Arc::new(self.event_bus.on_async::( + "channel_updated", + move |payload| { + let sender = sender.clone(); + async move { + let event = GatewayEvent { + namespace: "Channel", + action: "update", + content: channel_model_to_channel_response(payload), + }; + if let Ok(json) = serde_json::to_string(&event) { + let _ = sender.send(Message::Text(json.into())); + } + } + }, + ))); + + let sender = self.sender.clone(); + handles.push(Arc::new(self.event_bus.on_async::( + "channel_deleted", + move |payload| { + let sender = sender.clone(); + async move { + let event = GatewayEvent { + namespace: "Channel", + action: "remove", + content: payload, + }; + if let Ok(json) = serde_json::to_string(&event) { + let _ = sender.send(Message::Text(json.into())); + } + } + }, + ))); + + // ── Category ────────────────────────────────────────────────────────── + let sender = self.sender.clone(); + handles.push(Arc::new(self.event_bus.on_async::( + "category_created", + move |payload| { + let sender = sender.clone(); + async move { + let event = GatewayEvent { + namespace: "Category", + action: "add", + content: category_model_to_category_response(payload), + }; + if let Ok(json) = serde_json::to_string(&event) { + let _ = sender.send(Message::Text(json.into())); + } + } + }, + ))); + + let sender = self.sender.clone(); + handles.push(Arc::new(self.event_bus.on_async::( + "category_updated", + move |payload| { + let sender = sender.clone(); + async move { + let event = GatewayEvent { + namespace: "Category", + action: "update", + content: category_model_to_category_response(payload), + }; + if let Ok(json) = serde_json::to_string(&event) { + let _ = sender.send(Message::Text(json.into())); + } + } + }, + ))); + + let sender = self.sender.clone(); + handles.push(Arc::new(self.event_bus.on_async::( + "category_deleted", + move |payload| { + let sender = sender.clone(); + async move { + let event = GatewayEvent { + namespace: "Category", + action: "remove", + content: payload, + }; + if let Ok(json) = serde_json::to_string(&event) { + let _ = sender.send(Message::Text(json.into())); + } + } + }, + ))); + + // ── Server ──────────────────────────────────────────────────────────── + let sender = self.sender.clone(); + handles.push(Arc::new(self.event_bus.on_async::( + "server_created", + move |payload| { + let sender = sender.clone(); + async move { + let event = GatewayEvent { + namespace: "Server", + action: "add", + content: server_model_to_server_response(payload), + }; + if let Ok(json) = serde_json::to_string(&event) { + let _ = sender.send(Message::Text(json.into())); + } + } + }, + ))); + + let sender = self.sender.clone(); + handles.push(Arc::new(self.event_bus.on_async::( + "server_updated", + move |payload| { + let sender = sender.clone(); + async move { + let event = GatewayEvent { + namespace: "Server", + action: "update", + content: server_model_to_server_response(payload), + }; + if let Ok(json) = serde_json::to_string(&event) { + let _ = sender.send(Message::Text(json.into())); + } + } + }, + ))); + + let sender = self.sender.clone(); + handles.push(Arc::new(self.event_bus.on_async::( + "server_deleted", + move |payload| { + let sender = sender.clone(); + async move { + let event = GatewayEvent { + namespace: "Server", + action: "remove", + content: payload, + }; + if let Ok(json) = serde_json::to_string(&event) { + let _ = sender.send(Message::Text(json.into())); + } + } + }, + ))); + + self._event_handles = handles; + } + async fn on_connect(&self) { tracing::info!("Client connected: {:?}", self.user); }