use crate::core::AppState; use crate::domain::events::category::{ CategoryCreatedEvent, CategoryDeletedEvent, CategoryUpdatedEvent, }; use crate::domain::events::channel::{ ChannelCreatedEvent, ChannelDeletedEvent, ChannelUpdatedEvent, }; use crate::domain::events::emoji::{EmojiCreatedEvent, EmojiDeletedEvent, EmojiUpdatedEvent}; use crate::domain::events::message::{ MessageCreatedEvent, MessageDeletedEvent, MessageReactionAddedEvent, MessageReactionRemovedEvent, MessageUpdatedEvent, }; use crate::domain::events::server::{ServerCreatedEvent, ServerDeletedEvent, ServerUpdatedEvent}; use crate::domain::events::server_tree::ServerTreeInvalidatedEvent; use crate::domain::events::voice_presence::VoicePresenceEvent; use crate::models::{server_user, user}; use crate::repositories::Repositories; 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_with_data, message_model_to_message_response_with_reactions, reaction_model_to_response, }; use crate::routes::server::mapper::server_model_to_server_response; use crate::services::Services; use axum::extract::ws::Message; use event_bus::EventBus; use events::GatewayEvent; use parking_lot::RwLock; use sea_orm::{ColumnTrait, EntityTrait, QueryFilter, QuerySelect}; use serde::Serialize; use std::collections::{HashMap, HashSet}; use std::sync::Arc; use tokio::sync::mpsc; use uuid::Uuid; pub mod events; pub mod handlers; pub mod routes; /// Couche de transport WebSocket. Elle ne connaît pas les permissions. #[derive(Debug)] pub struct GatewayManager { pub clients: RwLock>, } #[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)] pub struct ConnectionKey { pub user_id: Uuid, pub connection_id: Uuid, } #[derive(Debug, Clone)] pub struct GatewayClient { user: crate::models::user::Model, connection_id: Uuid, pub sender: mpsc::UnboundedSender, } impl GatewayManager { pub fn new() -> Self { Self { clients: RwLock::new(HashMap::new()), } } pub(crate) fn add_client(&self, client: GatewayClient) { self.clients.write().insert(client.key(), client); } pub(crate) fn remove_client(&self, client: &GatewayClient) { self.clients.write().remove(&client.key()); } pub(crate) fn send_to_users( &self, users: impl IntoIterator, namespace: &'static str, action: &'static str, content: T, ) { let Ok(json) = serde_json::to_string(&GatewayEvent { namespace, action, content, }) else { return; }; let users: HashSet<_> = users.into_iter().collect(); for (key, client) in self.clients.read().iter() { if users.contains(&key.user_id) { let _ = client.sender.send(Message::Text(json.clone().into())); } } } } /// Résout les audiences puis délègue l'envoi au transport WebSocket. pub struct RealtimeRouter { gateway: Arc, services: Arc, repositories: Arc, } impl RealtimeRouter { pub fn new( gateway: Arc, services: Arc, repositories: Arc, ) -> Self { Self { gateway, services, repositories, } } pub fn start(self: &Arc, event_bus: Arc) { let router = Arc::clone(self); event_bus.on_async::(move |event| { let router = Arc::clone(&router); async move { let users = router.server_users(event.server_id).await; let action = if event.joined { "joined" } else { "left" }; router.gateway.send_to_users(users, "VoicePresence", action, event); } }); let router = Arc::clone(self); event_bus.on_async::(move |event| { let router = Arc::clone(&router); async move { let attachments = router .repositories .message .attachments_for_messages(&[event.message.id]) .await .ok() .and_then(|mut items| items.remove(&event.message.id)) .unwrap_or_default(); let users = router .services .realtime_registry .users_for_channel(event.channel_id); router.gateway.send_to_users( users, "Message", "add", message_model_to_message_response_with_data( event.message, event.server_id, Vec::new(), attachments, ), ); } }); let router = Arc::clone(self); event_bus.on_async::(move |event| { let router = Arc::clone(&router); async move { let reactions = router .services .message_reaction .grouped_for_messages(&[event.message.id]) .await .ok() .and_then(|mut groups| groups.remove(&event.message.id)) .unwrap_or_default(); let users = router .services .realtime_registry .users_for_channel(event.channel_id); router.gateway.send_to_users( users, "Message", "update", message_model_to_message_response_with_reactions( event.message, event.server_id, reactions, ), ); } }); let router = Arc::clone(self); event_bus.on_async::(move |event| { let router = Arc::clone(&router); async move { let users = router .services .realtime_registry .users_for_channel(event.channel_id); router .gateway .send_to_users(users, "Message", "remove", event.message.id); } }); let router = Arc::clone(self); event_bus.on_async::(move |event| { let router = Arc::clone(&router); async move { let users = router .services .realtime_registry .users_for_channel(event.channel_id); router.gateway.send_to_users( users, "Reaction", "add", reaction_model_to_response(event.reaction), ); } }); let router = Arc::clone(self); event_bus.on_async::(move |event| { let router = Arc::clone(&router); async move { let users = router .services .realtime_registry .users_for_channel(event.channel_id); router.gateway.send_to_users( users, "Reaction", "remove", reaction_model_to_response(event.reaction), ); } }); let router = Arc::clone(self); event_bus.on_async::(move |event| { let router = Arc::clone(&router); async move { if let Err(error) = router .services .realtime_registry .refresh_channel(&router.repositories, event.channel.id) .await { tracing::error!(channel_id = %event.channel.id, ?error, "Unable to refresh channel audience"); return; } let users = router .services .realtime_registry .users_for_channel(event.channel.id); router.gateway.send_to_users( users, "Channel", "add", channel_model_to_channel_response(event.channel), ); } }); let router = Arc::clone(self); event_bus.on_async::(move |event| { let router = Arc::clone(&router); async move { let mut users = router .services .realtime_registry .users_for_channel(event.channel.id); users.extend( router .services .realtime_registry .users_for_channel(event.previous.id), ); router.gateway.send_to_users( users, "Channel", "update", channel_model_to_channel_response(event.channel), ); } }); let router = Arc::clone(self); event_bus.on_async::(move |event| { let router = Arc::clone(&router); async move { let users = router .services .realtime_registry .users_for_channel(event.channel.id); router .gateway .send_to_users(users, "Channel", "remove", event.channel.id); router .services .realtime_registry .remove_channel(event.channel.id); } }); let router = Arc::clone(self); event_bus.on_async::(move |event| { let router = Arc::clone(&router); async move { let category = event.category; let users = router.server_users(category.server_id).await; router.gateway.send_to_users( users, "Category", "add", category_model_to_category_response(category), ); } }); let router = Arc::clone(self); event_bus.on_async::(move |event| { let router = Arc::clone(&router); async move { let category = event.category; let users = router.server_users(category.server_id).await; router.gateway.send_to_users( users, "Category", "update", category_model_to_category_response(category), ); } }); let router = Arc::clone(self); event_bus.on_async::(move |event| { let router = Arc::clone(&router); async move { let category = event.category; let users = router.server_users(category.server_id).await; router .gateway .send_to_users(users, "Category", "remove", category.id); } }); let router = Arc::clone(self); event_bus.on_async::(move |event| { let router = Arc::clone(&router); async move { let server = event.server; let users = router.server_users(server.id).await; router.gateway.send_to_users( users, "Server", "add", server_model_to_server_response(server), ); } }); let router = Arc::clone(self); event_bus.on_async::(move |event| { let router = Arc::clone(&router); async move { let server = event.server; let users = router.server_users(server.id).await; router.gateway.send_to_users( users, "Server", "update", server_model_to_server_response(server), ); } }); let router = Arc::clone(self); event_bus.on_async::(move |event| { let router = Arc::clone(&router); async move { router .gateway .send_to_users(event.audience, "Server", "remove", event.server.id); } }); let router = Arc::clone(self); event_bus.on_async::(move |event| { let router = Arc::clone(&router); async move { let users = match event.user_ids { Some(users) => users, None => router.server_users(event.server_id).await, }; router .gateway .send_to_users(users, "ServerTree", "refresh", event.server_id); } }); let router = Arc::clone(self); event_bus.on_async::(move |event| { let router = Arc::clone(&router); async move { let users = router.emoji_users(event.emoji.server_id).await; router.gateway.send_to_users( users, "Emoji", "add", crate::routes::emoji::mapper::response(event.emoji), ); } }); let router = Arc::clone(self); event_bus.on_async::(move |event| { let router = Arc::clone(&router); async move { let mut users = router.emoji_users(event.emoji.server_id).await; users.extend(router.emoji_users(event.previous.server_id).await); router.gateway.send_to_users( users, "Emoji", "update", crate::routes::emoji::mapper::response(event.emoji), ); } }); let router = Arc::clone(self); event_bus.on_async::(move |event| { let router = Arc::clone(&router); async move { let users = router.emoji_users(event.emoji.server_id).await; router.gateway.send_to_users( users, "Emoji", "remove", crate::routes::emoji::mapper::response(event.emoji), ); } }); } async fn server_users(&self, server_id: Uuid) -> Vec { server_user::Entity::find() .filter(server_user::Column::ServerId.eq(server_id)) .select_only() .column(server_user::Column::UserId) .into_tuple::() .all(&self.repositories.server.context.db) .await .unwrap_or_default() } async fn emoji_users(&self, server_id: Option) -> Vec { match server_id { Some(server_id) => self.server_users(server_id).await, None => user::Entity::find() .select_only() .column(user::Column::Id) .into_tuple::() .all(&self.repositories.server.context.db) .await .unwrap_or_default(), } } } impl GatewayClient { pub fn new(user: crate::models::user::Model, sender: mpsc::UnboundedSender) -> Self { Self { user, connection_id: Uuid::new_v4(), sender, } } pub fn key(&self) -> ConnectionKey { ConnectionKey { user_id: self.user.id, connection_id: self.connection_id, } } pub async fn on_connect(&mut self) { tracing::info!(user_id = %self.user.id, "Client connected"); } pub async fn on_disconnect(&mut self, state: &AppState) { tracing::info!(user_id = %self.user.id, "Client disconnected"); state.voice.leave_all(self.user.id).await; } pub async fn on_message(&self, message: Message, _state: &AppState) { if let Message::Text(content) = message { tracing::debug!(user_id = %self.user.id, "Received text message: {}", content); } } }