From a0d3a71ba8b0beba68c007ec09c8a7602936d0c0 Mon Sep 17 00:00:00 2001 From: kralverde <80051564+kralverde@users.noreply.github.com> Date: Tue, 18 Feb 2025 05:10:09 -1000 Subject: [PATCH] Gracefully exit on stop, new logger and rustyline (#557) * initial commit * fix linter * save remaining chunks in memory * fix colors not working when terminal prompt is enabled * fix typo --- pumpkin-world/src/level.rs | 8 + pumpkin/Cargo.toml | 9 +- pumpkin/src/command/commands/stop.rs | 14 +- pumpkin/src/entity/player.rs | 10 +- pumpkin/src/lib.rs | 237 +++++++++++++++++++++------ pumpkin/src/main.rs | 6 +- pumpkin/src/net/mod.rs | 34 +++- pumpkin/src/net/packet/status.rs | 2 +- pumpkin/src/server/mod.rs | 14 +- pumpkin/src/server/ticker.rs | 5 +- 10 files changed, 254 insertions(+), 85 deletions(-) diff --git a/pumpkin-world/src/level.rs b/pumpkin-world/src/level.rs index 7d8f09e79..fc5b1dfaa 100644 --- a/pumpkin-world/src/level.rs +++ b/pumpkin-world/src/level.rs @@ -97,6 +97,14 @@ impl Level { log::info!("Saving level..."); // chunks are automatically saved when all players get removed + // TODO: Await chunks that have been called by this ^ + + // save all stragling chunks + for chunk in self.loaded_chunks.iter() { + let pos = chunk.key(); + let chunk = chunk.value(); + self.write_chunk((*pos, chunk.clone())).await; + } // then lets save the world info self.world_info_writer diff --git a/pumpkin/Cargo.toml b/pumpkin/Cargo.toml index 07ce10914..e513da7b8 100644 --- a/pumpkin/Cargo.toml +++ b/pumpkin/Cargo.toml @@ -46,7 +46,7 @@ rand = "0.8" num-bigint = "0.4" # Console line reading -rustyline = "15.0" +rustyline-async = "0.4.5" # encryption rsa = "0.9" @@ -72,11 +72,12 @@ base64 = "0.22" png = "0.17" # logging -simple_logger = { version = "5.0", features = ["threads"] } +simplelog = { version = "0.12.2", features = ["ansi_term"] } + # Remove time in favor of chrono? time = "0.3" -chrono = { version = "0.4", features = ["serde"]} +chrono = { version = "0.4", features = ["serde"] } # plugins libloading = "0.8" @@ -85,4 +86,4 @@ libloading = "0.8" git-version = "0.3" # This makes it so the entire project doesn't recompile on each build on linux. [target.'cfg(target_os = "windows")'.build-dependencies] -tauri-winres= "0.3" +tauri-winres = "0.3" diff --git a/pumpkin/src/command/commands/stop.rs b/pumpkin/src/command/commands/stop.rs index 994ae1e12..c8891b1e1 100644 --- a/pumpkin/src/command/commands/stop.rs +++ b/pumpkin/src/command/commands/stop.rs @@ -5,6 +5,7 @@ use pumpkin_util::text::TextComponent; use crate::command::args::ConsumedArgs; use crate::command::tree::CommandTree; use crate::command::{CommandError, CommandExecutor, CommandSender}; +use crate::stop_server; const NAMES: [&str; 1] = ["stop"]; @@ -17,7 +18,7 @@ impl CommandExecutor for StopExecutor { async fn execute<'a>( &self, sender: &mut CommandSender<'a>, - server: &crate::server::Server, + _server: &crate::server::Server, _args: &ConsumedArgs<'a>, ) -> Result<(), CommandError> { sender @@ -25,15 +26,8 @@ impl CommandExecutor for StopExecutor { TextComponent::translate("commands.stop.stopping", []).color_named(NamedColor::Red), ) .await; - - // TODO: Gracefully stop - - let kick_message = TextComponent::text("Server stopped"); - for player in server.get_all_players().await { - player.kick(kick_message.clone()).await; - } - server.save().await; - std::process::exit(0) + stop_server(); + Ok(()) } } diff --git a/pumpkin/src/entity/player.rs b/pumpkin/src/entity/player.rs index 81221c3aa..f05772530 100644 --- a/pumpkin/src/entity/player.rs +++ b/pumpkin/src/entity/player.rs @@ -815,17 +815,19 @@ impl Player { return; } - self.client + let _ = self + .client .try_send_packet(&CPlayDisconnect::new(&reason)) - .await - .unwrap_or_else(|_| self.client.close()); + .await; + log::info!( "Kicked Player {} ({}) for {}", self.gameprofile.name, self.client.id, reason.to_pretty_console() ); - self.client.close(); + + self.client.close().await; } pub fn can_food_heal(&self) -> bool { diff --git a/pumpkin/src/lib.rs b/pumpkin/src/lib.rs index 4525e049e..d30ce5503 100644 --- a/pumpkin/src/lib.rs +++ b/pumpkin/src/lib.rs @@ -3,16 +3,23 @@ use crate::net::{lan_broadcast, query, rcon::RCONServer, Client}; use crate::server::{ticker::Ticker, Server}; -use log::{LevelFilter, Log}; +use log::{logger, Level, LevelFilter, Log}; +use net::PacketHandlerState; use plugin::PluginManager; use pumpkin_config::{ADVANCED_CONFIG, BASIC_CONFIG}; use pumpkin_util::text::TextComponent; -use rustyline::DefaultEditor; -use simple_logger::SimpleLogger; +use rustyline_async::{Readline, ReadlineEvent}; +use std::collections::HashMap; +use std::str::FromStr; +use std::sync::atomic::AtomicBool; +use std::sync::OnceLock; use std::{ net::SocketAddr, sync::{Arc, LazyLock}, }; +use tokio::select; +use tokio::sync::Notify; +use tokio::task::JoinHandle; use tokio::{ io::{AsyncReadExt, AsyncWriteExt}, net::{tcp::OwnedReadHalf, TcpListener}, @@ -35,28 +42,53 @@ const GIT_VERSION: &str = env!("GIT_VERSION"); pub static PLUGIN_MANAGER: LazyLock> = LazyLock::new(|| Mutex::new(PluginManager::new())); +// Yucky, is there a way to do this better? revisit our static LOGGER_IMPL? +static _INPUT_HOLDER: OnceLock>> = OnceLock::new(); + pub static LOGGER_IMPL: LazyLock, LevelFilter)>> = LazyLock::new(|| { if ADVANCED_CONFIG.logging.enabled { - let mut logger = SimpleLogger::new(); + let mut config = simplelog::ConfigBuilder::new(); if ADVANCED_CONFIG.logging.timestamp { - logger = logger.with_timestamp_format(time::macros::format_description!( + config.set_time_format_custom(time::macros::format_description!( "[year]-[month]-[day] [hour]:[minute]:[second]" )); + config.set_time_level(LevelFilter::Trace); } else { - logger = logger.without_timestamps(); + config.set_time_level(LevelFilter::Off); } - let logger = logger - .with_level(LevelFilter::Info) - .with_colors(ADVANCED_CONFIG.logging.color) - .with_threads(ADVANCED_CONFIG.logging.threads) - .env(); + if !ADVANCED_CONFIG.logging.color { + for level in Level::iter() { + config.set_level_color(level, None); + } + } else if ADVANCED_CONFIG.commands.use_console { + // We are technically logging to a file like object + config.set_write_log_enable_colors(true); + } - // Incase environment variables change it - let max_level = logger.max_level(); + if !ADVANCED_CONFIG.logging.threads { + config.set_thread_level(LevelFilter::Off); + } else { + config.set_thread_level(LevelFilter::Info); + } - Some((Box::new(logger), max_level)) + let level = std::env::var("RUST_LOG") + .ok() + .as_deref() + .map(LevelFilter::from_str) + .and_then(Result::ok) + .unwrap_or(LevelFilter::Info); + + if ADVANCED_CONFIG.commands.use_console { + let (rl, stdout) = Readline::new("$ ".to_owned()).unwrap(); + let logger = simplelog::WriteLogger::new(level, config.build(), stdout); + let _ = _INPUT_HOLDER.set(Mutex::new(Some(rl))); + Some((Box::new(logger), level)) + } else { + let logger = simplelog::SimpleLogger::new(level, config.build()); + Some((Box::new(logger), level)) + } } else { None } @@ -72,10 +104,20 @@ macro_rules! init_log { }; } +pub static SHOULD_STOP: AtomicBool = AtomicBool::new(false); +pub static STOP_INTERRUPT: LazyLock = LazyLock::new(Notify::new); + +pub fn stop_server() { + SHOULD_STOP.store(true, std::sync::atomic::Ordering::Relaxed); + STOP_INTERRUPT.notify_waiters(); +} + pub struct PumpkinServer { pub server: Arc, pub listener: TcpListener, pub server_addr: SocketAddr, + readline_handle: Option>, + tasks_to_await: Vec>, } impl PumpkinServer { @@ -91,19 +133,16 @@ impl PumpkinServer { .local_addr() .expect("Unable to get the address of server!"); - let pumpkin_server = Self { - server: server.clone(), - listener, - server_addr: addr, - }; - - let use_console = ADVANCED_CONFIG.commands.use_console; let rcon = ADVANCED_CONFIG.networking.rcon.clone(); let mut ticker = Ticker::new(BASIC_CONFIG.tps); - if use_console { - setup_console(server.clone()); + let mut readline = None; + if let Some(rt) = _INPUT_HOLDER.get() { + let mut rt = rt.lock().await; + let rt = rt.take().unwrap(); + let handle = setup_console(rt, server.clone()); + readline = Some(handle); } if rcon.enabled { @@ -123,15 +162,23 @@ impl PumpkinServer { tokio::spawn(lan_broadcast::start_lan_broadcast(addr)); } + let mut tasks_to_await = Vec::new(); // Ticker { let server = server.clone(); - tokio::spawn(async move { + let handle = tokio::spawn(async move { ticker.run(&server).await; - }) + }); + tasks_to_await.push(handle); }; - pumpkin_server + Self { + server: server.clone(), + listener, + server_addr: addr, + readline_handle: readline, + tasks_to_await, + } } pub async fn init_plugins(&self) { @@ -142,11 +189,25 @@ impl PumpkinServer { }; } - pub async fn start(&self) { - let mut master_client_id: u16 = 0; - loop { + pub async fn start(self) { + let mut master_client_id: usize = 0; + let tasks = Arc::new(Mutex::new(HashMap::new())); + + while !SHOULD_STOP.load(std::sync::atomic::Ordering::Relaxed) { + let await_new_client = || async { + let t1 = self.listener.accept(); + let t2 = STOP_INTERRUPT.notified(); + + select! { + client = t1 => Some(client.unwrap()), + () = t2 => None, + } + }; + // Asynchronously wait for an inbound socket. - let (connection, client_addr) = self.listener.accept().await.unwrap(); + let Some((connection, client_addr)) = await_new_client().await else { + break; + }; if let Err(e) = connection.set_nodelay(true) { log::warn!("failed to set TCP_NODELAY {e}"); @@ -174,20 +235,30 @@ impl PumpkinServer { let client = Arc::new(Client::new(tx, client_addr, id)); let client_clone = client.clone(); + // This task will be cleaned up on its own tokio::spawn(async move { - while (rx.recv().await).is_some() { - let mut enc = client_clone.enc.lock().await; - let buf = enc.take(); - if let Err(e) = connection_writer.lock().await.write_all(&buf).await { - log::warn!("Failed to write packet to client: {e}"); - client_clone.close(); - break; + // We clone ownership of a tx into here thru the client so this will never drop + // since there is always a tx in memory. We need to explicitly tell the recv to stop + while let Some(notif) = rx.recv().await { + match notif { + PacketHandlerState::PacketReady => { + let mut enc = client_clone.enc.lock().await; + let buf = enc.take(); + if let Err(e) = connection_writer.lock().await.write_all(&buf).await { + log::warn!("Failed to write packet to client: {e}"); + client_clone.close().await; + break; + } + } + PacketHandlerState::Stop => break, } } }); let server = self.server.clone(); - tokio::spawn(async move { + let tasks_clone = tasks.clone(); + // We need to await these to verify all cleanup code is complete + let handle = tokio::spawn(async move { while !client.closed.load(std::sync::atomic::Ordering::Relaxed) && !client .make_player @@ -221,34 +292,96 @@ impl PumpkinServer { log::debug!("Cleaning up player for id {}", id); player.remove().await; server.remove_player().await; + tasks_clone.lock().await.remove(&id); } }); + tasks.lock().await.insert(id, Some(handle)); + } + // Keep it in scope for logging + let rl = match self.readline_handle { + Some(rl) => Some(rl.await.unwrap()), + None => None, + }; + + log::info!("Stopped accepting incoming connections"); + + let kick_message = TextComponent::text("Server stopped"); + for player in self.server.get_all_players().await { + player.kick(kick_message.clone()).await; + } + + log::info!("Ending server tasks"); + + for handle in self.tasks_to_await.into_iter() { + if let Err(err) = handle.await { + log::error!("Failed to join server task: {}", err.to_string()); + } + } + + let handles: Vec>> = tasks + .lock() + .await + .values_mut() + .map(|val| val.take()) + .collect(); + + log::info!("Ending player tasks"); + + for handle in handles.into_iter().flatten() { + if let Err(err) = handle.await { + log::error!("Failed to join player task: {}", err.to_string()); + } + } + + self.server.save().await; + + log::info!("Completed save!"); + logger().flush(); + if let Some(mut rl) = rl { + let _ = rl.flush(); } } } -fn setup_console(server: Arc) { +fn setup_console(rl: Readline, server: Arc) -> JoinHandle { + // This needs to be async or it will hog a thread tokio::spawn(async move { - let mut rl = DefaultEditor::new().unwrap(); - loop { - // maybe put this into config ? - let readline = rl.readline("$ "); + let mut rl = rl; + while !SHOULD_STOP.load(std::sync::atomic::Ordering::Relaxed) { + let t1 = rl.readline(); + let t2 = STOP_INTERRUPT.notified(); - match readline { - Ok(line) => { - rl.add_history_entry(line.as_str()).unwrap(); + let result = select! { + line = t1 => Some(line), + () = t2 => None, + }; + + let Some(result) = result else { break }; + + match result { + Ok(ReadlineEvent::Line(line)) => { let dispatcher = server.command_dispatcher.read().await; + dispatcher .handle_command(&mut command::CommandSender::Console, &server, &line) .await; + rl.add_history_entry(line).unwrap(); } - Err(_) => { - // TODO: we can handle CTRL+C and stuff here + Ok(ReadlineEvent::Interrupted) => { + stop_server(); + break; + } + err => { + log::error!("Console command loop failed!"); + log::error!("{:?}", err); break; } } } - }); + log::debug!("Stopped console commands task"); + let _ = rl.flush(); + rl + }) } async fn poll(client: &Client, connection_reader: Arc>) -> bool { @@ -268,7 +401,7 @@ async fn poll(client: &Client, connection_reader: Arc>) -> Ok(None) => (), //log::debug!("Waiting for more data to complete packet..."), Err(err) => { log::warn!("Failed to decode packet for: {}", err.to_string()); - client.close(); + client.close().await; return false; // return to avoid reserving additional bytes } } @@ -281,13 +414,13 @@ async fn poll(client: &Client, connection_reader: Arc>) -> Ok(cnt) => { //log::debug!("Read {} bytes", cnt); if cnt == 0 { - client.close(); + client.close().await; return false; } } Err(error) => { log::error!("Error while reading incoming packet {}", error); - client.close(); + client.close().await; return false; } }; diff --git a/pumpkin/src/main.rs b/pumpkin/src/main.rs index 5aea5ae37..4908a7bce 100644 --- a/pumpkin/src/main.rs +++ b/pumpkin/src/main.rs @@ -48,7 +48,7 @@ use tokio::signal::unix::{signal, SignalKind}; use tokio::sync::Mutex; use crate::server::CURRENT_MC_VERSION; -use pumpkin::{init_log, PumpkinServer}; +use pumpkin::{init_log, stop_server, PumpkinServer, SHOULD_STOP}; use pumpkin_protocol::CURRENT_MC_PROTOCOL; use pumpkin_util::text::{color::NamedColor, TextComponent}; use std::time::Instant; @@ -84,6 +84,7 @@ async fn main() { std::panic::set_hook(Box::new(move |info| { default_panic(info); // TODO: Gracefully exit? + // we need to abide by the panic rules here std::process::exit(1); })); @@ -121,6 +122,7 @@ async fn main() { ); pumpkin_server.start().await; + log::info!("The server has stopped."); } fn handle_interrupt() { @@ -130,7 +132,7 @@ fn handle_interrupt() { .color_named(NamedColor::Red) .to_pretty_console() ); - std::process::exit(0); + stop_server(); } // Non-UNIX Ctrl-C handling diff --git a/pumpkin/src/net/mod.rs b/pumpkin/src/net/mod.rs index 541b8f536..6ad287e42 100644 --- a/pumpkin/src/net/mod.rs +++ b/pumpkin/src/net/mod.rs @@ -108,12 +108,17 @@ impl Default for PlayerConfig { } } +pub enum PacketHandlerState { + PacketReady, + Stop, +} + /// Everything which makes a Connection with our Server is a `Client`. /// Client will become Players when they reach the `Play` state pub struct Client { /// The client id. This is good for coorelating a connection with a player /// Only used for logging purposes - pub id: u16, + pub id: usize, /// The client's game profile information. pub gameprofile: Mutex>, /// The client's configuration settings, Optional @@ -135,7 +140,7 @@ pub struct Client { /// The packet decoder for incoming packets. pub dec: Arc>, /// A channel for sending packets to the client. - pub server_packets_channel: mpsc::Sender<()>, + pub server_packets_channel: mpsc::Sender, /// A queue of raw packets received from the client, waiting to be processed. pub client_packets_queue: Arc>>, /// Indicates whether the client should be converted into a player. @@ -144,7 +149,11 @@ pub struct Client { impl Client { #[must_use] - pub fn new(server_packets_channel: mpsc::Sender<()>, address: SocketAddr, id: u16) -> Self { + pub fn new( + server_packets_channel: mpsc::Sender, + address: SocketAddr, + id: usize, + ) -> Self { Self { id, protocol_version: AtomicI32::new(0), @@ -249,7 +258,10 @@ impl Client { return; } - let _ = self.server_packets_channel.send(()).await; + let _ = self + .server_packets_channel + .send(PacketHandlerState::PacketReady) + .await; /* let mut writer = self.connection_writer.lock().await; if let Err(error) = writer.write_all(&enc.take()).await { @@ -296,7 +308,10 @@ impl Client { let mut enc = self.enc.lock().await; enc.append_packet(packet)?; - let _ = self.server_packets_channel.send(()).await; + let _ = self + .server_packets_channel + .send(PacketHandlerState::PacketReady) + .await; /* let mut writer = self.connection_writer.lock().await; let _ = writer.write_all(&enc.take()).await; */ @@ -538,7 +553,7 @@ impl Client { log::warn!("Failed to kick {}: {}", self.id, err.to_string()); } log::debug!("Closing connection for {}", self.id); - self.close(); + self.close().await; } /// Checks if the client can join the server. @@ -596,9 +611,14 @@ impl Client { /// # Notes /// /// This function does not attempt to send any disconnect packets to the client. - pub fn close(&self) { + pub async fn close(&self) { self.closed .store(true, std::sync::atomic::Ordering::Relaxed); + // We dont care if this fails because if it doesn that means the task has already stopped + let _ = self + .server_packets_channel + .send(PacketHandlerState::Stop) + .await; log::debug!("Closed connection for {}", self.id); } } diff --git a/pumpkin/src/net/packet/status.rs b/pumpkin/src/net/packet/status.rs index 0e9790c29..d92719bde 100644 --- a/pumpkin/src/net/packet/status.rs +++ b/pumpkin/src/net/packet/status.rs @@ -13,6 +13,6 @@ impl Client { log::debug!("Handling ping request"); self.send_packet(&CPingResponse::new(ping_request.payload)) .await; - self.close(); + self.close().await; } } diff --git a/pumpkin/src/server/mod.rs b/pumpkin/src/server/mod.rs index ab9d58398..044045388 100644 --- a/pumpkin/src/server/mod.rs +++ b/pumpkin/src/server/mod.rs @@ -112,10 +112,8 @@ impl Server { ); // Spawn chunks are never unloaded - for x in -1..=1 { - for z in -1..=1 { - world.level.mark_chunk_as_newly_watched(Vector2::new(x, z)); - } + for chunk in Self::spawn_chunks() { + world.level.mark_chunk_as_newly_watched(chunk); } Self { @@ -144,6 +142,14 @@ impl Server { } } + const SPAWN_CHUNK_RADIUS: i32 = 1; + + pub fn spawn_chunks() -> impl Iterator> { + (-Self::SPAWN_CHUNK_RADIUS..=Self::SPAWN_CHUNK_RADIUS).flat_map(|x| { + (-Self::SPAWN_CHUNK_RADIUS..=Self::SPAWN_CHUNK_RADIUS).map(move |z| Vector2::new(x, z)) + }) + } + /// Adds a new player to the server. /// /// This function takes an `Arc` representing the connected client and performs the following actions: diff --git a/pumpkin/src/server/ticker.rs b/pumpkin/src/server/ticker.rs index ea7d32847..773710f94 100644 --- a/pumpkin/src/server/ticker.rs +++ b/pumpkin/src/server/ticker.rs @@ -2,6 +2,8 @@ use std::time::{Duration, Instant}; use tokio::time::sleep; +use crate::SHOULD_STOP; + use super::Server; pub struct Ticker { @@ -20,7 +22,7 @@ impl Ticker { /// IMPORTANT: Run this in a new thread/tokio task pub async fn run(&mut self, server: &Server) { - loop { + while !SHOULD_STOP.load(std::sync::atomic::Ordering::Relaxed) { let now = Instant::now(); let elapsed = now - self.last_tick; @@ -33,5 +35,6 @@ impl Ticker { sleep(sleep_time).await; } } + log::debug!("Ticker stopped"); } }