1use std::{
4 cmp::Reverse,
5 collections::{BinaryHeap, VecDeque},
6 hash::Hash,
7 sync::Arc,
8};
9
10use parking_lot::Condvar;
11use rustc_hash::{FxHashMap, FxHashSet};
12use steel_utils::{locks::SyncMutex, translations};
13use tokio::sync::Notify;
14use tokio::task::yield_now;
15use uuid::Uuid;
16
17use crate::{
18 entity::Entity,
19 player::{
20 Player,
21 connection::NetworkConnection,
22 connection::{ScheduledPacketExecution, ScheduledPlayPacket},
23 },
24};
25
26use super::Server;
27
28const MAX_OUTSTANDING_PACKETS_PER_PLAYER: usize = 8_192;
32const MAX_OUTSTANDING_BYTES_PER_PLAYER: usize = 32 * 1024 * 1024;
33const PACKET_ADMISSION_OVERHEAD: usize = 256;
35
36#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
37struct PlayerPacketLaneKey {
38 player_id: Uuid,
39 entity_id: i32,
40}
41
42impl PlayerPacketLaneKey {
43 fn new(player: &Player) -> Self {
44 Self {
45 player_id: player.gameprofile.id,
46 entity_id: player.id(),
47 }
48 }
49}
50
51struct PendingPlayPacket {
52 player: Arc<Player>,
53 packet: ScheduledPlayPacket,
54}
55
56pub(super) struct PacketProcessor {
61 queued: PacketQueue<PlayerPacketLaneKey, PendingPlayPacket>,
62}
63
64impl PacketProcessor {
65 pub(super) fn new() -> Self {
66 Self {
67 queued: PacketQueue::new(),
68 }
69 }
70
71 pub(super) fn schedule(
72 &self,
73 player: Arc<Player>,
74 packet: ScheduledPlayPacket,
75 payload_bytes: usize,
76 ) {
77 let player_id = player.gameprofile.id;
78 let lane_key = PlayerPacketLaneKey::new(&player);
79 if player.connection.closed() {
80 self.queued.discard_lane(lane_key);
81 return;
82 }
83
84 let execution = packet.execution();
85 let admission_bytes = payload_bytes.saturating_add(PACKET_ADMISSION_OVERHEAD);
86 let result = self.queued.try_submit(
87 lane_key,
88 execution,
89 admission_bytes,
90 PendingPlayPacket {
91 player: Arc::clone(&player),
92 packet,
93 },
94 );
95 let Err(error) = result else {
96 return;
97 };
98 if error == PacketAdmissionError::Stopped {
99 return;
100 }
101
102 tracing::warn!(
103 player_id = %player_id,
104 entity_id = lane_key.entity_id,
105 ?error,
106 "Disconnecting player after inbound packet admission limit"
107 );
108 player.disconnect(translations::DISCONNECT_EXCEEDED_PACKET_RATE.msg());
109 self.queued.discard_lane(lane_key);
110 }
111
112 pub(super) fn run(&self, server: &Arc<Server>) {
114 while let Some(mut work) = self.queued.next() {
115 let Some(pending) = work.take() else {
116 continue;
117 };
118 if !Self::packet_is_runnable(
119 &pending.player,
120 &pending.packet,
121 server.cancel_token.is_cancelled(),
122 ) {
123 continue;
124 }
125
126 pending.packet.handle(pending.player, server);
127 }
128 }
129
130 fn packet_is_runnable(
131 player: &Player,
132 packet: &ScheduledPlayPacket,
133 server_cancelled: bool,
134 ) -> bool {
135 !player.connection.closed()
136 && !server_cancelled
137 && player.gate_domain_switch_packet(
138 packet.is_domain_handshake_packet(),
139 packet.is_perform_respawn(),
140 )
141 }
142
143 pub(super) fn open_after_tick(&self) {
145 self.queued.open();
146 }
147
148 pub(super) async fn wait_for_overload_progress(&self) {
150 let Some(completed) = self.queued.progress_baseline() else {
151 return;
152 };
153 yield_now().await;
154 self.queued.wait_for_progress_since(completed).await;
155 }
156
157 pub(super) async fn close_for_tick(&self) {
159 self.queued.drain_for_tick().await;
160 }
161
162 pub(super) fn stop(&self) {
164 self.queued.stop();
165 }
166}
167
168impl Default for PacketProcessor {
169 fn default() -> Self {
170 Self::new()
171 }
172}
173
174#[derive(Clone, Copy, Eq, PartialEq)]
175enum PacketPhase {
176 Closed,
177 Open,
178 Draining(u64),
179 Stopped,
180}
181
182struct SequencedPacket<T> {
183 sequence: u64,
184 execution: ScheduledPacketExecution,
185 admission_bytes: usize,
186 value: T,
187}
188
189struct PacketLane<T> {
190 queued: VecDeque<SequencedPacket<T>>,
191 active: bool,
192 outstanding_packets: usize,
193 outstanding_bytes: usize,
194}
195
196impl<T> PacketLane<T> {
197 const fn new() -> Self {
198 Self {
199 queued: VecDeque::new(),
200 active: false,
201 outstanding_packets: 0,
202 outstanding_bytes: 0,
203 }
204 }
205}
206
207#[derive(Clone, Copy)]
208struct PacketAdmissionLimits {
209 per_player_packets: usize,
210 per_player_bytes: usize,
211}
212
213impl PacketAdmissionLimits {
214 const PRODUCTION: Self = Self {
215 per_player_packets: MAX_OUTSTANDING_PACKETS_PER_PLAYER,
216 per_player_bytes: MAX_OUTSTANDING_BYTES_PER_PLAYER,
217 };
218}
219
220#[derive(Clone, Copy, Debug, Eq, PartialEq)]
221enum PacketAdmissionError {
222 Stopped,
223 PlayerPacketLimit,
224 PlayerByteLimit,
225}
226
227struct PacketQueueState<K, T> {
228 phase: PacketPhase,
229 lanes: FxHashMap<K, PacketLane<T>>,
230 player_local_ready: BinaryHeap<Reverse<(u64, K)>>,
231 serialized_ready: BinaryHeap<Reverse<(u64, K)>>,
232 exclusive_ready: BinaryHeap<Reverse<(u64, K)>>,
233 serialized: BinaryHeap<Reverse<u64>>,
234 exclusive: BinaryHeap<Reverse<u64>>,
235 next_sequence: u64,
236 active: usize,
237 serialized_active: bool,
238 exclusive_active: bool,
239 completed: u64,
240}
241
242struct PacketQueue<K, T> {
244 state: SyncMutex<PacketQueueState<K, T>>,
245 limits: PacketAdmissionLimits,
246 work_available: Condvar,
247 idle: Notify,
248 progress: Notify,
249}
250
251impl<K, T> PacketQueue<K, T>
252where
253 K: Copy + Eq + Hash + Ord,
254{
255 fn new() -> Self {
256 Self::with_limits(PacketAdmissionLimits::PRODUCTION)
257 }
258
259 fn with_limits(limits: PacketAdmissionLimits) -> Self {
260 Self {
261 state: SyncMutex::new(PacketQueueState {
262 phase: PacketPhase::Closed,
263 lanes: FxHashMap::default(),
264 player_local_ready: BinaryHeap::new(),
265 serialized_ready: BinaryHeap::new(),
266 exclusive_ready: BinaryHeap::new(),
267 serialized: BinaryHeap::new(),
268 exclusive: BinaryHeap::new(),
269 next_sequence: 0,
270 active: 0,
271 serialized_active: false,
272 exclusive_active: false,
273 completed: 0,
274 }),
275 limits,
276 work_available: Condvar::new(),
277 idle: Notify::new(),
278 progress: Notify::new(),
279 }
280 }
281
282 fn try_submit(
283 &self,
284 key: K,
285 execution: ScheduledPacketExecution,
286 admission_bytes: usize,
287 value: T,
288 ) -> Result<(), PacketAdmissionError> {
289 let mut state = self.state.lock();
290 if state.phase == PacketPhase::Stopped {
291 return Err(PacketAdmissionError::Stopped);
292 }
293
294 let (player_packets, player_bytes) = state
295 .lanes
296 .get(&key)
297 .map(|lane| (lane.outstanding_packets, lane.outstanding_bytes))
298 .unwrap_or_default();
299 if player_packets >= self.limits.per_player_packets {
300 return Err(PacketAdmissionError::PlayerPacketLimit);
301 }
302 if admission_bytes > self.limits.per_player_bytes.saturating_sub(player_bytes) {
303 return Err(PacketAdmissionError::PlayerByteLimit);
304 }
305 let sequence = state.next_sequence;
306 assert!(sequence != u64::MAX, "packet submission sequence exhausted");
307 state.next_sequence = sequence + 1;
308
309 let lane = state.lanes.entry(key).or_insert_with(PacketLane::new);
310 let became_ready = !lane.active && lane.queued.is_empty();
311 lane.outstanding_packets += 1;
312 lane.outstanding_bytes += admission_bytes;
313 lane.queued.push_back(SequencedPacket {
314 sequence,
315 execution,
316 admission_bytes,
317 value,
318 });
319 if became_ready {
320 Self::mark_ready(&mut state, sequence, key, execution);
321 }
322 match execution {
323 ScheduledPacketExecution::PlayerLocal => {}
324 ScheduledPacketExecution::Serialized => state.serialized.push(Reverse(sequence)),
325 ScheduledPacketExecution::Exclusive => state.exclusive.push(Reverse(sequence)),
326 }
327 let should_wake = became_ready && state.phase == PacketPhase::Open;
328 drop(state);
329 if should_wake {
330 self.work_available.notify_one();
331 }
332 Ok(())
333 }
334
335 #[cfg(test)]
336 fn submit(&self, key: K, execution: ScheduledPacketExecution, value: T) {
337 let result = self.try_submit(key, execution, 1, value);
338 assert!(
339 matches!(result, Ok(()) | Err(PacketAdmissionError::Stopped)),
340 "test packet exceeded admission limits: {result:?}"
341 );
342 }
343
344 fn open(&self) {
345 let mut state = self.state.lock();
346 if state.phase == PacketPhase::Stopped {
347 return;
348 }
349 state.phase = PacketPhase::Open;
350 let should_wake = Self::select_next(&state, None).is_some();
351 drop(state);
352 if should_wake {
353 self.work_available.notify_all();
354 }
355 }
356
357 #[cfg(test)]
358 fn close(&self) {
359 let mut state = self.state.lock();
360 if state.phase != PacketPhase::Stopped {
361 state.phase = PacketPhase::Closed;
362 }
363 }
364
365 async fn drain_for_tick(&self) {
366 let Some((before_sequence, should_wake)) = self.begin_tick_drain() else {
367 return;
368 };
369 if should_wake {
370 self.work_available.notify_all();
371 }
372
373 loop {
374 let idle = self.idle.notified();
375 if self.tick_drain_complete(before_sequence) {
376 break;
377 }
378 idle.await;
379 }
380
381 self.finish_tick_drain(before_sequence);
382 }
383
384 fn finish_tick_drain(&self, before_sequence: u64) {
385 let mut state = self.state.lock();
386 if state.phase == PacketPhase::Draining(before_sequence) {
387 state.phase = PacketPhase::Closed;
388 }
389 }
390
391 fn begin_tick_drain(&self) -> Option<(u64, bool)> {
392 let mut state = self.state.lock();
393 let before_sequence = match state.phase {
394 PacketPhase::Stopped => None,
395 PacketPhase::Draining(before_sequence) => Some(before_sequence),
396 PacketPhase::Closed | PacketPhase::Open => {
397 let before_sequence = state.next_sequence;
398 state.phase = PacketPhase::Draining(before_sequence);
399 Some(before_sequence)
400 }
401 }?;
402 let should_wake = Self::select_next(&state, Some(before_sequence)).is_some();
403 Some((before_sequence, should_wake))
404 }
405
406 fn tick_drain_complete(&self, before_sequence: u64) -> bool {
407 let state = self.state.lock();
408 if state.phase == PacketPhase::Stopped {
409 return true;
410 }
411 if state.active != 0 {
412 return false;
413 }
414 Self::next_ready_sequence(&state).is_none_or(|sequence| sequence >= before_sequence)
415 }
416
417 async fn wait_for_progress_since(&self, completed: u64) {
418 loop {
419 let progress = self.progress.notified();
420 if self.has_progress_since(completed) {
421 return;
422 }
423 progress.await;
424 }
425 }
426
427 fn progress_baseline(&self) -> Option<u64> {
428 let state = self.state.lock();
429 Self::has_work(&state).then_some(state.completed)
430 }
431
432 fn has_progress_since(&self, completed: u64) -> bool {
433 let state = self.state.lock();
434 state.completed != completed
435 || !Self::has_work(&state)
436 || state.phase == PacketPhase::Stopped
437 }
438
439 fn has_work(state: &PacketQueueState<K, T>) -> bool {
440 state.active != 0 || state.lanes.values().any(|lane| !lane.queued.is_empty())
441 }
442
443 fn discard_lane(&self, key: K) {
444 let mut state = self.state.lock();
445 let Some(lane) = state.lanes.get_mut(&key) else {
446 return;
447 };
448 let discarded_packets = lane.queued.len();
449 if discarded_packets == 0 {
450 return;
451 }
452 let discarded_bytes = lane
453 .queued
454 .iter()
455 .map(|packet| packet.admission_bytes)
456 .sum::<usize>();
457 let discarded_sequences = lane
458 .queued
459 .iter()
460 .map(|packet| packet.sequence)
461 .collect::<FxHashSet<_>>();
462 lane.queued.clear();
463 assert!(
464 lane.outstanding_packets >= discarded_packets,
465 "session packet admission accounting underflow while discarding"
466 );
467 lane.outstanding_packets -= discarded_packets;
468 assert!(
469 lane.outstanding_bytes >= discarded_bytes,
470 "session byte admission accounting underflow while discarding"
471 );
472 lane.outstanding_bytes -= discarded_bytes;
473 let remove_lane = !lane.active;
474
475 state.player_local_ready.retain(|entry| entry.0.1 != key);
476 state.serialized_ready.retain(|entry| entry.0.1 != key);
477 state.exclusive_ready.retain(|entry| entry.0.1 != key);
478 state
479 .serialized
480 .retain(|entry| !discarded_sequences.contains(&entry.0));
481 state
482 .exclusive
483 .retain(|entry| !discarded_sequences.contains(&entry.0));
484 if remove_lane {
485 state.lanes.remove(&key);
486 }
487
488 let is_idle = state.active == 0;
489 let should_wake = matches!(state.phase, PacketPhase::Open | PacketPhase::Draining(_))
490 && Self::next_ready_sequence(&state).is_some();
491 drop(state);
492 self.progress.notify_one();
493 if should_wake {
494 self.work_available.notify_all();
495 }
496 if is_idle {
497 self.idle.notify_one();
498 }
499 }
500
501 fn stop(&self) {
502 let mut state = self.state.lock();
503 state.phase = PacketPhase::Stopped;
504 state.player_local_ready.clear();
505 state.serialized_ready.clear();
506 state.exclusive_ready.clear();
507 state.serialized.clear();
508 state.exclusive.clear();
509 state.lanes.retain(|_, lane| {
510 let lane_discarded_packets = lane.queued.len();
511 let lane_discarded_bytes = lane
512 .queued
513 .iter()
514 .map(|packet| packet.admission_bytes)
515 .sum::<usize>();
516 lane.queued.clear();
517 assert!(
518 lane.outstanding_packets >= lane_discarded_packets,
519 "session packet admission accounting underflow while stopping"
520 );
521 lane.outstanding_packets -= lane_discarded_packets;
522 assert!(
523 lane.outstanding_bytes >= lane_discarded_bytes,
524 "session byte admission accounting underflow while stopping"
525 );
526 lane.outstanding_bytes -= lane_discarded_bytes;
527 lane.active
528 });
529 let is_idle = state.active == 0;
530 drop(state);
531 self.work_available.notify_all();
532 self.progress.notify_waiters();
533 if is_idle {
534 self.idle.notify_one();
535 }
536 }
537
538 fn next(&self) -> Option<PacketWork<'_, K, T>> {
539 let mut state = self.state.lock();
540 loop {
541 match state.phase {
542 PacketPhase::Stopped => return None,
543 PacketPhase::Open => {
544 if let Some((key, execution, admission_bytes, value)) =
545 Self::start_next(&mut state, None)
546 {
547 state.active += 1;
548 drop(state);
549 return Some(PacketWork {
550 value: Some(value),
551 key,
552 execution,
553 admission_bytes,
554 queue: self,
555 });
556 }
557 }
558 PacketPhase::Draining(before_sequence) => {
559 if let Some((key, execution, admission_bytes, value)) =
560 Self::start_next(&mut state, Some(before_sequence))
561 {
562 state.active += 1;
563 drop(state);
564 return Some(PacketWork {
565 value: Some(value),
566 key,
567 execution,
568 admission_bytes,
569 queue: self,
570 });
571 }
572 }
573 PacketPhase::Closed => {}
574 }
575 self.work_available.wait(&mut state);
576 }
577 }
578
579 #[cfg(test)]
580 fn try_next(&self) -> Option<PacketWork<'_, K, T>> {
581 let mut state = self.state.lock();
582 let before_sequence = match state.phase {
583 PacketPhase::Open => None,
584 PacketPhase::Draining(before_sequence) => Some(before_sequence),
585 PacketPhase::Closed | PacketPhase::Stopped => return None,
586 };
587 let (key, execution, admission_bytes, value) =
588 Self::start_next(&mut state, before_sequence)?;
589 state.active += 1;
590 drop(state);
591 Some(PacketWork {
592 value: Some(value),
593 key,
594 execution,
595 admission_bytes,
596 queue: self,
597 })
598 }
599
600 fn start_next(
601 state: &mut PacketQueueState<K, T>,
602 before_sequence: Option<u64>,
603 ) -> Option<(K, ScheduledPacketExecution, usize, T)> {
604 let (ready_sequence, key, execution) = Self::select_next(state, before_sequence)?;
605 let Some(lane) = state.lanes.get(&key) else {
606 panic!("ready packet lane disappeared before starting");
607 };
608 assert!(!lane.active, "ready packet lane is already active");
609 let Some(packet) = lane.queued.front() else {
610 panic!("ready packet lane has no queued packet");
611 };
612 assert_eq!(
613 ready_sequence, packet.sequence,
614 "ready packet sequence does not match lane front"
615 );
616 assert_eq!(
617 execution, packet.execution,
618 "packet is registered in the wrong ready queue"
619 );
620
621 match execution {
622 ScheduledPacketExecution::PlayerLocal => {
623 assert!(
624 state.player_local_ready.pop() == Some(Reverse((ready_sequence, key))),
625 "player-local ready queue changed while the queue lock was held"
626 );
627 }
628 ScheduledPacketExecution::Serialized => {
629 assert!(
630 state.serialized_ready.pop() == Some(Reverse((ready_sequence, key))),
631 "serialized ready queue changed while the queue lock was held"
632 );
633 assert_eq!(
634 state.serialized.pop(),
635 Some(Reverse(ready_sequence)),
636 "serialized packet order changed while the queue lock was held"
637 );
638 assert!(
639 !state.serialized_active,
640 "serialized packet started while another was active"
641 );
642 state.serialized_active = true;
643 }
644 ScheduledPacketExecution::Exclusive => {
645 assert!(
646 state.exclusive_ready.pop() == Some(Reverse((ready_sequence, key))),
647 "exclusive ready queue changed while the queue lock was held"
648 );
649 assert_eq!(
650 state.exclusive.pop(),
651 Some(Reverse(ready_sequence)),
652 "exclusive packet barrier changed while the queue lock was held"
653 );
654 state.exclusive_active = true;
655 }
656 }
657
658 let Some(lane) = state.lanes.get_mut(&key) else {
659 panic!("ready packet lane disappeared before removal");
660 };
661 let Some(packet) = lane.queued.pop_front() else {
662 panic!("ready packet lane has no queued packet during removal");
663 };
664 lane.active = true;
665 Some((key, execution, packet.admission_bytes, packet.value))
666 }
667
668 fn select_next(
669 state: &PacketQueueState<K, T>,
670 before_sequence: Option<u64>,
671 ) -> Option<(u64, K, ScheduledPacketExecution)> {
672 if state.exclusive_active {
673 return None;
674 }
675
676 let next_exclusive = state.exclusive.peek().map(|entry| entry.0);
677 let next_serialized = state.serialized.peek().map(|entry| entry.0);
678 let mut selected = None;
679 let mut consider = |entry: Option<&Reverse<(u64, K)>>,
680 execution: ScheduledPacketExecution| {
681 let Some(Reverse((sequence, key))) = entry.copied() else {
682 return;
683 };
684 if before_sequence.is_some_and(|cutoff| sequence >= cutoff)
685 || next_exclusive.is_some_and(|exclusive| exclusive < sequence)
686 {
687 return;
688 }
689 if selected.is_none_or(|(selected_sequence, _, _)| sequence < selected_sequence) {
690 selected = Some((sequence, key, execution));
691 }
692 };
693
694 consider(
695 state.player_local_ready.peek(),
696 ScheduledPacketExecution::PlayerLocal,
697 );
698 if !state.serialized_active
699 && state
700 .serialized_ready
701 .peek()
702 .is_some_and(|entry| Some(entry.0.0) == next_serialized)
703 {
704 consider(
705 state.serialized_ready.peek(),
706 ScheduledPacketExecution::Serialized,
707 );
708 }
709 if state.active == 0
710 && state
711 .exclusive_ready
712 .peek()
713 .is_some_and(|entry| Some(entry.0.0) == next_exclusive)
714 {
715 consider(
716 state.exclusive_ready.peek(),
717 ScheduledPacketExecution::Exclusive,
718 );
719 }
720 selected
721 }
722
723 fn finish_one(&self, key: K, execution: ScheduledPacketExecution, admission_bytes: usize) {
724 let mut state = self.state.lock();
725 assert!(state.active > 0, "packet work accounting underflow");
726 state.active -= 1;
727 state.completed = state.completed.wrapping_add(1);
728 match execution {
729 ScheduledPacketExecution::PlayerLocal => {}
730 ScheduledPacketExecution::Serialized => {
731 assert!(state.serialized_active, "serialized packet is not active");
732 state.serialized_active = false;
733 }
734 ScheduledPacketExecution::Exclusive => {
735 assert!(state.exclusive_active, "exclusive packet is not active");
736 state.exclusive_active = false;
737 }
738 }
739
740 let next_sequence = {
741 let Some(lane) = state.lanes.get_mut(&key) else {
742 panic!("active packet lane disappeared before completion");
743 };
744 assert!(lane.active, "completed packet lane is not active");
745 lane.active = false;
746 assert!(
747 lane.outstanding_packets > 0,
748 "session packet admission accounting underflow on completion"
749 );
750 lane.outstanding_packets -= 1;
751 assert!(
752 lane.outstanding_bytes >= admission_bytes,
753 "session byte admission accounting underflow on completion"
754 );
755 lane.outstanding_bytes -= admission_bytes;
756 lane.queued.front().map(|packet| packet.sequence)
757 };
758 if let Some(sequence) = next_sequence {
759 let Some(lane) = state.lanes.get(&key) else {
760 panic!("packet lane disappeared before its next packet became ready");
761 };
762 let Some(packet) = lane.queued.front() else {
763 panic!("packet lane has no next packet after reporting its sequence");
764 };
765 let next_execution = packet.execution;
766 Self::mark_ready(&mut state, sequence, key, next_execution);
767 } else {
768 let Some(lane) = state.lanes.get(&key) else {
769 panic!("completed packet lane disappeared before removal");
770 };
771 assert_eq!(lane.outstanding_packets, 0);
772 assert_eq!(lane.outstanding_bytes, 0);
773 state.lanes.remove(&key);
774 }
775
776 let is_idle = state.active == 0;
777 let should_wake = matches!(state.phase, PacketPhase::Open | PacketPhase::Draining(_))
778 && Self::next_ready_sequence(&state).is_some();
779 drop(state);
780 self.progress.notify_one();
781 if should_wake {
782 match execution {
783 ScheduledPacketExecution::Exclusive => {
784 self.work_available.notify_all();
785 }
786 ScheduledPacketExecution::PlayerLocal | ScheduledPacketExecution::Serialized => {
787 self.work_available.notify_one();
788 }
789 }
790 }
791 if is_idle {
792 self.idle.notify_one();
793 }
794 }
795
796 fn mark_ready(
797 state: &mut PacketQueueState<K, T>,
798 sequence: u64,
799 key: K,
800 execution: ScheduledPacketExecution,
801 ) {
802 let entry = Reverse((sequence, key));
803 match execution {
804 ScheduledPacketExecution::PlayerLocal => state.player_local_ready.push(entry),
805 ScheduledPacketExecution::Serialized => state.serialized_ready.push(entry),
806 ScheduledPacketExecution::Exclusive => state.exclusive_ready.push(entry),
807 }
808 }
809
810 fn next_ready_sequence(state: &PacketQueueState<K, T>) -> Option<u64> {
811 [
812 state.player_local_ready.peek().map(|entry| entry.0.0),
813 state.serialized_ready.peek().map(|entry| entry.0.0),
814 state.exclusive_ready.peek().map(|entry| entry.0.0),
815 ]
816 .into_iter()
817 .flatten()
818 .min()
819 }
820}
821
822struct PacketWork<'a, K, T>
823where
824 K: Copy + Eq + Hash + Ord,
825{
826 value: Option<T>,
827 key: K,
828 execution: ScheduledPacketExecution,
829 admission_bytes: usize,
830 queue: &'a PacketQueue<K, T>,
831}
832
833impl<K, T> PacketWork<'_, K, T>
834where
835 K: Copy + Eq + Hash + Ord,
836{
837 const fn take(&mut self) -> Option<T> {
838 self.value.take()
839 }
840}
841
842impl<K, T> Drop for PacketWork<'_, K, T>
843where
844 K: Copy + Eq + Hash + Ord,
845{
846 fn drop(&mut self) {
847 self.queue
848 .finish_one(self.key, self.execution, self.admission_bytes);
849 }
850}
851
852#[cfg(test)]
853mod tests {
854 use std::{
855 sync::{Arc, mpsc},
856 thread,
857 time::Duration,
858 };
859
860 use tokio::time::timeout;
861 use uuid::Uuid;
862
863 use crate::{
864 entity::{Entity as _, LivingEntity as _},
865 player::connection::ScheduledPlayPacket,
866 test_support::{TestPlayerBuilder, fresh_test_world},
867 };
868
869 use super::{
870 PACKET_ADMISSION_OVERHEAD, PacketAdmissionError, PacketAdmissionLimits, PacketProcessor,
871 PacketQueue, PlayerPacketLaneKey, ScheduledPacketExecution,
872 };
873
874 const fn limits(per_player_packets: usize, per_player_bytes: usize) -> PacketAdmissionLimits {
875 PacketAdmissionLimits {
876 per_player_packets,
877 per_player_bytes,
878 }
879 }
880
881 #[test]
882 fn per_player_packet_limit_rejects_before_assigning_a_sequence() {
883 let queue = PacketQueue::with_limits(limits(2, 100));
884 assert_eq!(
885 queue.try_submit(1, ScheduledPacketExecution::PlayerLocal, 1, "first"),
886 Ok(())
887 );
888 assert_eq!(
889 queue.try_submit(1, ScheduledPacketExecution::Serialized, 1, "second"),
890 Ok(())
891 );
892
893 assert_eq!(
894 queue.try_submit(1, ScheduledPacketExecution::Exclusive, 1, "rejected"),
895 Err(PacketAdmissionError::PlayerPacketLimit)
896 );
897 assert_eq!(queue.state.lock().next_sequence, 2);
898 assert_eq!(
899 queue.try_submit(2, ScheduledPacketExecution::PlayerLocal, 1, "other player"),
900 Ok(())
901 );
902 }
903
904 #[test]
905 fn scheduled_respawn_is_retained_if_domain_switch_queues_before_worker_gate() {
906 let world = fresh_test_world("scheduled_domain_switch_respawn_packet");
907 let player = TestPlayerBuilder::new(world, Uuid::from_u128(1), "RespawnTester", 1).build();
908 let packet = ScheduledPlayPacket::perform_respawn_for_test();
909 let Some(token) = player.begin_pending_world_change() else {
910 panic!("test player should acquire a world-change token");
911 };
912 assert!(player.begin_domain_switch(token));
913 player.set_health(0.0);
914
915 assert!(!PacketProcessor::packet_is_runnable(
916 &player, &packet, false
917 ));
918 assert!(player.has_deferred_death_respawn_for_test());
919
920 assert!(player.finish_domain_switch(token));
921 assert!(player.finish_pending_world_change(token));
922 }
923
924 #[test]
925 fn per_player_byte_limit_is_independent_of_packet_count() {
926 let queue = PacketQueue::with_limits(limits(10, 6));
927 assert_eq!(
928 queue.try_submit(1, ScheduledPacketExecution::PlayerLocal, 4, "first"),
929 Ok(())
930 );
931
932 assert_eq!(
933 queue.try_submit(1, ScheduledPacketExecution::PlayerLocal, 3, "rejected"),
934 Err(PacketAdmissionError::PlayerByteLimit)
935 );
936 assert_eq!(
937 queue.try_submit(2, ScheduledPacketExecution::PlayerLocal, 6, "other player"),
938 Ok(())
939 );
940 }
941
942 #[test]
943 fn production_limit_allows_a_watchdog_window_of_mounted_client_traffic() {
944 const VANILLA_WATCHDOG_SECONDS: usize = 60;
945 const CLIENT_TICKS_PER_SECOND: usize = 20;
946 const MOUNTED_PACKETS_PER_CLIENT_TICK: usize = 4;
947 const EXPECTED_PACKETS: usize =
948 VANILLA_WATCHDOG_SECONDS * CLIENT_TICKS_PER_SECOND * MOUNTED_PACKETS_PER_CLIENT_TICK;
949
950 let key = PlayerPacketLaneKey {
951 player_id: Uuid::nil(),
952 entity_id: 1,
953 };
954 let queue = PacketQueue::new();
955 for packet in 0..EXPECTED_PACKETS {
956 assert!(
957 queue
958 .try_submit(
959 key,
960 ScheduledPacketExecution::PlayerLocal,
961 PACKET_ADMISSION_OVERHEAD,
962 packet,
963 )
964 .is_ok()
965 );
966 }
967
968 let state = queue.state.lock();
969 let lane = state.lanes.get(&key).expect("session lane should exist");
970 assert_eq!(lane.outstanding_packets, EXPECTED_PACKETS);
971 }
972
973 #[test]
974 fn independent_session_limits_do_not_accumulate_globally() {
975 const SESSION_COUNT: usize = 64;
976 const PACKETS_PER_SESSION: usize = 33;
977
978 let queue = PacketQueue::new();
979 for entity_id in 1_i32..=64 {
980 let key = PlayerPacketLaneKey {
981 player_id: Uuid::from_u128(u128::from(entity_id.unsigned_abs())),
982 entity_id,
983 };
984 for packet in 0..PACKETS_PER_SESSION {
985 assert!(
986 queue
987 .try_submit(
988 key,
989 ScheduledPacketExecution::PlayerLocal,
990 PACKET_ADMISSION_OVERHEAD,
991 packet,
992 )
993 .is_ok()
994 );
995 }
996 }
997
998 let state = queue.state.lock();
999 assert_eq!(state.lanes.len(), SESSION_COUNT);
1000 assert_eq!(
1001 usize::try_from(state.next_sequence),
1002 Ok(SESSION_COUNT * PACKETS_PER_SESSION)
1003 );
1004 }
1005
1006 #[test]
1007 fn discarding_stale_session_keeps_replacement_session_work() {
1008 let player_id = Uuid::nil();
1009 let stale_key = PlayerPacketLaneKey {
1010 player_id,
1011 entity_id: 1,
1012 };
1013 let replacement_key = PlayerPacketLaneKey {
1014 player_id,
1015 entity_id: 2,
1016 };
1017 let queue = PacketQueue::with_limits(limits(10, 100));
1018 queue.submit(
1019 stale_key,
1020 ScheduledPacketExecution::Exclusive,
1021 "stale barrier",
1022 );
1023 queue.submit(
1024 stale_key,
1025 ScheduledPacketExecution::Serialized,
1026 "stale serialized",
1027 );
1028 queue.submit(
1029 replacement_key,
1030 ScheduledPacketExecution::PlayerLocal,
1031 "replacement",
1032 );
1033
1034 queue.discard_lane(stale_key);
1035
1036 {
1037 let queue_state = queue.state.lock();
1038 assert!(!queue_state.lanes.contains_key(&stale_key));
1039 let replacement_lane = queue_state
1040 .lanes
1041 .get(&replacement_key)
1042 .expect("replacement session lane should remain");
1043 assert_eq!(replacement_lane.outstanding_packets, 1);
1044 assert_eq!(replacement_lane.outstanding_bytes, 1);
1045 assert!(queue_state.serialized.is_empty());
1046 assert!(queue_state.exclusive.is_empty());
1047 }
1048
1049 queue.open();
1050 let Some(mut work) = queue.try_next() else {
1051 panic!("replacement session packet should remain runnable");
1052 };
1053 assert_eq!(work.take(), Some("replacement"));
1054 }
1055
1056 #[test]
1057 fn active_packet_remains_charged_until_completion() {
1058 let queue = PacketQueue::with_limits(limits(1, 4));
1059 assert_eq!(
1060 queue.try_submit(1, ScheduledPacketExecution::PlayerLocal, 4, "active"),
1061 Ok(())
1062 );
1063 queue.open();
1064 let Some(work) = queue.try_next() else {
1065 panic!("accepted packet should start");
1066 };
1067
1068 assert_eq!(
1069 queue.try_submit(1, ScheduledPacketExecution::PlayerLocal, 1, "rejected"),
1070 Err(PacketAdmissionError::PlayerPacketLimit)
1071 );
1072 {
1073 let state = queue.state.lock();
1074 let lane = state.lanes.get(&1).expect("active lane should remain");
1075 assert_eq!(lane.outstanding_packets, 1);
1076 assert_eq!(lane.outstanding_bytes, 4);
1077 }
1078
1079 drop(work);
1080 {
1081 let state = queue.state.lock();
1082 assert!(state.lanes.is_empty());
1083 }
1084 assert_eq!(
1085 queue.try_submit(1, ScheduledPacketExecution::PlayerLocal, 4, "next"),
1086 Ok(())
1087 );
1088 }
1089
1090 #[test]
1091 fn discarding_lane_removes_hidden_barriers_and_keeps_active_work_charged() {
1092 let queue = PacketQueue::with_limits(limits(10, 100));
1093 queue.submit(1, ScheduledPacketExecution::PlayerLocal, "active attacker");
1094 queue.submit(1, ScheduledPacketExecution::Serialized, "hidden serialized");
1095 queue.submit(1, ScheduledPacketExecution::Exclusive, "hidden exclusive");
1096 queue.submit(2, ScheduledPacketExecution::PlayerLocal, "other player");
1097 queue.open();
1098 let Some(mut active) = queue.try_next() else {
1099 panic!("attacker's first packet should start");
1100 };
1101 assert_eq!(active.take(), Some("active attacker"));
1102
1103 queue.discard_lane(1);
1104
1105 {
1106 let state = queue.state.lock();
1107 assert!(state.serialized.is_empty());
1108 assert!(state.exclusive.is_empty());
1109 assert!(state.serialized_ready.is_empty());
1110 assert!(state.exclusive_ready.is_empty());
1111 let Some(lane) = state.lanes.get(&1) else {
1112 panic!("active attacker lane should remain");
1113 };
1114 assert!(lane.active);
1115 assert!(lane.queued.is_empty());
1116 assert_eq!(lane.outstanding_packets, 1);
1117 assert_eq!(lane.outstanding_bytes, 1);
1118 }
1119
1120 let Some(mut other) = queue.try_next() else {
1121 panic!("purged barriers should no longer block another player");
1122 };
1123 assert_eq!(other.take(), Some("other player"));
1124 drop(other);
1125 drop(active);
1126
1127 let state = queue.state.lock();
1128 assert!(state.lanes.is_empty());
1129 }
1130
1131 #[test]
1132 fn discarding_idle_lane_removes_ready_exclusive_barrier() {
1133 let queue = PacketQueue::with_limits(limits(10, 100));
1134 queue.submit(1, ScheduledPacketExecution::Exclusive, "attacker barrier");
1135 queue.submit(1, ScheduledPacketExecution::Serialized, "hidden serialized");
1136 queue.submit(2, ScheduledPacketExecution::PlayerLocal, "other player");
1137
1138 queue.discard_lane(1);
1139
1140 {
1141 let state = queue.state.lock();
1142 assert!(!state.lanes.contains_key(&1));
1143 let other = state.lanes.get(&2).expect("other lane should remain");
1144 assert_eq!(other.outstanding_packets, 1);
1145 assert_eq!(other.outstanding_bytes, 1);
1146 assert!(state.exclusive_ready.is_empty());
1147 assert!(state.exclusive.is_empty());
1148 assert!(state.serialized.is_empty());
1149 }
1150
1151 queue.open();
1152 let Some(mut other) = queue.try_next() else {
1153 panic!("discarded exclusive barrier should not block another player");
1154 };
1155 assert_eq!(other.take(), Some("other player"));
1156 }
1157
1158 #[test]
1159 fn queued_packets_start_in_submission_order_when_opened() {
1160 let queue = PacketQueue::new();
1161 queue.submit(1, ScheduledPacketExecution::PlayerLocal, 1);
1162 queue.submit(2, ScheduledPacketExecution::PlayerLocal, 2);
1163 queue.submit(1, ScheduledPacketExecution::PlayerLocal, 3);
1164 assert!(queue.try_next().is_none());
1165
1166 queue.open();
1167 let mut processed = Vec::new();
1168 while let Some(mut work) = queue.try_next() {
1169 if let Some(value) = work.take() {
1170 processed.push(value);
1171 }
1172 }
1173
1174 assert_eq!(processed, [1, 2, 3]);
1175 }
1176
1177 #[test]
1178 fn packet_lane_only_allows_one_active_handler() {
1179 let queue = PacketQueue::new();
1180 queue.submit(
1181 1,
1182 ScheduledPacketExecution::PlayerLocal,
1183 "first player packet",
1184 );
1185 queue.submit(
1186 1,
1187 ScheduledPacketExecution::PlayerLocal,
1188 "second player packet",
1189 );
1190 queue.submit(
1191 2,
1192 ScheduledPacketExecution::PlayerLocal,
1193 "other player packet",
1194 );
1195 queue.open();
1196
1197 let Some(mut first) = queue.try_next() else {
1198 panic!("first packet should start");
1199 };
1200 assert_eq!(first.take(), Some("first player packet"));
1201
1202 let Some(mut other_player) = queue.try_next() else {
1203 panic!("another player's packet should be able to start");
1204 };
1205 assert_eq!(other_player.take(), Some("other player packet"));
1206 assert!(queue.try_next().is_none());
1207
1208 drop(first);
1209 let Some(mut second) = queue.try_next() else {
1210 panic!("the next packet should start after its lane becomes idle");
1211 };
1212 assert_eq!(second.take(), Some("second player packet"));
1213 }
1214
1215 #[test]
1216 fn serialized_packet_overlaps_player_local_work_but_not_another_serialized_packet() {
1217 let queue = PacketQueue::new();
1218 queue.submit(1, ScheduledPacketExecution::PlayerLocal, "first local");
1219 queue.submit(2, ScheduledPacketExecution::Serialized, "first serialized");
1220 queue.submit(3, ScheduledPacketExecution::Serialized, "second serialized");
1221 queue.submit(4, ScheduledPacketExecution::PlayerLocal, "later local");
1222 queue.open();
1223
1224 let Some(mut first_local) = queue.try_next() else {
1225 panic!("first player-local packet should start");
1226 };
1227 assert_eq!(first_local.take(), Some("first local"));
1228
1229 let Some(mut first_serialized) = queue.try_next() else {
1230 panic!("serialized packet should overlap player-local work");
1231 };
1232 assert_eq!(first_serialized.take(), Some("first serialized"));
1233
1234 let Some(mut later_local) = queue.try_next() else {
1235 panic!("player-local work should bypass a blocked serialized packet");
1236 };
1237 assert_eq!(later_local.take(), Some("later local"));
1238 assert!(queue.try_next().is_none());
1239
1240 drop(first_serialized);
1241 let Some(mut second_serialized) = queue.try_next() else {
1242 panic!("next serialized packet should start after its predecessor finishes");
1243 };
1244 assert_eq!(second_serialized.take(), Some("second serialized"));
1245 }
1246
1247 #[test]
1248 fn serialized_packet_hidden_in_an_active_lane_preserves_serialized_order() {
1249 let queue = PacketQueue::new();
1250 queue.submit(1, ScheduledPacketExecution::PlayerLocal, "active lane");
1251 queue.submit(1, ScheduledPacketExecution::Serialized, "hidden serialized");
1252 queue.submit(2, ScheduledPacketExecution::Serialized, "later serialized");
1253 queue.submit(
1254 3,
1255 ScheduledPacketExecution::PlayerLocal,
1256 "independent local",
1257 );
1258 queue.open();
1259
1260 let Some(active_lane) = queue.try_next() else {
1261 panic!("first lane packet should start");
1262 };
1263 let Some(mut independent_local) = queue.try_next() else {
1264 panic!("player-local work should bypass serialized ordering contention");
1265 };
1266 assert_eq!(independent_local.take(), Some("independent local"));
1267 assert!(queue.try_next().is_none());
1268
1269 drop(active_lane);
1270 let Some(mut hidden_serialized) = queue.try_next() else {
1271 panic!("earliest serialized packet should start when its lane becomes idle");
1272 };
1273 assert_eq!(hidden_serialized.take(), Some("hidden serialized"));
1274 assert!(queue.try_next().is_none());
1275
1276 drop(hidden_serialized);
1277 let Some(mut later_serialized) = queue.try_next() else {
1278 panic!("later serialized packet should preserve global submission order");
1279 };
1280 assert_eq!(later_serialized.take(), Some("later serialized"));
1281 }
1282
1283 #[test]
1284 fn blocking_worker_bypasses_active_serialized_work_for_player_local_work() {
1285 let queue = Arc::new(PacketQueue::new());
1286 queue.submit(1, ScheduledPacketExecution::Serialized, "active serialized");
1287 queue.submit(2, ScheduledPacketExecution::Serialized, "queued serialized");
1288 queue.submit(3, ScheduledPacketExecution::PlayerLocal, "player local");
1289 queue.open();
1290
1291 let Some(active_serialized) = queue.try_next() else {
1292 panic!("first serialized packet should start");
1293 };
1294 let worker_queue = Arc::clone(&queue);
1295 let (sender, receiver) = mpsc::channel();
1296 let worker = thread::spawn(move || {
1297 for _ in 0..2 {
1298 let Some(mut work) = worker_queue.next() else {
1299 return;
1300 };
1301 if let Some(value) = work.take() {
1302 let _ = sender.send(value);
1303 }
1304 }
1305 });
1306
1307 assert_eq!(
1308 receiver.recv_timeout(Duration::from_secs(1)),
1309 Ok("player local")
1310 );
1311 assert!(receiver.recv_timeout(Duration::from_millis(10)).is_err());
1312
1313 drop(active_serialized);
1314 assert_eq!(
1315 receiver.recv_timeout(Duration::from_secs(1)),
1316 Ok("queued serialized")
1317 );
1318 assert!(worker.join().is_ok());
1319 }
1320
1321 #[test]
1322 fn exclusive_packet_waits_for_active_work_and_blocks_later_packets() {
1323 let queue = PacketQueue::new();
1324 queue.submit(1, ScheduledPacketExecution::PlayerLocal, "before barrier");
1325 queue.submit(2, ScheduledPacketExecution::Exclusive, "barrier");
1326 queue.submit(3, ScheduledPacketExecution::PlayerLocal, "after barrier");
1327 queue.open();
1328
1329 let Some(mut before) = queue.try_next() else {
1330 panic!("packet before the barrier should start");
1331 };
1332 assert_eq!(before.take(), Some("before barrier"));
1333 assert!(queue.try_next().is_none());
1334
1335 drop(before);
1336 let Some(mut barrier) = queue.try_next() else {
1337 panic!("exclusive packet should start after active work finishes");
1338 };
1339 assert_eq!(barrier.take(), Some("barrier"));
1340 assert!(queue.try_next().is_none());
1341
1342 drop(barrier);
1343 let Some(mut after) = queue.try_next() else {
1344 panic!("packet after the barrier should start after it finishes");
1345 };
1346 assert_eq!(after.take(), Some("after barrier"));
1347 }
1348
1349 #[test]
1350 fn exclusive_packet_waits_for_serialized_work_and_blocks_player_local_work() {
1351 let queue = PacketQueue::new();
1352 queue.submit(1, ScheduledPacketExecution::Serialized, "before barrier");
1353 queue.submit(2, ScheduledPacketExecution::Exclusive, "barrier");
1354 queue.submit(3, ScheduledPacketExecution::PlayerLocal, "after barrier");
1355 queue.open();
1356
1357 let Some(mut before) = queue.try_next() else {
1358 panic!("serialized packet before the barrier should start");
1359 };
1360 assert_eq!(before.take(), Some("before barrier"));
1361 assert!(queue.try_next().is_none());
1362
1363 drop(before);
1364 let Some(mut barrier) = queue.try_next() else {
1365 panic!("exclusive packet should wait for serialized work");
1366 };
1367 assert_eq!(barrier.take(), Some("barrier"));
1368 assert!(queue.try_next().is_none());
1369
1370 drop(barrier);
1371 let Some(mut after) = queue.try_next() else {
1372 panic!("player-local packet should wait for the exclusive barrier");
1373 };
1374 assert_eq!(after.take(), Some("after barrier"));
1375 }
1376
1377 #[test]
1378 fn exclusive_packet_hidden_in_an_active_lane_still_blocks_later_lanes() {
1379 let queue = PacketQueue::new();
1380 queue.submit(1, ScheduledPacketExecution::PlayerLocal, "active");
1381 queue.submit(1, ScheduledPacketExecution::Exclusive, "barrier");
1382 queue.submit(2, ScheduledPacketExecution::PlayerLocal, "later lane");
1383 queue.open();
1384
1385 let Some(active) = queue.try_next() else {
1386 panic!("first packet should start");
1387 };
1388 assert!(queue.try_next().is_none());
1389
1390 drop(active);
1391 let Some(mut barrier) = queue.try_next() else {
1392 panic!("hidden exclusive packet should become runnable");
1393 };
1394 assert_eq!(barrier.take(), Some("barrier"));
1395 }
1396
1397 #[test]
1398 fn closed_phase_retains_new_packets_for_the_next_open_phase() {
1399 let queue = PacketQueue::new();
1400 queue.open();
1401 queue.close();
1402 queue.submit(1, ScheduledPacketExecution::PlayerLocal, 1);
1403
1404 assert!(queue.try_next().is_none());
1405 queue.open();
1406 let Some(mut work) = queue.try_next() else {
1407 panic!("queued packet should become available when the packet phase opens");
1408 };
1409 assert_eq!(work.take(), Some(1));
1410 }
1411
1412 #[test]
1413 fn tick_drain_only_processes_packets_submitted_before_its_cutoff() {
1414 let queue = PacketQueue::new();
1415 queue.submit(1, ScheduledPacketExecution::PlayerLocal, "before cutoff");
1416 queue.open();
1417 let Some((before_sequence, should_wake)) = queue.begin_tick_drain() else {
1418 panic!("running queue should begin a tick drain");
1419 };
1420 assert!(should_wake);
1421 queue.submit(2, ScheduledPacketExecution::PlayerLocal, "after cutoff");
1422
1423 let Some(mut before) = queue.try_next() else {
1424 panic!("packet before the cutoff should drain");
1425 };
1426 assert_eq!(before.take(), Some("before cutoff"));
1427 drop(before);
1428
1429 assert!(queue.try_next().is_none());
1430 assert!(queue.tick_drain_complete(before_sequence));
1431 queue.finish_tick_drain(before_sequence);
1432 queue.open();
1433 let Some(mut after) = queue.try_next() else {
1434 panic!("packet after the cutoff should wait for the next open phase");
1435 };
1436 assert_eq!(after.take(), Some("after cutoff"));
1437 }
1438
1439 #[test]
1440 fn tick_drain_orders_hidden_serialized_and_exclusive_work_before_its_cutoff() {
1441 let queue = PacketQueue::new();
1442 queue.submit(1, ScheduledPacketExecution::PlayerLocal, "active local");
1443 queue.submit(1, ScheduledPacketExecution::Serialized, "hidden serialized");
1444 queue.submit(2, ScheduledPacketExecution::Exclusive, "exclusive");
1445 queue.open();
1446
1447 let Some(active_local) = queue.try_next() else {
1448 panic!("player-local packet should start before draining");
1449 };
1450 let Some((before_sequence, should_wake)) = queue.begin_tick_drain() else {
1451 panic!("running queue should begin a tick drain");
1452 };
1453 assert!(!should_wake);
1454 queue.submit(
1455 3,
1456 ScheduledPacketExecution::Serialized,
1457 "post-cutoff serialized",
1458 );
1459 queue.submit(
1460 4,
1461 ScheduledPacketExecution::PlayerLocal,
1462 "post-cutoff local",
1463 );
1464
1465 assert!(!queue.tick_drain_complete(before_sequence));
1466 assert!(queue.try_next().is_none());
1467 drop(active_local);
1468
1469 let Some(mut hidden_serialized) = queue.try_next() else {
1470 panic!("hidden pre-cutoff serialized packet should drain");
1471 };
1472 assert_eq!(hidden_serialized.take(), Some("hidden serialized"));
1473 assert!(!queue.tick_drain_complete(before_sequence));
1474 drop(hidden_serialized);
1475
1476 let Some(mut exclusive) = queue.try_next() else {
1477 panic!("pre-cutoff exclusive packet should drain after earlier work");
1478 };
1479 assert_eq!(exclusive.take(), Some("exclusive"));
1480 assert!(!queue.tick_drain_complete(before_sequence));
1481 drop(exclusive);
1482
1483 assert!(queue.try_next().is_none());
1484 assert!(queue.tick_drain_complete(before_sequence));
1485 queue.finish_tick_drain(before_sequence);
1486 queue.open();
1487
1488 let Some(mut post_cutoff_serialized) = queue.try_next() else {
1489 panic!("post-cutoff serialized packet should remain for the next phase");
1490 };
1491 assert_eq!(
1492 post_cutoff_serialized.take(),
1493 Some("post-cutoff serialized")
1494 );
1495 let Some(mut post_cutoff_local) = queue.try_next() else {
1496 panic!("post-cutoff player-local packet should remain for the next phase");
1497 };
1498 assert_eq!(post_cutoff_local.take(), Some("post-cutoff local"));
1499 }
1500
1501 #[tokio::test]
1502 async fn tick_drain_waits_for_active_packet_work() {
1503 let queue = Arc::new(PacketQueue::new());
1504 queue.open();
1505 queue.submit(1, ScheduledPacketExecution::PlayerLocal, 1);
1506 let Some(work) = queue.try_next() else {
1507 panic!("open packet phase should start queued work");
1508 };
1509 let drain = queue.drain_for_tick();
1510 tokio::pin!(drain);
1511
1512 assert!(
1513 timeout(Duration::from_millis(10), drain.as_mut())
1514 .await
1515 .is_err()
1516 );
1517 drop(work);
1518 assert!(
1519 timeout(Duration::from_secs(1), drain.as_mut())
1520 .await
1521 .is_ok()
1522 );
1523 }
1524
1525 #[tokio::test]
1526 async fn overload_progress_waits_for_one_active_packet_to_finish() {
1527 let queue = Arc::new(PacketQueue::new());
1528 queue.open();
1529 queue.submit(1, ScheduledPacketExecution::PlayerLocal, 1);
1530 let Some(work) = queue.try_next() else {
1531 panic!("open packet phase should start queued work");
1532 };
1533 let Some(completed) = queue.progress_baseline() else {
1534 panic!("active packet work should require progress");
1535 };
1536 let progress = queue.wait_for_progress_since(completed);
1537 tokio::pin!(progress);
1538
1539 assert!(
1540 timeout(Duration::from_millis(10), progress.as_mut())
1541 .await
1542 .is_err()
1543 );
1544 drop(work);
1545 assert!(
1546 timeout(Duration::from_secs(1), progress.as_mut())
1547 .await
1548 .is_ok()
1549 );
1550 }
1551
1552 #[test]
1553 fn stopped_queue_discards_pending_and_future_work() {
1554 let queue = PacketQueue::new();
1555 queue.open();
1556 queue.submit(1, ScheduledPacketExecution::PlayerLocal, 1);
1557 queue.stop();
1558 queue.submit(1, ScheduledPacketExecution::PlayerLocal, 2);
1559
1560 assert!(queue.try_next().is_none());
1561 let state = queue.state.lock();
1562 assert!(state.lanes.is_empty());
1563 }
1564
1565 #[test]
1566 fn stopping_with_active_work_keeps_completion_accounting_valid() {
1567 let queue = PacketQueue::new();
1568 queue.submit(1, ScheduledPacketExecution::Exclusive, 1);
1569 queue.submit(1, ScheduledPacketExecution::PlayerLocal, 2);
1570 queue.open();
1571 let Some(work) = queue.try_next() else {
1572 panic!("open packet phase should start queued work");
1573 };
1574
1575 queue.stop();
1576 {
1577 let state = queue.state.lock();
1578 let lane = state.lanes.get(&1).expect("active lane should remain");
1579 assert!(lane.active);
1580 assert!(lane.queued.is_empty());
1581 assert_eq!(lane.outstanding_packets, 1);
1582 assert_eq!(lane.outstanding_bytes, 1);
1583 }
1584 drop(work);
1585
1586 assert!(queue.try_next().is_none());
1587 let state = queue.state.lock();
1588 assert_eq!(state.active, 0);
1589 assert!(state.lanes.is_empty());
1590 }
1591
1592 #[tokio::test]
1593 async fn stopping_with_active_serialized_and_local_work_clears_all_queue_state() {
1594 let queue = PacketQueue::new();
1595 queue.submit(1, ScheduledPacketExecution::Serialized, "active serialized");
1596 queue.submit(2, ScheduledPacketExecution::PlayerLocal, "active local");
1597 queue.submit(3, ScheduledPacketExecution::Serialized, "queued serialized");
1598 queue.submit(4, ScheduledPacketExecution::Exclusive, "queued exclusive");
1599 queue.open();
1600
1601 let Some(active_serialized) = queue.try_next() else {
1602 panic!("serialized packet should start");
1603 };
1604 let Some(active_local) = queue.try_next() else {
1605 panic!("player-local packet should overlap serialized work");
1606 };
1607 queue.stop();
1608 queue.drain_for_tick().await;
1609
1610 {
1611 let state = queue.state.lock();
1612 assert_eq!(state.active, 2);
1613 assert!(state.serialized_active);
1614 assert!(state.player_local_ready.is_empty());
1615 assert!(state.serialized_ready.is_empty());
1616 assert!(state.exclusive_ready.is_empty());
1617 assert!(state.serialized.is_empty());
1618 assert!(state.exclusive.is_empty());
1619 assert_eq!(state.lanes.len(), 2);
1620 for key in [1, 2] {
1621 let lane = state.lanes.get(&key).expect("active lane should remain");
1622 assert!(lane.active);
1623 assert!(lane.queued.is_empty());
1624 assert_eq!(lane.outstanding_packets, 1);
1625 assert_eq!(lane.outstanding_bytes, 1);
1626 }
1627 }
1628
1629 drop(active_serialized);
1630 drop(active_local);
1631
1632 let state = queue.state.lock();
1633 assert_eq!(state.active, 0);
1634 assert!(!state.serialized_active);
1635 assert!(state.lanes.is_empty());
1636 }
1637
1638 #[test]
1639 fn blocking_worker_only_starts_work_during_the_open_phase() {
1640 let queue = Arc::new(PacketQueue::new());
1641 let worker_queue = Arc::clone(&queue);
1642 let (sender, receiver) = mpsc::channel();
1643 let worker = thread::spawn(move || {
1644 while let Some(mut work) = worker_queue.next() {
1645 if let Some(value) = work.take() {
1646 let _ = sender.send(value);
1647 }
1648 }
1649 });
1650
1651 queue.submit(1, ScheduledPacketExecution::PlayerLocal, 1);
1652 assert!(receiver.recv_timeout(Duration::from_millis(10)).is_err());
1653 queue.open();
1654 assert_eq!(receiver.recv_timeout(Duration::from_secs(1)), Ok(1));
1655 queue.close();
1656 queue.submit(1, ScheduledPacketExecution::PlayerLocal, 2);
1657 assert!(receiver.recv_timeout(Duration::from_millis(10)).is_err());
1658 queue.stop();
1659
1660 assert!(worker.join().is_ok());
1661 }
1662}