8.4 KiB
8.4 KiB
sessionId
| 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
{"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 surHashMap<UserId, HashMap<ConnId, GatewayClient>>) +GatewayClientavec unmpsc::UnboundedSender<axum::ws::Message>src/routes/gateway/handlers.rs:handle_socketcrée le client, spawn deux tâches (send/recv), mais n'écoute pas encore l'EventBussrc/core/state.rs:AppStateexposegateway: Arc<GatewayManager>etevent_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 lesModelSeaORM (ouUuidpour les deletes) - EventBus : API
on_async<T, _, _>(topic, |payload| async { ... })avecT: 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 :
#[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é depuisAppStateà 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 viaself.event_bus.on_async, cloneself.senderdans chaque callback, et stocke les handles dansself._event_handles
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
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
// 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
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.rsavec uniquementGatewayEvent<T: Serialize>(namespace, action, content) - Déclarer
pub mod events;danssrc/routes/gateway/mod.rs - Ajouter à
GatewayClient:event_bus: Arc<EventBus>et_event_handles: Vec<SubscriptionHandle> - Mettre à jour
GatewayClient::newpour accepterevent_bus: Arc<EventBus>en paramètre - Implémenter
pub fn subscribe_to_events(&mut self)surGatewayClient:- S'abonner via
self.event_bus.on_asyncaux 12 topics (message/channel/category/server × created/updated/deleted) - Dans chaque callback, cloner
self.senderet envoyer le JSONGatewayEvent - Pour
_created/_updated: mapperModel→ DTO via les mappers existants - Pour
_deleted: content =Uuidbrut - Stocker les handles dans
self._event_handles
- S'abonner via
✓ 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()avantstate.gateway.add_client() - Les abonnements vivent le temps de la connexion via
_event_handlesdans 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}