Files
oxspeak_server/.junie/plans/gateway-event-bus-triggers.md
T
2026-07-04 16:19:43 +02:00

179 lines
8.4 KiB
Markdown
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
---
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": "<uuid>"}
```
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<UserId, HashMap<ConnId, GatewayClient>>`) + `GatewayClient` avec un `mpsc::UnboundedSender<axum::ws::Message>`
- `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<GatewayManager>` et `event_bus: Arc<EventBus>`
- **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<T, _, _>(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<T>` sérialisable en JSON :
```rust
#[derive(Serialize)]
pub struct GatewayEvent<T: Serialize> {
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<EventBus>` (copié depuis `AppState` à la création)
- Un champ `_event_handles: Vec<SubscriptionHandle>` (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<Message>,
pub event_bus: Arc<EventBus>,
_event_handles: Vec<SubscriptionHandle>,
}
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<T: Serialize> {
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<T>
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<T: Serialize>` (namespace, action, content)
- Déclarer `pub mod events;` dans `src/routes/gateway/mod.rs`
- Ajouter à `GatewayClient` : `event_bus: Arc<EventBus>` et `_event_handles: Vec<SubscriptionHandle>`
- Mettre à jour `GatewayClient::new` pour accepter `event_bus: Arc<EventBus>` 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}`