This commit is contained in:
2026-07-07 01:39:38 +02:00
parent 30d928ddde
commit 4fef4f5162
2 changed files with 98 additions and 215 deletions
-1
View File
@@ -42,7 +42,6 @@ async fn handle_socket(socket: WebSocket, state: AppState, user: User) {
let (tx, mut rx) = mpsc::unbounded_channel::<Message>(); let (tx, mut rx) = mpsc::unbounded_channel::<Message>();
let mut client = GatewayClient::new(user, tx, state.event_bus.clone()); let mut client = GatewayClient::new(user, tx, state.event_bus.clone());
client.subscribe_to_events();
client.on_connect().await; client.on_connect().await;
state.gateway.add_client(client.clone()); state.gateway.add_client(client.clone());
+98 -214
View File
@@ -1,4 +1,3 @@
use event_bus::EventBus;
use crate::models::category; use crate::models::category;
use crate::models::channel; use crate::models::channel;
use crate::models::message; 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::message::mapper::message_model_to_message_response;
use crate::routes::server::mapper::server_model_to_server_response; use crate::routes::server::mapper::server_model_to_server_response;
use axum::extract::ws::Message; use axum::extract::ws::Message;
use event_bus::EventBus;
use events::GatewayEvent; use events::GatewayEvent;
use parking_lot::RwLock; use parking_lot::RwLock;
use std::collections::HashMap; use std::collections::HashMap;
@@ -55,7 +55,11 @@ impl GatewayManager {
} }
impl GatewayClient { impl GatewayClient {
pub fn new(user: User, sender: mpsc::UnboundedSender<Message>, event_bus: Arc<EventBus>) -> Self { pub fn new(
user: User,
sender: mpsc::UnboundedSender<Message>,
event_bus: Arc<EventBus>,
) -> Self {
let connection_id = Uuid::new_v4(); let connection_id = Uuid::new_v4();
Self { Self {
user, user,
@@ -66,238 +70,118 @@ impl GatewayClient {
} }
} }
fn subscribe_event<T, F, R>(
&self,
event_name: &'static str,
namespace: &'static str,
action: &'static str,
mapper: F,
) -> Arc<JoinHandle<()>>
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::<T, _, _>(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) { pub fn subscribe_to_events(&mut self) {
let mut handles: Vec<Arc<JoinHandle<()>>> = Vec::new(); let mut handles = Vec::new();
// ── Message ────────────────────────────────────────────────────────── // Message
let sender = self.sender.clone(); handles.push(self.subscribe_event(
handles.push(Arc::new(self.event_bus.on_async::<message::Model, _, _>(
"message_created", "message_created",
move |payload| { "Message",
let sender = sender.clone(); "add",
async move { message_model_to_message_response,
let event = GatewayEvent { ));
namespace: "Message", handles.push(self.subscribe_event(
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::Model, _, _>(
"message_updated", "message_updated",
move |payload| { "Message",
let sender = sender.clone(); "update",
async move { message_model_to_message_response,
let event = GatewayEvent { ));
namespace: "Message", handles.push(self.subscribe_event("message_deleted", "Message", "remove", |id: Uuid| id));
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()));
}
}
},
)));
let sender = self.sender.clone(); // Channel
handles.push(Arc::new(self.event_bus.on_async::<Uuid, _, _>( handles.push(self.subscribe_event(
"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::Model, _, _>(
"channel_created", "channel_created",
move |payload| { "Channel",
let sender = sender.clone(); "add",
async move { channel_model_to_channel_response,
let event = GatewayEvent { ));
namespace: "Channel", handles.push(self.subscribe_event(
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::Model, _, _>(
"channel_updated", "channel_updated",
move |payload| { "Channel",
let sender = sender.clone(); "update",
async move { channel_model_to_channel_response,
let event = GatewayEvent { ));
namespace: "Channel", handles.push(self.subscribe_event("channel_deleted", "Channel", "remove", |id: Uuid| id));
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()));
}
}
},
)));
let sender = self.sender.clone(); // Category
handles.push(Arc::new(self.event_bus.on_async::<Uuid, _, _>( handles.push(self.subscribe_event(
"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::Model, _, _>(
"category_created", "category_created",
move |payload| { "Category",
let sender = sender.clone(); "add",
async move { category_model_to_category_response,
let event = GatewayEvent { ));
namespace: "Category", handles.push(self.subscribe_event(
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::Model, _, _>(
"category_updated", "category_updated",
move |payload| { "Category",
let sender = sender.clone(); "update",
async move { category_model_to_category_response,
let event = GatewayEvent { ));
namespace: "Category", handles.push(self.subscribe_event("category_deleted", "Category", "remove", |id: Uuid| id));
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()));
}
}
},
)));
let sender = self.sender.clone(); // Server
handles.push(Arc::new(self.event_bus.on_async::<Uuid, _, _>( handles.push(self.subscribe_event(
"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::Model, _, _>(
"server_created", "server_created",
move |payload| { "Server",
let sender = sender.clone(); "add",
async move { server_model_to_server_response,
let event = GatewayEvent { ));
namespace: "Server", handles.push(self.subscribe_event(
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::Model, _, _>(
"server_updated", "server_updated",
move |payload| { "Server",
let sender = sender.clone(); "update",
async move { server_model_to_server_response,
let event = GatewayEvent { ));
namespace: "Server", handles.push(self.subscribe_event("server_deleted", "Server", "remove", |id: Uuid| id));
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::<Uuid, _, _>(
"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()));
}
}
},
)));
self._event_handles = handles; self._event_handles = handles;
} }
async fn on_connect(&self) { pub fn unsubscribe_all(&mut self) {
tracing::info!("Client connected: {:?}", self.user); 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); tracing::info!("Client disconnected: {:?}", self.user);
self.unsubscribe_all();
} }
async fn on_message(&self, message: Message) { async fn on_message(&self, message: Message) {