From ca4a0ce6759b7acbf3a86430d40a562fca6c4183 Mon Sep 17 00:00:00 2001 From: "NGnius (Graham)" Date: Sat, 8 Feb 2025 11:57:23 -0500 Subject: [PATCH] Allow concurrent connections on the load balancer --- assets/robocraft/service_packets.md | 6 +++++- rc_services/src/main.rs | 19 ++++++++++++++++--- rc_services/src/state.rs | 2 +- 3 files changed, 22 insertions(+), 5 deletions(-) diff --git a/assets/robocraft/service_packets.md b/assets/robocraft/service_packets.md index 233f247..bcbe8cd 100644 --- a/assets/robocraft/service_packets.md +++ b/assets/robocraft/service_packets.md @@ -32,7 +32,7 @@ Connect Data 41 bytes 0..1 =240 1..5 current tick count -### Modes +### Operation Request & Response 220 Region info (in) 222 Friends list info (in) @@ -60,6 +60,10 @@ Connect Data 41 bytes #### Op code always??? 221 Auth token (string) +### Event + +14 Concurrent user check passed, i.e. this room still has space for you + diff --git a/rc_services/src/main.rs b/rc_services/src/main.rs index dba4ede..26c515b 100644 --- a/rc_services/src/main.rs +++ b/rc_services/src/main.rs @@ -18,11 +18,15 @@ async fn main() -> std::io::Result<()> { log::debug!("Got cli args {:?}", args); let ip_addr: std::net::IpAddr = args.ip.parse().expect("Invalid IP address"); + // memory leak, but only once (so not a big deal) + let redirect_static = Box::leak(Box::new(args.redirect.clone())); + let room_name_static = Box::leak(Box::new(args.room_name.clone())); + let listener = net::TcpListener::bind(std::net::SocketAddr::new(ip_addr, args.port)).await?; loop { let (socket, address) = listener.accept().await?; - process_socket(socket, address, NonZero::new(args.retries), &args.redirect, &args.room_name).await; + tokio::spawn(process_socket(socket, address, NonZero::new(args.retries), redirect_static, room_name_static)); } } @@ -37,10 +41,19 @@ async fn process_socket(mut socket: net::TcpStream, address: std::net::SocketAdd return; } }; - buf.clear(); + let sock_state = state::State::new(enc); + while let Ok(packet) = receive_packet(&mut buf, &mut socket, retries, sock_state.binrw_args()).await { + match packet { + Packet::Ping(ping) => { + handle_ping(ping, &mut buf, &mut socket).await; + }, + Packet::Packet(packet) => log::warn!("Not handling packet {:?}", packet), + } + } } async fn handle_ping(ping: Ping, buf: &mut Vec, socket: &mut net::TcpStream) { + buf.clear(); let resp = Packet::Ping(polariton_auth::ping_pong(ping)); resp.to_buf(buf, None).unwrap(); let write_count = socket.write(buf).await.unwrap(); @@ -57,7 +70,7 @@ async fn read_more(buf: &mut Vec, socket: &mut net::TcpStream) -> Result, socket: &mut net::TcpStream, max_retries: Option>, args: Option>>) -> Result { +async fn receive_packet(buf: &mut Vec, socket: &mut net::TcpStream, max_retries: Option>, args: Option>>) -> Result { buf.clear(); let read_count = read_more(buf, socket).await?; if read_count == 0 { return Err(std::io::Error::new(std::io::ErrorKind::InvalidData, "socket did not read any bytes")); } // bad packet diff --git a/rc_services/src/state.rs b/rc_services/src/state.rs index 2bd18ab..cf7f73d 100644 --- a/rc_services/src/state.rs +++ b/rc_services/src/state.rs @@ -5,7 +5,7 @@ pub struct State { } impl State { - pub fn new(c: Box>) -> Self { + pub fn new(c: Box>) -> Self { Self { crypto: c, }