Skip to main content

steel_core/server/
packet_processor.rs

1//! Serverbound gameplay packet scheduling between game ticks.
2
3use 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
28// Per-session count limits bound tick-drain work; byte limits bound retained decoded payloads.
29// Vanilla's optional connection rate limit is disabled by default. This higher safety ceiling is
30// intended to catch unbounded backlogs without treating ordinary traffic during a long tick as spam.
31const MAX_OUTSTANDING_PACKETS_PER_PLAYER: usize = 8_192;
32const MAX_OUTSTANDING_BYTES_PER_PLAYER: usize = 32 * 1024 * 1024;
33// Approximate fixed queue, lane, and scheduling-index storage per admitted packet.
34const 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
56/// Gameplay packets submitted by network tasks.
57///
58/// The processor runs while the game tick is idle. At each tick boundary it drains every packet
59/// submitted before that boundary, while retaining later submissions for the next packet phase.
60pub(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    /// Runs the blocking packet worker until the processor is stopped.
113    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    /// Opens the inter-tick packet phase and wakes the worker.
144    pub(super) fn open_after_tick(&self) {
145        self.queued.open();
146    }
147
148    /// Guarantees packet progress when a late tick leaves no normal inter-tick window.
149    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    /// Drains packets submitted before this tick boundary, then closes packet admission.
158    pub(super) async fn close_for_tick(&self) {
159        self.queued.drain_for_tick().await;
160    }
161
162    /// Stops the packet worker and discards queued work during server shutdown.
163    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
242/// Bounded ordered multi-producer lanes with a finite snapshot drain at each tick boundary.
243struct 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}