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

8.4 KiB
Raw Blame History

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 : createdadd, updatedupdate, deletedremove

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 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 :

#[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
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.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}