1use 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
44pub type JavaTransportRead = Box<dyn AsyncRead + Send + Unpin>;
46pub type JavaTransportWrite = Box<dyn AsyncWrite + Send + Unpin>;
48pub type JavaNetworkReader = TCPNetworkDecoder<BufReader<JavaTransportRead>>;
50pub type JavaNetworkWriter =
52 Arc<AsyncMutex<Option<TCPNetworkEncoder<BufWriter<JavaTransportWrite>>>>>;
53
54const DISCONNECT_FLUSH_TIMEOUT: Duration = Duration::from_secs(1);
55
56pub enum OutboundPacket {
58 Packet(EncodedPacket),
60 Disconnect(EncodedPacket),
62}
63
64pub(crate) struct ScheduledPlayPacket(ScheduledPlayPacketKind);
66
67#[derive(Clone, Copy, Debug, Eq, PartialEq)]
69pub(crate) enum ScheduledPacketExecution {
70 PlayerLocal,
75 Serialized,
78 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 pub(crate) const fn is_domain_handshake_packet(&self) -> bool {
135 matches!(
136 self.0,
137 ScheduledPlayPacketKind::AcceptTeleportation(_) | ScheduledPlayPacketKind::PlayerLoaded
138 )
139 }
140
141 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 pub(crate) const fn execution(&self) -> ScheduledPacketExecution {
163 match &self.0 {
164 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 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 PlayerAction::ReleaseUseItem | PlayerAction::Stab => {
202 ScheduledPacketExecution::Exclusive
203 }
204 },
205 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 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
357pub struct BundleBuilder {
361 packets: Vec<EncodedPacket>,
362 compression: Option<CompressionInfo>,
363}
364
365impl BundleBuilder {
366 #[must_use]
368 pub const fn new(compression: Option<CompressionInfo>) -> Self {
369 Self {
370 packets: Vec::new(),
371 compression,
372 }
373 }
374
375 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 #[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
402pub 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 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 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 #[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 #[must_use]
523 pub fn latency(&self) -> i32 {
524 *self.latency.lock() as i32
525 }
526
527 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 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 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 pub fn close(&self) {
588 self.cancel_token.cancel();
589 }
590
591 #[must_use]
593 pub fn closed(&self) -> bool {
594 self.cancel_token.is_cancelled()
595 }
596
597 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 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 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 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 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}