From e1117ce4a68e3bffc26c8ce5da95c1c97ddf27af Mon Sep 17 00:00:00 2001 From: Alexander Medvedev Date: Mon, 3 Aug 2026 21:16:07 +0200 Subject: [PATCH] chore: nicer event firing --- pumpkin-macros/src/lib.rs | 136 +++++++------ pumpkin-protocol/src/java/packet_decoder.rs | 3 - .../src/cylindrical_chunk_iterator.rs | 31 ++- pumpkin/src/block/blocks/plant/crop/mod.rs | 4 +- .../blocks/redstone/pressure_plate/mod.rs | 4 +- pumpkin/src/block/registry.rs | 8 +- pumpkin/src/command/commands/me.rs | 11 +- pumpkin/src/command/commands/plugin.rs | 22 +- pumpkin/src/command/commands/say.rs | 11 +- pumpkin/src/command/commands/setworldspawn.rs | 9 +- pumpkin/src/command/dispatcher.rs | 11 +- pumpkin/src/entity/player.rs | 46 ++--- pumpkin/src/entity/projectile/egg.rs | 4 +- pumpkin/src/lib.rs | 70 ++++--- pumpkin/src/net/bedrock/login.rs | 2 +- pumpkin/src/net/bedrock/mod.rs | 8 +- pumpkin/src/net/bedrock/play.rs | 39 ++-- pumpkin/src/net/java/mod.rs | 26 ++- pumpkin/src/net/java/play.rs | 71 ++++--- pumpkin/src/plugin/api/context.rs | 2 +- pumpkin/src/plugin/mod.rs | 192 ++++++++---------- pumpkin/src/server/mod.rs | 7 +- pumpkin/src/server/ticker.rs | 11 +- pumpkin/src/world/chunker.rs | 19 +- pumpkin/src/world/mod.rs | 90 ++++---- 25 files changed, 437 insertions(+), 400 deletions(-) diff --git a/pumpkin-macros/src/lib.rs b/pumpkin-macros/src/lib.rs index 8462cf465..971275384 100644 --- a/pumpkin-macros/src/lib.rs +++ b/pumpkin-macros/src/lib.rs @@ -1,12 +1,12 @@ use heck::ToShoutySnakeCase; use proc_macro::TokenStream; -use proc_macro_error2::{abort, abort_call_site, proc_macro_error}; +use proc_macro_error2::{abort, abort_call_site}; use pumpkin_data::tag::{RegistryKey, get_tag_ids}; use pumpkin_data::{Block, BlockId}; use quote::{format_ident, quote}; use syn::spanned::Spanned; use syn::{self, Attribute, DeriveInput, LitStr, Type, parse_quote}; -use syn::{Expr, Field, Fields, ItemStruct, Stmt, parse_macro_input}; +use syn::{Block as SynBlock, Expr, Field, Fields, ItemStruct, Stmt, parse_macro_input}; /// Derives the `Payload` trait for an event struct, enabling it to be used in the plugin system. /// @@ -89,87 +89,97 @@ pub fn cancellable(_args: TokenStream, input: TokenStream) -> TokenStream { .into() } -/// Sends a cancellable event asynchronously. +/// Sends a cancellable event through the plugin manager. /// -/// # Arguments -/// - `input` – The input `TokenStream` representing the labelled blocks and event expressions. -#[proc_macro_error] +/// # Syntax +/// ```ignore +/// send_cancellable! {{ +/// ; +/// ; +/// 'after: { } +/// 'cancelled: { } +/// }} +/// ``` #[proc_macro] pub fn send_cancellable(input: TokenStream) -> TokenStream { - let block = parse_macro_input!(input as syn::Block); + let input = parse_macro_input!(input as SynBlock); + + let mut stmts_iter = input.stmts.into_iter(); + + let Some(Stmt::Expr(server_stmt, _)) = stmts_iter.next() else { + abort_call_site!("expected server expression as first statement") + }; + + let Some(Stmt::Expr(event_stmt, _)) = stmts_iter.next() else { + abort_call_site!("expected event expression as second statement") + }; + + let event_expr = if let Expr::Reference(syn::ExprReference { + expr, + mutability: Some(_), + .. + }) = event_stmt + { + *expr + } else { + event_stmt + }; - let mut server_expr = None; - let mut event_expr = None; let mut after_block = None; let mut cancelled_block = None; - for stmt in block.stmts { - match stmt { - Stmt::Expr(expr, _) => { - // Check if it is a labelled block first. - let mut is_special_block = false; - if let Expr::Block(ref b) = expr - && let Some(ref label) = b.label - { - let label_name = label.name.ident.to_string(); - if label_name == "after" { - after_block = Some(b.clone()); // Clone strictly necessary here as we split AST. - is_special_block = true; - } else if label_name == "cancelled" { - cancelled_block = Some(b.clone()); - is_special_block = true; - } - } - - // If it wasn't a special block, it must be the event expression. - if !is_special_block { - if server_expr.is_none() { - server_expr = Some(expr); - } else if event_expr.is_none() { - event_expr = Some(expr); - } else { - abort!( - expr.span(), - "Multiple event expressions found. Only one event expression allowed." - ); - } - } + for stmt in stmts_iter { + if let Stmt::Expr(Expr::Block(b), _) = stmt + && let Some(ref label) = b.label + { + if label.name.ident == "after" { + after_block = Some(b.block); + } else if label.name.ident == "cancelled" { + cancelled_block = Some(b.block); } - // Abort on other statements (like `let x = ...`) if strictness is desired. - _ => abort!( - stmt.span(), - "Only event expressions and labeled blocks allowed in `send_cancellable!`" - ), } } - let Some(server) = server_expr else { - abort_call_site!("Server expression must be specified"); - }; - - let Some(event) = event_expr else { - abort_call_site!("Event expression must be specified"); - }; - - // Construct the if/else logic. - let logic = match (after_block, cancelled_block) { + let execution = match (after_block, cancelled_block) { (Some(after), Some(cancelled)) => quote! { - if !event.cancelled { #after } else { #cancelled } + if !is_cancelled { + #after + } else { + #cancelled + } }, (Some(after), None) => quote! { - if !event.cancelled { #after } + if !is_cancelled { + #after + } }, (None, Some(cancelled)) => quote! { - if event.cancelled { #cancelled } + if is_cancelled { + #cancelled + } }, (None, None) => quote! {}, }; - quote! { - let event = #server.plugin_manager.fire(#event).await; - #logic - } - .into() + let expanded = quote! { + { + let mut event = #event_expr; + let server_ref: &std::sync::Arc = { + use std::borrow::Borrow; + (#server_stmt).borrow() + }; + server_ref.plugin_manager.fire(server_ref, &mut event).await; + + let is_cancelled = { + use crate::plugin::Cancellable; + event.cancelled() + }; + + #execution + } + }; + + expanded.into() } /// Attaches a fixed packet ID to a struct implementing `Packet`. diff --git a/pumpkin-protocol/src/java/packet_decoder.rs b/pumpkin-protocol/src/java/packet_decoder.rs index ae4330de8..5dc3a3fbc 100644 --- a/pumpkin-protocol/src/java/packet_decoder.rs +++ b/pumpkin-protocol/src/java/packet_decoder.rs @@ -157,9 +157,6 @@ impl TCPNetworkDecoder { DecompressionReader::None(bounded_reader) }; - // TODO: Serde is sync so we need to write to a buffer here :( - // Is there a way to deserialize in an asynchronous manner? - let packet_id = VarInt::decode_async(&mut reader) .await .map_err(|_| PacketDecodeError::DecodeID)? diff --git a/pumpkin-world/src/cylindrical_chunk_iterator.rs b/pumpkin-world/src/cylindrical_chunk_iterator.rs index f80266ab2..ddf4b9ffa 100644 --- a/pumpkin-world/src/cylindrical_chunk_iterator.rs +++ b/pumpkin-world/src/cylindrical_chunk_iterator.rs @@ -17,25 +17,22 @@ impl Cylindrical { } } - pub fn for_each_changed_chunk( - old_cylindrical: Self, - new_cylindrical: Self, - newly_included: &mut Vec>, - just_removed: &mut Vec>, + pub fn changed_chunks( + old: Self, + new: Self, + ) -> ( + impl Iterator>, + impl Iterator>, ) { - for new_cylindrical_chunk in new_cylindrical.all_chunks_within() { - if !old_cylindrical.is_within_distance(new_cylindrical_chunk.x, new_cylindrical_chunk.y) - { - newly_included.push(new_cylindrical_chunk); - } - } + let loading = new + .all_chunks_within() + .filter(move |c| !old.is_within_distance(c.x, c.y)); - for old_cylindrical_chunk in old_cylindrical.all_chunks_within() { - if !new_cylindrical.is_within_distance(old_cylindrical_chunk.x, old_cylindrical_chunk.y) - { - just_removed.push(old_cylindrical_chunk); - } - } + let unloading = old + .all_chunks_within() + .filter(move |c| !new.is_within_distance(c.x, c.y)); + + (loading, unloading) } #[allow(dead_code)] diff --git a/pumpkin/src/block/blocks/plant/crop/mod.rs b/pumpkin/src/block/blocks/plant/crop/mod.rs index 71a874659..664d5773f 100644 --- a/pumpkin/src/block/blocks/plant/crop/mod.rs +++ b/pumpkin/src/block/blocks/plant/crop/mod.rs @@ -56,7 +56,7 @@ trait CropBlockBase: PlantBlockBase { if rand::rng().random_range(0..=(25.0 / f).floor() as i64) == 0 { let mut new_state_id = self.state_with_age(block, state, age + 1); if let Some(server) = world.server.upgrade() { - let event = BlockGrowEvent::new( + let mut event = BlockGrowEvent::new( world.clone(), block, state, @@ -64,7 +64,7 @@ trait CropBlockBase: PlantBlockBase { new_state_id, *pos, ); - let event = server.plugin_manager.fire(event).await; + server.plugin_manager.fire(&server, &mut event).await; if event.cancelled { return; } diff --git a/pumpkin/src/block/blocks/redstone/pressure_plate/mod.rs b/pumpkin/src/block/blocks/redstone/pressure_plate/mod.rs index b260f3839..9023e1f6b 100644 --- a/pumpkin/src/block/blocks/redstone/pressure_plate/mod.rs +++ b/pumpkin/src/block/blocks/redstone/pressure_plate/mod.rs @@ -64,14 +64,14 @@ pub(crate) trait PressurePlate { let has_output = calc_output > 0; if calc_output != output { let next_output = if let Some(server) = world.server.upgrade() { - let event = crate::plugin::block::block_redstone::BlockRedstoneEvent::new( + let mut event = crate::plugin::block::block_redstone::BlockRedstoneEvent::new( world.clone(), state.id, *pos, i32::from(output), i32::from(calc_output), ); - let event = server.plugin_manager.fire(event).await; + server.plugin_manager.fire(&server, &mut event).await; if event.cancelled { return; } diff --git a/pumpkin/src/block/registry.rs b/pumpkin/src/block/registry.rs index f1250cdd1..0da221f81 100644 --- a/pumpkin/src/block/registry.rs +++ b/pumpkin/src/block/registry.rs @@ -460,7 +460,7 @@ impl BlockRegistry { &self, player: &Arc, placed_block: &'static Block, - server: &Server, + server: &Arc, use_item_on: &SUseItemOn, location: BlockPos, face: BlockDirection, @@ -605,16 +605,16 @@ impl BlockRegistry { } } - let event = crate::plugin::block::block_place::BlockPlaceEvent::new( + let mut event = crate::plugin::block::block_place::BlockPlaceEvent::new( player.clone(), placed_block, clicked_block, final_block_pos, true, ); - let event = server + server .plugin_manager - .fire::(event) + .fire::(server, &mut event) .await; if event.cancelled { return Ok(None); diff --git a/pumpkin/src/command/commands/me.rs b/pumpkin/src/command/commands/me.rs index fb3639e2f..f7b8f21bb 100644 --- a/pumpkin/src/command/commands/me.rs +++ b/pumpkin/src/command/commands/me.rs @@ -28,7 +28,16 @@ impl CommandExecutor for Executor { return Err(InvalidConsumption(Some(ARG_MESSAGE.into()))); }; - server + let Some(server_arc) = sender + .world_or_first(server) + .and_then(|w| w.server.upgrade()) + else { + return Err(CommandError::CommandFailed(TextComponent::text( + "Failed to get server instance", + ))); + }; + + server_arc .broadcast_message( &TextComponent::text(msg.clone()), &TextComponent::text(format!("{sender}")), diff --git a/pumpkin/src/command/commands/plugin.rs b/pumpkin/src/command/commands/plugin.rs index 54f26e1d4..6a703c220 100644 --- a/pumpkin/src/command/commands/plugin.rs +++ b/pumpkin/src/command/commands/plugin.rs @@ -90,9 +90,18 @@ impl CommandExecutor for LoadExecutor { )))); } + let Some(server_arc) = sender + .world_or_first(server) + .and_then(|w| w.server.upgrade()) + else { + return Err(CommandError::CommandFailed(TextComponent::text( + "Failed to get server instance", + ))); + }; + let result = server .plugin_manager - .try_load_plugin(Path::new(plugin_name)) + .try_load_plugin(&server_arc, Path::new(plugin_name)) .await; match result { @@ -171,7 +180,16 @@ impl CommandExecutor for HotReloadExecutor { let enabled = self.0; if enabled { - if let Err(e) = server.plugin_manager.start_watcher().await { + let Some(server_arc) = sender + .world_or_first(server) + .and_then(|w| w.server.upgrade()) + else { + return Err(CommandError::CommandFailed(TextComponent::text( + "Failed to get server instance", + ))); + }; + + if let Err(e) = server.plugin_manager.start_watcher(&server_arc).await { return Err(CommandError::CommandFailed(TextComponent::text(format!( "Failed to start plugin watcher: {e}" )))); diff --git a/pumpkin/src/command/commands/say.rs b/pumpkin/src/command/commands/say.rs index 75b7af29d..cb17b76c4 100644 --- a/pumpkin/src/command/commands/say.rs +++ b/pumpkin/src/command/commands/say.rs @@ -28,7 +28,16 @@ impl CommandExecutor for Executor { return Err(InvalidConsumption(Some(ARG_MESSAGE.into()))); }; - server + let Some(server_arc) = sender + .world_or_first(server) + .and_then(|w| w.server.upgrade()) + else { + return Err(CommandError::CommandFailed(TextComponent::text( + "Failed to get server instance", + ))); + }; + + server_arc .broadcast_message( &TextComponent::text(msg.clone()), &TextComponent::text(format!("{sender}")), diff --git a/pumpkin/src/command/commands/setworldspawn.rs b/pumpkin/src/command/commands/setworldspawn.rs index a1f9ede41..6ac21323c 100644 --- a/pumpkin/src/command/commands/setworldspawn.rs +++ b/pumpkin/src/command/commands/setworldspawn.rs @@ -126,7 +126,7 @@ async fn setworldspawn( let previous_pitch = current_info.spawn_pitch; let mut new_yaw = yaw; let mut new_pitch = pitch; - let event = SpawnChangeEvent::new( + let mut event = SpawnChangeEvent::new( world.clone(), previous_position, previous_yaw, @@ -135,7 +135,12 @@ async fn setworldspawn( new_yaw, new_pitch, ); - let event = server.plugin_manager.fire(event).await; + if let Some(server_arc) = world.server.upgrade() { + server_arc + .plugin_manager + .fire(&server_arc, &mut event) + .await; + } new_position = event.new_position; new_yaw = event.new_yaw; new_pitch = event.new_pitch; diff --git a/pumpkin/src/command/dispatcher.rs b/pumpkin/src/command/dispatcher.rs index 01ee76b65..301df2046 100644 --- a/pumpkin/src/command/dispatcher.rs +++ b/pumpkin/src/command/dispatcher.rs @@ -436,7 +436,8 @@ impl CommandDispatcher { return Err(CommandFailed(TextComponent::text("Empty Command"))); } let key = args.remove(0).value; - Ok((key, args.into_iter().rev().collect())) + args.reverse(); + Ok((key, args)) } /// Execute a command using its corresponding [`CommandTree`]. @@ -448,9 +449,9 @@ impl CommandDispatcher { ) -> Result<(), CommandError> { let (key, raw_args) = Self::split_parts(cmd)?; - if !self.commands.contains_key(key) { - return Err(SyntaxError(unknown_command_syntax_error(cmd, 0))); - } + let tree = self + .get_tree(key) + .map_err(|_| SyntaxError(unknown_command_syntax_error(cmd, 0)))?; let Some(permission) = self.permissions.get(key) else { return Err(CommandFailed(TextComponent::text( @@ -462,8 +463,6 @@ impl CommandDispatcher { return Err(PermissionDenied); } - let tree = self.get_tree(key)?; - let mut path_failures = Vec::new(); // try paths until fitting path is found diff --git a/pumpkin/src/entity/player.rs b/pumpkin/src/entity/player.rs index 08192a71f..469e398ca 100644 --- a/pumpkin/src/entity/player.rs +++ b/pumpkin/src/entity/player.rs @@ -2174,8 +2174,9 @@ impl Player { payload: Bytes, ) -> bool { if let Some(server) = self.world().server.upgrade() { - let event = PacketSentEvent::new(self.clone(), packet_id, payload, Arc::new(packet)); - let event = server.plugin_manager.fire(event).await; + let mut event = + PacketSentEvent::new(self.clone(), packet_id, payload, Arc::new(packet)); + server.plugin_manager.fire(&server, &mut event).await; return event.cancelled; } false @@ -2186,8 +2187,9 @@ impl Player { // This is a dummy object to satisfy the non-optional requirement in WIT // In the future we should make all packets 'static or have a way to represent raw packets in WIT struct RawPacket; - let event = PacketSentEvent::new(self.clone(), packet_id, payload, Arc::new(RawPacket)); - let event = server.plugin_manager.fire(event).await; + let mut event = + PacketSentEvent::new(self.clone(), packet_id, payload, Arc::new(RawPacket)); + server.plugin_manager.fire(&server, &mut event).await; return event.cancelled; } false @@ -3501,8 +3503,8 @@ impl Player { /// Add experience points to the player. pub async fn add_experience_points(self: &Arc, mut added_points: i32) { if let Some(server) = self.world().server.upgrade() { - let event = PlayerExpChangeEvent::new(self.clone(), added_points); - let event = server.plugin_manager.fire(event).await; + let mut event = PlayerExpChangeEvent::new(self.clone(), added_points); + server.plugin_manager.fire(&server, &mut event).await; added_points = event.amount; } @@ -3637,15 +3639,12 @@ impl Player { }; if let Some(server) = self.living_entity.entity.world.load().server.upgrade() { - server - .plugin_manager - .fire( - crate::plugin::api::events::player::inventory_close::InventoryCloseEvent::new( - self, - window_type, - ), - ) - .await; + let mut event = + crate::plugin::api::events::player::inventory_close::InventoryCloseEvent::new( + self, + window_type, + ); + server.plugin_manager.fire(&server, &mut event).await; } let player_screen_handler: Arc> = @@ -3829,7 +3828,7 @@ impl Player { } #[allow(clippy::too_many_lines)] - pub async fn on_slot_click(self: &Arc, packet: SClickSlot, server: &Server) { + pub async fn on_slot_click(self: &Arc, packet: SClickSlot, server: &Arc) { self.update_last_action_time(); let screen_handler_arc = self.current_screen_handler.lock().await.clone(); let mut screen_handler = screen_handler_arc.lock().await; @@ -4060,14 +4059,13 @@ impl Player { .await; drop(perm_manager); - let event = server - .plugin_manager - .fire(PlayerPermissionCheckEvent::new( - self.clone(), - node.to_string(), - result, - )) - .await; + let mut event = PlayerPermissionCheckEvent::new(self.clone(), node.to_string(), result); + if let Some(server_arc) = self.world().server.upgrade() { + server_arc + .plugin_manager + .fire(&server_arc, &mut event) + .await; + } event.result } diff --git a/pumpkin/src/entity/projectile/egg.rs b/pumpkin/src/entity/projectile/egg.rs index c5b44c727..159b38339 100644 --- a/pumpkin/src/entity/projectile/egg.rs +++ b/pumpkin/src/entity/projectile/egg.rs @@ -129,14 +129,14 @@ impl EntityBase for EggEntity { && let Some(player) = world.get_player_by_id(owner_id) && let Some(server) = world.server.upgrade() { - let event = PlayerEggThrowEvent::new( + let mut event = PlayerEggThrowEvent::new( player, self.get_entity().entity_uuid, hatching, to_spawn as u8, hatching_type, ); - let event = server.plugin_manager.fire(event).await; + server.plugin_manager.fire(&server, &mut event).await; if event.cancelled { hatching = false; } else { diff --git a/pumpkin/src/lib.rs b/pumpkin/src/lib.rs index 9ea5b2cf9..8b88942da 100644 --- a/pumpkin/src/lib.rs +++ b/pumpkin/src/lib.rs @@ -1,6 +1,9 @@ // Not warn event sending macros #![allow(unused_labels)] +#[macro_use] +extern crate pumpkin_macros; + use crate::crash::CrashReport; use crate::data::VanillaData; use crate::logging::{GzipRollingLogger, PumpkinCommandCompleter, ReadlineLogWrapper}; @@ -8,11 +11,10 @@ use crate::net::bedrock::BedrockClient; use crate::net::java::JavaClient; use crate::net::{ClientPlatform, DisconnectReason, PacketHandlerResult}; use crate::net::{lan_broadcast::LANBroadcast, query, rcon::RCONServer}; +use crate::plugin::server::server_command::ServerCommandEvent; use crate::server::{Server, ticker::Ticker}; -use plugin::server::server_command::ServerCommandEvent; use plugin::server::server_load::{LoadType, ServerLoadEvent}; use pumpkin_config::{AdvancedConfiguration, BasicConfiguration}; -use pumpkin_macros::send_cancellable; use pumpkin_util::text::TextComponent; use pumpkin_util::text::color::{Color, NamedColor}; use rustyline::Editor; @@ -310,15 +312,7 @@ impl PumpkinServer { } pub async fn init_plugins(&self) -> std::time::Duration { - self.server - .plugin_manager - .set_self_ref(self.server.plugin_manager.clone()) - .await; - self.server - .plugin_manager - .set_server(self.server.clone()) - .await; - match self.server.plugin_manager.load_plugins().await { + match self.server.plugin_manager.load_plugins(&self.server).await { Ok(duration) => duration, Err(err) => { error!("{err}"); @@ -355,10 +349,9 @@ impl PumpkinServer { let mut master_client_id: u64 = 0; let bedrock_clients = Arc::new(Mutex::new(HashMap::new())); - let _ = self - .server + self.server .plugin_manager - .fire(ServerLoadEvent::new(LoadType::Startup)) + .fire(&self.server, &mut ServerLoadEvent::new(LoadType::Startup)) .await; while !SHOULD_STOP.load(Ordering::Relaxed) { @@ -645,16 +638,19 @@ fn setup_stdin_console(server: Arc) { tokio::spawn(async move { while !SHOULD_STOP.load(Ordering::Relaxed) { if let Some(command) = rx.recv().await { - send_cancellable! {{ - &server; - ServerCommandEvent::new(command.clone()); - - 'after: { - server.command_dispatcher.read().await - .handle_command(&command::CommandSender::Console.into_source(&server).await, command.as_str()) - .await; - }; - }} + let mut event = ServerCommandEvent::new(command.clone()); + server.plugin_manager.fire(&server, &mut event).await; + if !event.cancelled { + server + .command_dispatcher + .read() + .await + .handle_command( + &command::CommandSender::Console.into_source(&server).await, + command.as_str(), + ) + .await; + } } } }); @@ -716,18 +712,20 @@ fn setup_console(mut rl: Editor, server: A }; if let Some(line) = result { - send_cancellable! {{ - &server; - ServerCommandEvent::new(line.clone()); - - 'after: { - server.command_dispatcher.read().await - .handle_command(&command::CommandSender::Console.into_source(&server).await, &line) - .await; - - let _ = tx_reply.send(1).await; - } - }} + let mut event = ServerCommandEvent::new(line.clone()); + server.plugin_manager.fire(&server, &mut event).await; + if !event.cancelled { + server + .command_dispatcher + .read() + .await + .handle_command( + &command::CommandSender::Console.into_source(&server).await, + &line, + ) + .await; + let _ = tx_reply.send(1).await; + } } else { break; } diff --git a/pumpkin/src/net/bedrock/login.rs b/pumpkin/src/net/bedrock/login.rs index 7b2ced3dd..eea353c44 100644 --- a/pumpkin/src/net/bedrock/login.rs +++ b/pumpkin/src/net/bedrock/login.rs @@ -309,7 +309,7 @@ impl BedrockClient { pub async fn handle_resource_pack_response( &self, packet: SResourcePackResponse, - server: &Server, + server: &Arc, ) { // TODO: warn & ignore if the player is already spawned in diff --git a/pumpkin/src/net/bedrock/mod.rs b/pumpkin/src/net/bedrock/mod.rs index 5cb62b540..322b2165b 100644 --- a/pumpkin/src/net/bedrock/mod.rs +++ b/pumpkin/src/net/bedrock/mod.rs @@ -372,8 +372,8 @@ impl BedrockClient { let mut valid_chunks = Vec::with_capacity(chunks.len()); for chunk in chunks { - let event = ChunkSend::new(player.world(), chunk.clone()); - let event = server.plugin_manager.fire(event).await; + let mut event = ChunkSend::new(player.world(), chunk.clone()); + server.plugin_manager.fire(&server, &mut event).await; if !event.cancelled { valid_chunks.push(chunk.clone()); } @@ -1106,7 +1106,7 @@ impl BedrockClient { packet.id, packet.payload.clone(), ); - event = server.plugin_manager.fire(event).await; + server.plugin_manager.fire(server, &mut event).await; if event.cancelled { continue; } @@ -1208,7 +1208,7 @@ impl BedrockClient { .await; } SMobEquipment::PACKET_ID => { - self.handle_mob_equipment(player, SMobEquipment::read(reader)?) + self.handle_mob_equipment(server, player, SMobEquipment::read(reader)?) .await; } _ => { diff --git a/pumpkin/src/net/bedrock/play.rs b/pumpkin/src/net/bedrock/play.rs index acb9c73dd..ab0517289 100644 --- a/pumpkin/src/net/bedrock/play.rs +++ b/pumpkin/src/net/bedrock/play.rs @@ -10,7 +10,6 @@ use pumpkin_data::{ item_stack::ItemStack, }; use pumpkin_inventory::screen_handler::{InventoryPlayer, ScreenHandler}; -use pumpkin_macros::send_cancellable; use pumpkin_protocol::bedrock::{ client::inventory_content::CInventoryContent, network_item::{ @@ -209,7 +208,7 @@ impl BedrockClient { &self, player: &Arc, packet: SPlayerAuthInput, - server: &Server, + server: &Arc, ) { if !player.has_client_loaded() { return; @@ -839,7 +838,7 @@ impl BedrockClient { } send_cancellable! {{ - server; + &server; event; 'after: { server.item_registry.on_use(&stack_for_use, player).await; @@ -1012,7 +1011,7 @@ impl BedrockClient { pub async fn handle_chat_message( &self, - server: &Server, + server: &Arc, player: &Arc, packet: SText<'_>, ) { @@ -1236,15 +1235,16 @@ impl BedrockClient { pub async fn handle_modal_form_response( &self, player: &Arc, - server: &Server, + server: &Arc, packet: pumpkin_protocol::bedrock::server::modal_form_response::SModalFormResponse<'_>, ) { - let event = crate::plugin::api::events::player::bedrock_form_response::BedrockFormResponseEvent::new( - player.clone(), - packet.form_id.0 as u32, - packet.form_data.map(std::borrow::Cow::into_owned), - ); - let _ = server.plugin_manager.fire(event).await; + let mut event = + crate::plugin::api::events::player::bedrock_form_response::BedrockFormResponseEvent::new( + player.clone(), + packet.form_id.0 as u32, + packet.form_data.map(std::borrow::Cow::into_owned), + ); + server.plugin_manager.fire(server, &mut event).await; } #[allow(clippy::too_many_lines)] @@ -1797,7 +1797,12 @@ impl BedrockClient { .await; } - pub async fn handle_mob_equipment(&self, player: &Arc, packet: SMobEquipment) { + pub async fn handle_mob_equipment( + &self, + _server: &Arc, + player: &Arc, + packet: SMobEquipment, + ) { player.update_last_action_time(); let slot = packet.hotbar_slot; if slot >= 9 { @@ -1805,9 +1810,13 @@ impl BedrockClient { } let previous_slot = player.inventory.get_selected_slot(); if let Some(server) = player.world().server.upgrade() { - let event = PlayerItemHeldEvent::new(player.clone(), previous_slot, slot); - let event = server.plugin_manager.fire(event).await; - if event.cancelled { + let mut event = PlayerItemHeldEvent::new(player.clone(), previous_slot, slot); + server.plugin_manager.fire(&server, &mut event).await; + let is_cancelled = { + use crate::plugin::Cancellable; + event.cancelled() + }; + if is_cancelled { self.enqueue_packet(&CPlayerHotbar { selected_slot: VarUInt(previous_slot as u32), container_id: 0, diff --git a/pumpkin/src/net/java/mod.rs b/pumpkin/src/net/java/mod.rs index 96a627acb..b0a338372 100644 --- a/pumpkin/src/net/java/mod.rs +++ b/pumpkin/src/net/java/mod.rs @@ -354,8 +354,8 @@ impl JavaClient { self.send_packet_now(&CChunkBatchStart).await; for chunk in chunks { - let event = ChunkSend::new(player.world(), chunk.clone()); - let event = server.plugin_manager.fire(event).await; + let mut event = ChunkSend::new(player.world(), chunk.clone()); + server.plugin_manager.fire(&server, &mut event).await; if event.cancelled { continue; } @@ -393,7 +393,6 @@ impl JavaClient { } else { false }; - if !cancelled { self.enqueue_packet_data(payload).await; } @@ -907,7 +906,7 @@ impl JavaClient { packet.id, packet.payload.clone(), ); - event = server.plugin_manager.fire(event).await; + server.plugin_manager.fire(server, &mut event).await; if event.cancelled { return Ok(()); } @@ -945,6 +944,7 @@ impl JavaClient { } id if id == SClientInformationPlay::to_id(version) => { self.handle_client_information( + server, player, SClientInformationPlay::read(&mut payload, &version)?, ) @@ -1126,8 +1126,12 @@ impl JavaClient { .await; } id if id == SSetHeldItem::to_id(version) => { - self.handle_set_held_item(player, SSetHeldItem::read(&mut payload, &version)?) - .await; + self.handle_set_held_item( + server, + player, + SSetHeldItem::read(&mut payload, &version)?, + ) + .await; } id if id == SSetCreativeSlot::to_id(version) => { self.handle_set_creative_slot( @@ -1137,7 +1141,7 @@ impl JavaClient { .await?; } id if id == SSwingArm::to_id(version) => { - self.handle_swing_arm(player, SSwingArm::read(&mut payload, &version)?) + self.handle_swing_arm(server, player, SSwingArm::read(&mut payload, &version)?) .await; } id if id == SUpdateSign::to_id(version) => { @@ -1185,12 +1189,12 @@ impl JavaClient { } id if id == SCustomPayload::to_id(version) => { let payload = SCustomPayload::read(&mut payload, &version)?; - let event = PlayerCustomPayloadEvent::new( + let mut event = PlayerCustomPayloadEvent::new( player.clone(), payload.channel.to_string(), Bytes::copy_from_slice(payload.data), ); - server.plugin_manager.fire(event).await; + server.plugin_manager.fire(server, &mut event).await; } id if id == SRecipeBookChangeSettings::to_id(version) => { self.handle_recipe_book_change_settings( @@ -1222,12 +1226,12 @@ impl JavaClient { &mut payload, &version, )?; - let event = crate::plugin::api::events::player::custom_click_action::CustomClickActionEvent::new( + let mut event = crate::plugin::api::events::player::custom_click_action::CustomClickActionEvent::new( player.clone(), packet.action_id.to_string(), packet.payload.map(Bytes::copy_from_slice), ); - server.plugin_manager.fire(event).await; + server.plugin_manager.fire(server, &mut event).await; } id if id == SSelectTrade::to_id(version) => { self.handle_select_trade(player, SSelectTrade::read(&mut payload, &version)?) diff --git a/pumpkin/src/net/java/play.rs b/pumpkin/src/net/java/play.rs index f390a29fe..ebd2c62dc 100644 --- a/pumpkin/src/net/java/play.rs +++ b/pumpkin/src/net/java/play.rs @@ -51,7 +51,6 @@ use pumpkin_inventory::InventoryError; use pumpkin_inventory::merchant::merchant_screen_handler::MerchantScreenHandler; use pumpkin_inventory::player::player_inventory::PlayerInventory; use pumpkin_inventory::screen_handler::{InventoryPlayer, ScreenHandler}; -use pumpkin_macros::send_cancellable; use pumpkin_protocol::bedrock::client::CMovePlayer; use pumpkin_protocol::codec::var_int::VarInt; use pumpkin_protocol::codec::var_ulong::VarULong; @@ -986,7 +985,7 @@ impl JavaClient { &self, player: &Arc, command: SPlayerCommand, - server: &Server, + server: &Arc, ) { if command.entity_id != player.entity_id().into() { return; @@ -1049,7 +1048,7 @@ impl JavaClient { &self, player: &Arc, input: SPlayerInput, - server: &Server, + server: &Arc, ) { player.last_input.store(input.input, Ordering::Relaxed); @@ -1345,7 +1344,12 @@ impl JavaClient { screen_handler_arc.lock().await.send_content_updates().await; } - pub async fn handle_swing_arm(&self, player: &Arc, swing_arm: SSwingArm) { + pub async fn handle_swing_arm( + &self, + server: &Arc, + player: &Arc, + swing_arm: SSwingArm, + ) { player.update_last_action_time(); let Ok(hand) = Hand::try_from(swing_arm.hand.0) else { self.kick(TextComponent::text("Invalid hand")).await; @@ -1378,12 +1382,8 @@ impl JavaClient { PlayerInteractEvent::new(player, InteractAction::LeftClickAir, &Block::AIR, None) }; - let Some(server) = player.world().server.upgrade() else { - return; - }; - send_cancellable! {{ - server; + &server; event; 'after: { player.swing_hand(hand, false).await; @@ -1393,7 +1393,7 @@ impl JavaClient { pub async fn handle_chat_message( &self, - server: &Server, + server: &Arc, player: &Arc, chat_message: SChatMessage<'_>, ) { @@ -1609,6 +1609,7 @@ impl JavaClient { pub async fn handle_client_information( &self, + server: &Arc, player: &Arc, client_information: SClientInformationPlay<'_>, ) { @@ -1673,9 +1674,9 @@ impl JavaClient { chunker::update_position(player).await; } - if main_hand_changed && let Some(server) = player.world().server.upgrade() { - let event = PlayerChangedMainHandEvent::new(player.clone(), main_hand); - let _ = server.plugin_manager.fire(event).await; + if main_hand_changed { + let mut event = PlayerChangedMainHandEvent::new(player.clone(), main_hand); + server.plugin_manager.fire(server, &mut event).await; } if update_settings { @@ -2178,7 +2179,7 @@ impl JavaClient { &self, player: &Arc, player_abilities: SPlayerAbilities, - server: &Server, + server: &Arc, ) { let (flying, allow_flying) = { let abilities = player.abilities.lock().await; @@ -2467,7 +2468,7 @@ impl JavaClient { &self, player: &Arc, use_item: &SUseItem, - server: &Server, + server: &Arc, ) { if !player.has_client_loaded() { return; @@ -2619,7 +2620,7 @@ impl JavaClient { async fn should_continue_use_after_fish_event( &self, - server: &Server, + server: &Arc, player: &Arc, hand: Hand, item_for_use: &Item, @@ -2629,7 +2630,7 @@ impl JavaClient { } // TODO: Apply fishing rod durability on retrieval based on catch type. - let fish_event = PlayerFishEvent::new( + let mut fish_event = PlayerFishEvent::new( player.clone(), None, uuid::Uuid::nil(), @@ -2638,11 +2639,17 @@ impl JavaClient { hand, 0, ); - let fish_event = server.plugin_manager.fire(fish_event).await; + server.plugin_manager.fire(server, &mut fish_event).await; !fish_event.cancelled } - pub async fn handle_set_held_item(&self, player: &Player, held: SSetHeldItem) { + pub async fn handle_set_held_item( + &self, + server: &Arc, + + player: &Player, + held: SSetHeldItem, + ) { player.update_last_action_time(); let slot = held.slot; if !(0..=8).contains(&slot) { @@ -2651,19 +2658,17 @@ impl JavaClient { } let slot = slot as u8; let previous_slot = player.inventory.get_selected_slot(); - if let Some(server) = player.world().server.upgrade() { - let Some(player_arc) = player.world().get_player_by_uuid(player.gameprofile.id) else { - return; - }; - let event = PlayerItemHeldEvent::new(player_arc, previous_slot, slot); - let event = server.plugin_manager.fire(event).await; - if event.cancelled { - player - .client - .enqueue_packet(&CSetSelectedSlot::new(previous_slot as i8)) - .await; - return; - } + let Some(player_arc) = player.world().get_player_by_uuid(player.gameprofile.id) else { + return; + }; + let mut event = PlayerItemHeldEvent::new(player_arc, previous_slot, slot); + server.plugin_manager.fire(server, &mut event).await; + if event.cancelled { + player + .client + .enqueue_packet(&CSetSelectedSlot::new(previous_slot as i8)) + .await; + return; } let inv = player.inventory(); @@ -2805,7 +2810,7 @@ impl JavaClient { &self, player: &Arc, block: &'static Block, - server: &Server, + server: &Arc, use_item_on: SUseItemOn, location: BlockPos, face: BlockDirection, diff --git a/pumpkin/src/plugin/api/context.rs b/pumpkin/src/plugin/api/context.rs index 96ef6268f..eeabae836 100644 --- a/pumpkin/src/plugin/api/context.rs +++ b/pumpkin/src/plugin/api/context.rs @@ -319,7 +319,7 @@ impl Context { loader: Arc, ) -> bool { let before_count = self.plugin_manager.loaded_plugins().await.len(); - self.plugin_manager.add_loader(loader).await; + self.plugin_manager.add_loader(&self.server, loader).await; let after_count = self.plugin_manager.loaded_plugins().await.len(); // Return true if any new plugins were loaded diff --git a/pumpkin/src/plugin/mod.rs b/pumpkin/src/plugin/mod.rs index 22b0c7ca4..f4150ecc0 100644 --- a/pumpkin/src/plugin/mod.rs +++ b/pumpkin/src/plugin/mod.rs @@ -172,11 +172,8 @@ pub enum PluginState { pub struct PluginManager { plugins: RwLock>, loaders: RwLock>>, - server: RwLock>>, handlers: Arc>, unloaded_files: RwLock>, - // Self-reference for sharing with contexts - self_ref: RwLock>>, services: Arc>>>, // Plugin state tracking plugin_states: RwLock>, @@ -204,43 +201,19 @@ struct LoadedPlugin { /// Error types for plugin management #[derive(Error, Debug)] pub enum ManagerError { - #[error("Server not initialized")] - ServerNotInitialized, - #[error("Plugin not found: {0}")] PluginNotFound(String), - #[error("Loader error: {0}")] LoaderError(#[from] LoaderError), - #[error("IO error: {0}")] IoError(#[from] std::io::Error), - - #[error("Plugin manager not initialized properly")] - ManagerNotInitialized, - #[error("Dependency error: {0}")] DependencyError(String), } impl Default for PluginManager { fn default() -> Self { - Self { - plugins: RwLock::new(Vec::new()), - loaders: RwLock::new(vec![ - Arc::new(NativePluginLoader), - Arc::new(WasmPluginLoader), - ]), - server: RwLock::new(None), - handlers: Arc::new(RwLock::new(HashMap::new())), - unloaded_files: RwLock::new(HashSet::new()), - self_ref: RwLock::new(None), - services: Arc::new(RwLock::new(HashMap::new())), - plugin_states: RwLock::new(HashMap::new()), - state_notify: Arc::new(Notify::new()), - hot_reload_task: RwLock::new(None), - hot_reload_enabled: AtomicBool::new(false), - } + Self::new() } } @@ -248,7 +221,20 @@ impl PluginManager { /// Create a new plugin manager with default loaders #[must_use] pub fn new() -> Self { - Self::default() + Self { + plugins: RwLock::new(Vec::new()), + loaders: RwLock::new(vec![ + Arc::new(NativePluginLoader), + Arc::new(WasmPluginLoader), + ]), + handlers: Arc::new(RwLock::new(HashMap::new())), + unloaded_files: RwLock::new(HashSet::new()), + services: Arc::new(RwLock::new(HashMap::new())), + plugin_states: RwLock::new(HashMap::new()), + state_notify: Arc::new(Notify::new()), + hot_reload_task: RwLock::new(None), + hot_reload_enabled: AtomicBool::new(false), + } } /// Unload all loaded plugins @@ -272,15 +258,15 @@ impl PluginManager { } /// Add a new plugin loader implementation - pub async fn add_loader(&self, loader: Arc) { + pub async fn add_loader(self: &Arc, server: &Arc, loader: Arc) { self.loaders.write().await.push(loader); // Try to load previously unloaded files with the new loader - self.retry_unloaded_files().await; + self.retry_unloaded_files(server).await; } /// Start watching the plugins directory for changes - pub async fn start_watcher(&self) -> Result<(), ManagerError> { + pub async fn start_watcher(self: &Arc, server: &Arc) -> Result<(), ManagerError> { if self.hot_reload_task.read().await.is_some() { return Ok(()); } @@ -302,19 +288,14 @@ impl PluginManager { .watch(plugin_dir, RecursiveMode::NonRecursive) .map_err(|e| ManagerError::IoError(std::io::Error::other(e)))?; - let self_ref = self - .self_ref - .read() - .await - .clone() - .ok_or(ManagerError::ManagerNotInitialized)?; - + let manager = self.clone(); + let server = server.clone(); let task = tokio::spawn(async move { // Keep watcher alive by moving it into the task let _watcher = watcher; while let Some(event) = rx.recv().await { - if !self_ref + if !manager .hot_reload_enabled .load(std::sync::atomic::Ordering::Relaxed) { @@ -331,7 +312,7 @@ impl PluginManager { // We need to find if this plugin is already loaded to unload it first let plugin_name = { - let plugins = self_ref.plugins.read().await; + let plugins = manager.plugins.read().await; plugins .iter() .find(|p| p.path == path) @@ -340,13 +321,13 @@ impl PluginManager { if let Some(name) = plugin_name { info!("Hot-reloading plugin: {}", name); - let _ = self_ref.unload_plugin(&name).await; + let _ = manager.unload_plugin(&name).await; } // For now, we just try to load it. If it's already loaded, // the loader might handle it or we might get a duplicate. // Most WASM loaders will just create a new instance. - if let Err(e) = self_ref.start_loading_plugin(&path).await { + if let Err(e) = manager.start_loading_plugin(&server, &path).await { error!("Failed to hot-reload plugin {:?}: {}", path, e); } } @@ -382,13 +363,13 @@ impl PluginManager { } /// Retry loading files that couldn't be loaded previously - async fn retry_unloaded_files(&self) { + async fn retry_unloaded_files(self: &Arc, server: &Arc) { let files_to_retry: Vec = { self.unloaded_files.read().await.iter().cloned().collect() }; let mut retry_tasks = Vec::new(); for path in files_to_retry { - if let Ok(task) = self.start_loading_plugin(&path).await { + if let Ok(task) = self.start_loading_plugin(server, &path).await { retry_tasks.push(task); } } @@ -397,18 +378,6 @@ impl PluginManager { join_all(retry_tasks).await; } - /// Set server reference for plugin context - pub async fn set_server(&self, server: Arc) { - let mut srv = self.server.write().await; - srv.replace(server); - } - - /// Set self reference for creating contexts - pub async fn set_self_ref(&self, self_ref: Arc) { - let mut sref = self.self_ref.write().await; - sref.replace(self_ref); - } - /// Get a clone of the loaders for context use #[must_use] pub async fn get_loaders(&self) -> Vec> { @@ -530,9 +499,10 @@ impl PluginManager { } /// Spawn initialization for a single plugin - #[expect(clippy::too_many_lines)] + #[allow(clippy::too_many_lines)] async fn spawn_plugin_initialization( - &self, + self: &Arc, + server: Arc, instance: Arc, metadata: PluginMetadata, loader_data: Box, @@ -545,25 +515,11 @@ impl PluginManager { .await .insert(metadata.name.clone(), PluginState::Loading); - let self_ref = self - .self_ref - .read() - .await - .clone() - .ok_or(ManagerError::ServerNotInitialized)?; - let context = Arc::new(Context::new( metadata.clone(), - Arc::clone( - &self - .server - .read() - .await - .clone() - .ok_or(ManagerError::ServerNotInitialized)?, - ), + server, Arc::clone(&self.handlers), - Arc::clone(&self_ref), + Arc::clone(self), Arc::clone(&LOGGER_IMPL), )); @@ -585,7 +541,7 @@ impl PluginManager { }; // Spawn async task for plugin initialization - let self_ref_clone = Arc::clone(&self_ref); + let self_ref_clone = Arc::clone(self); let state_notify = Arc::clone(&self.state_notify); let plugin_name = metadata.name.clone(); let loader_clone = loader.clone(); @@ -662,7 +618,11 @@ impl PluginManager { } /// Load all plugins from the plugin directory - pub async fn load_plugins(&self) -> Result { + #[allow(clippy::too_many_lines)] + pub async fn load_plugins( + self: &Arc, + server: &Arc, + ) -> Result { let path = Path::new(PLUGIN_DIR); if !path.exists() { @@ -748,6 +708,7 @@ impl PluginManager { if let Some((instance, metadata, loader_data, loader, path)) = plugins_map.remove(&name) { let (allowed, wait_time) = self + .clone() .check_permissions_cached(&path, &metadata, &mut cache, &cache_path) .await; @@ -762,7 +723,14 @@ impl PluginManager { } match self - .spawn_plugin_initialization(instance, metadata, loader_data, loader, path) + .spawn_plugin_initialization( + server.clone(), + instance, + metadata, + loader_data, + loader, + path, + ) .await { Ok(task) => { @@ -812,7 +780,8 @@ impl PluginManager { /// Start loading a plugin asynchronously async fn start_loading_plugin( - &self, + self: &Arc, + server: &Arc, path: &Path, ) -> Result, ManagerError> { for loader in self.loaders.read().await.iter() { @@ -838,6 +807,7 @@ impl PluginManager { return self .spawn_plugin_initialization( + server.clone(), instance, metadata, loader_data, @@ -857,12 +827,19 @@ impl PluginManager { } /// Attempt to load a single plugin file - pub async fn try_load_plugin(&self, path: &Path) -> Result<(), ManagerError> { - self.start_loading_plugin(path).await?.await.map_err(|e| { - ManagerError::LoaderError(LoaderError::InitializationFailed(format!( - "Task join error: {e}" - ))) - }) + pub async fn try_load_plugin( + self: &Arc, + server: &Arc, + path: &Path, + ) -> Result<(), ManagerError> { + self.start_loading_plugin(server, path) + .await? + .await + .map_err(|e| { + ManagerError::LoaderError(LoaderError::InitializationFailed(format!( + "Task join error: {e}" + ))) + }) } /// Wait for a plugin to finish loading @@ -1022,28 +999,37 @@ impl PluginManager { } /// Fire an event to all registered handlers - pub async fn fire(&self, mut event: E) -> E { - if let Some(server) = self.server.read().await.as_ref() { - let handlers = self.handlers.read().await; - if let Some(handlers) = handlers.get(&E::get_name_static()) { - let (blocking, non_blocking): (Vec<_>, Vec<_>) = - handlers.iter().partition(|h| h.is_blocking()); + pub async fn fire( + &self, + server: &Arc, + event: &mut E, + ) { + let handlers_lock = self.handlers.read().await; + if handlers_lock.is_empty() { + return; + } - // Process blocking handlers first - for handler in blocking { - handler.handle_blocking_dyn(server, &mut event).await; - } + let Some(handlers) = handlers_lock.get(&E::get_name_static()) else { + return; + }; - // Process non-blocking handlers - join_all( - non_blocking - .into_iter() - .map(|h| h.handle_dyn(server, &event)), - ) - .await; + if handlers.is_empty() { + return; + } + + // Process blocking handlers first + for handler in handlers { + if handler.is_blocking() { + handler.handle_blocking_dyn(server, event).await; + } + } + + // Process non-blocking handlers + for handler in handlers { + if !handler.is_blocking() { + handler.handle_dyn(server, event).await; } } - event } } diff --git a/pumpkin/src/server/mod.rs b/pumpkin/src/server/mod.rs index b9898e584..dfcfae5b4 100644 --- a/pumpkin/src/server/mod.rs +++ b/pumpkin/src/server/mod.rs @@ -30,7 +30,6 @@ use pumpkin_world::world::WorldPortalExt; use tracing::{debug, error, info, warn}; use crate::command::CommandSender; -use pumpkin_macros::send_cancellable; use pumpkin_protocol::java::client::login::CEncryptionRequest; use pumpkin_protocol::java::client::play::{CChangeDifficulty, CTabList}; use pumpkin_protocol::{ClientPacket, java::client::config::CPluginMessage}; @@ -481,7 +480,7 @@ impl Server { /// /// You still have to spawn the `Player` in a `World` to let them join and make them visible. pub async fn add_player( - &self, + self: &Arc, client: Arc, profile: GameProfile, config: Option, @@ -550,7 +549,7 @@ impl Server { send_cancellable! {{ self; - PlayerLoginEvent::new(player.clone(), TextComponent::text("You have been kicked from the server")); + &mut PlayerLoginEvent::new(player.clone(), TextComponent::text("You have been kicked from the server")); 'after: { player.screen_handler_sync_handler.store_player(player.clone()).await; if world @@ -643,7 +642,7 @@ impl Server { } pub async fn broadcast_message( - &self, + self: &Arc, message: &TextComponent, sender_name: &TextComponent, chat_type: u8, diff --git a/pumpkin/src/server/ticker.rs b/pumpkin/src/server/ticker.rs index 26fcc606c..ba324b54d 100644 --- a/pumpkin/src/server/ticker.rs +++ b/pumpkin/src/server/ticker.rs @@ -25,9 +25,9 @@ impl Ticker { manager.tick(); let tick_number = server.tick_count.load(Ordering::Relaxed); - let _ = server + server .plugin_manager - .fire(ServerTickStartEvent::new(tick_number)) + .fire(server, &mut ServerTickStartEvent::new(tick_number)) .await; if manager.is_sprinting() { @@ -44,9 +44,12 @@ impl Ticker { let tick_duration_nanos = tick_start_time.elapsed().as_nanos() as i64; let tick_number = server.tick_count.load(Ordering::Relaxed); - let _ = server + server .plugin_manager - .fire(ServerTickEndEvent::new(tick_number, tick_duration_nanos)) + .fire( + server, + &mut ServerTickEndEvent::new(tick_number, tick_duration_nanos), + ) .await; server.update_tick_times(tick_duration_nanos).await; diff --git a/pumpkin/src/world/chunker.rs b/pumpkin/src/world/chunker.rs index 6c0ef9e42..e8e4c20c9 100644 --- a/pumpkin/src/world/chunker.rs +++ b/pumpkin/src/world/chunker.rs @@ -72,14 +72,11 @@ pub async fn update_position(player: &Arc) { .await; } } - let mut loading_chunks = Vec::new(); - let mut unloading_chunks = Vec::new(); - Cylindrical::for_each_changed_chunk( - old_cylindrical, - new_cylindrical, - &mut loading_chunks, - &mut unloading_chunks, - ); + + let (loading_iter, unloading_iter) = + Cylindrical::changed_chunks(old_cylindrical, new_cylindrical); + let loading_chunks: Vec<_> = loading_iter.collect(); + let unloading_chunks: Vec<_> = unloading_iter.collect(); // Use the chunk_manager's world reference, which is updated on dimension change. // This ensures we load chunks from the correct world after portal teleportation. @@ -95,13 +92,11 @@ pub async fn update_position(player: &Arc) { ); world }; - player.watched_section.store(new_cylindrical); - if let ClientPlatform::Java(_) = player.client.as_ref() { + if let ClientPlatform::Java(client) = player.client.as_ref() { for chunk in &unloading_chunks { - player - .client + client .enqueue_packet(&CUnloadChunk::new(chunk.x, chunk.y)) .await; } diff --git a/pumpkin/src/world/mod.rs b/pumpkin/src/world/mod.rs index d858eb240..b98ce485a 100644 --- a/pumpkin/src/world/mod.rs +++ b/pumpkin/src/world/mod.rs @@ -1818,7 +1818,7 @@ impl World { &self, base_config: &BasicConfiguration, player: Arc, - server: &Server, + server: &Arc, ) { static CREATIVE_CONTENT: std::sync::OnceLock<(Vec, Vec)> = std::sync::OnceLock::new(); @@ -2557,8 +2557,8 @@ impl World { ) .color_named(NamedColor::Yellow); - let event = PlayerJoinEvent::new(player.clone(), msg_comp); - let event = server.plugin_manager.fire(event).await; + let mut event = PlayerJoinEvent::new(player.clone(), msg_comp); + server.plugin_manager.fire(server, &mut event).await; if !event.cancelled { self.broadcast_system_message(&event.join_message, false) @@ -3231,9 +3231,9 @@ impl World { [TextComponent::text(player.gameprofile.name.clone())], ) .color_named(NamedColor::Yellow); - let event = PlayerJoinEvent::new(player.clone(), msg_comp); + let mut event = PlayerJoinEvent::new(player.clone(), msg_comp); - let event = server.plugin_manager.fire(event).await; + server.plugin_manager.fire(server, &mut event).await; if !event.cancelled { self.broadcast_system_message(&event.join_message, false) @@ -3425,18 +3425,16 @@ impl World { // the non-cancellable PlayerRespawnEvent, which observes the resolved world. let (resolved_world, position, yaw, pitch) = if let Some(new_world) = candidate_world { if let Some(server) = self.server.upgrade() { - let event = server - .plugin_manager - .fire(PlayerChangeWorldEvent { - player: player.clone(), - previous_world: self.clone(), - new_world: new_world.clone(), - position, - yaw, - pitch, - cancelled: false, - }) - .await; + let mut event = PlayerChangeWorldEvent { + player: player.clone(), + previous_world: self.clone(), + new_world: new_world.clone(), + position, + yaw, + pitch, + cancelled: false, + }; + server.plugin_manager.fire(&server, &mut event).await; if event.cancelled { (None, position, yaw, pitch) @@ -3508,17 +3506,20 @@ impl World { // Notify plugins that the player has respawned (non-cancellable). if let Some(server) = self.server.upgrade() { - let _ = server + server .plugin_manager - .fire(PlayerRespawnEvent::new( - player.clone(), - self.clone(), - target_world.clone(), - position, - yaw, - pitch, - alive, - )) + .fire( + &server, + &mut PlayerRespawnEvent::new( + player.clone(), + self.clone(), + target_world.clone(), + position, + yaw, + pitch, + alive, + ), + ) .await; } @@ -4167,21 +4168,17 @@ impl World { [TextComponent::text(player.gameprofile.name.clone())], ) .color_named(NamedColor::Yellow); - let event = PlayerLeaveEvent::new(player.clone(), msg_comp); + let mut event = PlayerLeaveEvent::new(player.clone(), msg_comp); - let event = self - .server - .upgrade() - .unwrap() - .plugin_manager - .fire(event) - .await; + if let Some(server) = self.server.upgrade() { + server.plugin_manager.fire(&server, &mut event).await; - if !event.cancelled { - for player in self.players.load().iter() { - player.send_system_message(&event.leave_message).await; + if !event.cancelled { + for player in self.players.load().iter() { + player.send_system_message(&event.leave_message).await; + } + info!("{}", event.leave_message.to_pretty_console()); } - info!("{}", event.leave_message.to_pretty_console()); } } } @@ -4571,7 +4568,7 @@ impl World { if is_air(broken_block_state) { return None; } - let event = BlockBreakEvent::new( + let mut event = BlockBreakEvent::new( cause.clone(), broken_block, *position, @@ -4579,13 +4576,12 @@ impl World { !flags.contains(BlockFlags::SKIP_DROPS), ); - let event = self - .server - .upgrade() - .unwrap() - .plugin_manager - .fire::(event) - .await; + if let Some(server) = self.server.upgrade() { + server + .plugin_manager + .fire::(&server, &mut event) + .await; + } if !event.cancelled { let mut flags = flags;