aboutsummaryrefslogtreecommitdiff
path: root/server/src/main.rs
diff options
context:
space:
mode:
authormetamuffin <metamuffin@disroot.org>2025-10-24 00:24:24 +0200
committermetamuffin <metamuffin@disroot.org>2025-10-24 00:24:26 +0200
commit1e5dc0dee2fed17d6cc5c0e98edbb9b72daa6345 (patch)
tree2309491b2f205f7bacac42c0760fc12b3092676f /server/src/main.rs
parentd53a9660c38e2d7a7b896675f4295b48b902506a (diff)
downloadhurrycurry-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.rs168
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,