Skip to main content

steel_core/player/connection/
java.rs

1//! This module contains the `JavaConnection` struct, which is used to represent a connection to a Java client.
2use std::io::Cursor;
3use std::sync::{Arc, Weak};
4use std::time::{SystemTime, UNIX_EPOCH};
5
6use steel_protocol::packet_reader::TCPNetworkDecoder;
7use steel_protocol::packet_traits::{ClientPacket, CompressionInfo, EncodedPacket, ServerPacket};
8use steel_protocol::packet_writer::TCPNetworkEncoder;
9use steel_protocol::packets::common::{
10    CDisconnect, CKeepAlive, CPongResponse, SClientInformation, SCustomPayload, SKeepAlive,
11    SPingRequest,
12};
13use steel_protocol::packets::game::{
14    CBundleDelimiter, CCommandSuggestions, ClientCommandAction, PlayerAction, PlayerCommandAction,
15    SAcceptTeleportation, SAttack, SChangeDifficulty, SChangeGameMode, SChat, SChatAck,
16    SChatCommand, SChatSessionUpdate, SChunkBatchReceived, SClientCommand, SClientTickEnd,
17    SCommandSuggestion, SContainerButtonClick, SContainerClick, SContainerClose,
18    SContainerSlotStateChanged, SInteract, SMovePlayer, SMovePlayerPos, SMovePlayerPosRot,
19    SMovePlayerRot, SMovePlayerStatusOnly, SMoveVehicle, SPickItemFromBlock, SPlayerAbilities,
20    SPlayerAction, SPlayerCommand, SPlayerInput, SPlayerLoad, SRenameItem, SSetCarriedItem,
21    SSetCreativeModeSlot, SSignUpdate, SSpectatorAction, SSwing, SUseItem, SUseItemOn,
22};
23
24use steel_protocol::utils::{ConnectionProtocol, PacketError, RawPacket};
25use steel_registry::packets::play;
26use steel_utils::locks::{AsyncMutex, SyncMutex};
27use steel_utils::translations;
28use text_components::content::Resolvable;
29use text_components::custom::CustomData;
30use text_components::resolving::TextResolutor;
31use text_components::{Modifier, TextComponent, format::Color};
32use tokio::io::{BufReader, BufWriter};
33use tokio::net::tcp::{OwnedReadHalf, OwnedWriteHalf};
34use tokio::select;
35use tokio::sync::mpsc::{UnboundedReceiver, UnboundedSender, error::TryRecvError};
36use tokio_util::sync::CancellationToken;
37
38use crate::command::{handle_client_request, sender::CommandSender};
39use crate::player::Player;
40use crate::player::connection::NetworkConnection;
41use crate::server::Server;
42
43/// Shared Java socket writer.
44pub type JavaNetworkWriter = Arc<AsyncMutex<Option<TCPNetworkEncoder<BufWriter<OwnedWriteHalf>>>>>;
45
46/// Outbound packet queue message for Java connections.
47pub enum OutboundPacket {
48    /// Normal packet write that may be interrupted by connection shutdown.
49    Packet(EncodedPacket),
50    /// Final disconnect packet that must be flushed before closing the socket.
51    Disconnect(EncodedPacket),
52}
53
54/// A decoded play packet whose handler runs in the server's inter-tick packet phase.
55pub(crate) struct ScheduledPlayPacket(ScheduledPlayPacketKind);
56
57/// Cross-player concurrency permitted for a scheduled packet handler.
58#[derive(Clone, Copy, Debug, Eq, PartialEq)]
59pub(crate) enum ScheduledPacketExecution {
60    /// The handler may overlap handlers for other players, but never its own player lane.
61    ///
62    /// Shared mutations must be fully linearized by their resource locks, and the handler must
63    /// tolerate cross-player execution order differing from packet submission order.
64    PlayerLocal,
65    /// The handler may overlap player-local work, but not another serialized handler. Serialized
66    /// handlers start in global packet submission order.
67    Serialized,
68    /// The handler is a global submission-order barrier and must not overlap scheduled work.
69    Exclusive,
70}
71
72enum ScheduledPlayPacketKind {
73    AcceptTeleportation(SAcceptTeleportation),
74    Attack(SAttack),
75    Interact(SInteract),
76    CustomPayload(SCustomPayload),
77    Chat(Box<SChat>),
78    ChatAck(SChatAck),
79    ChatSessionUpdate(SChatSessionUpdate),
80    ClientInformation(SClientInformation),
81    ClientTickEnd,
82    MovePlayer(SMovePlayer),
83    MoveVehicle(SMoveVehicle),
84    PlayerLoaded,
85    ChatCommand(SChatCommand),
86    CommandSuggestion(SCommandSuggestion),
87    ContainerButtonClick(SContainerButtonClick),
88    ContainerClick(SContainerClick),
89    ContainerClose(SContainerClose),
90    ContainerSlotStateChanged(SContainerSlotStateChanged),
91    SetCreativeModeSlot(SSetCreativeModeSlot),
92    PlayerInput(SPlayerInput),
93    PlayerCommand(SPlayerCommand),
94    PlayerAbilities(SPlayerAbilities),
95    RenameItem(SRenameItem),
96    UseItemOn(SUseItemOn),
97    UseItem(SUseItem),
98    SetCarriedItem(SSetCarriedItem),
99    Swing(SSwing),
100    PlayerAction(SPlayerAction),
101    PickItemFromBlock(SPickItemFromBlock),
102    SignUpdate(SSignUpdate),
103    SpectatorAction(SSpectatorAction),
104    ClientCommand(SClientCommand),
105    ChangeGameMode(SChangeGameMode),
106    ChangeDifficulty(SChangeDifficulty),
107}
108
109enum ImmediatePlayPacket {
110    KeepAlive(SKeepAlive),
111    PingRequest(SPingRequest),
112    ChunkBatchReceived(SChunkBatchReceived),
113    Unknown(i32),
114}
115
116enum DecodedPlayPacket {
117    Scheduled(ScheduledPlayPacket),
118    Immediate(ImmediatePlayPacket),
119}
120
121impl ScheduledPlayPacket {
122    /// Returns whether this packet acknowledges target-world synchronization.
123    pub(crate) const fn is_domain_handshake_packet(&self) -> bool {
124        matches!(
125            self.0,
126            ScheduledPlayPacketKind::AcceptTeleportation(_) | ScheduledPlayPacketKind::PlayerLoaded
127        )
128    }
129
130    /// Returns whether this is the death screen's one-shot respawn request.
131    pub(crate) const fn is_perform_respawn(&self) -> bool {
132        matches!(
133            self.0,
134            ScheduledPlayPacketKind::ClientCommand(SClientCommand {
135                action: ClientCommandAction::PerformRespawn,
136            })
137        )
138    }
139
140    #[cfg(test)]
141    pub(crate) const fn perform_respawn_for_test() -> Self {
142        Self(ScheduledPlayPacketKind::ClientCommand(SClientCommand {
143            action: ClientCommandAction::PerformRespawn,
144        }))
145    }
146
147    /// Returns the handler's audited cross-player concurrency class.
148    ///
149    /// This match is intentionally exhaustive so every newly implemented packet requires an
150    /// explicit concurrency decision.
151    pub(crate) const fn execution(&self) -> ScheduledPacketExecution {
152        match &self.0 {
153            // These handlers touch player-owned state or use an individually linearizable shared
154            // operation. The player lane preserves same-player order.
155            ScheduledPlayPacketKind::ChatSessionUpdate(_)
156            | ScheduledPlayPacketKind::ClientInformation(_)
157            | ScheduledPlayPacketKind::ClientTickEnd
158            | ScheduledPlayPacketKind::PlayerLoaded
159            | ScheduledPlayPacketKind::ChatCommand(_)
160            | ScheduledPlayPacketKind::CommandSuggestion(_)
161            | ScheduledPlayPacketKind::ContainerClose(_)
162            | ScheduledPlayPacketKind::SetCreativeModeSlot(_)
163            | ScheduledPlayPacketKind::PlayerInput(_)
164            | ScheduledPlayPacketKind::PlayerAbilities(_)
165            | ScheduledPlayPacketKind::SetCarriedItem(_)
166            | ScheduledPlayPacketKind::Swing(_)
167            | ScheduledPlayPacketKind::PickItemFromBlock(_)
168            | ScheduledPlayPacketKind::ClientCommand(_) => ScheduledPacketExecution::PlayerLocal,
169            ScheduledPlayPacketKind::PlayerCommand(packet) => match packet.action {
170                PlayerCommandAction::StartSprinting
171                | PlayerCommandAction::StopSprinting
172                | PlayerCommandAction::StartFallFlying => ScheduledPacketExecution::PlayerLocal,
173                PlayerCommandAction::LeaveBed => ScheduledPacketExecution::Serialized,
174                // These handlers are not implemented, so their eventual vehicle transaction
175                // cannot yet be audited against concurrently player-local work.
176                PlayerCommandAction::StartRidingJump
177                | PlayerCommandAction::StopRidingJump
178                | PlayerCommandAction::OpenVehicleInventory => ScheduledPacketExecution::Exclusive,
179            },
180            ScheduledPlayPacketKind::PlayerAction(packet) => match packet.action {
181                PlayerAction::AbortDestroyBlock | PlayerAction::SwapItemWithOffhand => {
182                    ScheduledPacketExecution::PlayerLocal
183                }
184                PlayerAction::StartDestroyBlock
185                | PlayerAction::StopDestroyBlock
186                | PlayerAction::DropAllItems
187                | PlayerAction::DropItem => ScheduledPacketExecution::Serialized,
188                // Release-use is not implemented, while stab spans independently locked targets;
189                // neither can yet overlap player-local work safely.
190                PlayerAction::ReleaseUseItem | PlayerAction::Stab => {
191                    ScheduledPacketExecution::Exclusive
192                }
193            },
194            // Position, world, menu, chat, and domain mutations may overlap player-local work but
195            // retain one global mutation order matching the packet submission order.
196            ScheduledPlayPacketKind::AcceptTeleportation(_)
197            | ScheduledPlayPacketKind::Chat(_)
198            | ScheduledPlayPacketKind::ChatAck(_)
199            | ScheduledPlayPacketKind::MovePlayer(_)
200            | ScheduledPlayPacketKind::MoveVehicle(_)
201            | ScheduledPlayPacketKind::ContainerClick(_)
202            | ScheduledPlayPacketKind::RenameItem(_)
203            | ScheduledPlayPacketKind::UseItemOn(_)
204            | ScheduledPlayPacketKind::UseItem(_)
205            | ScheduledPlayPacketKind::SignUpdate(_)
206            | ScheduledPlayPacketKind::SpectatorAction(_)
207            | ScheduledPlayPacketKind::ChangeGameMode(_)
208            | ScheduledPlayPacketKind::ChangeDifficulty(_) => ScheduledPacketExecution::Serialized,
209            // Combat spans source and target state, custom payloads have no constrained resource
210            // contract, and the unimplemented menu handlers have no auditable transaction yet.
211            ScheduledPlayPacketKind::Attack(_)
212            | ScheduledPlayPacketKind::Interact(_)
213            | ScheduledPlayPacketKind::CustomPayload(_)
214            | ScheduledPlayPacketKind::ContainerButtonClick(_)
215            | ScheduledPlayPacketKind::ContainerSlotStateChanged(_) => {
216                ScheduledPacketExecution::Exclusive
217            }
218        }
219    }
220
221    pub(crate) const fn can_process_before_join(&self) -> bool {
222        matches!(
223            &self.0,
224            ScheduledPlayPacketKind::AcceptTeleportation(_)
225                | ScheduledPlayPacketKind::ClientInformation(_)
226                | ScheduledPlayPacketKind::ClientTickEnd
227                | ScheduledPlayPacketKind::CustomPayload(_)
228                | ScheduledPlayPacketKind::ChatAck(_)
229                | ScheduledPlayPacketKind::ChatSessionUpdate(_)
230                | ScheduledPlayPacketKind::PlayerLoaded
231        )
232    }
233
234    pub(crate) fn handle(self, player: Arc<Player>, server: &Arc<Server>) {
235        if !player.has_joined_world() && !self.can_process_before_join() {
236            return;
237        }
238
239        match self.0 {
240            ScheduledPlayPacketKind::AcceptTeleportation(packet) => {
241                player.handle_accept_teleportation(packet);
242            }
243            ScheduledPlayPacketKind::Attack(packet) => player.handle_attack(packet),
244            ScheduledPlayPacketKind::Interact(packet) => player.handle_interact(packet),
245            ScheduledPlayPacketKind::CustomPayload(packet) => {
246                player.handle_custom_payload(packet);
247            }
248            ScheduledPlayPacketKind::Chat(packet) => {
249                player.handle_chat(*packet, Arc::clone(&player));
250            }
251            ScheduledPlayPacketKind::ChatAck(packet) => player.handle_chat_ack(packet),
252            ScheduledPlayPacketKind::ChatSessionUpdate(packet) => {
253                player.handle_chat_session_update(packet);
254            }
255            ScheduledPlayPacketKind::ClientInformation(packet) => {
256                player.handle_client_information(packet);
257            }
258            ScheduledPlayPacketKind::ClientTickEnd => player.handle_client_tick_end(),
259            ScheduledPlayPacketKind::MovePlayer(packet) => player.handle_move_player(packet),
260            ScheduledPlayPacketKind::MoveVehicle(packet) => player.handle_move_vehicle(packet),
261            ScheduledPlayPacketKind::PlayerLoaded => {
262                if player.mark_client_loaded_from_network() {
263                    player.send_inventory_to_remote();
264                }
265            }
266            ScheduledPlayPacketKind::ChatCommand(packet) => {
267                if server
268                    .submit_command(CommandSender::Player(Arc::clone(&player)), packet.command)
269                    .is_err()
270                {
271                    player.send_message(
272                        &TextComponent::const_plain("Command queue is full").color(Color::Red),
273                    );
274                }
275                player.detect_command_rate_spam();
276            }
277            ScheduledPlayPacketKind::CommandSuggestion(packet) => {
278                if server
279                    .submit_command_suggestions(Arc::clone(&player), packet.id, packet.command)
280                    .is_err()
281                {
282                    player.send_packet(CCommandSuggestions::new(packet.id, 0, 0, Vec::new()));
283                }
284            }
285            ScheduledPlayPacketKind::ContainerButtonClick(packet) => {
286                player.handle_container_button_click(packet);
287            }
288            ScheduledPlayPacketKind::ContainerClick(packet) => {
289                player.handle_container_click(packet);
290            }
291            ScheduledPlayPacketKind::ContainerClose(packet) => {
292                player.handle_container_close(packet);
293            }
294            ScheduledPlayPacketKind::ContainerSlotStateChanged(packet) => {
295                player.handle_container_slot_state_changed(packet);
296            }
297            ScheduledPlayPacketKind::SetCreativeModeSlot(packet) => {
298                player.handle_set_creative_mode_slot(packet);
299            }
300            ScheduledPlayPacketKind::PlayerInput(packet) => player.handle_player_input(packet),
301            ScheduledPlayPacketKind::PlayerCommand(packet) => {
302                player.handle_player_command(packet);
303            }
304            ScheduledPlayPacketKind::PlayerAbilities(packet) => {
305                player.handle_player_abilities(packet);
306            }
307            ScheduledPlayPacketKind::RenameItem(packet) => player.handle_rename_item(packet),
308            ScheduledPlayPacketKind::UseItemOn(packet) => player.handle_use_item_on(packet),
309            ScheduledPlayPacketKind::UseItem(packet) => player.handle_use_item(packet),
310            ScheduledPlayPacketKind::SetCarriedItem(packet) => {
311                player.handle_set_carried_item(packet);
312            }
313            ScheduledPlayPacketKind::Swing(packet) => player.swing(packet.hand, false),
314            ScheduledPlayPacketKind::PlayerAction(packet) => {
315                player.handle_player_action(packet);
316            }
317            ScheduledPlayPacketKind::PickItemFromBlock(packet) => {
318                player.handle_pick_item_from_block(packet);
319            }
320            ScheduledPlayPacketKind::SignUpdate(packet) => player.handle_sign_update(packet),
321            ScheduledPlayPacketKind::SpectatorAction(packet) => {
322                player.handle_spectator_action(packet);
323            }
324            ScheduledPlayPacketKind::ClientCommand(packet) => {
325                player.handle_client_command(packet.action);
326            }
327            ScheduledPlayPacketKind::ChangeGameMode(packet) => {
328                handle_client_request(&player, server, packet.gamemode);
329            }
330            ScheduledPlayPacketKind::ChangeDifficulty(packet) => {
331                player.handle_change_difficulty(packet.difficulty);
332            }
333        }
334    }
335}
336
337/// Builder for creating packet bundles.
338///
339/// Used with [`JavaConnection::send_bundle`] to send multiple packets atomically.
340pub struct BundleBuilder {
341    packets: Vec<EncodedPacket>,
342    compression: Option<CompressionInfo>,
343}
344
345impl BundleBuilder {
346    /// Creates a new `BundleBuilder` with the given compression settings.
347    #[must_use]
348    pub const fn new(compression: Option<CompressionInfo>) -> Self {
349        Self {
350            packets: Vec::new(),
351            compression,
352        }
353    }
354
355    /// Adds a packet to the bundle.
356    ///
357    /// # Panics
358    /// Panics if the packet fails to encode.
359    pub fn add<P: ClientPacket>(&mut self, packet: P) {
360        let encoded = EncodedPacket::from_bare(packet, self.compression, ConnectionProtocol::Play)
361            .expect("Failed to encode packet");
362        self.packets.push(encoded);
363    }
364
365    /// Consumes the builder and returns the collected encoded packets.
366    #[must_use]
367    pub fn into_packets(self) -> Vec<EncodedPacket> {
368        self.packets
369    }
370}
371
372#[expect(
373    clippy::struct_field_names,
374    reason = "alive_ prefix is intentional to group related keep-alive fields"
375)]
376struct KeepAliveTracker {
377    alive_time: u64,
378    alive_pending: bool,
379    alive_id: u64,
380}
381
382/// A connection to a Java client.
383pub struct JavaConnection {
384    outgoing_packets: UnboundedSender<OutboundPacket>,
385    cancel_token: CancellationToken,
386    compression: Option<CompressionInfo>,
387    network_writer: JavaNetworkWriter,
388    id: u64,
389
390    player: Weak<Player>,
391    keep_alive_tracker: SyncMutex<KeepAliveTracker>,
392    latency: SyncMutex<u32>,
393}
394
395impl JavaConnection {
396    /// Creates a new `JavaConnection`.
397    pub const fn new(
398        outgoing_packets: UnboundedSender<OutboundPacket>,
399        cancel_token: CancellationToken,
400        compression: Option<CompressionInfo>,
401        network_writer: JavaNetworkWriter,
402        id: u64,
403        player: Weak<Player>,
404    ) -> Self {
405        Self {
406            outgoing_packets,
407            cancel_token,
408            compression,
409            network_writer,
410            id,
411            player,
412            keep_alive_tracker: SyncMutex::new(KeepAliveTracker {
413                alive_time: 0,
414                alive_pending: false,
415                alive_id: 0,
416            }),
417            latency: SyncMutex::new(0),
418        }
419    }
420
421    async fn write_packet_now(&self, packet: &EncodedPacket) -> Result<(), PacketError> {
422        let mut network_writer = self.network_writer.lock().await;
423        let Some(network_writer) = network_writer.as_mut() else {
424            return Err(PacketError::ConnectionClosed);
425        };
426        network_writer.write_packet(packet).await
427    }
428
429    async fn release_network_writer(&self) {
430        self.network_writer.lock().await.take();
431    }
432
433    /// Ticks the connection.
434    pub fn tick(&self) {
435        self.keep_connection_alive();
436    }
437
438    fn keep_connection_alive(&self) {
439        let mut tracker = self.keep_alive_tracker.lock();
440        let now = SystemTime::now()
441            .duration_since(UNIX_EPOCH)
442            .expect("System time before UNIX EPOCH")
443            .as_millis() as u64;
444
445        if now - tracker.alive_time >= 15000 {
446            if tracker.alive_pending {
447                self.disconnect(translations::DISCONNECT_TIMEOUT.msg());
448            } else {
449                tracker.alive_pending = true;
450                tracker.alive_id = now;
451                tracker.alive_time = now;
452                self.send_packet(CKeepAlive::new(tracker.alive_id as i64));
453            }
454        }
455    }
456
457    /// Handles a keep alive packet.
458    #[expect(
459        clippy::cast_possible_truncation,
460        reason = "latency saturates at u32::MAX ms (~49 days), which is unreachable in practice"
461    )]
462    fn handle_keep_alive(&self, packet: SKeepAlive) {
463        let mut tracker = self.keep_alive_tracker.lock();
464        if tracker.alive_pending && packet.id as u64 == tracker.alive_id {
465            let now = SystemTime::now()
466                .duration_since(UNIX_EPOCH)
467                .expect("System time before UNIX EPOCH")
468                .as_millis() as u64;
469
470            let time = now.saturating_sub(tracker.alive_time) as u32;
471            tracker.alive_pending = false;
472            drop(tracker);
473            let mut latency = self.latency.lock();
474            *latency = (*latency * 3 + time) / 4;
475        } else {
476            self.disconnect(translations::DISCONNECT_TIMEOUT.msg());
477        }
478    }
479
480    /// Returns the current latency in milliseconds.
481    /// This is a smoothed average calculated from keep-alive round-trip times.
482    #[must_use]
483    pub fn latency(&self) -> i32 {
484        *self.latency.lock() as i32
485    }
486
487    /// Disconnects the client.
488    pub fn disconnect(&self, reason: impl Into<TextComponent>) {
489        let packet = match EncodedPacket::from_bare(
490            CDisconnect::new(&reason.into(), self),
491            self.compression,
492            ConnectionProtocol::Play,
493        ) {
494            Ok(packet) => packet,
495            Err(err) => {
496                log::warn!(
497                    "Failed to encode disconnect packet for client {}: {err}",
498                    self.id
499                );
500                self.close();
501                return;
502            }
503        };
504        if self
505            .outgoing_packets
506            .send(OutboundPacket::Disconnect(packet))
507            .is_err()
508        {
509            self.close();
510            return;
511        }
512        self.close();
513    }
514
515    /// Sends a packet to the client.
516    ///
517    /// # Panics
518    /// - If the packet fails to be encoded.
519    /// - If the packet fails to be sent through the channel.
520    pub fn send_packet<P: ClientPacket>(&self, packet: P) {
521        let packet = EncodedPacket::from_bare(packet, self.compression, ConnectionProtocol::Play)
522            .expect("Failed to encode packet");
523        if self
524            .outgoing_packets
525            .send(OutboundPacket::Packet(packet))
526            .is_err()
527        {
528            self.close();
529        }
530    }
531
532    /// Sends an encoded packet to the client.
533    ///
534    /// # Panics
535    /// - If the packet fails to be sent through the channel.
536    pub fn send_encoded_packet(&self, packet: EncodedPacket) {
537        if self
538            .outgoing_packets
539            .send(OutboundPacket::Packet(packet))
540            .is_err()
541        {
542            self.close();
543        }
544    }
545
546    /// Closes the connection.
547    pub fn close(&self) {
548        self.cancel_token.cancel();
549    }
550
551    /// Returns whether the connection is closed.
552    #[must_use]
553    pub fn closed(&self) -> bool {
554        self.cancel_token.is_cancelled()
555    }
556
557    /// Waits for the connection to be closed.
558    pub async fn wait_for_close(&self) {
559        self.cancel_token.cancelled().await;
560    }
561
562    const fn can_process_before_join(packet_id: i32) -> bool {
563        matches!(
564            packet_id,
565            play::S_ACCEPT_TELEPORTATION
566                | play::S_KEEP_ALIVE
567                | play::S_PING_REQUEST
568                | play::S_CLIENT_INFORMATION
569                | play::S_CUSTOM_PAYLOAD
570                | play::S_CHUNK_BATCH_RECEIVED
571                | play::S_CHAT_SESSION_UPDATE
572                | play::S_CHAT_ACK
573                | play::S_CLIENT_TICK_END
574                | play::S_PLAYER_LOADED
575        )
576    }
577
578    const fn can_process_during_domain_handshake(packet_id: i32) -> bool {
579        matches!(
580            packet_id,
581            play::S_ACCEPT_TELEPORTATION | play::S_CHUNK_BATCH_RECEIVED | play::S_PLAYER_LOADED
582        )
583    }
584
585    /// Decodes and dispatches one packet received from the client.
586    fn process_packet(
587        &self,
588        packet: RawPacket,
589        player: Arc<Player>,
590        server: &Server,
591    ) -> Result<(), PacketError> {
592        if !player.has_joined_world() && !Self::can_process_before_join(packet.id) {
593            return Ok(());
594        }
595
596        let payload_bytes = packet.payload.len();
597        let Some(packet) = Self::decode_domain_gated_packet(packet, &player)? else {
598            return Ok(());
599        };
600
601        match packet {
602            DecodedPlayPacket::Scheduled(packet) => {
603                server.schedule_play_packet(player, packet, payload_bytes);
604            }
605            DecodedPlayPacket::Immediate(packet) => {
606                self.handle_immediate_packet(packet, &player);
607            }
608        }
609        Ok(())
610    }
611
612    fn decode_domain_gated_packet(
613        packet: RawPacket,
614        player: &Player,
615    ) -> Result<Option<DecodedPlayPacket>, PacketError> {
616        let maintenance_packet = matches!(packet.id, play::S_KEEP_ALIVE | play::S_PING_REQUEST);
617        let handshake_packet = Self::can_process_during_domain_handshake(packet.id);
618        if packet.id == play::S_CLIENT_COMMAND {
619            let decoded = Self::decode_play_packet(packet)?;
620            let perform_respawn = matches!(
621                &decoded,
622                DecodedPlayPacket::Scheduled(packet) if packet.is_perform_respawn()
623            );
624            if player.gate_domain_switch_packet(false, perform_respawn) {
625                return Ok(Some(decoded));
626            }
627            return Ok(None);
628        }
629
630        if !maintenance_packet && !player.gate_domain_switch_packet(handshake_packet, false) {
631            return Ok(None);
632        }
633
634        Self::decode_play_packet(packet).map(Some)
635    }
636
637    #[expect(
638        clippy::too_many_lines,
639        reason = "single match decode over all implemented play packets keeps protocol routing auditable"
640    )]
641    fn decode_play_packet(packet: RawPacket) -> Result<DecodedPlayPacket, PacketError> {
642        let data = &mut Cursor::new(packet.payload.as_slice());
643        let scheduled = |packet| DecodedPlayPacket::Scheduled(ScheduledPlayPacket(packet));
644
645        Ok(match packet.id {
646            play::S_ACCEPT_TELEPORTATION => {
647                scheduled(ScheduledPlayPacketKind::AcceptTeleportation(
648                    SAcceptTeleportation::read_packet(data)?,
649                ))
650            }
651            play::S_ATTACK => {
652                scheduled(ScheduledPlayPacketKind::Attack(SAttack::read_packet(data)?))
653            }
654            play::S_INTERACT => scheduled(ScheduledPlayPacketKind::Interact(
655                SInteract::read_packet(data)?,
656            )),
657            play::S_CUSTOM_PAYLOAD => scheduled(ScheduledPlayPacketKind::CustomPayload(
658                SCustomPayload::read_packet(data)?,
659            )),
660            play::S_CHAT => scheduled(ScheduledPlayPacketKind::Chat(Box::new(SChat::read_packet(
661                data,
662            )?))),
663            play::S_CHAT_SESSION_UPDATE => scheduled(ScheduledPlayPacketKind::ChatSessionUpdate(
664                SChatSessionUpdate::read_packet(data)?,
665            )),
666            play::S_CHAT_ACK => scheduled(ScheduledPlayPacketKind::ChatAck(SChatAck::read_packet(
667                data,
668            )?)),
669            play::S_CLIENT_INFORMATION => scheduled(ScheduledPlayPacketKind::ClientInformation(
670                SClientInformation::read_packet(data)?,
671            )),
672            play::S_CLIENT_TICK_END => {
673                let _ = SClientTickEnd::read_packet(data)?;
674                scheduled(ScheduledPlayPacketKind::ClientTickEnd)
675            }
676            play::S_CHUNK_BATCH_RECEIVED => DecodedPlayPacket::Immediate(
677                ImmediatePlayPacket::ChunkBatchReceived(SChunkBatchReceived::read_packet(data)?),
678            ),
679            play::S_KEEP_ALIVE => DecodedPlayPacket::Immediate(ImmediatePlayPacket::KeepAlive(
680                SKeepAlive::read_packet(data)?,
681            )),
682            play::S_MOVE_PLAYER_POS => scheduled(ScheduledPlayPacketKind::MovePlayer(
683                SMovePlayerPos::read_packet(data)?.into(),
684            )),
685            play::S_MOVE_PLAYER_POS_ROT => scheduled(ScheduledPlayPacketKind::MovePlayer(
686                SMovePlayerPosRot::read_packet(data)?.into(),
687            )),
688            play::S_MOVE_PLAYER_ROT => scheduled(ScheduledPlayPacketKind::MovePlayer(
689                SMovePlayerRot::read_packet(data)?.into(),
690            )),
691            play::S_MOVE_PLAYER_STATUS_ONLY => scheduled(ScheduledPlayPacketKind::MovePlayer(
692                SMovePlayerStatusOnly::read_packet(data)?.into(),
693            )),
694            play::S_MOVE_VEHICLE => scheduled(ScheduledPlayPacketKind::MoveVehicle(
695                SMoveVehicle::read_packet(data)?,
696            )),
697            play::S_PLAYER_LOADED => {
698                let _ = SPlayerLoad::read_packet(data)?;
699                scheduled(ScheduledPlayPacketKind::PlayerLoaded)
700            }
701            play::S_CHAT_COMMAND => scheduled(ScheduledPlayPacketKind::ChatCommand(
702                SChatCommand::read_packet(data)?,
703            )),
704            play::S_COMMAND_SUGGESTION => scheduled(ScheduledPlayPacketKind::CommandSuggestion(
705                SCommandSuggestion::read_packet(data)?,
706            )),
707            play::S_CONTAINER_BUTTON_CLICK => {
708                scheduled(ScheduledPlayPacketKind::ContainerButtonClick(
709                    SContainerButtonClick::read_packet(data)?,
710                ))
711            }
712            play::S_CONTAINER_CLICK => scheduled(ScheduledPlayPacketKind::ContainerClick(
713                SContainerClick::read_packet(data)?,
714            )),
715            play::S_CONTAINER_CLOSE => scheduled(ScheduledPlayPacketKind::ContainerClose(
716                SContainerClose::read_packet(data)?,
717            )),
718            play::S_CONTAINER_SLOT_STATE_CHANGED => {
719                scheduled(ScheduledPlayPacketKind::ContainerSlotStateChanged(
720                    SContainerSlotStateChanged::read_packet(data)?,
721                ))
722            }
723            play::S_SET_CREATIVE_MODE_SLOT => {
724                scheduled(ScheduledPlayPacketKind::SetCreativeModeSlot(
725                    SSetCreativeModeSlot::read_packet(data)?,
726                ))
727            }
728            play::S_PLAYER_INPUT => scheduled(ScheduledPlayPacketKind::PlayerInput(
729                SPlayerInput::read_packet(data)?,
730            )),
731            play::S_PLAYER_COMMAND => scheduled(ScheduledPlayPacketKind::PlayerCommand(
732                SPlayerCommand::read_packet(data)?,
733            )),
734            play::S_PLAYER_ABILITIES => scheduled(ScheduledPlayPacketKind::PlayerAbilities(
735                SPlayerAbilities::read_packet(data)?,
736            )),
737            play::S_RENAME_ITEM => scheduled(ScheduledPlayPacketKind::RenameItem(
738                SRenameItem::read_packet(data)?,
739            )),
740            play::S_USE_ITEM_ON => scheduled(ScheduledPlayPacketKind::UseItemOn(
741                SUseItemOn::read_packet(data)?,
742            )),
743            play::S_USE_ITEM => scheduled(ScheduledPlayPacketKind::UseItem(SUseItem::read_packet(
744                data,
745            )?)),
746            play::S_SET_CARRIED_ITEM => scheduled(ScheduledPlayPacketKind::SetCarriedItem(
747                SSetCarriedItem::read_packet(data)?,
748            )),
749            play::S_SWING => scheduled(ScheduledPlayPacketKind::Swing(SSwing::read_packet(data)?)),
750            play::S_PLAYER_ACTION => scheduled(ScheduledPlayPacketKind::PlayerAction(
751                SPlayerAction::read_packet(data)?,
752            )),
753            play::S_PICK_ITEM_FROM_BLOCK => scheduled(ScheduledPlayPacketKind::PickItemFromBlock(
754                SPickItemFromBlock::read_packet(data)?,
755            )),
756            play::S_SIGN_UPDATE => scheduled(ScheduledPlayPacketKind::SignUpdate(
757                SSignUpdate::read_packet(data)?,
758            )),
759            play::S_SPECTATOR_ACTION => scheduled(ScheduledPlayPacketKind::SpectatorAction(
760                SSpectatorAction::read_packet(data)?,
761            )),
762            play::S_CLIENT_COMMAND => scheduled(ScheduledPlayPacketKind::ClientCommand(
763                SClientCommand::read_packet(data)?,
764            )),
765            play::S_PING_REQUEST => DecodedPlayPacket::Immediate(ImmediatePlayPacket::PingRequest(
766                SPingRequest::read_packet(data)?,
767            )),
768            play::S_CHANGE_GAME_MODE => scheduled(ScheduledPlayPacketKind::ChangeGameMode(
769                SChangeGameMode::read_packet(data)?,
770            )),
771            play::S_CHANGE_DIFFICULTY => scheduled(ScheduledPlayPacketKind::ChangeDifficulty(
772                SChangeDifficulty::read_packet(data)?,
773            )),
774            id => DecodedPlayPacket::Immediate(ImmediatePlayPacket::Unknown(id)),
775        })
776    }
777
778    fn handle_immediate_packet(&self, packet: ImmediatePlayPacket, player: &Player) {
779        match packet {
780            ImmediatePlayPacket::KeepAlive(packet) => self.handle_keep_alive(packet),
781            ImmediatePlayPacket::PingRequest(packet) => {
782                player.send_packet(CPongResponse::new(packet.time));
783            }
784            ImmediatePlayPacket::ChunkBatchReceived(packet) => {
785                player
786                    .chunk_sender
787                    .lock()
788                    .on_chunk_batch_received_by_client(packet.desired_chunks_per_tick);
789            }
790            ImmediatePlayPacket::Unknown(id) => log::info!("play packet id {id} is not known"),
791        }
792    }
793
794    /// Listens for packets from the client.
795    pub async fn listener(
796        &self,
797        mut reader: TCPNetworkDecoder<BufReader<OwnedReadHalf>>,
798        server: Arc<Server>,
799    ) {
800        loop {
801            select! {
802                () = self.wait_for_close() => {
803                    break;
804                }
805                packet = reader.get_raw_packet() => {
806                    match packet {
807                        Ok(packet) => {
808                            if let Some(player) = self.player.upgrade()
809                                && let Err(err) = self.process_packet(packet, player, &server) {
810                                log::warn!(
811                                    "Failed to get packet from client {}: {err}",
812                                    self.id
813                                );
814                            }
815                        }
816                        Err(err) => {
817                            log::debug!("Failed to get raw packet from client {}: {err}", self.id);
818                            self.close();
819                        }
820                    }
821                }
822            }
823        }
824    }
825
826    /// Sends packets to the client.
827    ///
828    pub async fn sender(&self, mut sender_recv: UnboundedReceiver<OutboundPacket>) {
829        loop {
830            select! {
831                biased;
832                () = self.wait_for_close() => {
833                    self.write_queued_disconnect(&mut sender_recv).await;
834                    break;
835                }
836                outbound = sender_recv.recv() => {
837                    if let Some(outbound) = outbound {
838                        let (packet, close_after_write) = match outbound {
839                            OutboundPacket::Packet(packet) => (packet, false),
840                            OutboundPacket::Disconnect(packet) => (packet, true),
841                        };
842
843                        if close_after_write {
844                            if let Err(err) = self.write_packet_now(&packet).await {
845                                log::warn!("Failed to send disconnect packet to client {}: {err}", self.id);
846                            }
847                            self.close();
848                            break;
849                        }
850
851                        let write_result = self.write_packet_now(&packet);
852                        select! {
853                            biased;
854                            () = self.wait_for_close() => {
855                                self.write_queued_disconnect(&mut sender_recv).await;
856                                break;
857                            },
858                            result = write_result => {
859                                if let Err(err) = result {
860                                    log::warn!("Failed to send packet to client {}: {err}", self.id);
861                                    self.close();
862                                    break;
863                                }
864                            }
865                        }
866                    } else {
867                        //log::warn!(
868                        //    "Internal packet_sender_recv channel closed for client {}",
869                        //    self.id
870                        //);
871                        self.close();
872                    }
873                }
874            }
875        }
876
877        self.release_network_writer().await;
878
879        let Some(player) = self.player.upgrade() else {
880            return;
881        };
882        if !player.has_joined_world() || player.server().cancel_token.is_cancelled() {
883            return;
884        }
885        player.server().queue_player_disconnect(player);
886    }
887
888    async fn write_queued_disconnect(&self, sender_recv: &mut UnboundedReceiver<OutboundPacket>) {
889        let mut disconnect_packet = None;
890        loop {
891            match sender_recv.try_recv() {
892                Ok(OutboundPacket::Packet(_)) => {}
893                Ok(OutboundPacket::Disconnect(packet)) => disconnect_packet = Some(packet),
894                Err(TryRecvError::Empty | TryRecvError::Disconnected) => break,
895            }
896        }
897
898        let Some(packet) = disconnect_packet else {
899            return;
900        };
901        if let Err(err) = self.write_packet_now(&packet).await {
902            log::warn!(
903                "Failed to send disconnect packet to client {} during close: {err}",
904                self.id
905            );
906        }
907    }
908}
909
910impl TextResolutor for JavaConnection {
911    fn resolve_content(&self, _resolvable: &Resolvable) -> TextComponent {
912        TextComponent::new()
913    }
914
915    fn resolve_custom(&self, _data: &CustomData) -> Option<TextComponent> {
916        None
917    }
918
919    fn translate(&self, _key: &str) -> Option<String> {
920        None
921    }
922}
923
924impl NetworkConnection for JavaConnection {
925    fn compression(&self) -> Option<CompressionInfo> {
926        self.compression
927    }
928
929    fn send_encoded(&self, packet: EncodedPacket) {
930        self.send_encoded_packet(packet);
931    }
932
933    fn send_encoded_bundle(&self, packets: Vec<EncodedPacket>) {
934        self.send_packet(CBundleDelimiter);
935        for packet in packets {
936            self.send_encoded_packet(packet);
937        }
938        self.send_packet(CBundleDelimiter);
939    }
940
941    fn disconnect_with_reason(&self, reason: TextComponent) {
942        self.disconnect(reason);
943    }
944
945    fn tick(&self) {
946        self.keep_connection_alive();
947    }
948
949    fn latency(&self) -> i32 {
950        *self.latency.lock() as i32
951    }
952
953    fn close(&self) {
954        self.cancel_token.cancel();
955    }
956
957    fn closed(&self) -> bool {
958        self.cancel_token.is_cancelled()
959    }
960}
961
962#[cfg(test)]
963mod tests {
964    use std::array;
965
966    use crate::{
967        entity::{Entity as _, LivingEntity as _},
968        test_support::{TestPlayerBuilder, fresh_test_world},
969    };
970    use rustc_hash::FxHashMap;
971    use steel_protocol::packets::common::{ChatVisibility, HumanoidArm, ParticleStatus};
972    use steel_protocol::packets::game::{ClickType, ClientCommandAction, HashedStack};
973    use steel_registry::{blocks::properties::Direction, item_stack::ItemStack};
974    use steel_utils::{BlockPos, codec::VarInt, types::InteractionHand};
975    use uuid::Uuid;
976
977    use super::*;
978
979    fn decode(packet: RawPacket) -> DecodedPlayPacket {
980        let Ok(decoded) = JavaConnection::decode_play_packet(packet) else {
981            panic!("test play packet should decode");
982        };
983        decoded
984    }
985
986    fn execution(kind: ScheduledPlayPacketKind) -> ScheduledPacketExecution {
987        ScheduledPlayPacket(kind).execution()
988    }
989
990    #[test]
991    fn pre_join_custom_payload_uses_serverbound_play_packet_id() {
992        assert!(JavaConnection::can_process_before_join(
993            play::S_CUSTOM_PAYLOAD
994        ));
995        assert!(!JavaConnection::can_process_before_join(
996            play::C_CUSTOM_PAYLOAD
997        ));
998    }
999
1000    #[test]
1001    fn queued_domain_switch_records_only_perform_respawn_at_connection_gate() {
1002        let world = fresh_test_world("queued_domain_switch_respawn_packet");
1003        let player = TestPlayerBuilder::new(world, Uuid::from_u128(1), "RespawnTester", 1).build();
1004        let Some(token) = player.begin_pending_world_change() else {
1005            panic!("test player should acquire a world-change token");
1006        };
1007        assert!(player.begin_domain_switch(token));
1008        player.set_health(0.0);
1009
1010        let request_stats = JavaConnection::decode_domain_gated_packet(
1011            RawPacket {
1012                id: play::S_CLIENT_COMMAND,
1013                payload: vec![ClientCommandAction::RequestStats as u8],
1014            },
1015            &player,
1016        );
1017        assert!(matches!(request_stats, Ok(None)));
1018        assert!(!player.has_deferred_death_respawn_for_test());
1019
1020        let perform_respawn = JavaConnection::decode_domain_gated_packet(
1021            RawPacket {
1022                id: play::S_CLIENT_COMMAND,
1023                payload: vec![ClientCommandAction::PerformRespawn as u8],
1024            },
1025            &player,
1026        );
1027        assert!(matches!(perform_respawn, Ok(None)));
1028        assert!(player.has_deferred_death_respawn_for_test());
1029
1030        assert!(player.finish_domain_switch(token));
1031        assert!(player.finish_pending_world_change(token));
1032    }
1033
1034    #[test]
1035    fn custom_payload_defaults_to_global_exclusive_scheduling() {
1036        let channel = b"minecraft:brand";
1037        let mut payload = vec![channel.len() as u8];
1038        payload.extend_from_slice(channel);
1039        payload.extend_from_slice(b"steel");
1040        let decoded = decode(RawPacket {
1041            id: play::S_CUSTOM_PAYLOAD,
1042            payload,
1043        });
1044        let DecodedPlayPacket::Scheduled(
1045            packet @ ScheduledPlayPacket(ScheduledPlayPacketKind::CustomPayload(_)),
1046        ) = decoded
1047        else {
1048            panic!("custom payload should use the scheduled packet path");
1049        };
1050
1051        assert_eq!(packet.execution(), ScheduledPacketExecution::Exclusive);
1052    }
1053
1054    #[test]
1055    fn pre_join_allows_initial_play_acknowledgements() {
1056        assert!(JavaConnection::can_process_before_join(
1057            play::S_ACCEPT_TELEPORTATION
1058        ));
1059        assert!(JavaConnection::can_process_before_join(
1060            play::S_CHUNK_BATCH_RECEIVED
1061        ));
1062        assert!(JavaConnection::can_process_before_join(
1063            play::S_PLAYER_LOADED
1064        ));
1065        assert!(JavaConnection::can_process_during_domain_handshake(
1066            play::S_ACCEPT_TELEPORTATION
1067        ));
1068        assert!(JavaConnection::can_process_during_domain_handshake(
1069            play::S_CHUNK_BATCH_RECEIVED
1070        ));
1071        assert!(JavaConnection::can_process_during_domain_handshake(
1072            play::S_PLAYER_LOADED
1073        ));
1074        assert!(!JavaConnection::can_process_during_domain_handshake(
1075            play::S_MOVE_PLAYER_POS
1076        ));
1077    }
1078
1079    #[test]
1080    fn scheduled_domain_handshake_classification_is_narrow() {
1081        let accept = decode(RawPacket {
1082            id: play::S_ACCEPT_TELEPORTATION,
1083            payload: vec![0],
1084        });
1085        let DecodedPlayPacket::Scheduled(accept) = accept else {
1086            panic!("teleport acknowledgement should be scheduled");
1087        };
1088        assert!(accept.is_domain_handshake_packet());
1089
1090        let client_tick_end = decode(RawPacket {
1091            id: play::S_CLIENT_TICK_END,
1092            payload: Vec::new(),
1093        });
1094        let DecodedPlayPacket::Scheduled(client_tick_end) = client_tick_end else {
1095            panic!("client tick end should be scheduled");
1096        };
1097        assert!(!client_tick_end.is_domain_handshake_packet());
1098    }
1099
1100    #[test]
1101    fn client_tick_end_is_scheduled_for_the_inter_tick_phase() {
1102        let decoded = decode(RawPacket {
1103            id: play::S_CLIENT_TICK_END,
1104            payload: Vec::new(),
1105        });
1106
1107        assert!(matches!(
1108            decoded,
1109            DecodedPlayPacket::Scheduled(ScheduledPlayPacket(
1110                ScheduledPlayPacketKind::ClientTickEnd
1111            ))
1112        ));
1113    }
1114
1115    #[test]
1116    fn packet_execution_classification_separates_local_and_serialized_work() {
1117        assert_eq!(
1118            execution(ScheduledPlayPacketKind::PlayerAbilities(SPlayerAbilities {
1119                flags: 0
1120            },)),
1121            ScheduledPacketExecution::PlayerLocal
1122        );
1123        assert_eq!(
1124            execution(ScheduledPlayPacketKind::MovePlayer(
1125                SMovePlayerStatusOnly { packed_byte: 0 }.into(),
1126            )),
1127            ScheduledPacketExecution::Serialized
1128        );
1129    }
1130
1131    #[test]
1132    fn inventory_execution_reflects_complete_transaction_boundaries() {
1133        let click = SContainerClick {
1134            container_id: 0,
1135            state_id: 0,
1136            slot_num: 0,
1137            button_num: 0,
1138            click_type: ClickType::Pickup,
1139            changed_slots: FxHashMap::default(),
1140            carried_item: HashedStack::Empty,
1141        };
1142
1143        assert_eq!(
1144            execution(ScheduledPlayPacketKind::ContainerClick(click)),
1145            ScheduledPacketExecution::Serialized
1146        );
1147        assert_eq!(
1148            execution(ScheduledPlayPacketKind::ContainerClose(SContainerClose {
1149                container_id: 0,
1150            })),
1151            ScheduledPacketExecution::PlayerLocal
1152        );
1153        assert_eq!(
1154            execution(ScheduledPlayPacketKind::SetCreativeModeSlot(
1155                SSetCreativeModeSlot {
1156                    slot_num: 1,
1157                    item_stack: ItemStack::empty(),
1158                },
1159            )),
1160            ScheduledPacketExecution::PlayerLocal
1161        );
1162    }
1163
1164    #[test]
1165    fn player_command_execution_is_action_sensitive() {
1166        let command = |action| {
1167            execution(ScheduledPlayPacketKind::PlayerCommand(SPlayerCommand {
1168                entity_id: 1,
1169                action,
1170                data: 0,
1171            }))
1172        };
1173
1174        assert_eq!(
1175            command(PlayerCommandAction::StartSprinting),
1176            ScheduledPacketExecution::PlayerLocal
1177        );
1178        assert_eq!(
1179            command(PlayerCommandAction::StartFallFlying),
1180            ScheduledPacketExecution::PlayerLocal
1181        );
1182        assert_eq!(
1183            command(PlayerCommandAction::LeaveBed),
1184            ScheduledPacketExecution::Serialized
1185        );
1186        assert_eq!(
1187            command(PlayerCommandAction::OpenVehicleInventory),
1188            ScheduledPacketExecution::Exclusive
1189        );
1190    }
1191
1192    #[test]
1193    fn player_action_execution_is_action_sensitive() {
1194        let action = |action| {
1195            execution(ScheduledPlayPacketKind::PlayerAction(SPlayerAction {
1196                action,
1197                pos: BlockPos::new(0, 64, 0),
1198                direction: Direction::Down,
1199                sequence: 0,
1200            }))
1201        };
1202
1203        assert_eq!(
1204            action(PlayerAction::AbortDestroyBlock),
1205            ScheduledPacketExecution::PlayerLocal
1206        );
1207        assert_eq!(
1208            action(PlayerAction::SwapItemWithOffhand),
1209            ScheduledPacketExecution::PlayerLocal
1210        );
1211        assert_eq!(
1212            action(PlayerAction::StartDestroyBlock),
1213            ScheduledPacketExecution::Serialized
1214        );
1215        assert_eq!(
1216            action(PlayerAction::Stab),
1217            ScheduledPacketExecution::Exclusive
1218        );
1219    }
1220
1221    #[test]
1222    fn chat_message_and_ack_share_the_serialized_commit_lane() {
1223        assert_eq!(
1224            execution(ScheduledPlayPacketKind::Chat(Box::new(SChat {
1225                message: "hello".to_owned(),
1226                timestamp: 0,
1227                salt: 0,
1228                signature: None,
1229                offset: 0,
1230                acknowledged: [0; 3],
1231                checksum: 0,
1232            }))),
1233            ScheduledPacketExecution::Serialized
1234        );
1235        assert_eq!(
1236            execution(ScheduledPlayPacketKind::ChatAck(SChatAck {
1237                offset: VarInt(0),
1238            })),
1239            ScheduledPacketExecution::Serialized
1240        );
1241    }
1242
1243    #[test]
1244    fn cross_player_and_unimplemented_handlers_remain_global_barriers() {
1245        assert_eq!(
1246            execution(ScheduledPlayPacketKind::Attack(SAttack { entity_id: 1 })),
1247            ScheduledPacketExecution::Exclusive
1248        );
1249        assert_eq!(
1250            execution(ScheduledPlayPacketKind::ContainerButtonClick(
1251                SContainerButtonClick {
1252                    container_id: 1,
1253                    button_id: 0,
1254                },
1255            )),
1256            ScheduledPacketExecution::Exclusive
1257        );
1258        assert_eq!(
1259            execution(ScheduledPlayPacketKind::PlayerAction(SPlayerAction {
1260                action: PlayerAction::ReleaseUseItem,
1261                pos: BlockPos::new(0, 64, 0),
1262                direction: Direction::Down,
1263                sequence: 0,
1264            })),
1265            ScheduledPacketExecution::Exclusive
1266        );
1267    }
1268
1269    #[test]
1270    fn audited_handlers_use_the_narrowest_safe_execution_class() {
1271        assert_eq!(
1272            execution(ScheduledPlayPacketKind::AcceptTeleportation(
1273                SAcceptTeleportation { teleport_id: 1 },
1274            )),
1275            ScheduledPacketExecution::Serialized
1276        );
1277        assert_eq!(
1278            execution(ScheduledPlayPacketKind::PlayerInput(SPlayerInput {
1279                flags: 0,
1280            })),
1281            ScheduledPacketExecution::PlayerLocal
1282        );
1283        assert_eq!(
1284            execution(ScheduledPlayPacketKind::ChatSessionUpdate(
1285                SChatSessionUpdate {
1286                    session_id: Uuid::nil(),
1287                    expires_at: 0,
1288                    public_key: Vec::new(),
1289                    key_signature: Vec::new(),
1290                },
1291            )),
1292            ScheduledPacketExecution::PlayerLocal
1293        );
1294        assert_eq!(
1295            execution(ScheduledPlayPacketKind::ClientInformation(
1296                SClientInformation {
1297                    language: "en_us".to_owned(),
1298                    view_distance: 8,
1299                    chat_visibility: ChatVisibility::Full,
1300                    chat_colors: true,
1301                    model_customization: 0,
1302                    main_hand: HumanoidArm::Right,
1303                    text_filtering_enabled: false,
1304                    allows_listing: true,
1305                    particle_status: ParticleStatus::All,
1306                },
1307            )),
1308            ScheduledPacketExecution::PlayerLocal
1309        );
1310        assert_eq!(
1311            execution(ScheduledPlayPacketKind::ChatCommand(SChatCommand {
1312                command: "help".to_owned(),
1313            })),
1314            ScheduledPacketExecution::PlayerLocal
1315        );
1316        assert_eq!(
1317            execution(ScheduledPlayPacketKind::PickItemFromBlock(
1318                SPickItemFromBlock {
1319                    pos: BlockPos::new(0, 64, 0),
1320                    include_data: false,
1321                },
1322            )),
1323            ScheduledPacketExecution::PlayerLocal
1324        );
1325        assert_eq!(
1326            execution(ScheduledPlayPacketKind::SignUpdate(SSignUpdate {
1327                pos: BlockPos::new(0, 64, 0),
1328                is_front_text: true,
1329                lines: array::from_fn(|_| String::new()),
1330            })),
1331            ScheduledPacketExecution::Serialized
1332        );
1333        assert_eq!(
1334            execution(ScheduledPlayPacketKind::Swing(SSwing {
1335                hand: InteractionHand::MainHand,
1336            })),
1337            ScheduledPacketExecution::PlayerLocal
1338        );
1339        assert_eq!(
1340            execution(ScheduledPlayPacketKind::ClientCommand(SClientCommand {
1341                action: ClientCommandAction::PerformRespawn,
1342            })),
1343            ScheduledPacketExecution::PlayerLocal
1344        );
1345    }
1346
1347    #[test]
1348    fn keep_alive_remains_on_the_immediate_connection_path() {
1349        let decoded = decode(RawPacket {
1350            id: play::S_KEEP_ALIVE,
1351            payload: 42_i64.to_be_bytes().to_vec(),
1352        });
1353
1354        assert!(matches!(
1355            decoded,
1356            DecodedPlayPacket::Immediate(ImmediatePlayPacket::KeepAlive(SKeepAlive { id: 42 }))
1357        ));
1358    }
1359
1360    #[test]
1361    fn chunk_batch_ack_uses_the_immediate_connection_path() {
1362        let decoded = decode(RawPacket {
1363            id: play::S_CHUNK_BATCH_RECEIVED,
1364            payload: 12.5_f32.to_be_bytes().to_vec(),
1365        });
1366
1367        assert!(matches!(
1368            decoded,
1369            DecodedPlayPacket::Immediate(ImmediatePlayPacket::ChunkBatchReceived(
1370                SChunkBatchReceived {
1371                    desired_chunks_per_tick: 12.5
1372                }
1373            ))
1374        ));
1375    }
1376}