1use std::{
4 cmp::Reverse,
5 collections::{BinaryHeap, VecDeque},
6 hash::Hash,
7 mem,
8 sync::Arc,
9};
10
11use crate::{
12 entity::Entity,
13 player::{
14 Player, PlayerSession, PlayerSessionId,
15 connection::NetworkConnection,
16 connection::{ScheduledPacketExecution, ScheduledPlayPacket},
17 },
18};
19use parking_lot::Condvar;
20use rustc_hash::{FxHashMap, FxHashSet};
21use steel_utils::{locks::SyncMutex, translations};
22use tokio::sync::Notify;
23use tokio::task::yield_now;
24
25use super::Server;
26
27const MAX_OUTSTANDING_PACKETS_PER_PLAYER: usize = 8_192;
31const MAX_OUTSTANDING_BYTES_PER_PLAYER: usize = 32 * 1024 * 1024;
32const PACKET_ADMISSION_OVERHEAD: usize = 256;
34
35#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
36struct PlayerPacketLaneKey(PlayerSessionId);
37
38impl PlayerPacketLaneKey {
39 const fn new(session: &PlayerSession) -> Self {
40 Self(session.id())
41 }
42}
43
44#[derive(Clone, Copy)]
46pub(crate) struct PlayerPacketTransition(PacketLanePause<PlayerPacketLaneKey>);
47
48struct PendingPlayPacket {
49 session: Arc<PlayerSession>,
50 packet: ScheduledPlayPacket,
51}
52
53pub(super) struct PacketProcessor {
60 queued: PacketQueue<PlayerPacketLaneKey, PendingPlayPacket>,
61}
62
63impl PacketProcessor {
64 pub(super) fn new() -> Self {
65 Self {
66 queued: PacketQueue::new(),
67 }
68 }
69
70 pub(super) fn schedule(
71 &self,
72 player: Arc<Player>,
73 packet: ScheduledPlayPacket,
74 payload_bytes: usize,
75 ) {
76 let player_id = player.gameprofile.id;
77 let lane_key = PlayerPacketLaneKey::new(&player.session);
78 if player.connection.closed() {
79 self.queued.discard_lane(lane_key);
80 return;
81 }
82
83 let execution = packet.execution();
84 let admission_bytes = payload_bytes.saturating_add(PACKET_ADMISSION_OVERHEAD);
85 let result = self.queued.try_submit(
86 lane_key,
87 execution,
88 admission_bytes,
89 PendingPlayPacket {
90 session: Arc::clone(&player.session),
91 packet,
92 },
93 );
94 let Err(error) = result else {
95 return;
96 };
97 if error == PacketAdmissionError::Stopped {
98 return;
99 }
100
101 tracing::warn!(
102 player_id = %player_id,
103 entity_id = player.id(),
104 ?error,
105 "Disconnecting player after inbound packet admission limit"
106 );
107 player.disconnect(translations::DISCONNECT_EXCEEDED_PACKET_RATE.msg());
108 self.queued.discard_lane(lane_key);
109 }
110
111 pub(super) fn run(&self, server: &Arc<Server>) {
113 while let Some(mut work) = self.queued.next() {
114 let Some(pending) = work.take() else {
115 continue;
116 };
117 let Some(player) = pending.session.current_player() else {
118 continue;
119 };
120 if !Self::packet_is_runnable(
121 &player,
122 &pending.packet,
123 server.cancel_token.is_cancelled(),
124 ) {
125 continue;
126 }
127
128 pending.packet.handle(player, server);
129 }
130 }
131
132 pub(super) fn pause_player_session(
134 &self,
135 session: &PlayerSession,
136 ) -> Option<PlayerPacketTransition> {
137 self.queued
138 .pause_lane(PlayerPacketLaneKey::new(session))
139 .map(PlayerPacketTransition)
140 }
141
142 pub(super) fn resume_player_session(&self, transition: PlayerPacketTransition) -> bool {
144 self.queued.resume_lane(transition.0)
145 }
146
147 pub(super) fn discard_player_session(&self, session: &PlayerSession) {
149 self.queued.discard_lane(PlayerPacketLaneKey::new(session));
150 }
151
152 fn packet_is_runnable(
153 player: &Player,
154 packet: &ScheduledPlayPacket,
155 server_cancelled: bool,
156 ) -> bool {
157 !player.connection.closed()
158 && !server_cancelled
159 && player.gate_domain_switch_packet(
160 packet.is_domain_handshake_packet(),
161 packet.is_perform_respawn(),
162 )
163 }
164
165 pub(super) fn open_after_tick(&self) {
167 self.queued.open();
168 }
169
170 pub(super) async fn wait_for_overload_progress(&self) {
172 let Some(completed) = self.queued.progress_baseline() else {
173 return;
174 };
175 yield_now().await;
176 self.queued.wait_for_progress_since(completed).await;
177 }
178
179 pub(super) async fn close_for_tick(&self) {
181 self.queued.drain_for_tick().await;
182 }
183
184 pub(super) fn stop(&self) {
186 self.queued.stop();
187 }
188}
189
190impl Default for PacketProcessor {
191 fn default() -> Self {
192 Self::new()
193 }
194}
195
196#[derive(Clone, Copy, Eq, PartialEq)]
197enum PacketPhase {
198 Closed,
199 Open,
200 Draining(u64),
201 Stopped,
202}
203
204struct SequencedPacket<T> {
205 sequence: u64,
206 execution: ScheduledPacketExecution,
207 admission_bytes: usize,
208 value: T,
209}
210
211struct DeferredPacket<T> {
212 execution: ScheduledPacketExecution,
213 admission_bytes: usize,
214 value: T,
215}
216
217struct PacketLane<T> {
218 queued: VecDeque<SequencedPacket<T>>,
219 deferred: VecDeque<DeferredPacket<T>>,
220 active: bool,
221 pause_generation: Option<u64>,
222 outstanding_packets: usize,
223 outstanding_bytes: usize,
224}
225
226impl<T> PacketLane<T> {
227 const fn new() -> Self {
228 Self {
229 queued: VecDeque::new(),
230 deferred: VecDeque::new(),
231 active: false,
232 pause_generation: None,
233 outstanding_packets: 0,
234 outstanding_bytes: 0,
235 }
236 }
237}
238
239#[derive(Clone, Copy)]
240struct PacketAdmissionLimits {
241 per_player_packets: usize,
242 per_player_bytes: usize,
243}
244
245impl PacketAdmissionLimits {
246 const PRODUCTION: Self = Self {
247 per_player_packets: MAX_OUTSTANDING_PACKETS_PER_PLAYER,
248 per_player_bytes: MAX_OUTSTANDING_BYTES_PER_PLAYER,
249 };
250}
251
252#[derive(Clone, Copy, Debug, Eq, PartialEq)]
253enum PacketAdmissionError {
254 Stopped,
255 PlayerPacketLimit,
256 PlayerByteLimit,
257}
258
259struct PacketQueueState<K, T> {
260 phase: PacketPhase,
261 lanes: FxHashMap<K, PacketLane<T>>,
262 player_local_ready: BinaryHeap<Reverse<(u64, K)>>,
263 serialized_ready: BinaryHeap<Reverse<(u64, K)>>,
264 exclusive_ready: BinaryHeap<Reverse<(u64, K)>>,
265 serialized: BinaryHeap<Reverse<u64>>,
266 exclusive: BinaryHeap<Reverse<u64>>,
267 next_sequence: u64,
268 next_pause_generation: u64,
269 active: usize,
270 serialized_active: bool,
271 exclusive_active: bool,
272 completed: u64,
273}
274
275struct PacketQueue<K, T> {
277 state: SyncMutex<PacketQueueState<K, T>>,
278 limits: PacketAdmissionLimits,
279 work_available: Condvar,
280 idle: Notify,
281 progress: Notify,
282}
283
284#[derive(Clone, Copy)]
285struct PacketLanePause<K> {
286 key: K,
287 generation: u64,
288}
289
290impl<K, T> PacketQueue<K, T>
291where
292 K: Copy + Eq + Hash + Ord,
293{
294 fn new() -> Self {
295 Self::with_limits(PacketAdmissionLimits::PRODUCTION)
296 }
297
298 fn with_limits(limits: PacketAdmissionLimits) -> Self {
299 Self {
300 state: SyncMutex::new(PacketQueueState {
301 phase: PacketPhase::Closed,
302 lanes: FxHashMap::default(),
303 player_local_ready: BinaryHeap::new(),
304 serialized_ready: BinaryHeap::new(),
305 exclusive_ready: BinaryHeap::new(),
306 serialized: BinaryHeap::new(),
307 exclusive: BinaryHeap::new(),
308 next_sequence: 0,
309 next_pause_generation: 0,
310 active: 0,
311 serialized_active: false,
312 exclusive_active: false,
313 completed: 0,
314 }),
315 limits,
316 work_available: Condvar::new(),
317 idle: Notify::new(),
318 progress: Notify::new(),
319 }
320 }
321
322 fn try_submit(
323 &self,
324 key: K,
325 execution: ScheduledPacketExecution,
326 admission_bytes: usize,
327 value: T,
328 ) -> Result<(), PacketAdmissionError> {
329 let mut state = self.state.lock();
330 if state.phase == PacketPhase::Stopped {
331 return Err(PacketAdmissionError::Stopped);
332 }
333
334 let (player_packets, player_bytes) = state
335 .lanes
336 .get(&key)
337 .map(|lane| (lane.outstanding_packets, lane.outstanding_bytes))
338 .unwrap_or_default();
339 if player_packets >= self.limits.per_player_packets {
340 return Err(PacketAdmissionError::PlayerPacketLimit);
341 }
342 if admission_bytes > self.limits.per_player_bytes.saturating_sub(player_bytes) {
343 return Err(PacketAdmissionError::PlayerByteLimit);
344 }
345 let paused = state
346 .lanes
347 .get(&key)
348 .is_some_and(|lane| lane.pause_generation.is_some());
349 if paused {
350 let lane = state.lanes.entry(key).or_insert_with(PacketLane::new);
351 lane.outstanding_packets += 1;
352 lane.outstanding_bytes += admission_bytes;
353 lane.deferred.push_back(DeferredPacket {
354 execution,
355 admission_bytes,
356 value,
357 });
358 return Ok(());
359 }
360
361 let sequence = state.next_sequence;
362 assert!(sequence != u64::MAX, "packet submission sequence exhausted");
363 state.next_sequence = sequence + 1;
364
365 let lane = state.lanes.entry(key).or_insert_with(PacketLane::new);
366 let became_ready = !lane.active && lane.queued.is_empty();
367 lane.outstanding_packets += 1;
368 lane.outstanding_bytes += admission_bytes;
369 lane.queued.push_back(SequencedPacket {
370 sequence,
371 execution,
372 admission_bytes,
373 value,
374 });
375 if became_ready {
376 Self::mark_ready(&mut state, sequence, key, execution);
377 }
378 Self::mark_global_order(&mut state, sequence, execution);
379 let should_wake = became_ready && state.phase == PacketPhase::Open;
380 drop(state);
381 if should_wake {
382 self.work_available.notify_one();
383 }
384 Ok(())
385 }
386
387 #[cfg(test)]
388 fn submit(&self, key: K, execution: ScheduledPacketExecution, value: T) {
389 let result = self.try_submit(key, execution, 1, value);
390 assert!(
391 matches!(result, Ok(()) | Err(PacketAdmissionError::Stopped)),
392 "test packet exceeded admission limits: {result:?}"
393 );
394 }
395
396 fn open(&self) {
397 let mut state = self.state.lock();
398 if state.phase == PacketPhase::Stopped {
399 return;
400 }
401 state.phase = PacketPhase::Open;
402 let should_wake = Self::select_next(&state, None).is_some();
403 drop(state);
404 if should_wake {
405 self.work_available.notify_all();
406 }
407 }
408
409 #[cfg(test)]
410 fn close(&self) {
411 let mut state = self.state.lock();
412 if state.phase != PacketPhase::Stopped {
413 state.phase = PacketPhase::Closed;
414 }
415 }
416
417 async fn drain_for_tick(&self) {
418 let Some((before_sequence, should_wake)) = self.begin_tick_drain() else {
419 return;
420 };
421 if should_wake {
422 self.work_available.notify_all();
423 }
424
425 loop {
426 let idle = self.idle.notified();
427 if self.tick_drain_complete(before_sequence) {
428 break;
429 }
430 idle.await;
431 }
432
433 self.finish_tick_drain(before_sequence);
434 }
435
436 fn finish_tick_drain(&self, before_sequence: u64) {
437 let mut state = self.state.lock();
438 if state.phase == PacketPhase::Draining(before_sequence) {
439 state.phase = PacketPhase::Closed;
440 }
441 }
442
443 fn begin_tick_drain(&self) -> Option<(u64, bool)> {
444 let mut state = self.state.lock();
445 let before_sequence = match state.phase {
446 PacketPhase::Stopped => None,
447 PacketPhase::Draining(before_sequence) => Some(before_sequence),
448 PacketPhase::Closed | PacketPhase::Open => {
449 let before_sequence = state.next_sequence;
450 state.phase = PacketPhase::Draining(before_sequence);
451 Some(before_sequence)
452 }
453 }?;
454 let should_wake = Self::select_next(&state, Some(before_sequence)).is_some();
455 Some((before_sequence, should_wake))
456 }
457
458 fn tick_drain_complete(&self, before_sequence: u64) -> bool {
459 let state = self.state.lock();
460 if state.phase == PacketPhase::Stopped {
461 return true;
462 }
463 if state.active != 0 {
464 return false;
465 }
466 Self::next_ready_sequence(&state).is_none_or(|sequence| sequence >= before_sequence)
467 }
468
469 async fn wait_for_progress_since(&self, completed: u64) {
470 loop {
471 let progress = self.progress.notified();
472 if self.has_progress_since(completed) {
473 return;
474 }
475 progress.await;
476 }
477 }
478
479 fn progress_baseline(&self) -> Option<u64> {
480 let state = self.state.lock();
481 Self::has_work(&state).then_some(state.completed)
482 }
483
484 fn has_progress_since(&self, completed: u64) -> bool {
485 let state = self.state.lock();
486 state.completed != completed
487 || !Self::has_work(&state)
488 || state.phase == PacketPhase::Stopped
489 }
490
491 fn has_work(state: &PacketQueueState<K, T>) -> bool {
492 state.active != 0
493 || state
494 .lanes
495 .values()
496 .any(|lane| lane.pause_generation.is_none() && !lane.queued.is_empty())
497 }
498
499 fn pause_lane(&self, key: K) -> Option<PacketLanePause<K>> {
500 let mut state = self.state.lock();
501 if state.phase == PacketPhase::Stopped {
502 return None;
503 }
504 if state
505 .lanes
506 .get(&key)
507 .is_some_and(|lane| lane.pause_generation.is_some())
508 {
509 return None;
510 }
511
512 let generation = state.next_pause_generation;
513 assert!(
514 generation != u64::MAX,
515 "packet lane pause generation exhausted"
516 );
517 state.next_pause_generation = generation + 1;
518
519 let queued = {
520 let lane = state.lanes.entry(key).or_insert_with(PacketLane::new);
521 lane.pause_generation = Some(generation);
522 lane.queued.drain(..).collect::<Vec<_>>()
523 };
524
525 state.player_local_ready.retain(|entry| entry.0.1 != key);
526 state.serialized_ready.retain(|entry| entry.0.1 != key);
527 state.exclusive_ready.retain(|entry| entry.0.1 != key);
528 let queued_sequences = queued
529 .iter()
530 .map(|packet| packet.sequence)
531 .collect::<FxHashSet<_>>();
532 state
533 .serialized
534 .retain(|entry| !queued_sequences.contains(&entry.0));
535 state
536 .exclusive
537 .retain(|entry| !queued_sequences.contains(&entry.0));
538 let Some(lane) = state.lanes.get_mut(&key) else {
539 panic!("paused packet lane disappeared while deferring its packets");
540 };
541 lane.deferred
542 .extend(queued.into_iter().map(|packet| DeferredPacket {
543 execution: packet.execution,
544 admission_bytes: packet.admission_bytes,
545 value: packet.value,
546 }));
547
548 let is_idle = state.active == 0;
549 let should_wake = matches!(state.phase, PacketPhase::Open | PacketPhase::Draining(_))
550 && Self::next_ready_sequence(&state).is_some();
551 drop(state);
552 self.progress.notify_one();
553 if should_wake {
554 self.work_available.notify_all();
555 }
556 if is_idle {
557 self.idle.notify_one();
558 }
559 Some(PacketLanePause { key, generation })
560 }
561
562 fn resume_lane(&self, pause: PacketLanePause<K>) -> bool {
563 let mut state = self.state.lock();
564 if state.phase == PacketPhase::Stopped {
565 return false;
566 }
567
568 let Some(lane) = state.lanes.get_mut(&pause.key) else {
569 return false;
570 };
571 if lane.pause_generation != Some(pause.generation) {
572 return false;
573 }
574 lane.pause_generation = None;
575 let active = lane.active;
576 let deferred = mem::take(&mut lane.deferred);
577
578 let Ok(deferred_count) = u64::try_from(deferred.len()) else {
579 panic!("deferred packet count does not fit in the sequence space");
580 };
581 let first_sequence = state.next_sequence;
582 let Some(next_sequence) = first_sequence.checked_add(deferred_count) else {
583 panic!("packet submission sequence exhausted while resuming a lane");
584 };
585 state.next_sequence = next_sequence;
586
587 let mut queued = VecDeque::with_capacity(deferred.len());
591 for (sequence, packet) in (first_sequence..next_sequence).zip(deferred) {
592 Self::mark_global_order(&mut state, sequence, packet.execution);
593 queued.push_back(SequencedPacket {
594 sequence,
595 execution: packet.execution,
596 admission_bytes: packet.admission_bytes,
597 value: packet.value,
598 });
599 }
600 let first = queued
601 .front()
602 .map(|packet| (packet.sequence, packet.execution));
603 let Some(lane) = state.lanes.get_mut(&pause.key) else {
604 panic!("resumed packet lane disappeared before packet re-admission");
605 };
606 lane.queued.extend(queued);
607 if !active && let Some((sequence, execution)) = first {
608 Self::mark_ready(&mut state, sequence, pause.key, execution);
609 }
610 if !active && first.is_none() {
611 state.lanes.remove(&pause.key);
612 }
613
614 let should_wake = matches!(state.phase, PacketPhase::Open | PacketPhase::Draining(_))
615 && Self::next_ready_sequence(&state).is_some();
616 drop(state);
617 self.progress.notify_one();
618 if should_wake {
619 self.work_available.notify_all();
620 }
621 true
622 }
623
624 fn discard_lane(&self, key: K) {
625 let mut state = self.state.lock();
626 let (discarded_sequences, remove_lane) = {
627 let Some(lane) = state.lanes.get_mut(&key) else {
628 return;
629 };
630 let discarded_packets = lane.queued.len() + lane.deferred.len();
631 let discarded_bytes = lane
632 .queued
633 .iter()
634 .map(|packet| packet.admission_bytes)
635 .chain(lane.deferred.iter().map(|packet| packet.admission_bytes))
636 .sum::<usize>();
637 let discarded_sequences = lane
638 .queued
639 .iter()
640 .map(|packet| packet.sequence)
641 .collect::<FxHashSet<_>>();
642 lane.queued.clear();
643 lane.deferred.clear();
644 lane.pause_generation = None;
645 assert!(
646 lane.outstanding_packets >= discarded_packets,
647 "session packet admission accounting underflow while discarding"
648 );
649 lane.outstanding_packets -= discarded_packets;
650 assert!(
651 lane.outstanding_bytes >= discarded_bytes,
652 "session byte admission accounting underflow while discarding"
653 );
654 lane.outstanding_bytes -= discarded_bytes;
655 (discarded_sequences, !lane.active)
656 };
657
658 state.player_local_ready.retain(|entry| entry.0.1 != key);
659 state.serialized_ready.retain(|entry| entry.0.1 != key);
660 state.exclusive_ready.retain(|entry| entry.0.1 != key);
661 state
662 .serialized
663 .retain(|entry| !discarded_sequences.contains(&entry.0));
664 state
665 .exclusive
666 .retain(|entry| !discarded_sequences.contains(&entry.0));
667 if remove_lane {
668 state.lanes.remove(&key);
669 }
670
671 let is_idle = state.active == 0;
672 let should_wake = matches!(state.phase, PacketPhase::Open | PacketPhase::Draining(_))
673 && Self::next_ready_sequence(&state).is_some();
674 drop(state);
675 self.progress.notify_one();
676 if should_wake {
677 self.work_available.notify_all();
678 }
679 if is_idle {
680 self.idle.notify_one();
681 }
682 }
683
684 fn stop(&self) {
685 let mut state = self.state.lock();
686 state.phase = PacketPhase::Stopped;
687 state.player_local_ready.clear();
688 state.serialized_ready.clear();
689 state.exclusive_ready.clear();
690 state.serialized.clear();
691 state.exclusive.clear();
692 state.lanes.retain(|_, lane| {
693 let lane_discarded_packets = lane.queued.len() + lane.deferred.len();
694 let lane_discarded_bytes = lane
695 .queued
696 .iter()
697 .map(|packet| packet.admission_bytes)
698 .chain(lane.deferred.iter().map(|packet| packet.admission_bytes))
699 .sum::<usize>();
700 lane.queued.clear();
701 lane.deferred.clear();
702 assert!(
703 lane.outstanding_packets >= lane_discarded_packets,
704 "session packet admission accounting underflow while stopping"
705 );
706 lane.outstanding_packets -= lane_discarded_packets;
707 assert!(
708 lane.outstanding_bytes >= lane_discarded_bytes,
709 "session byte admission accounting underflow while stopping"
710 );
711 lane.outstanding_bytes -= lane_discarded_bytes;
712 lane.pause_generation = None;
713 lane.active
714 });
715 let is_idle = state.active == 0;
716 drop(state);
717 self.work_available.notify_all();
718 self.progress.notify_waiters();
719 if is_idle {
720 self.idle.notify_one();
721 }
722 }
723
724 fn next(&self) -> Option<PacketWork<'_, K, T>> {
725 let mut state = self.state.lock();
726 loop {
727 match state.phase {
728 PacketPhase::Stopped => return None,
729 PacketPhase::Open => {
730 if let Some((key, execution, admission_bytes, value)) =
731 Self::start_next(&mut state, None)
732 {
733 state.active += 1;
734 drop(state);
735 return Some(PacketWork {
736 value: Some(value),
737 key,
738 execution,
739 admission_bytes,
740 queue: self,
741 });
742 }
743 }
744 PacketPhase::Draining(before_sequence) => {
745 if let Some((key, execution, admission_bytes, value)) =
746 Self::start_next(&mut state, Some(before_sequence))
747 {
748 state.active += 1;
749 drop(state);
750 return Some(PacketWork {
751 value: Some(value),
752 key,
753 execution,
754 admission_bytes,
755 queue: self,
756 });
757 }
758 }
759 PacketPhase::Closed => {}
760 }
761 self.work_available.wait(&mut state);
762 }
763 }
764
765 #[cfg(test)]
766 fn try_next(&self) -> Option<PacketWork<'_, K, T>> {
767 let mut state = self.state.lock();
768 let before_sequence = match state.phase {
769 PacketPhase::Open => None,
770 PacketPhase::Draining(before_sequence) => Some(before_sequence),
771 PacketPhase::Closed | PacketPhase::Stopped => return None,
772 };
773 let (key, execution, admission_bytes, value) =
774 Self::start_next(&mut state, before_sequence)?;
775 state.active += 1;
776 drop(state);
777 Some(PacketWork {
778 value: Some(value),
779 key,
780 execution,
781 admission_bytes,
782 queue: self,
783 })
784 }
785
786 fn start_next(
787 state: &mut PacketQueueState<K, T>,
788 before_sequence: Option<u64>,
789 ) -> Option<(K, ScheduledPacketExecution, usize, T)> {
790 let (ready_sequence, key, execution) = Self::select_next(state, before_sequence)?;
791 let Some(lane) = state.lanes.get(&key) else {
792 panic!("ready packet lane disappeared before starting");
793 };
794 assert!(!lane.active, "ready packet lane is already active");
795 let Some(packet) = lane.queued.front() else {
796 panic!("ready packet lane has no queued packet");
797 };
798 assert_eq!(
799 ready_sequence, packet.sequence,
800 "ready packet sequence does not match lane front"
801 );
802 assert_eq!(
803 execution, packet.execution,
804 "packet is registered in the wrong ready queue"
805 );
806
807 match execution {
808 ScheduledPacketExecution::PlayerLocal => {
809 assert!(
810 state.player_local_ready.pop() == Some(Reverse((ready_sequence, key))),
811 "player-local ready queue changed while the queue lock was held"
812 );
813 }
814 ScheduledPacketExecution::Serialized => {
815 assert!(
816 state.serialized_ready.pop() == Some(Reverse((ready_sequence, key))),
817 "serialized ready queue changed while the queue lock was held"
818 );
819 assert_eq!(
820 state.serialized.pop(),
821 Some(Reverse(ready_sequence)),
822 "serialized packet order changed while the queue lock was held"
823 );
824 assert!(
825 !state.serialized_active,
826 "serialized packet started while another was active"
827 );
828 state.serialized_active = true;
829 }
830 ScheduledPacketExecution::Exclusive => {
831 assert!(
832 state.exclusive_ready.pop() == Some(Reverse((ready_sequence, key))),
833 "exclusive ready queue changed while the queue lock was held"
834 );
835 assert_eq!(
836 state.exclusive.pop(),
837 Some(Reverse(ready_sequence)),
838 "exclusive packet barrier changed while the queue lock was held"
839 );
840 state.exclusive_active = true;
841 }
842 }
843
844 let Some(lane) = state.lanes.get_mut(&key) else {
845 panic!("ready packet lane disappeared before removal");
846 };
847 let Some(packet) = lane.queued.pop_front() else {
848 panic!("ready packet lane has no queued packet during removal");
849 };
850 lane.active = true;
851 Some((key, execution, packet.admission_bytes, packet.value))
852 }
853
854 fn select_next(
855 state: &PacketQueueState<K, T>,
856 before_sequence: Option<u64>,
857 ) -> Option<(u64, K, ScheduledPacketExecution)> {
858 if state.exclusive_active {
859 return None;
860 }
861
862 let next_exclusive = state.exclusive.peek().map(|entry| entry.0);
863 let next_serialized = state.serialized.peek().map(|entry| entry.0);
864 let mut selected = None;
865 let mut consider = |entry: Option<&Reverse<(u64, K)>>,
866 execution: ScheduledPacketExecution| {
867 let Some(Reverse((sequence, key))) = entry.copied() else {
868 return;
869 };
870 if before_sequence.is_some_and(|cutoff| sequence >= cutoff)
871 || next_exclusive.is_some_and(|exclusive| exclusive < sequence)
872 {
873 return;
874 }
875 if selected.is_none_or(|(selected_sequence, _, _)| sequence < selected_sequence) {
876 selected = Some((sequence, key, execution));
877 }
878 };
879
880 consider(
881 state.player_local_ready.peek(),
882 ScheduledPacketExecution::PlayerLocal,
883 );
884 if !state.serialized_active
885 && state
886 .serialized_ready
887 .peek()
888 .is_some_and(|entry| Some(entry.0.0) == next_serialized)
889 {
890 consider(
891 state.serialized_ready.peek(),
892 ScheduledPacketExecution::Serialized,
893 );
894 }
895 if state.active == 0
896 && state
897 .exclusive_ready
898 .peek()
899 .is_some_and(|entry| Some(entry.0.0) == next_exclusive)
900 {
901 consider(
902 state.exclusive_ready.peek(),
903 ScheduledPacketExecution::Exclusive,
904 );
905 }
906 selected
907 }
908
909 fn finish_one(&self, key: K, execution: ScheduledPacketExecution, admission_bytes: usize) {
910 let mut state = self.state.lock();
911 assert!(state.active > 0, "packet work accounting underflow");
912 state.active -= 1;
913 state.completed = state.completed.wrapping_add(1);
914 match execution {
915 ScheduledPacketExecution::PlayerLocal => {}
916 ScheduledPacketExecution::Serialized => {
917 assert!(state.serialized_active, "serialized packet is not active");
918 state.serialized_active = false;
919 }
920 ScheduledPacketExecution::Exclusive => {
921 assert!(state.exclusive_active, "exclusive packet is not active");
922 state.exclusive_active = false;
923 }
924 }
925
926 let (next_packet, paused) = {
927 let Some(lane) = state.lanes.get_mut(&key) else {
928 panic!("active packet lane disappeared before completion");
929 };
930 assert!(lane.active, "completed packet lane is not active");
931 lane.active = false;
932 assert!(
933 lane.outstanding_packets > 0,
934 "session packet admission accounting underflow on completion"
935 );
936 lane.outstanding_packets -= 1;
937 assert!(
938 lane.outstanding_bytes >= admission_bytes,
939 "session byte admission accounting underflow on completion"
940 );
941 lane.outstanding_bytes -= admission_bytes;
942 (
943 lane.queued
944 .front()
945 .map(|packet| (packet.sequence, packet.execution)),
946 lane.pause_generation.is_some(),
947 )
948 };
949 if let Some((sequence, next_execution)) = next_packet {
950 if !paused {
951 Self::mark_ready(&mut state, sequence, key, next_execution);
952 }
953 } else if !paused {
954 let Some(lane) = state.lanes.get(&key) else {
955 panic!("completed packet lane disappeared before removal");
956 };
957 assert_eq!(lane.outstanding_packets, 0);
958 assert_eq!(lane.outstanding_bytes, 0);
959 state.lanes.remove(&key);
960 }
961
962 let is_idle = state.active == 0;
963 let should_wake = matches!(state.phase, PacketPhase::Open | PacketPhase::Draining(_))
964 && Self::next_ready_sequence(&state).is_some();
965 drop(state);
966 self.progress.notify_one();
967 if should_wake {
968 match execution {
969 ScheduledPacketExecution::Exclusive => {
970 self.work_available.notify_all();
971 }
972 ScheduledPacketExecution::PlayerLocal | ScheduledPacketExecution::Serialized => {
973 self.work_available.notify_one();
974 }
975 }
976 }
977 if is_idle {
978 self.idle.notify_one();
979 }
980 }
981
982 fn mark_ready(
983 state: &mut PacketQueueState<K, T>,
984 sequence: u64,
985 key: K,
986 execution: ScheduledPacketExecution,
987 ) {
988 let entry = Reverse((sequence, key));
989 match execution {
990 ScheduledPacketExecution::PlayerLocal => state.player_local_ready.push(entry),
991 ScheduledPacketExecution::Serialized => state.serialized_ready.push(entry),
992 ScheduledPacketExecution::Exclusive => state.exclusive_ready.push(entry),
993 }
994 }
995
996 fn mark_global_order(
997 state: &mut PacketQueueState<K, T>,
998 sequence: u64,
999 execution: ScheduledPacketExecution,
1000 ) {
1001 match execution {
1002 ScheduledPacketExecution::PlayerLocal => {}
1003 ScheduledPacketExecution::Serialized => state.serialized.push(Reverse(sequence)),
1004 ScheduledPacketExecution::Exclusive => state.exclusive.push(Reverse(sequence)),
1005 }
1006 }
1007
1008 fn next_ready_sequence(state: &PacketQueueState<K, T>) -> Option<u64> {
1009 [
1010 state.player_local_ready.peek().map(|entry| entry.0.0),
1011 state.serialized_ready.peek().map(|entry| entry.0.0),
1012 state.exclusive_ready.peek().map(|entry| entry.0.0),
1013 ]
1014 .into_iter()
1015 .flatten()
1016 .min()
1017 }
1018}
1019
1020struct PacketWork<'a, K, T>
1021where
1022 K: Copy + Eq + Hash + Ord,
1023{
1024 value: Option<T>,
1025 key: K,
1026 execution: ScheduledPacketExecution,
1027 admission_bytes: usize,
1028 queue: &'a PacketQueue<K, T>,
1029}
1030
1031impl<K, T> PacketWork<'_, K, T>
1032where
1033 K: Copy + Eq + Hash + Ord,
1034{
1035 const fn take(&mut self) -> Option<T> {
1036 self.value.take()
1037 }
1038}
1039
1040impl<K, T> Drop for PacketWork<'_, K, T>
1041where
1042 K: Copy + Eq + Hash + Ord,
1043{
1044 fn drop(&mut self) {
1045 self.queue
1046 .finish_one(self.key, self.execution, self.admission_bytes);
1047 }
1048}
1049
1050#[cfg(test)]
1051mod tests {
1052 use std::{
1053 sync::{Arc, mpsc},
1054 thread,
1055 time::Duration,
1056 };
1057
1058 use tokio::time::timeout;
1059
1060 use crate::{
1061 entity::{Entity as _, LivingEntity as _},
1062 player::{ClientInformation, Player, connection::ScheduledPlayPacket},
1063 test_support::{TestPlayerBuilder, fresh_test_world},
1064 };
1065
1066 use super::{
1067 PACKET_ADMISSION_OVERHEAD, PacketAdmissionError, PacketAdmissionLimits, PacketProcessor,
1068 PacketQueue, PendingPlayPacket, ScheduledPacketExecution,
1069 };
1070
1071 const fn limits(per_player_packets: usize, per_player_bytes: usize) -> PacketAdmissionLimits {
1072 PacketAdmissionLimits {
1073 per_player_packets,
1074 per_player_bytes,
1075 }
1076 }
1077
1078 fn replacement_for(player: &Arc<Player>) -> Arc<Player> {
1079 Arc::new(Player::new(
1080 player.gameprofile.clone(),
1081 Arc::clone(&player.connection),
1082 Arc::clone(&player.session),
1083 player.get_world(),
1084 player.server.clone(),
1085 Arc::clone(&player.config),
1086 player.id(),
1087 ClientInformation::default(),
1088 ))
1089 }
1090
1091 #[test]
1092 fn per_player_packet_limit_rejects_before_assigning_a_sequence() {
1093 let queue = PacketQueue::with_limits(limits(2, 100));
1094 assert_eq!(
1095 queue.try_submit(1, ScheduledPacketExecution::PlayerLocal, 1, "first"),
1096 Ok(())
1097 );
1098 assert_eq!(
1099 queue.try_submit(1, ScheduledPacketExecution::Serialized, 1, "second"),
1100 Ok(())
1101 );
1102
1103 assert_eq!(
1104 queue.try_submit(1, ScheduledPacketExecution::Exclusive, 1, "rejected"),
1105 Err(PacketAdmissionError::PlayerPacketLimit)
1106 );
1107 assert_eq!(queue.state.lock().next_sequence, 2);
1108 assert_eq!(
1109 queue.try_submit(2, ScheduledPacketExecution::PlayerLocal, 1, "other player"),
1110 Ok(())
1111 );
1112 }
1113
1114 #[test]
1115 fn scheduled_respawn_is_retained_if_domain_switch_queues_before_worker_gate() {
1116 let world = fresh_test_world("scheduled_domain_switch_respawn_packet");
1117 let player = TestPlayerBuilder::new(world, "RespawnTester", 1).build();
1118 let packet = ScheduledPlayPacket::perform_respawn_for_test();
1119 let Some(token) = player.begin_pending_world_change() else {
1120 panic!("test player should acquire a world-change token");
1121 };
1122 assert!(player.begin_domain_switch(token));
1123 player.set_health(0.0);
1124
1125 assert!(!PacketProcessor::packet_is_runnable(
1126 &player, &packet, false
1127 ));
1128 assert!(player.has_deferred_death_respawn_for_test());
1129
1130 assert!(player.finish_domain_switch(token));
1131 assert!(player.finish_pending_world_change(token));
1132 }
1133
1134 #[test]
1135 fn queued_packets_resolve_the_player_bound_when_they_start() {
1136 let world = fresh_test_world("packet_session_replacement_resolution");
1137 let original = TestPlayerBuilder::new(world, "Original", 1).build();
1138 let replacement = replacement_for(&original);
1139 let session = Arc::clone(&original.session);
1140 let processor = PacketProcessor::new();
1141 let packet = ScheduledPlayPacket::perform_respawn_for_test();
1142
1143 processor.schedule(Arc::clone(&original), packet, 1);
1144 let Some(transition) = processor.pause_player_session(&session) else {
1145 panic!("session packet lane should pause");
1146 };
1147 assert!(session.replace_player(&original, &replacement));
1148
1149 processor.schedule(original, ScheduledPlayPacket::perform_respawn_for_test(), 1);
1152 assert!(processor.resume_player_session(transition));
1153 processor.queued.open();
1154
1155 for _ in 0..2 {
1156 let Some(mut work) = processor.queued.try_next() else {
1157 panic!("retained session packet should become runnable");
1158 };
1159 let Some(PendingPlayPacket { session, .. }) = work.take() else {
1160 panic!("retained packet should still carry its session");
1161 };
1162 let Some(current) = session.current_player() else {
1163 panic!("replacement should be bound before packet work resumes");
1164 };
1165 assert!(Arc::ptr_eq(¤t, &replacement));
1166 }
1167 }
1168
1169 #[test]
1170 fn per_player_byte_limit_is_independent_of_packet_count() {
1171 let queue = PacketQueue::with_limits(limits(10, 6));
1172 assert_eq!(
1173 queue.try_submit(1, ScheduledPacketExecution::PlayerLocal, 4, "first"),
1174 Ok(())
1175 );
1176
1177 assert_eq!(
1178 queue.try_submit(1, ScheduledPacketExecution::PlayerLocal, 3, "rejected"),
1179 Err(PacketAdmissionError::PlayerByteLimit)
1180 );
1181 assert_eq!(
1182 queue.try_submit(2, ScheduledPacketExecution::PlayerLocal, 6, "other player"),
1183 Ok(())
1184 );
1185 }
1186
1187 #[test]
1188 fn production_limit_allows_a_watchdog_window_of_mounted_client_traffic() {
1189 const VANILLA_WATCHDOG_SECONDS: usize = 60;
1190 const CLIENT_TICKS_PER_SECOND: usize = 20;
1191 const MOUNTED_PACKETS_PER_CLIENT_TICK: usize = 4;
1192 const EXPECTED_PACKETS: usize =
1193 VANILLA_WATCHDOG_SECONDS * CLIENT_TICKS_PER_SECOND * MOUNTED_PACKETS_PER_CLIENT_TICK;
1194
1195 let queue = PacketQueue::new();
1196 for packet in 0..EXPECTED_PACKETS {
1197 assert!(
1198 queue
1199 .try_submit(
1200 1,
1201 ScheduledPacketExecution::PlayerLocal,
1202 PACKET_ADMISSION_OVERHEAD,
1203 packet,
1204 )
1205 .is_ok()
1206 );
1207 }
1208
1209 let state = queue.state.lock();
1210 let lane = state.lanes.get(&1).expect("session lane should exist");
1211 assert_eq!(lane.outstanding_packets, EXPECTED_PACKETS);
1212 }
1213
1214 #[test]
1215 fn independent_session_limits_do_not_accumulate_globally() {
1216 const SESSION_COUNT: usize = 64;
1217 const PACKETS_PER_SESSION: usize = 33;
1218
1219 let queue = PacketQueue::new();
1220 for session_id in 1_u64..=64 {
1221 for packet in 0..PACKETS_PER_SESSION {
1222 assert!(
1223 queue
1224 .try_submit(
1225 session_id,
1226 ScheduledPacketExecution::PlayerLocal,
1227 PACKET_ADMISSION_OVERHEAD,
1228 packet,
1229 )
1230 .is_ok()
1231 );
1232 }
1233 }
1234
1235 let state = queue.state.lock();
1236 assert_eq!(state.lanes.len(), SESSION_COUNT);
1237 assert_eq!(
1238 usize::try_from(state.next_sequence),
1239 Ok(SESSION_COUNT * PACKETS_PER_SESSION)
1240 );
1241 }
1242
1243 #[test]
1244 fn discarding_stale_session_keeps_replacement_session_work() {
1245 let stale_key = 1;
1246 let replacement_key = 2;
1247 let queue = PacketQueue::with_limits(limits(10, 100));
1248 queue.submit(
1249 stale_key,
1250 ScheduledPacketExecution::Exclusive,
1251 "stale barrier",
1252 );
1253 queue.submit(
1254 stale_key,
1255 ScheduledPacketExecution::Serialized,
1256 "stale serialized",
1257 );
1258 queue.submit(
1259 replacement_key,
1260 ScheduledPacketExecution::PlayerLocal,
1261 "replacement",
1262 );
1263
1264 queue.discard_lane(stale_key);
1265
1266 {
1267 let queue_state = queue.state.lock();
1268 assert!(!queue_state.lanes.contains_key(&stale_key));
1269 let replacement_lane = queue_state
1270 .lanes
1271 .get(&replacement_key)
1272 .expect("replacement session lane should remain");
1273 assert_eq!(replacement_lane.outstanding_packets, 1);
1274 assert_eq!(replacement_lane.outstanding_bytes, 1);
1275 assert!(queue_state.serialized.is_empty());
1276 assert!(queue_state.exclusive.is_empty());
1277 }
1278
1279 queue.open();
1280 let Some(mut work) = queue.try_next() else {
1281 panic!("replacement session packet should remain runnable");
1282 };
1283 assert_eq!(work.take(), Some("replacement"));
1284 }
1285
1286 #[test]
1287 fn paused_fifo_is_re_admitted_after_unrelated_global_work() {
1288 let queue = PacketQueue::with_limits(limits(10, 100));
1289 queue.submit(
1290 1,
1291 ScheduledPacketExecution::Serialized,
1292 "deferred serialized",
1293 );
1294 queue.submit(1, ScheduledPacketExecution::Exclusive, "deferred exclusive");
1295 let Some(pause) = queue.pause_lane(1) else {
1296 panic!("session lane should pause");
1297 };
1298
1299 queue.submit(
1300 2,
1301 ScheduledPacketExecution::Serialized,
1302 "unrelated serialized",
1303 );
1304 queue.submit(
1305 3,
1306 ScheduledPacketExecution::Exclusive,
1307 "unrelated exclusive",
1308 );
1309 queue.open();
1310
1311 let Some(mut unrelated_serialized) = queue.try_next() else {
1312 panic!("paused serialized work must not block another session");
1313 };
1314 assert_eq!(unrelated_serialized.take(), Some("unrelated serialized"));
1315 drop(unrelated_serialized);
1316
1317 let Some(mut unrelated_exclusive) = queue.try_next() else {
1318 panic!("paused exclusive work must not remain a global barrier");
1319 };
1320 assert_eq!(unrelated_exclusive.take(), Some("unrelated exclusive"));
1321 drop(unrelated_exclusive);
1322 assert!(queue.try_next().is_none());
1323
1324 assert!(queue.resume_lane(pause));
1325 let Some(mut deferred_serialized) = queue.try_next() else {
1326 panic!("resumed FIFO front should become runnable");
1327 };
1328 assert_eq!(deferred_serialized.take(), Some("deferred serialized"));
1329 drop(deferred_serialized);
1330
1331 let Some(mut deferred_exclusive) = queue.try_next() else {
1332 panic!("resumed FIFO tail should follow its front");
1333 };
1334 assert_eq!(deferred_exclusive.take(), Some("deferred exclusive"));
1335 }
1336
1337 #[test]
1338 fn pausing_active_player_session_defers_its_tail_until_exact_resume() {
1339 let world = fresh_test_world("packet_active_session_pause");
1340 let player = TestPlayerBuilder::new(Arc::clone(&world), "Paused", 1).build();
1341 let unrelated_player = TestPlayerBuilder::new(world, "Unrelated", 2).build();
1342 let processor = PacketProcessor::new();
1343 processor.schedule(
1344 Arc::clone(&player),
1345 ScheduledPlayPacket::perform_respawn_for_test(),
1346 1,
1347 );
1348 processor.schedule(
1349 Arc::clone(&player),
1350 ScheduledPlayPacket::perform_respawn_for_test(),
1351 2,
1352 );
1353 processor.open_after_tick();
1354
1355 let Some(mut active) = processor.queued.try_next() else {
1356 panic!("respawn packet should be active before its lane pauses");
1357 };
1358 assert_eq!(active.admission_bytes, PACKET_ADMISSION_OVERHEAD + 1);
1359 let Some(PendingPlayPacket {
1360 session: active_session,
1361 ..
1362 }) = active.take()
1363 else {
1364 panic!("active respawn should carry its player session");
1365 };
1366 assert!(Arc::ptr_eq(&active_session, &player.session));
1367
1368 let Some(transition) = processor.pause_player_session(&player.session) else {
1369 panic!("active session lane should pause");
1370 };
1371 processor.schedule(
1372 Arc::clone(&player),
1373 ScheduledPlayPacket::perform_respawn_for_test(),
1374 3,
1375 );
1376 processor.schedule(
1377 Arc::clone(&unrelated_player),
1378 ScheduledPlayPacket::perform_respawn_for_test(),
1379 4,
1380 );
1381
1382 let Some(mut unrelated) = processor.queued.try_next() else {
1383 panic!("paused lane must not block unrelated player work");
1384 };
1385 assert_eq!(unrelated.admission_bytes, PACKET_ADMISSION_OVERHEAD + 4);
1386 let Some(PendingPlayPacket {
1387 session: unrelated_session,
1388 ..
1389 }) = unrelated.take()
1390 else {
1391 panic!("unrelated packet should carry its player session");
1392 };
1393 assert!(Arc::ptr_eq(&unrelated_session, &unrelated_player.session));
1394 drop(unrelated);
1395 assert!(processor.queued.try_next().is_none());
1396
1397 assert!(processor.resume_player_session(transition));
1398 assert!(processor.queued.try_next().is_none());
1399 drop(active);
1400
1401 let Some(mut queued_before_pause) = processor.queued.try_next() else {
1402 panic!("exact pause owner should re-admit the original lane tail");
1403 };
1404 assert_eq!(
1405 queued_before_pause.admission_bytes,
1406 PACKET_ADMISSION_OVERHEAD + 2
1407 );
1408 let Some(PendingPlayPacket {
1409 session: queued_before_pause_session,
1410 ..
1411 }) = queued_before_pause.take()
1412 else {
1413 panic!("packet queued before the pause should retain its session");
1414 };
1415 assert!(Arc::ptr_eq(&queued_before_pause_session, &player.session));
1416 drop(queued_before_pause);
1417
1418 let Some(mut queued_during_pause) = processor.queued.try_next() else {
1419 panic!("packets submitted while paused should retain session FIFO order");
1420 };
1421 assert_eq!(
1422 queued_during_pause.admission_bytes,
1423 PACKET_ADMISSION_OVERHEAD + 3
1424 );
1425 let Some(PendingPlayPacket {
1426 session: queued_during_pause_session,
1427 ..
1428 }) = queued_during_pause.take()
1429 else {
1430 panic!("packet queued during the pause should retain its session");
1431 };
1432 assert!(Arc::ptr_eq(&queued_during_pause_session, &player.session));
1433 }
1434
1435 #[test]
1436 fn stale_pause_token_cannot_resume_a_later_pause() {
1437 let queue = PacketQueue::with_limits(limits(10, 100));
1438 let Some(first_pause) = queue.pause_lane(1) else {
1439 panic!("first lane pause should succeed");
1440 };
1441 assert!(queue.resume_lane(first_pause));
1442
1443 let Some(second_pause) = queue.pause_lane(1) else {
1444 panic!("second lane pause should succeed");
1445 };
1446 queue.submit(1, ScheduledPacketExecution::PlayerLocal, "deferred");
1447 queue.open();
1448
1449 assert!(!queue.resume_lane(first_pause));
1450 assert!(queue.try_next().is_none());
1451 assert!(queue.resume_lane(second_pause));
1452
1453 let Some(mut resumed) = queue.try_next() else {
1454 panic!("current pause owner should resume the lane");
1455 };
1456 assert_eq!(resumed.take(), Some("deferred"));
1457 }
1458
1459 #[tokio::test]
1460 async fn paused_packets_do_not_hold_tick_drain_open() {
1461 let queue = PacketQueue::new();
1462 queue.submit(1, ScheduledPacketExecution::Exclusive, "deferred");
1463 let Some(pause) = queue.pause_lane(1) else {
1464 panic!("session lane should pause");
1465 };
1466 queue.open();
1467
1468 assert!(
1469 timeout(Duration::from_secs(1), queue.drain_for_tick())
1470 .await
1471 .is_ok()
1472 );
1473 assert!(queue.resume_lane(pause));
1474 queue.open();
1475
1476 let Some(mut resumed) = queue.try_next() else {
1477 panic!("deferred packet should remain retained after the tick drain");
1478 };
1479 assert_eq!(resumed.take(), Some("deferred"));
1480 }
1481
1482 #[test]
1483 fn active_packet_remains_charged_until_completion() {
1484 let queue = PacketQueue::with_limits(limits(1, 4));
1485 assert_eq!(
1486 queue.try_submit(1, ScheduledPacketExecution::PlayerLocal, 4, "active"),
1487 Ok(())
1488 );
1489 queue.open();
1490 let Some(work) = queue.try_next() else {
1491 panic!("accepted packet should start");
1492 };
1493
1494 assert_eq!(
1495 queue.try_submit(1, ScheduledPacketExecution::PlayerLocal, 1, "rejected"),
1496 Err(PacketAdmissionError::PlayerPacketLimit)
1497 );
1498 {
1499 let state = queue.state.lock();
1500 let lane = state.lanes.get(&1).expect("active lane should remain");
1501 assert_eq!(lane.outstanding_packets, 1);
1502 assert_eq!(lane.outstanding_bytes, 4);
1503 }
1504
1505 drop(work);
1506 {
1507 let state = queue.state.lock();
1508 assert!(state.lanes.is_empty());
1509 }
1510 assert_eq!(
1511 queue.try_submit(1, ScheduledPacketExecution::PlayerLocal, 4, "next"),
1512 Ok(())
1513 );
1514 }
1515
1516 #[test]
1517 fn discarding_lane_removes_hidden_barriers_and_keeps_active_work_charged() {
1518 let queue = PacketQueue::with_limits(limits(10, 100));
1519 queue.submit(1, ScheduledPacketExecution::PlayerLocal, "active attacker");
1520 queue.submit(1, ScheduledPacketExecution::Serialized, "hidden serialized");
1521 queue.submit(1, ScheduledPacketExecution::Exclusive, "hidden exclusive");
1522 queue.submit(2, ScheduledPacketExecution::PlayerLocal, "other player");
1523 queue.open();
1524 let Some(mut active) = queue.try_next() else {
1525 panic!("attacker's first packet should start");
1526 };
1527 assert_eq!(active.take(), Some("active attacker"));
1528
1529 queue.discard_lane(1);
1530
1531 {
1532 let state = queue.state.lock();
1533 assert!(state.serialized.is_empty());
1534 assert!(state.exclusive.is_empty());
1535 assert!(state.serialized_ready.is_empty());
1536 assert!(state.exclusive_ready.is_empty());
1537 let Some(lane) = state.lanes.get(&1) else {
1538 panic!("active attacker lane should remain");
1539 };
1540 assert!(lane.active);
1541 assert!(lane.queued.is_empty());
1542 assert_eq!(lane.outstanding_packets, 1);
1543 assert_eq!(lane.outstanding_bytes, 1);
1544 }
1545
1546 let Some(mut other) = queue.try_next() else {
1547 panic!("purged barriers should no longer block another player");
1548 };
1549 assert_eq!(other.take(), Some("other player"));
1550 drop(other);
1551 drop(active);
1552
1553 let state = queue.state.lock();
1554 assert!(state.lanes.is_empty());
1555 }
1556
1557 #[test]
1558 fn discarding_idle_lane_removes_ready_exclusive_barrier() {
1559 let queue = PacketQueue::with_limits(limits(10, 100));
1560 queue.submit(1, ScheduledPacketExecution::Exclusive, "attacker barrier");
1561 queue.submit(1, ScheduledPacketExecution::Serialized, "hidden serialized");
1562 queue.submit(2, ScheduledPacketExecution::PlayerLocal, "other player");
1563
1564 queue.discard_lane(1);
1565
1566 {
1567 let state = queue.state.lock();
1568 assert!(!state.lanes.contains_key(&1));
1569 let other = state.lanes.get(&2).expect("other lane should remain");
1570 assert_eq!(other.outstanding_packets, 1);
1571 assert_eq!(other.outstanding_bytes, 1);
1572 assert!(state.exclusive_ready.is_empty());
1573 assert!(state.exclusive.is_empty());
1574 assert!(state.serialized.is_empty());
1575 }
1576
1577 queue.open();
1578 let Some(mut other) = queue.try_next() else {
1579 panic!("discarded exclusive barrier should not block another player");
1580 };
1581 assert_eq!(other.take(), Some("other player"));
1582 }
1583
1584 #[test]
1585 fn queued_packets_start_in_submission_order_when_opened() {
1586 let queue = PacketQueue::new();
1587 queue.submit(1, ScheduledPacketExecution::PlayerLocal, 1);
1588 queue.submit(2, ScheduledPacketExecution::PlayerLocal, 2);
1589 queue.submit(1, ScheduledPacketExecution::PlayerLocal, 3);
1590 assert!(queue.try_next().is_none());
1591
1592 queue.open();
1593 let mut processed = Vec::new();
1594 while let Some(mut work) = queue.try_next() {
1595 if let Some(value) = work.take() {
1596 processed.push(value);
1597 }
1598 }
1599
1600 assert_eq!(processed, [1, 2, 3]);
1601 }
1602
1603 #[test]
1604 fn packet_lane_only_allows_one_active_handler() {
1605 let queue = PacketQueue::new();
1606 queue.submit(
1607 1,
1608 ScheduledPacketExecution::PlayerLocal,
1609 "first player packet",
1610 );
1611 queue.submit(
1612 1,
1613 ScheduledPacketExecution::PlayerLocal,
1614 "second player packet",
1615 );
1616 queue.submit(
1617 2,
1618 ScheduledPacketExecution::PlayerLocal,
1619 "other player packet",
1620 );
1621 queue.open();
1622
1623 let Some(mut first) = queue.try_next() else {
1624 panic!("first packet should start");
1625 };
1626 assert_eq!(first.take(), Some("first player packet"));
1627
1628 let Some(mut other_player) = queue.try_next() else {
1629 panic!("another player's packet should be able to start");
1630 };
1631 assert_eq!(other_player.take(), Some("other player packet"));
1632 assert!(queue.try_next().is_none());
1633
1634 drop(first);
1635 let Some(mut second) = queue.try_next() else {
1636 panic!("the next packet should start after its lane becomes idle");
1637 };
1638 assert_eq!(second.take(), Some("second player packet"));
1639 }
1640
1641 #[test]
1642 fn serialized_packet_overlaps_player_local_work_but_not_another_serialized_packet() {
1643 let queue = PacketQueue::new();
1644 queue.submit(1, ScheduledPacketExecution::PlayerLocal, "first local");
1645 queue.submit(2, ScheduledPacketExecution::Serialized, "first serialized");
1646 queue.submit(3, ScheduledPacketExecution::Serialized, "second serialized");
1647 queue.submit(4, ScheduledPacketExecution::PlayerLocal, "later local");
1648 queue.open();
1649
1650 let Some(mut first_local) = queue.try_next() else {
1651 panic!("first player-local packet should start");
1652 };
1653 assert_eq!(first_local.take(), Some("first local"));
1654
1655 let Some(mut first_serialized) = queue.try_next() else {
1656 panic!("serialized packet should overlap player-local work");
1657 };
1658 assert_eq!(first_serialized.take(), Some("first serialized"));
1659
1660 let Some(mut later_local) = queue.try_next() else {
1661 panic!("player-local work should bypass a blocked serialized packet");
1662 };
1663 assert_eq!(later_local.take(), Some("later local"));
1664 assert!(queue.try_next().is_none());
1665
1666 drop(first_serialized);
1667 let Some(mut second_serialized) = queue.try_next() else {
1668 panic!("next serialized packet should start after its predecessor finishes");
1669 };
1670 assert_eq!(second_serialized.take(), Some("second serialized"));
1671 }
1672
1673 #[test]
1674 fn serialized_packet_hidden_in_an_active_lane_preserves_serialized_order() {
1675 let queue = PacketQueue::new();
1676 queue.submit(1, ScheduledPacketExecution::PlayerLocal, "active lane");
1677 queue.submit(1, ScheduledPacketExecution::Serialized, "hidden serialized");
1678 queue.submit(2, ScheduledPacketExecution::Serialized, "later serialized");
1679 queue.submit(
1680 3,
1681 ScheduledPacketExecution::PlayerLocal,
1682 "independent local",
1683 );
1684 queue.open();
1685
1686 let Some(active_lane) = queue.try_next() else {
1687 panic!("first lane packet should start");
1688 };
1689 let Some(mut independent_local) = queue.try_next() else {
1690 panic!("player-local work should bypass serialized ordering contention");
1691 };
1692 assert_eq!(independent_local.take(), Some("independent local"));
1693 assert!(queue.try_next().is_none());
1694
1695 drop(active_lane);
1696 let Some(mut hidden_serialized) = queue.try_next() else {
1697 panic!("earliest serialized packet should start when its lane becomes idle");
1698 };
1699 assert_eq!(hidden_serialized.take(), Some("hidden serialized"));
1700 assert!(queue.try_next().is_none());
1701
1702 drop(hidden_serialized);
1703 let Some(mut later_serialized) = queue.try_next() else {
1704 panic!("later serialized packet should preserve global submission order");
1705 };
1706 assert_eq!(later_serialized.take(), Some("later serialized"));
1707 }
1708
1709 #[test]
1710 fn blocking_worker_bypasses_active_serialized_work_for_player_local_work() {
1711 let queue = Arc::new(PacketQueue::new());
1712 queue.submit(1, ScheduledPacketExecution::Serialized, "active serialized");
1713 queue.submit(2, ScheduledPacketExecution::Serialized, "queued serialized");
1714 queue.submit(3, ScheduledPacketExecution::PlayerLocal, "player local");
1715 queue.open();
1716
1717 let Some(active_serialized) = queue.try_next() else {
1718 panic!("first serialized packet should start");
1719 };
1720 let worker_queue = Arc::clone(&queue);
1721 let (sender, receiver) = mpsc::channel();
1722 let worker = thread::spawn(move || {
1723 for _ in 0..2 {
1724 let Some(mut work) = worker_queue.next() else {
1725 return;
1726 };
1727 if let Some(value) = work.take() {
1728 let _ = sender.send(value);
1729 }
1730 }
1731 });
1732
1733 assert_eq!(
1734 receiver.recv_timeout(Duration::from_secs(1)),
1735 Ok("player local")
1736 );
1737 assert!(receiver.recv_timeout(Duration::from_millis(10)).is_err());
1738
1739 drop(active_serialized);
1740 assert_eq!(
1741 receiver.recv_timeout(Duration::from_secs(1)),
1742 Ok("queued serialized")
1743 );
1744 assert!(worker.join().is_ok());
1745 }
1746
1747 #[test]
1748 fn exclusive_packet_waits_for_active_work_and_blocks_later_packets() {
1749 let queue = PacketQueue::new();
1750 queue.submit(1, ScheduledPacketExecution::PlayerLocal, "before barrier");
1751 queue.submit(2, ScheduledPacketExecution::Exclusive, "barrier");
1752 queue.submit(3, ScheduledPacketExecution::PlayerLocal, "after barrier");
1753 queue.open();
1754
1755 let Some(mut before) = queue.try_next() else {
1756 panic!("packet before the barrier should start");
1757 };
1758 assert_eq!(before.take(), Some("before barrier"));
1759 assert!(queue.try_next().is_none());
1760
1761 drop(before);
1762 let Some(mut barrier) = queue.try_next() else {
1763 panic!("exclusive packet should start after active work finishes");
1764 };
1765 assert_eq!(barrier.take(), Some("barrier"));
1766 assert!(queue.try_next().is_none());
1767
1768 drop(barrier);
1769 let Some(mut after) = queue.try_next() else {
1770 panic!("packet after the barrier should start after it finishes");
1771 };
1772 assert_eq!(after.take(), Some("after barrier"));
1773 }
1774
1775 #[test]
1776 fn exclusive_packet_waits_for_serialized_work_and_blocks_player_local_work() {
1777 let queue = PacketQueue::new();
1778 queue.submit(1, ScheduledPacketExecution::Serialized, "before barrier");
1779 queue.submit(2, ScheduledPacketExecution::Exclusive, "barrier");
1780 queue.submit(3, ScheduledPacketExecution::PlayerLocal, "after barrier");
1781 queue.open();
1782
1783 let Some(mut before) = queue.try_next() else {
1784 panic!("serialized packet before the barrier should start");
1785 };
1786 assert_eq!(before.take(), Some("before barrier"));
1787 assert!(queue.try_next().is_none());
1788
1789 drop(before);
1790 let Some(mut barrier) = queue.try_next() else {
1791 panic!("exclusive packet should wait for serialized work");
1792 };
1793 assert_eq!(barrier.take(), Some("barrier"));
1794 assert!(queue.try_next().is_none());
1795
1796 drop(barrier);
1797 let Some(mut after) = queue.try_next() else {
1798 panic!("player-local packet should wait for the exclusive barrier");
1799 };
1800 assert_eq!(after.take(), Some("after barrier"));
1801 }
1802
1803 #[test]
1804 fn exclusive_packet_hidden_in_an_active_lane_still_blocks_later_lanes() {
1805 let queue = PacketQueue::new();
1806 queue.submit(1, ScheduledPacketExecution::PlayerLocal, "active");
1807 queue.submit(1, ScheduledPacketExecution::Exclusive, "barrier");
1808 queue.submit(2, ScheduledPacketExecution::PlayerLocal, "later lane");
1809 queue.open();
1810
1811 let Some(active) = queue.try_next() else {
1812 panic!("first packet should start");
1813 };
1814 assert!(queue.try_next().is_none());
1815
1816 drop(active);
1817 let Some(mut barrier) = queue.try_next() else {
1818 panic!("hidden exclusive packet should become runnable");
1819 };
1820 assert_eq!(barrier.take(), Some("barrier"));
1821 }
1822
1823 #[test]
1824 fn closed_phase_retains_new_packets_for_the_next_open_phase() {
1825 let queue = PacketQueue::new();
1826 queue.open();
1827 queue.close();
1828 queue.submit(1, ScheduledPacketExecution::PlayerLocal, 1);
1829
1830 assert!(queue.try_next().is_none());
1831 queue.open();
1832 let Some(mut work) = queue.try_next() else {
1833 panic!("queued packet should become available when the packet phase opens");
1834 };
1835 assert_eq!(work.take(), Some(1));
1836 }
1837
1838 #[test]
1839 fn tick_drain_only_processes_packets_submitted_before_its_cutoff() {
1840 let queue = PacketQueue::new();
1841 queue.submit(1, ScheduledPacketExecution::PlayerLocal, "before cutoff");
1842 queue.open();
1843 let Some((before_sequence, should_wake)) = queue.begin_tick_drain() else {
1844 panic!("running queue should begin a tick drain");
1845 };
1846 assert!(should_wake);
1847 queue.submit(2, ScheduledPacketExecution::PlayerLocal, "after cutoff");
1848
1849 let Some(mut before) = queue.try_next() else {
1850 panic!("packet before the cutoff should drain");
1851 };
1852 assert_eq!(before.take(), Some("before cutoff"));
1853 drop(before);
1854
1855 assert!(queue.try_next().is_none());
1856 assert!(queue.tick_drain_complete(before_sequence));
1857 queue.finish_tick_drain(before_sequence);
1858 queue.open();
1859 let Some(mut after) = queue.try_next() else {
1860 panic!("packet after the cutoff should wait for the next open phase");
1861 };
1862 assert_eq!(after.take(), Some("after cutoff"));
1863 }
1864
1865 #[test]
1866 fn tick_drain_orders_hidden_serialized_and_exclusive_work_before_its_cutoff() {
1867 let queue = PacketQueue::new();
1868 queue.submit(1, ScheduledPacketExecution::PlayerLocal, "active local");
1869 queue.submit(1, ScheduledPacketExecution::Serialized, "hidden serialized");
1870 queue.submit(2, ScheduledPacketExecution::Exclusive, "exclusive");
1871 queue.open();
1872
1873 let Some(active_local) = queue.try_next() else {
1874 panic!("player-local packet should start before draining");
1875 };
1876 let Some((before_sequence, should_wake)) = queue.begin_tick_drain() else {
1877 panic!("running queue should begin a tick drain");
1878 };
1879 assert!(!should_wake);
1880 queue.submit(
1881 3,
1882 ScheduledPacketExecution::Serialized,
1883 "post-cutoff serialized",
1884 );
1885 queue.submit(
1886 4,
1887 ScheduledPacketExecution::PlayerLocal,
1888 "post-cutoff local",
1889 );
1890
1891 assert!(!queue.tick_drain_complete(before_sequence));
1892 assert!(queue.try_next().is_none());
1893 drop(active_local);
1894
1895 let Some(mut hidden_serialized) = queue.try_next() else {
1896 panic!("hidden pre-cutoff serialized packet should drain");
1897 };
1898 assert_eq!(hidden_serialized.take(), Some("hidden serialized"));
1899 assert!(!queue.tick_drain_complete(before_sequence));
1900 drop(hidden_serialized);
1901
1902 let Some(mut exclusive) = queue.try_next() else {
1903 panic!("pre-cutoff exclusive packet should drain after earlier work");
1904 };
1905 assert_eq!(exclusive.take(), Some("exclusive"));
1906 assert!(!queue.tick_drain_complete(before_sequence));
1907 drop(exclusive);
1908
1909 assert!(queue.try_next().is_none());
1910 assert!(queue.tick_drain_complete(before_sequence));
1911 queue.finish_tick_drain(before_sequence);
1912 queue.open();
1913
1914 let Some(mut post_cutoff_serialized) = queue.try_next() else {
1915 panic!("post-cutoff serialized packet should remain for the next phase");
1916 };
1917 assert_eq!(
1918 post_cutoff_serialized.take(),
1919 Some("post-cutoff serialized")
1920 );
1921 let Some(mut post_cutoff_local) = queue.try_next() else {
1922 panic!("post-cutoff player-local packet should remain for the next phase");
1923 };
1924 assert_eq!(post_cutoff_local.take(), Some("post-cutoff local"));
1925 }
1926
1927 #[tokio::test]
1928 async fn tick_drain_waits_for_active_packet_work() {
1929 let queue = Arc::new(PacketQueue::new());
1930 queue.open();
1931 queue.submit(1, ScheduledPacketExecution::PlayerLocal, 1);
1932 let Some(work) = queue.try_next() else {
1933 panic!("open packet phase should start queued work");
1934 };
1935 let drain = queue.drain_for_tick();
1936 tokio::pin!(drain);
1937
1938 assert!(
1939 timeout(Duration::from_millis(10), drain.as_mut())
1940 .await
1941 .is_err()
1942 );
1943 drop(work);
1944 assert!(
1945 timeout(Duration::from_secs(1), drain.as_mut())
1946 .await
1947 .is_ok()
1948 );
1949 }
1950
1951 #[tokio::test]
1952 async fn overload_progress_waits_for_one_active_packet_to_finish() {
1953 let queue = Arc::new(PacketQueue::new());
1954 queue.open();
1955 queue.submit(1, ScheduledPacketExecution::PlayerLocal, 1);
1956 let Some(work) = queue.try_next() else {
1957 panic!("open packet phase should start queued work");
1958 };
1959 let Some(completed) = queue.progress_baseline() else {
1960 panic!("active packet work should require progress");
1961 };
1962 let progress = queue.wait_for_progress_since(completed);
1963 tokio::pin!(progress);
1964
1965 assert!(
1966 timeout(Duration::from_millis(10), progress.as_mut())
1967 .await
1968 .is_err()
1969 );
1970 drop(work);
1971 assert!(
1972 timeout(Duration::from_secs(1), progress.as_mut())
1973 .await
1974 .is_ok()
1975 );
1976 }
1977
1978 #[test]
1979 fn stopped_queue_discards_pending_and_future_work() {
1980 let queue = PacketQueue::new();
1981 queue.open();
1982 queue.submit(1, ScheduledPacketExecution::PlayerLocal, 1);
1983 queue.stop();
1984 queue.submit(1, ScheduledPacketExecution::PlayerLocal, 2);
1985
1986 assert!(queue.try_next().is_none());
1987 let state = queue.state.lock();
1988 assert!(state.lanes.is_empty());
1989 }
1990
1991 #[test]
1992 fn stopping_with_active_work_keeps_completion_accounting_valid() {
1993 let queue = PacketQueue::new();
1994 queue.submit(1, ScheduledPacketExecution::Exclusive, 1);
1995 queue.submit(1, ScheduledPacketExecution::PlayerLocal, 2);
1996 queue.open();
1997 let Some(work) = queue.try_next() else {
1998 panic!("open packet phase should start queued work");
1999 };
2000
2001 queue.stop();
2002 {
2003 let state = queue.state.lock();
2004 let lane = state.lanes.get(&1).expect("active lane should remain");
2005 assert!(lane.active);
2006 assert!(lane.queued.is_empty());
2007 assert_eq!(lane.outstanding_packets, 1);
2008 assert_eq!(lane.outstanding_bytes, 1);
2009 }
2010 drop(work);
2011
2012 assert!(queue.try_next().is_none());
2013 let state = queue.state.lock();
2014 assert_eq!(state.active, 0);
2015 assert!(state.lanes.is_empty());
2016 }
2017
2018 #[tokio::test]
2019 async fn stopping_with_active_serialized_and_local_work_clears_all_queue_state() {
2020 let queue = PacketQueue::new();
2021 queue.submit(1, ScheduledPacketExecution::Serialized, "active serialized");
2022 queue.submit(2, ScheduledPacketExecution::PlayerLocal, "active local");
2023 queue.submit(3, ScheduledPacketExecution::Serialized, "queued serialized");
2024 queue.submit(4, ScheduledPacketExecution::Exclusive, "queued exclusive");
2025 queue.open();
2026
2027 let Some(active_serialized) = queue.try_next() else {
2028 panic!("serialized packet should start");
2029 };
2030 let Some(active_local) = queue.try_next() else {
2031 panic!("player-local packet should overlap serialized work");
2032 };
2033 queue.stop();
2034 queue.drain_for_tick().await;
2035
2036 {
2037 let state = queue.state.lock();
2038 assert_eq!(state.active, 2);
2039 assert!(state.serialized_active);
2040 assert!(state.player_local_ready.is_empty());
2041 assert!(state.serialized_ready.is_empty());
2042 assert!(state.exclusive_ready.is_empty());
2043 assert!(state.serialized.is_empty());
2044 assert!(state.exclusive.is_empty());
2045 assert_eq!(state.lanes.len(), 2);
2046 for key in [1, 2] {
2047 let lane = state.lanes.get(&key).expect("active lane should remain");
2048 assert!(lane.active);
2049 assert!(lane.queued.is_empty());
2050 assert_eq!(lane.outstanding_packets, 1);
2051 assert_eq!(lane.outstanding_bytes, 1);
2052 }
2053 }
2054
2055 drop(active_serialized);
2056 drop(active_local);
2057
2058 let state = queue.state.lock();
2059 assert_eq!(state.active, 0);
2060 assert!(!state.serialized_active);
2061 assert!(state.lanes.is_empty());
2062 }
2063
2064 #[test]
2065 fn blocking_worker_only_starts_work_during_the_open_phase() {
2066 let queue = Arc::new(PacketQueue::new());
2067 let worker_queue = Arc::clone(&queue);
2068 let (sender, receiver) = mpsc::channel();
2069 let worker = thread::spawn(move || {
2070 while let Some(mut work) = worker_queue.next() {
2071 if let Some(value) = work.take() {
2072 let _ = sender.send(value);
2073 }
2074 }
2075 });
2076
2077 queue.submit(1, ScheduledPacketExecution::PlayerLocal, 1);
2078 assert!(receiver.recv_timeout(Duration::from_millis(10)).is_err());
2079 queue.open();
2080 assert_eq!(receiver.recv_timeout(Duration::from_secs(1)), Ok(1));
2081 queue.close();
2082 queue.submit(1, ScheduledPacketExecution::PlayerLocal, 2);
2083 assert!(receiver.recv_timeout(Duration::from_millis(10)).is_err());
2084 queue.stop();
2085
2086 assert!(worker.join().is_ok());
2087 }
2088}