diff options
| author | metamuffin <metamuffin@disroot.org> | 2025-10-24 00:24:24 +0200 |
|---|---|---|
| committer | metamuffin <metamuffin@disroot.org> | 2025-10-24 00:24:26 +0200 |
| commit | 1e5dc0dee2fed17d6cc5c0e98edbb9b72daa6345 (patch) | |
| tree | 2309491b2f205f7bacac42c0760fc12b3092676f /server/src/main.rs | |
| parent | d53a9660c38e2d7a7b896675f4295b48b902506a (diff) | |
| download | hurrycurry-1e5dc0dee2fed17d6cc5c0e98edbb9b72daa6345.tar hurrycurry-1e5dc0dee2fed17d6cc5c0e98edbb9b72daa6345.tar.bz2 hurrycurry-1e5dc0dee2fed17d6cc5c0e98edbb9b72daa6345.tar.zst | |
Kick player with no inputs in last 60s; refactor server io
Diffstat (limited to 'server/src/main.rs')
| -rw-r--r-- | server/src/main.rs | 168 |
1 files changed, 79 insertions, 89 deletions
diff --git a/server/src/main.rs b/server/src/main.rs index e87d52b6..0236ebe4 100644 --- a/server/src/main.rs +++ b/server/src/main.rs @@ -27,12 +27,12 @@ use std::{ time::Duration, }; use tokio::{ - net::TcpListener, + net::{TcpListener, TcpStream}, spawn, - sync::{RwLock, broadcast, mpsc::channel}, + sync::{RwLock, broadcast}, time::interval, }; -use tokio_tungstenite::tungstenite::Message; +use tokio_tungstenite::{WebSocketStream, tungstenite::Message}; #[derive(Parser)] pub(crate) struct Args { @@ -150,7 +150,7 @@ async fn run(data_path: PathBuf, args: Args) -> anyhow::Result<()> { let ws_listener = TcpListener::bind(args.listen).await?; info!("Listening for websockets on {}", ws_listener.local_addr()?); - let (tx, rx) = broadcast::channel::<PacketC>(128 * 1024); + let (tx, _) = broadcast::channel::<PacketC>(128 * 1024); let mut state = Server::new(data_path, tx)?; state.load(state.index.generate_with_book("lobby")?, None); @@ -200,25 +200,18 @@ async fn run(data_path: PathBuf, args: Args) -> anyhow::Result<()> { for id in (1..).map(ConnectionID) { let (sock, addr) = ws_listener.accept().await?; - let Ok(sock) = tokio_tungstenite::accept_async(sock).await else { + let Ok(mut sock) = tokio_tungstenite::accept_async(sock).await else { warn!("Invalid ws handshake"); continue; }; info!("{id} Client connected ({addr})"); - let (mut write, mut read) = sock.split(); let state = state.clone(); - let mut rx = rx.resubscribe(); - let (error_tx, mut error_rx) = channel::<PacketC>(8); - - let init = state.write().await.connect(id).await; - - // let supports_binary = Arc::new(AtomicBool::new(false)); - // let supports_binary2 = supports_binary.clone(); + let (init, mut broadcast_rx, mut replies_rx) = state.write().await.connect(id).await; spawn(async move { for p in init { - if let Err(e) = write + if let Err(e) = sock .send(tokio_tungstenite::tungstenite::Message::Text( serde_json::to_string(&p).unwrap().into(), )) @@ -228,94 +221,93 @@ async fn run(data_path: PathBuf, args: Args) -> anyhow::Result<()> { return; } } + debug!("{id} client has caught up"); loop { - let Some(packet) = tokio::select!( - p = rx.recv() => Some(match p { - Ok(e) => e, + tokio::select! { + p = replies_rx.recv() => match p { + Some(packet) => { + if send_packet(id, &mut sock, packet).await { break; }; + }, + None => { + info!("{id} Server closed the connection"); + break; + } + }, + p = broadcast_rx.recv() => match p { + Ok(packet) => { + if send_packet(id, &mut sock, packet).await { break; }; + }, Err(e) => { - rx = rx.resubscribe(); + broadcast_rx = broadcast_rx.resubscribe(); warn!("{id} Client was lagging; resubscribed: {e}"); - PacketC::ServerMessage { + let packet = PacketC::ServerMessage { message: trm!("s.state.overflow_resubscribe"), error: true, - } - } - }), - p = error_rx.recv() => p - ) else { - info!("{id} Client outbound sender dropped. closing connection"); - break; - }; - // let message = if supports_binary.load(Ordering::Relaxed) { - // Message::Binary( - // bincode::encode_to_vec(&packet, BINCODE_CONFIG) - // .unwrap() - // .into(), - // ) - // } else { - let message = Message::Text(serde_json::to_string(&packet).unwrap().into()); - // }; - if let Err(e) = write.send(message).await { - warn!("{id} WebSocket error: {e}"); - break; - } - } - }); - - spawn(async move { - while let Some(Ok(message)) = read.next().await { - let packet = match message { - Message::Text(line) if line.len() < 8196 => match serde_json::from_str(&line) { - Ok(p) => p, - Err(e) => { - warn!("{id} Invalid json packet: {e}"); - break; + }; + if send_packet(id, &mut sock, packet).await { break; }; } }, - Message::Binary(_packet) => { - // supports_binary2.store(true, Ordering::Relaxed); - // match bincode::decode_from_slice::<PacketS, _>(&packet, BINCODE_CONFIG) { - // Ok((p, _size)) => p, - // Err(e) => { - // warn!("Invalid binary packet: {e}"); - // break; - // } - // } - continue; - } - Message::Close(_) => break, - _ => continue, - }; + Some(Ok(message)) = sock.next() => { + let packet = match message { + Message::Text(line) if line.len() < 8196 => match serde_json::from_str(&line) { + Ok(p) => p, + Err(e) => { + warn!("{id} Invalid json packet: {e}"); + break; + } + }, + Message::Binary(_packet) => continue, + Message::Close(_) => { + info!("{id} Client closed the connection"); + break + }, + _ => continue, + }; - if matches!( - packet, - PacketS::Movement { .. } | PacketS::ReplayTick { .. } - ) { - trace!("{id} <- {packet:?}"); - } else { - debug!("{id} <- {packet:?}"); - } - let packet_out = match state.write().await.packet_in_outer(id, packet).await { - Ok(packets) => packets, - Err(e) => { - warn!("Client error: {e}"); - vec![PacketC::ServerMessage { - message: e.into(), - error: true, - }] - } + if matches!( + packet, + PacketS::Movement { .. } | PacketS::ReplayTick { .. } + ) { + trace!("{id} <- {packet:?}"); + } else { + debug!("{id} <- {packet:?}"); + } + let packet_out = match state.write().await.packet_in_outer(id, packet) { + Ok(packets) => packets, + Err(e) => { + warn!("Client error: {e}"); + vec![PacketC::ServerMessage { + message: e.into(), + error: true, + }] + } + }; + for packet in packet_out { + if send_packet(id, &mut sock, packet).await { break; }; + } + } }; - for packet in packet_out { - let _ = error_tx.send(packet).await; - } } - info!("{id} Client disconnected"); - let _ = state.write().await.disconnect(id).await; + let _ = state.write().await.disconnect(id); }); } Ok(()) } +async fn send_packet( + id: ConnectionID, + sock: &mut WebSocketStream<TcpStream>, + packet: PacketC, +) -> bool { + let message = Message::Text(serde_json::to_string(&packet).unwrap().into()); + if let Err(e) = sock.send(message).await { + warn!("{id} WebSocket error: {e}"); + true + } else { + false + } +} + #[cfg(test)] mod test { use hurrycurry_protocol::{Character, PacketS, PlayerClass, PlayerID}; @@ -386,11 +378,9 @@ mod test { position: None, }, ) - .await .unwrap(); assert!( s.packet_in_outer(ConnectionID(conn.try_into().unwrap()), p) - .await .is_err(), "test {}", conn, |