Init
This commit is contained in:
@@ -42,38 +42,39 @@ 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();
|
||||
let (tx, mut rx) = mpsc::unbounded_channel::<Message>();
|
||||
let event_bus = state.event_bus.clone();
|
||||
|
||||
let client = GatewayClient::new(user);
|
||||
let client = GatewayClient::new(user, tx);
|
||||
client.on_connect().await;
|
||||
state.gateway.add_user(client);
|
||||
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 {
|
||||
// if sender.send(message).await.is_err() {
|
||||
// break;
|
||||
// }
|
||||
// }
|
||||
// });
|
||||
// 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 {
|
||||
if sender.send(message).await.is_err() {
|
||||
break;
|
||||
}
|
||||
}
|
||||
});
|
||||
//
|
||||
// // Task pour recevoir les messages du WebSocket
|
||||
// let state_clone = state.clone();
|
||||
// let mut recv_task = tokio::spawn(async move {
|
||||
// while let Some(Ok(message)) = receiver.next().await {
|
||||
// on_message(user_id, message, &state_clone).await;
|
||||
// }
|
||||
// });
|
||||
// Task pour recevoir les messages du WebSocket
|
||||
let client_clone = client.clone();
|
||||
let state_clone = state.clone();
|
||||
let mut recv_task = tokio::spawn(async move {
|
||||
while let Some(Ok(message)) = receiver.next().await {
|
||||
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(),
|
||||
// };
|
||||
// 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(),
|
||||
};
|
||||
//
|
||||
// // Déconnexion (Disconnect)
|
||||
// on_disconnect(user_id, &state).await;
|
||||
|
||||
Reference in New Issue
Block a user