This commit is contained in:
2026-06-21 19:11:23 +02:00
parent 3fd2b8ade8
commit e9d77fbd05
11 changed files with 96 additions and 77 deletions
+5 -50
View File
@@ -1,4 +1,3 @@
use crate::auth::token::verify_jwt;
use crate::core::AppState;
use crate::http::context::CurrentUser;
use crate::models::user::Model as User;
@@ -13,7 +12,6 @@ use axum::{
use futures_util::{sink::SinkExt, stream::StreamExt};
use serde::Deserialize;
use tokio::sync::mpsc;
use uuid::Uuid;
#[derive(Deserialize)]
pub struct WsQuery {
@@ -49,9 +47,6 @@ async fn handle_socket(socket: WebSocket, state: AppState, user: User) {
client.on_connect().await;
state.gateway.add_client(client.clone());
// // Enregistrement du client (Connect)
// on_connect(user_id, tx, &state).await;
//
// Task pour envoyer les messages du canal mpsc vers le WebSocket
let mut send_task = tokio::spawn(async move {
while let Some(message) = rx.recv().await {
@@ -60,7 +55,7 @@ async fn handle_socket(socket: WebSocket, state: AppState, user: User) {
}
}
});
//
// Task pour recevoir les messages du WebSocket
let client_clone = client.clone();
let state_clone = state.clone();
@@ -69,54 +64,14 @@ async fn handle_socket(socket: WebSocket, state: AppState, user: User) {
client_clone.on_message(message).await;
}
});
//
// Attente de la fin d'une des tâches (déconnexion)
tokio::select! {
_ = (&mut send_task) => recv_task.abort(),
_ = (&mut recv_task) => send_task.abort(),
};
//
state.gateway.remove_client(client.clone());
// // Déconnexion (Disconnect)
// on_disconnect(user_id, &state).await;
}
pub async fn on_connect(user_id: Uuid, tx: mpsc::UnboundedSender<Message>, state: &AppState) {
tracing::info!("Client connected: {}", user_id);
let mut clients = state
.gateway
.clients
.write()
.expect("Failed to lock clients for writing");
clients.insert(user_id, GatewayClient { user_id, tx });
}
pub async fn on_disconnect(user_id: Uuid, state: &AppState) {
tracing::info!("Client disconnected: {}", user_id);
let mut clients = state
.gateway
.clients
.write()
.expect("Failed to lock clients for writing");
clients.remove(&user_id);
}
pub async fn on_message(user_id: Uuid, message: Message, _state: &AppState) {
tracing::debug!("Message received from {}: {:?}", user_id, message);
// Exemple d'utilisation de l'état/repositories
// let user_opt = state.repositories.user.get_by_id(user_id).await.ok().flatten();
match message {
Message::Text(text) => {
tracing::info!("Received text from {}: {}", user_id, text);
// Logique de dispatch ou de traitement ici
}
Message::Binary(_) => {
tracing::info!("Received binary from {}", user_id);
}
Message::Close(_) => {
tracing::info!("Received close from {}", user_id);
}
_ => {}
}
client.on_disconnect().await;
}
+6 -2
View File
@@ -49,9 +49,13 @@ impl GatewayClient {
}
}
async fn on_connect(&self) {}
async fn on_connect(&self) {
tracing::info!("Client connected: {:?}", self.user);
}
async fn on_disconnect(&self) {}
async fn on_disconnect(&self) {
tracing::info!("Client disconnected: {:?}", self.user);
}
async fn on_message(&self, message: Message) {
match message {
+4 -2
View File
@@ -31,10 +31,12 @@ pub fn router() -> OxRouter {
let api_routes = Router::new()
.merge(secure_routes)
.merge(auth::routes::router())
.merge(core::routes::router())
.merge(gateway::routes::router());
.merge(core::routes::router());
let ws_routes = Router::new().merge(gateway::routes::router());
Router::new()
.nest("/api", api_routes)
.nest("/ws", ws_routes)
.merge(SwaggerUi::new("/swagger").url("/api-docs/openapi.json", openapi::ApiDoc::openapi()))
}