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    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
27// Per-session count limits bound tick-drain work; byte limits bound retained decoded payloads.
28// Vanilla's optional connection rate limit is disabled by default. This higher safety ceiling is
29// intended to catch unbounded backlogs without treating ordinary traffic during a long tick as spam.
30const MAX_OUTSTANDING_PACKETS_PER_PLAYER: usize = 8_192;
31const MAX_OUTSTANDING_BYTES_PER_PLAYER: usize = 32 * 1024 * 1024;
32// Approximate fixed queue, lane, and scheduling-index storage per admitted packet.
33const 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/// Exact ownership of one paused player-session packet lane.
45#[derive(Clone, Copy)]
46pub(crate) struct PlayerPacketTransition(PacketLanePause<PlayerPacketLaneKey>);
47
48struct PendingPlayPacket {
49    session: Arc<PlayerSession>,
50    packet: ScheduledPlayPacket,
51}
52
53/// Gameplay packets submitted by network tasks.
54///
55/// The processor runs while the game tick is idle. At each tick boundary it drains every packet
56/// submitted before that boundary, while retaining later submissions for the next packet phase.
57/// A player-replacement transition temporarily withdraws that session's queued packets from the
58/// runnable global order and re-admits them in session FIFO order after the replacement is bound.
59pub(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    /// Runs the blocking packet worker until the processor is stopped.
112    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    /// Suspends one connection's FIFO lane while its player entity is replaced.
133    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    /// Re-admits packets retained during player replacement into the runnable global order.
143    pub(super) fn resume_player_session(&self, transition: PlayerPacketTransition) -> bool {
144        self.queued.resume_lane(transition.0)
145    }
146
147    /// Drops all retained packets for a connection that can no longer own a player.
148    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    /// Opens the inter-tick packet phase and wakes the worker.
166    pub(super) fn open_after_tick(&self) {
167        self.queued.open();
168    }
169
170    /// Guarantees packet progress when a late tick leaves no normal inter-tick window.
171    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    /// Drains packets submitted before this tick boundary, then closes packet admission.
180    pub(super) async fn close_for_tick(&self) {
181        self.queued.drain_for_tick().await;
182    }
183
184    /// Stops the packet worker and discards queued work during server shutdown.
185    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
275/// Bounded ordered multi-producer lanes with a finite snapshot drain at each tick boundary.
276struct 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        // Deferred packets deliberately receive fresh sequence numbers. Keeping their original
588        // numbers would let a resumed serialized/exclusive packet jump ahead of unrelated global
589        // work that was allowed to run while this lane was paused.
590        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        // The listener may have captured the old player immediately before the bind. Scheduling
1150        // through that stale snapshot must still target the stable session and its replacement.
1151        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(&current, &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}