From f48334ef0c66b867d1e19720ab33c092287dec1f Mon Sep 17 00:00:00 2001 From: "NG (Graham)" Date: Thu, 24 Jul 2025 19:47:10 -0400 Subject: [PATCH] Remove user from queue on lobby disconnect --- rc_lobby_room/src/lobby.rs | 11 +++++++++++ rc_lobby_room/src/main.rs | 17 +++++++++++------ 2 files changed, 22 insertions(+), 6 deletions(-) diff --git a/rc_lobby_room/src/lobby.rs b/rc_lobby_room/src/lobby.rs index 100ce5a..4ce439a 100644 --- a/rc_lobby_room/src/lobby.rs +++ b/rc_lobby_room/src/lobby.rs @@ -143,4 +143,15 @@ impl QueueHandler { } } } + + pub async fn leave_queue(&self, user: &(dyn oj_rc_core::persist::user::LobbyUser + Send + Sync)) { + let user_id = user.user_id(); + for queue in self.users_in_queue.lock().await.values_mut() { + if let Some((i, _)) = queue.iter().enumerate().find(|(_, user)| user.user_id == user_id) { + queue.remove(i); + log::info!("User {} was removed from a queue", user_id); + break; + } + } + } } diff --git a/rc_lobby_room/src/main.rs b/rc_lobby_room/src/main.rs index 0cc208c..73cf7f6 100644 --- a/rc_lobby_room/src/main.rs +++ b/rc_lobby_room/src/main.rs @@ -20,7 +20,7 @@ pub struct InitConfig { pub queue: std::sync::Arc, } -pub type UserTy = oj_rc_core::UserState<()>; +pub type UserTy = std::sync::Arc>; #[tokio::main] async fn main() -> std::io::Result<()> { @@ -49,11 +49,11 @@ async fn main() -> std::io::Result<()> { if args.once { log::warn!("Handling first connection and then exiting"); let (socket, address) = listener.accept().await?; - process_socket(socket, address, server.clone(), init_ctx.users.clone()).await; + process_socket(socket, address, server.clone(), init_ctx.users.clone(), init_ctx.queue.clone()).await; } else { loop { let (socket, address) = listener.accept().await?; - tokio::spawn(process_socket(socket, address, server.clone(), init_ctx.users.clone())); + tokio::spawn(process_socket(socket, address, server.clone(), init_ctx.users.clone(), init_ctx.queue.clone())); } } server.join(); @@ -61,7 +61,7 @@ async fn main() -> std::io::Result<()> { Ok(()) } -async fn process_socket(mut socket: net::TcpStream, address: std::net::SocketAddr, server: std::sync::Arc>, users: std::sync::Arc) { +async fn process_socket(mut socket: net::TcpStream, address: std::net::SocketAddr, server: std::sync::Arc>, users: std::sync::Arc, queue: std::sync::Arc) { log::debug!("Accepting connection from address {}", address); let enc = match do_connect_handshake(&mut socket).await { Some(x) => x, @@ -72,9 +72,14 @@ async fn process_socket(mut socket: net::TcpStream, address: std::net::SocketAdd }; let (socket_r, socket_w) = socket.into_split(); let (chann_tx, chann_rx) = tokio::sync::mpsc::unbounded_channel(); - let user_state = oj_rc_core::UserState::<()>::new(users, chann_tx.clone()); + let user_state = std::sync::Arc::new(oj_rc_core::UserState::<()>::new(users, chann_tx.clone())); let ctx = polariton::packet::SerdesContext::from_boxed(Default::default(), enc); - server.handle_async_with_channel(socket_r, socket_w, user_state, ctx, chann_tx, chann_rx).await; + server.handle_async_with_channel_join(socket_r, socket_w, user_state.clone(), ctx, chann_tx, chann_rx).await; + if let Ok(user_info) = user_state.user() { + queue.leave_queue(user_info.as_ref().as_ref()).await; + } else { + log::info!("Unauthenticated user disconnected"); + } log::debug!("Goodbye connection from address {}", address); }