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