From 4fef4f516238ecb6ea775021042304020d795ca9 Mon Sep 17 00:00:00 2001 From: Nell Date: Tue, 7 Jul 2026 01:39:38 +0200 Subject: [PATCH] init --- src/routes/gateway/handlers.rs | 1 - src/routes/gateway/mod.rs | 312 +++++++++++---------------------- 2 files changed, 98 insertions(+), 215 deletions(-) diff --git a/src/routes/gateway/handlers.rs b/src/routes/gateway/handlers.rs index e338f07..530eb58 100644 --- a/src/routes/gateway/handlers.rs +++ b/src/routes/gateway/handlers.rs @@ -42,7 +42,6 @@ async fn handle_socket(socket: WebSocket, state: AppState, user: User) { let (tx, mut rx) = mpsc::unbounded_channel::(); 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 8071d99..4ddc941 100644 --- a/src/routes/gateway/mod.rs +++ b/src/routes/gateway/mod.rs @@ -1,4 +1,3 @@ -use event_bus::EventBus; use crate::models::category; use crate::models::channel; use crate::models::message; @@ -9,6 +8,7 @@ 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 event_bus::EventBus; use events::GatewayEvent; use parking_lot::RwLock; use std::collections::HashMap; @@ -55,7 +55,11 @@ impl GatewayManager { } impl GatewayClient { - pub fn new(user: User, sender: mpsc::UnboundedSender, event_bus: Arc) -> Self { + pub fn new( + user: User, + sender: mpsc::UnboundedSender, + event_bus: Arc, + ) -> Self { let connection_id = Uuid::new_v4(); Self { user, @@ -66,238 +70,118 @@ impl GatewayClient { } } + fn subscribe_event( + &self, + event_name: &'static str, + namespace: &'static str, + action: &'static str, + mapper: F, + ) -> Arc> + where + T: Clone + Send + Sync + 'static, + R: serde::Serialize + Send + 'static, + F: Fn(T) -> R + Send + Sync + 'static, + { + let sender = self.sender.clone(); + Arc::new( + self.event_bus + .on_async::(event_name, move |payload| { + let sender = sender.clone(); + let content = mapper(payload); + async move { + let event = GatewayEvent { + namespace, + action, + content, + }; + if let Ok(json) = serde_json::to_string(&event) { + let _ = sender.send(Message::Text(json.into())); + } + } + }), + ) + } + pub fn subscribe_to_events(&mut self) { - let mut handles: Vec>> = Vec::new(); + let mut handles = Vec::new(); - // ── Message ────────────────────────────────────────────────────────── - let sender = self.sender.clone(); - handles.push(Arc::new(self.event_bus.on_async::( + // Message + handles.push(self.subscribe_event( "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", + "add", + message_model_to_message_response, + )); + handles.push(self.subscribe_event( "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())); - } - } - }, - ))); + "Message", + "update", + message_model_to_message_response, + )); + handles.push(self.subscribe_event("message_deleted", "Message", "remove", |id: Uuid| id)); - 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 + handles.push(self.subscribe_event( "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", + "add", + channel_model_to_channel_response, + )); + handles.push(self.subscribe_event( "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())); - } - } - }, - ))); + "Channel", + "update", + channel_model_to_channel_response, + )); + handles.push(self.subscribe_event("channel_deleted", "Channel", "remove", |id: Uuid| id)); - 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 + handles.push(self.subscribe_event( "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", + "add", + category_model_to_category_response, + )); + handles.push(self.subscribe_event( "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())); - } - } - }, - ))); + "Category", + "update", + category_model_to_category_response, + )); + handles.push(self.subscribe_event("category_deleted", "Category", "remove", |id: Uuid| id)); - 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 + handles.push(self.subscribe_event( "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", + "add", + server_model_to_server_response, + )); + handles.push(self.subscribe_event( "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())); - } - } - }, - ))); + "Server", + "update", + server_model_to_server_response, + )); + handles.push(self.subscribe_event("server_deleted", "Server", "remove", |id: Uuid| id)); self._event_handles = handles; } - async fn on_connect(&self) { - tracing::info!("Client connected: {:?}", self.user); + pub fn unsubscribe_all(&mut self) { + for handle in self._event_handles.drain(..) { + handle.abort(); + } } - async fn on_disconnect(&self) { + async fn on_connect(&mut self) { + tracing::info!("Client connected: {:?}", self.user); + self.subscribe_to_events(); + } + + async fn on_disconnect(&mut self) { tracing::info!("Client disconnected: {:?}", self.user); + self.unsubscribe_all(); } async fn on_message(&self, message: Message) {