Implement event on gateway
This commit is contained in:
@@ -0,0 +1,8 @@
|
||||
use serde::Serialize;
|
||||
|
||||
#[derive(Serialize)]
|
||||
pub struct GatewayEvent<T: Serialize> {
|
||||
pub namespace: &'static str,
|
||||
pub action: &'static str,
|
||||
pub content: T,
|
||||
}
|
||||
@@ -40,9 +40,9 @@ pub async fn ws_handler(
|
||||
async fn handle_socket(socket: WebSocket, state: AppState, user: User) {
|
||||
let (mut sender, mut receiver) = socket.split();
|
||||
let (tx, mut rx) = mpsc::unbounded_channel::<Message>();
|
||||
let event_bus = state.event_bus.clone();
|
||||
|
||||
let client = GatewayClient::new(user, tx);
|
||||
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());
|
||||
|
||||
|
||||
+245
-2
@@ -1,10 +1,23 @@
|
||||
use event_bus::EventBus;
|
||||
use crate::models::category;
|
||||
use crate::models::channel;
|
||||
use crate::models::message;
|
||||
use crate::models::server;
|
||||
use crate::models::user::Model as User;
|
||||
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;
|
||||
use crate::routes::server::mapper::server_model_to_server_response;
|
||||
use axum::extract::ws::Message;
|
||||
use events::GatewayEvent;
|
||||
use parking_lot::RwLock;
|
||||
use std::collections::HashMap;
|
||||
use std::sync::Arc;
|
||||
use tokio::sync::mpsc;
|
||||
use tokio::task::JoinHandle;
|
||||
use uuid::Uuid;
|
||||
|
||||
pub mod events;
|
||||
pub mod handlers;
|
||||
pub mod routes;
|
||||
|
||||
@@ -18,7 +31,9 @@ pub struct GatewayManager {
|
||||
pub struct GatewayClient {
|
||||
user: User,
|
||||
connection_id: Uuid,
|
||||
sender: mpsc::UnboundedSender<Message>,
|
||||
pub sender: mpsc::UnboundedSender<Message>,
|
||||
pub event_bus: Arc<EventBus>,
|
||||
_event_handles: Vec<Arc<JoinHandle<()>>>,
|
||||
}
|
||||
|
||||
impl GatewayManager {
|
||||
@@ -40,15 +55,243 @@ impl GatewayManager {
|
||||
}
|
||||
|
||||
impl GatewayClient {
|
||||
pub fn new(user: User, sender: mpsc::UnboundedSender<Message>) -> Self {
|
||||
pub fn new(user: User, sender: mpsc::UnboundedSender<Message>, event_bus: Arc<EventBus>) -> Self {
|
||||
let connection_id = Uuid::new_v4();
|
||||
Self {
|
||||
user,
|
||||
connection_id,
|
||||
sender,
|
||||
event_bus,
|
||||
_event_handles: Vec::new(),
|
||||
}
|
||||
}
|
||||
|
||||
pub fn subscribe_to_events(&mut self) {
|
||||
let mut handles: Vec<Arc<JoinHandle<()>>> = Vec::new();
|
||||
|
||||
// ── Message ──────────────────────────────────────────────────────────
|
||||
let sender = self.sender.clone();
|
||||
handles.push(Arc::new(self.event_bus.on_async::<message::Model, _, _>(
|
||||
"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::Model, _, _>(
|
||||
"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()));
|
||||
}
|
||||
}
|
||||
},
|
||||
)));
|
||||
|
||||
let sender = self.sender.clone();
|
||||
handles.push(Arc::new(self.event_bus.on_async::<Uuid, _, _>(
|
||||
"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",
|
||||
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::Model, _, _>(
|
||||
"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()));
|
||||
}
|
||||
}
|
||||
},
|
||||
)));
|
||||
|
||||
let sender = self.sender.clone();
|
||||
handles.push(Arc::new(self.event_bus.on_async::<Uuid, _, _>(
|
||||
"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",
|
||||
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::Model, _, _>(
|
||||
"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()));
|
||||
}
|
||||
}
|
||||
},
|
||||
)));
|
||||
|
||||
let sender = self.sender.clone();
|
||||
handles.push(Arc::new(self.event_bus.on_async::<Uuid, _, _>(
|
||||
"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",
|
||||
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::Model, _, _>(
|
||||
"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::<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;
|
||||
}
|
||||
|
||||
async fn on_connect(&self) {
|
||||
tracing::info!("Client connected: {:?}", self.user);
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user