Skip to main content

steel_core/server/
player_admission.rs

1use super::{
2    Arc, CPlayerInfoUpdate, CRemovePlayerInfo, ClientPacket, ConnectionProtocol, DomainPlayerState,
3    EncodedPacket, Entity, GlobalPlayerData, Instant, JoinSet, NetworkConnection,
4    PendingWorldChangeToken, PersistentPlayerData, Player, ResetReason, SegQueue, Server,
5    SyncMutex, Uuid, mpsc,
6};
7
8pub(super) struct PendingPlayerJoin {
9    pub(super) player: Arc<Player>,
10    pub(super) state: Result<DomainPlayerState, String>,
11}
12
13pub(super) struct PendingPlayerDisconnect {
14    player: Arc<Player>,
15    domain: String,
16    player_data: Arc<PersistentPlayerData>,
17}
18
19#[derive(Clone, Copy, Debug, Eq, PartialEq)]
20pub(super) enum PlayerAdmissionState {
21    Joining,
22    Relocating,
23    Disconnecting,
24}
25
26pub(super) struct PlayerJoinQueue {
27    sender: mpsc::Sender<PendingPlayerJoin>,
28    receiver: SyncMutex<mpsc::Receiver<PendingPlayerJoin>>,
29}
30
31impl PlayerJoinQueue {
32    pub(super) fn new() -> Self {
33        let (sender, receiver) = mpsc::channel();
34        Self {
35            sender,
36            receiver: SyncMutex::new(receiver),
37        }
38    }
39
40    fn send(&self, join: PendingPlayerJoin) {
41        let _ = self.sender.send(join);
42    }
43
44    fn drain(&self) -> Vec<PendingPlayerJoin> {
45        let receiver = self.receiver.lock();
46        let mut joins = Vec::new();
47        while let Ok(join) = receiver.try_recv() {
48            joins.push(join);
49        }
50        joins
51    }
52}
53
54pub(super) struct PlayerDisconnectQueue {
55    queued: SegQueue<Arc<Player>>,
56    prepared: SegQueue<PendingPlayerDisconnect>,
57}
58
59impl PlayerDisconnectQueue {
60    pub(super) const fn new() -> Self {
61        Self {
62            queued: SegQueue::new(),
63            prepared: SegQueue::new(),
64        }
65    }
66
67    fn send(&self, player: Arc<Player>) {
68        self.queued.push(player);
69    }
70
71    fn pop(&self) -> Option<Arc<Player>> {
72        self.queued.pop()
73    }
74
75    fn send_prepared(&self, disconnect: PendingPlayerDisconnect) {
76        self.prepared.push(disconnect);
77    }
78
79    fn drain_prepared(&self) -> Vec<PendingPlayerDisconnect> {
80        let mut prepared = Vec::new();
81        while let Some(disconnect) = self.prepared.pop() {
82            prepared.push(disconnect);
83        }
84        prepared
85    }
86
87    pub(super) fn clear(&self) {
88        while self.queued.pop().is_some() {}
89        while self.prepared.pop().is_some() {}
90    }
91}
92
93impl Server {
94    /// Queues initial player join work.
95    ///
96    /// Persistent data is loaded asynchronously, then world insertion is finalized at the
97    /// game tick safe point so the socket reader can enter play immediately.
98    pub fn queue_player_join(self: &Arc<Self>, player: Arc<Player>) {
99        if player.connection.closed() {
100            return;
101        }
102        if !self.reserve_player_join(&player) {
103            player.disconnect("You are already connected to this server");
104            return;
105        }
106
107        let server = Arc::clone(self);
108        tokio::spawn(async move {
109            let state = server.prepare_player_join(&player).await;
110            server
111                .pending_player_joins
112                .send(PendingPlayerJoin { player, state });
113        });
114    }
115
116    async fn prepare_player_join(&self, player: &Player) -> Result<DomainPlayerState, String> {
117        let target_domain = self.load_join_domain(player).await?;
118        self.load_domain_player_state(player, &target_domain, None)
119            .await
120    }
121
122    pub(super) fn process_player_joins(self: &Arc<Self>) {
123        for join in self.pending_player_joins.drain() {
124            self.finish_prepared_player_join(join);
125        }
126    }
127
128    pub(super) fn finish_prepared_player_join(self: &Arc<Self>, join: PendingPlayerJoin) {
129        let PendingPlayerJoin { player, state } = join;
130        let uuid = player.gameprofile.id;
131        if player.connection.closed() {
132            self.release_player_admission(uuid, PlayerAdmissionState::Joining);
133            return;
134        }
135
136        let state = match state {
137            Ok(state) => state,
138            Err(error) => {
139                self.release_player_admission(uuid, PlayerAdmissionState::Joining);
140                log::error!(
141                    "Failed to load player data for {}: {error}",
142                    player.gameprofile.name
143                );
144                player.disconnect("Failed to load player data");
145                return;
146            }
147        };
148
149        if !self.admit_reserved_player(Arc::clone(&player)) {
150            player.disconnect("You are already connected to this server");
151            return;
152        }
153
154        self.apply_cached_or_default_permission_state(&player);
155        Self::apply_domain_player_state(&player, &state);
156        self.send_login_packet(&player, &state.world);
157
158        player.reset(Arc::clone(&state.world), ResetReason::InitialJoin);
159        Self::apply_domain_player_state(&player, &state);
160        let residence_token = player.domain_residence_token();
161        let restores = self.prepare_domain_restores(&player, &state);
162        if !Self::install_domain_restores(&player, residence_token, &restores) {
163            tracing::error!(
164                player = %player.gameprofile.name,
165                "Initial admission lost its domain residence before restore installation"
166            );
167            player.connection.close();
168            self.remove_online_player_sync(&player);
169            return;
170        }
171        let pos = player.position();
172        let rotation = player.rotation();
173        // The client drops a player entity spawn when it does not already know that
174        // player's profile. Vanilla publishes player info before adding the player to
175        // the level, which can immediately start entity tracking for existing players.
176        self.sync_tab_list(&player);
177        let admitted = player.spawn(pos, rotation, ResetReason::InitialJoin);
178        if !admitted {
179            self.remove_online_player_sync(&player);
180            self.broadcast_to_online(CRemovePlayerInfo { uuids: vec![uuid] });
181            return;
182        }
183        let previous_name = self.record_known_player(&player.gameprofile);
184        self.broadcast_player_join_message(&player, previous_name.as_deref());
185        if player.mark_joined_world() {
186            player.send_inventory_to_remote();
187        }
188        self.schedule_domain_restores(&player, residence_token, restores);
189        if player.connection.closed() {
190            self.queue_player_disconnect(player);
191        }
192    }
193
194    pub(super) fn reserve_player_join(&self, player: &Player) -> bool {
195        let uuid = player.gameprofile.id;
196        let mut admissions = self.player_admissions.lock();
197        if admissions.contains_key(&uuid) {
198            return false;
199        }
200        if self.online_players.get_by_uuid(&uuid).is_some() {
201            return false;
202        }
203        admissions
204            .insert(uuid, PlayerAdmissionState::Joining)
205            .is_none()
206    }
207
208    fn admit_reserved_player(&self, player: Arc<Player>) -> bool {
209        let uuid = player.gameprofile.id;
210        let mut admissions = self.player_admissions.lock();
211        if admissions.get(&uuid) != Some(&PlayerAdmissionState::Joining) {
212            return false;
213        }
214
215        let admitted = self.online_players.insert(player);
216        let _ = admissions.remove(&uuid);
217        admitted
218    }
219
220    fn reserve_player_disconnect(&self, player: &Arc<Player>) -> bool {
221        let uuid = player.gameprofile.id;
222        let mut admissions = self.player_admissions.lock();
223        if admissions.contains_key(&uuid) {
224            return false;
225        }
226        if !self
227            .online_players
228            .get_by_uuid(&uuid)
229            .is_some_and(|current| Arc::ptr_eq(&current, player))
230        {
231            return false;
232        }
233        admissions
234            .insert(uuid, PlayerAdmissionState::Disconnecting)
235            .is_none()
236    }
237
238    pub(super) fn reserve_player_relocation(&self, player: &Arc<Player>) -> bool {
239        let uuid = player.gameprofile.id;
240        let mut admissions = self.player_admissions.lock();
241        if admissions.contains_key(&uuid) {
242            return false;
243        }
244        if !self
245            .online_players
246            .get_by_uuid(&uuid)
247            .is_some_and(|current| Arc::ptr_eq(&current, player))
248        {
249            return false;
250        }
251        admissions
252            .insert(uuid, PlayerAdmissionState::Relocating)
253            .is_none()
254    }
255
256    fn transition_player_relocation_to_disconnect(&self, player: &Arc<Player>) -> bool {
257        let uuid = player.gameprofile.id;
258        let mut admissions = self.player_admissions.lock();
259        if admissions.get(&uuid) != Some(&PlayerAdmissionState::Relocating) {
260            return false;
261        }
262        if !self
263            .online_players
264            .get_by_uuid(&uuid)
265            .is_some_and(|current| Arc::ptr_eq(&current, player))
266        {
267            return false;
268        }
269        let _ = admissions.insert(uuid, PlayerAdmissionState::Disconnecting);
270        true
271    }
272
273    pub(super) fn release_player_admission(&self, uuid: Uuid, state: PlayerAdmissionState) {
274        let mut admissions = self.player_admissions.lock();
275        if admissions.get(&uuid) == Some(&state) {
276            let _ = admissions.remove(&uuid);
277        }
278    }
279
280    fn remove_online_player_sync(&self, player: &Arc<Player>) {
281        let _ = self.online_players.remove_player_sync(player);
282    }
283
284    pub(crate) fn queue_player_disconnect(&self, player: Arc<Player>) {
285        debug_assert!(
286            player.connection.closed(),
287            "only closed players may enter the disconnect queue"
288        );
289        self.pending_player_disconnects.send(player);
290    }
291
292    pub(crate) fn queue_detached_player_disconnect(
293        &self,
294        player: Arc<Player>,
295        domain: String,
296        player_data: Arc<PersistentPlayerData>,
297        pending_token: PendingWorldChangeToken,
298    ) {
299        let uuid = player.gameprofile.id;
300        let live_memberships = self
301            .worlds
302            .values()
303            .filter(|world| world.contains_player(&player))
304            .cloned()
305            .collect::<Vec<_>>();
306        if !live_memberships.is_empty() {
307            tracing::error!(
308                player = %player.gameprofile.name,
309                membership_count = live_memberships.len(),
310                "Cleaning live world membership after detached player admission failed"
311            );
312            for world in live_memberships {
313                world.remove_player_for_world_change(&player);
314            }
315        }
316
317        player.finish_player_transition(pending_token);
318        player.finish_pending_world_change(pending_token);
319        if !self.reserve_player_disconnect(&player) {
320            return;
321        }
322
323        self.broadcast_player_leave_message(&player);
324        self.remove_online_player_sync(&player);
325        self.broadcast_to_online(CRemovePlayerInfo { uuids: vec![uuid] });
326        self.pending_player_disconnects
327            .send_prepared(PendingPlayerDisconnect {
328                player,
329                domain,
330                player_data,
331            });
332    }
333
334    pub(crate) fn queue_relocating_player_disconnect(
335        &self,
336        player: Arc<Player>,
337        domain: String,
338        player_data: Arc<PersistentPlayerData>,
339        pending_token: PendingWorldChangeToken,
340    ) {
341        let uuid = player.gameprofile.id;
342        let live_memberships = self
343            .worlds
344            .values()
345            .filter(|world| world.contains_player(&player))
346            .cloned()
347            .collect::<Vec<_>>();
348        if !live_memberships.is_empty() {
349            tracing::error!(
350                player = %player.gameprofile.name,
351                membership_count = live_memberships.len(),
352                "Cleaning live world membership after relocating player admission failed"
353            );
354            for world in live_memberships {
355                world.remove_player_for_world_change(&player);
356            }
357        }
358
359        player.finish_player_transition(pending_token);
360        player.finish_pending_world_change(pending_token);
361        if !self.transition_player_relocation_to_disconnect(&player) {
362            tracing::error!(
363                player = %player.gameprofile.name,
364                "Relocating player lost its exclusive disconnect-save ownership"
365            );
366            return;
367        }
368
369        self.broadcast_player_leave_message(&player);
370        self.remove_online_player_sync(&player);
371        self.broadcast_to_online(CRemovePlayerInfo { uuids: vec![uuid] });
372        self.pending_player_disconnects
373            .send_prepared(PendingPlayerDisconnect {
374                player,
375                domain,
376                player_data,
377            });
378    }
379
380    pub(super) fn process_player_disconnects(&self) -> Vec<PendingPlayerDisconnect> {
381        let mut pending = Vec::new();
382        while let Some(player) = self.pending_player_disconnects.pop() {
383            if let Some(disconnect) = self.process_player_disconnect(player) {
384                pending.push(disconnect);
385            }
386        }
387
388        if !pending.is_empty() {
389            // Steel batches the protocol-supported UUID list to avoid quadratic broadcast work
390            // during mass disconnects; Vanilla normally emits one packet per player.
391            let uuids = pending
392                .iter()
393                .map(|disconnect| disconnect.player.gameprofile.id)
394                .collect();
395            self.broadcast_to_online(CRemovePlayerInfo { uuids });
396        }
397        pending
398    }
399
400    pub(super) fn process_player_disconnect(
401        &self,
402        player: Arc<Player>,
403    ) -> Option<PendingPlayerDisconnect> {
404        let uuid = player.gameprofile.id;
405        if !self.reserve_player_disconnect(&player) {
406            return None;
407        }
408
409        let world = player.get_world();
410        let (player, domain, player_data) = world.detach_player_for_disconnect(Arc::clone(&player));
411
412        // Vanilla broadcasts before removing the player from its global player list.
413        self.broadcast_player_leave_message(&player);
414        let player = self.online_players.remove_player_sync(&player);
415
416        let Some(player) = player else {
417            self.release_player_admission(uuid, PlayerAdmissionState::Disconnecting);
418            return None;
419        };
420
421        Some(PendingPlayerDisconnect {
422            player,
423            domain,
424            player_data: Arc::new(player_data),
425        })
426    }
427
428    async fn save_disconnected_player(&self, pending: PendingPlayerDisconnect) {
429        let PendingPlayerDisconnect {
430            player,
431            domain,
432            player_data,
433        } = pending;
434        let uuid = player.gameprofile.id;
435        let start = Instant::now();
436
437        if let Err(e) = self
438            .player_data_storage
439            .save_domain_data(&domain, uuid, player_data.as_ref())
440            .await
441        {
442            log::error!("Failed to save player domain data for {uuid}: {e}");
443        }
444        if let Err(e) = self
445            .player_data_storage
446            .save_global(
447                uuid,
448                &GlobalPlayerData {
449                    last_active_domain: domain,
450                },
451            )
452            .await
453        {
454            log::error!("Failed to save global player data for {uuid}: {e}");
455        }
456
457        player.cleanup();
458        self.release_player_admission(uuid, PlayerAdmissionState::Disconnecting);
459        log::info!("Player {uuid} removed in {:?}", start.elapsed());
460    }
461
462    pub(super) fn start_player_disconnect_saves(self: &Arc<Self>, saves: &mut JoinSet<()>) {
463        let mut pending = self.pending_player_disconnects.drain_prepared();
464        pending.extend(self.process_player_disconnects());
465        for pending in pending {
466            let server = Arc::clone(self);
467            saves.spawn(async move {
468                server.save_disconnected_player(pending).await;
469            });
470        }
471
472        while let Some(result) = saves.try_join_next() {
473            if let Err(error) = result {
474                log::error!("Player disconnect save task failed: {error}");
475            }
476        }
477    }
478
479    /// Broadcasts a packet to every online player, regardless of world membership.
480    pub fn broadcast_to_online<P: ClientPacket>(&self, packet: P) {
481        let Ok(encoded) =
482            EncodedPacket::from_bare(packet, self.config.compression, ConnectionProtocol::Play)
483        else {
484            return;
485        };
486        self.online_players.iter_players(|_, player| {
487            player.connection.send_encoded(encoded.clone());
488            true
489        });
490    }
491
492    pub(super) fn broadcast_to_online_with<P: ClientPacket, F: Fn(&Player) -> P>(&self, packet: F) {
493        self.online_players.iter_players(|_, player| {
494            player.send_packet(packet(player));
495            true
496        });
497    }
498
499    /// Sends full tab list synchronization for a newly joined player.
500    ///
501    /// Server membership mirrors vanilla `PlayerList`; world entity spawning remains
502    /// owned by the per-world entity tracker.
503    fn sync_tab_list(&self, player: &Arc<Player>) {
504        self.online_players.iter_players(|_, existing_player| {
505            if existing_player.gameprofile.id == player.gameprofile.id {
506                return true;
507            }
508
509            let add_existing = CPlayerInfoUpdate::create_player_initializing(
510                existing_player.gameprofile.id,
511                existing_player.gameprofile.name.clone(),
512                existing_player.gameprofile.properties.clone(),
513                existing_player.game_mode().into(),
514                existing_player.connection.latency(),
515                None,
516                true,
517            );
518            player.send_packet(add_existing);
519
520            if let Some(session) = existing_player.chat_session()
521                && let Ok(protocol_data) = session.as_data().to_protocol_data()
522            {
523                player.send_packet(CPlayerInfoUpdate::update_chat_session(
524                    existing_player.gameprofile.id,
525                    protocol_data,
526                ));
527            }
528
529            true
530        });
531
532        let player_info_packet = CPlayerInfoUpdate::create_player_initializing(
533            player.gameprofile.id,
534            player.gameprofile.name.clone(),
535            player.gameprofile.properties.clone(),
536            player.game_mode().into(),
537            player.connection.latency(),
538            None,
539            true,
540        );
541        self.broadcast_to_online(player_info_packet);
542    }
543
544    pub(super) fn broadcast_player_latency_updates(&self) {
545        let mut latency_entries = Vec::new();
546        self.online_players.iter_players(|uuid, player| {
547            latency_entries.push((*uuid, player.connection.latency()));
548            true
549        });
550
551        if !latency_entries.is_empty() {
552            self.broadcast_to_online(CPlayerInfoUpdate::update_latency(latency_entries));
553        }
554    }
555}