1use steel_utils::translations;
2
3use super::{
4 Arc, CPlayerInfoUpdate, CRemovePlayerInfo, CancellationToken, ClientPacket, ConnectionProtocol,
5 DomainPlayerState, EncodedPacket, Entity, GlobalPlayerData, Instant, JoinSet,
6 NetworkConnection, PendingWorldChangeToken, PersistentPlayerData, Player, ResetReason,
7 SegQueue, Server, SyncMutex, Uuid, mpsc,
8};
9
10pub(super) struct PendingPlayerJoin {
11 pub(super) player: Arc<Player>,
12 pub(super) state: Result<DomainPlayerState, String>,
13}
14
15pub(super) struct PendingPlayerDisconnect {
16 player: Arc<Player>,
17 domain: String,
18 player_data: Arc<PersistentPlayerData>,
19}
20
21#[derive(Clone, Copy, Debug, Eq, PartialEq)]
22pub(super) enum PlayerAdmissionState {
23 Joining,
24 Relocating,
25 Disconnecting,
26}
27
28#[derive(Debug, Eq, PartialEq)]
29pub(super) enum PlayerJoinError {
30 DuplicateLogin,
31 ServerFull,
32}
33
34#[derive(Clone, Copy, Debug, Eq, PartialEq, thiserror::Error)]
36pub enum DuplicatePlayerWaitError {
37 #[error("duplicate-player wait was cancelled")]
39 Cancelled,
40 #[error("duplicate-player wait exceeded the login deadline")]
42 TimedOut,
43}
44
45pub struct PlayerJoinReservation {
49 server: Arc<Server>,
50 uuid: Uuid,
51 queued: bool,
52}
53
54impl PlayerJoinReservation {
55 pub fn queue_player_join(mut self, player: Arc<Player>) {
57 if player.gameprofile.id != self.uuid {
58 log::error!(
59 "Player join reservation UUID mismatch: reserved {}, got {}",
60 self.uuid,
61 player.gameprofile.id
62 );
63 player.connection.close();
64 return;
65 }
66
67 self.queued = true;
68 self.server.queue_reserved_player_join(player);
69 }
70}
71
72impl Drop for PlayerJoinReservation {
73 fn drop(&mut self) {
74 if !self.queued {
75 self.server
76 .release_player_admission(self.uuid, PlayerAdmissionState::Joining);
77 }
78 }
79}
80
81pub(super) struct PlayerJoinQueue {
82 sender: mpsc::Sender<PendingPlayerJoin>,
83 receiver: SyncMutex<mpsc::Receiver<PendingPlayerJoin>>,
84}
85
86impl PlayerJoinQueue {
87 pub(super) fn new() -> Self {
88 let (sender, receiver) = mpsc::channel();
89 Self {
90 sender,
91 receiver: SyncMutex::new(receiver),
92 }
93 }
94
95 fn send(&self, join: PendingPlayerJoin) {
96 let _ = self.sender.send(join);
97 }
98
99 fn drain(&self) -> Vec<PendingPlayerJoin> {
100 let receiver = self.receiver.lock();
101 let mut joins = Vec::new();
102 while let Ok(join) = receiver.try_recv() {
103 joins.push(join);
104 }
105 joins
106 }
107}
108
109pub(super) struct PlayerDisconnectQueue {
110 queued: SegQueue<Arc<Player>>,
111 prepared: SegQueue<PendingPlayerDisconnect>,
112}
113
114impl PlayerDisconnectQueue {
115 pub(super) const fn new() -> Self {
116 Self {
117 queued: SegQueue::new(),
118 prepared: SegQueue::new(),
119 }
120 }
121
122 fn send(&self, player: Arc<Player>) {
123 self.queued.push(player);
124 }
125
126 fn pop(&self) -> Option<Arc<Player>> {
127 self.queued.pop()
128 }
129
130 fn send_prepared(&self, disconnect: PendingPlayerDisconnect) {
131 self.prepared.push(disconnect);
132 }
133
134 fn drain_prepared(&self) -> Vec<PendingPlayerDisconnect> {
135 let mut prepared = Vec::new();
136 while let Some(disconnect) = self.prepared.pop() {
137 prepared.push(disconnect);
138 }
139 prepared
140 }
141
142 pub(super) fn clear(&self) {
143 while self.queued.pop().is_some() {}
144 while self.prepared.pop().is_some() {}
145 }
146}
147
148impl Server {
149 pub async fn disconnect_duplicate_player_and_wait(
155 &self,
156 uuid: Uuid,
157 connection_cancel: &CancellationToken,
158 login_deadline_tick: u64,
159 ) -> Result<(), DuplicatePlayerWaitError> {
160 loop {
161 let admission_changed = self.player_admission_changed.notified();
162 tokio::pin!(admission_changed);
163 admission_changed.as_mut().enable();
164
165 let existing = {
166 let admissions = self.player_admissions.lock();
167 let existing = self.online_players.get_by_uuid(&uuid);
168 if existing.is_none() && !admissions.contains_key(&uuid) {
169 return Ok(());
170 }
171 existing
172 };
173
174 if let Some(existing) = existing
175 && !existing.connection.closed()
176 {
177 existing.disconnect(translations::MULTIPLAYER_DISCONNECT_DUPLICATE_LOGIN.msg());
178 }
179
180 tokio::select! {
181 biased;
182 () = connection_cancel.cancelled() => {
183 return Err(DuplicatePlayerWaitError::Cancelled);
184 }
185 () = self.cancel_token.cancelled() => {
186 return Err(DuplicatePlayerWaitError::Cancelled);
187 }
188 () = &mut admission_changed => {}
189 () = self.wait_until_tick(login_deadline_tick) => {
190 return Err(DuplicatePlayerWaitError::TimedOut);
191 }
192 }
193 }
194 }
195
196 pub fn queue_player_join(self: &Arc<Self>, player: Arc<Player>) {
201 if player.connection.closed() {
202 return;
203 }
204 let Some(reservation) = self.try_reserve_player_join(player.gameprofile.id) else {
205 player.disconnect(translations::MULTIPLAYER_DISCONNECT_DUPLICATE_LOGIN.msg());
206 return;
207 };
208
209 reservation.queue_player_join(player);
210 }
211
212 fn queue_reserved_player_join(self: &Arc<Self>, player: Arc<Player>) {
213 let uuid = player.gameprofile.id;
214 if player.connection.closed() {
215 self.release_player_admission(uuid, PlayerAdmissionState::Joining);
216 return;
217 }
218
219 let server = Arc::clone(self);
220 tokio::spawn(Self::prepare_and_queue_player_join(server, player));
221 }
222
223 async fn prepare_and_queue_player_join(server: Arc<Self>, player: Arc<Player>) {
224 let state = match server.load_join_domain(&player).await {
225 Ok(target_domain) => {
226 server
227 .load_domain_player_state(&player, &target_domain, None)
228 .await
229 }
230 Err(error) => Err(error),
231 };
232 server
233 .pending_player_joins
234 .send(PendingPlayerJoin { player, state });
235 }
236
237 pub(super) fn process_player_joins(self: &Arc<Self>) {
238 for join in self.pending_player_joins.drain() {
239 self.finish_prepared_player_join(join);
240 }
241 }
242
243 pub(super) fn finish_prepared_player_join(self: &Arc<Self>, join: PendingPlayerJoin) {
244 let PendingPlayerJoin { player, state } = join;
245 let uuid = player.gameprofile.id;
246 if player.connection.closed() {
247 self.release_player_admission(uuid, PlayerAdmissionState::Joining);
248 return;
249 }
250
251 let state = match state {
252 Ok(state) => state,
253 Err(error) => {
254 self.release_player_admission(uuid, PlayerAdmissionState::Joining);
255 log::error!(
256 "Failed to load player data for {}: {error}",
257 player.gameprofile.name
258 );
259 player.disconnect("Failed to load player data");
260 return;
261 }
262 };
263
264 if let Err(error) = self.admit_reserved_player(Arc::clone(&player)) {
265 let reason = match error {
266 PlayerJoinError::DuplicateLogin => {
267 translations::MULTIPLAYER_DISCONNECT_DUPLICATE_LOGIN.msg()
268 }
269 PlayerJoinError::ServerFull => {
270 translations::MULTIPLAYER_DISCONNECT_SERVER_FULL.msg()
271 }
272 };
273 player.disconnect(reason);
274 return;
275 }
276
277 self.apply_cached_or_default_permission_state(&player);
278 Self::apply_domain_player_state(&player, &state);
279 self.send_login_packet(&player, &state.world);
280
281 player.reset(Arc::clone(&state.world), ResetReason::InitialJoin);
282 Self::apply_domain_player_state(&player, &state);
283 let residence_token = player.domain_residence_token();
284 let restores = self.prepare_domain_restores(&player, &state);
285 if !Self::install_domain_restores(&player, residence_token, &restores) {
286 tracing::error!(
287 player = %player.gameprofile.name,
288 "Initial admission lost its domain residence before restore installation"
289 );
290 player.connection.close();
291 let _ = self.remove_online_player_sync(&player);
292 return;
293 }
294 let pos = player.position();
295 let rotation = player.rotation();
296 self.sync_tab_list(&player);
300 let admitted = player.spawn(pos, rotation, ResetReason::InitialJoin);
301 if !admitted {
302 let _ = self.remove_online_player_sync(&player);
303 self.broadcast_to_online(CRemovePlayerInfo { uuids: vec![uuid] });
304 return;
305 }
306 let previous_name = self.record_known_player(&player.gameprofile);
307 self.broadcast_player_join_message(&player, previous_name.as_deref());
308 if player.mark_joined_world() {
309 player.send_inventory_to_remote();
310 }
311 self.schedule_domain_restores(&player, residence_token, restores);
312 if player.connection.closed() {
313 self.queue_player_disconnect(player);
314 }
315 }
316
317 #[must_use]
322 pub fn is_player_limit_reached(&self, uuid: Uuid) -> bool {
323 self.online_players.len() >= self.config.max_players as usize
324 && !self.can_bypass_player_limit(uuid)
325 }
326
327 pub fn try_reserve_player_join(self: &Arc<Self>, uuid: Uuid) -> Option<PlayerJoinReservation> {
329 let mut admissions = self.player_admissions.lock();
330 if admissions.contains_key(&uuid) {
331 return None;
332 }
333 if self.online_players.get_by_uuid(&uuid).is_some() {
334 return None;
335 }
336 let previous = admissions.insert(uuid, PlayerAdmissionState::Joining);
337 debug_assert!(previous.is_none());
338 Some(PlayerJoinReservation {
339 server: Arc::clone(self),
340 uuid,
341 queued: false,
342 })
343 }
344
345 #[cfg(test)]
346 pub(super) fn reserve_player_join(&self, player: &Player) -> bool {
347 let uuid = player.gameprofile.id;
348 let mut admissions = self.player_admissions.lock();
349 if admissions.contains_key(&uuid) || self.online_players.get_by_uuid(&uuid).is_some() {
350 return false;
351 }
352 admissions
353 .insert(uuid, PlayerAdmissionState::Joining)
354 .is_none()
355 }
356
357 pub(super) fn admit_reserved_player(&self, player: Arc<Player>) -> Result<(), PlayerJoinError> {
358 let uuid = player.gameprofile.id;
359 let mut admissions = self.player_admissions.lock();
360 if admissions.get(&uuid) != Some(&PlayerAdmissionState::Joining) {
361 return Err(PlayerJoinError::DuplicateLogin);
362 }
363
364 let result = if self.is_player_limit_reached(uuid) {
366 Err(PlayerJoinError::ServerFull)
367 } else if self.online_players.insert(player) {
368 Ok(())
369 } else {
370 Err(PlayerJoinError::DuplicateLogin)
371 };
372 let _ = admissions.remove(&uuid);
373 drop(admissions);
374 self.player_admission_changed.notify_waiters();
375 result
376 }
377
378 pub(super) fn reserve_player_disconnect(&self, player: &Arc<Player>) -> bool {
379 let uuid = player.gameprofile.id;
380 let mut admissions = self.player_admissions.lock();
381 if admissions.contains_key(&uuid) {
382 return false;
383 }
384 if !self
385 .online_players
386 .get_by_uuid(&uuid)
387 .is_some_and(|current| Arc::ptr_eq(¤t, player))
388 {
389 return false;
390 }
391 admissions
392 .insert(uuid, PlayerAdmissionState::Disconnecting)
393 .is_none()
394 }
395
396 pub(super) fn reserve_player_relocation(&self, player: &Arc<Player>) -> bool {
397 let uuid = player.gameprofile.id;
398 let mut admissions = self.player_admissions.lock();
399 if admissions.contains_key(&uuid) {
400 return false;
401 }
402 if !self
403 .online_players
404 .get_by_uuid(&uuid)
405 .is_some_and(|current| Arc::ptr_eq(¤t, player))
406 {
407 return false;
408 }
409 admissions
410 .insert(uuid, PlayerAdmissionState::Relocating)
411 .is_none()
412 }
413
414 fn transition_player_relocation_to_disconnect(&self, player: &Arc<Player>) -> bool {
415 let uuid = player.gameprofile.id;
416 let mut admissions = self.player_admissions.lock();
417 if admissions.get(&uuid) != Some(&PlayerAdmissionState::Relocating) {
418 return false;
419 }
420 if !self
421 .online_players
422 .get_by_uuid(&uuid)
423 .is_some_and(|current| Arc::ptr_eq(¤t, player))
424 {
425 return false;
426 }
427 let _ = admissions.insert(uuid, PlayerAdmissionState::Disconnecting);
428 true
429 }
430
431 pub(super) fn release_player_admission(&self, uuid: Uuid, state: PlayerAdmissionState) {
432 let mut admissions = self.player_admissions.lock();
433 if admissions.get(&uuid) == Some(&state) {
434 let _ = admissions.remove(&uuid);
435 drop(admissions);
436 self.player_admission_changed.notify_waiters();
437 }
438 }
439
440 pub(crate) fn remove_online_player_sync(&self, player: &Arc<Player>) -> Option<Arc<Player>> {
441 let removed = self.online_players.remove_player_sync(player);
442 if let Some(removed) = &removed {
443 if removed.session.clear_player(removed) {
444 self.discard_player_packets(removed);
445 }
446 self.player_admission_changed.notify_waiters();
447 }
448 removed
449 }
450
451 pub(crate) fn replace_online_player(
457 &self,
458 expected: &Arc<Player>,
459 replacement: Arc<Player>,
460 ) -> bool {
461 if !Arc::ptr_eq(&expected.session, &replacement.session)
462 || !Arc::ptr_eq(&expected.connection, &replacement.connection)
463 || !expected.session.is_current_player(expected)
464 {
465 return false;
466 }
467
468 let uuid = expected.gameprofile.id;
469 let admissions = self.player_admissions.lock();
470 if admissions.contains_key(&uuid) {
471 return false;
472 }
473
474 self.online_players.replace_player(expected, replacement)
475 }
476
477 pub(crate) fn rollback_respawn_online_player(
479 &self,
480 failed_replacement: &Arc<Player>,
481 original: Arc<Player>,
482 ) -> bool {
483 if !Arc::ptr_eq(&failed_replacement.session, &original.session)
484 || !Arc::ptr_eq(&failed_replacement.connection, &original.connection)
485 || !original.session.is_current_player(&original)
486 {
487 return false;
488 }
489
490 let uuid = original.gameprofile.id;
491 let admissions = self.player_admissions.lock();
492 if admissions.contains_key(&uuid) {
493 return false;
494 }
495
496 self.online_players
497 .replace_player(failed_replacement, original)
498 }
499
500 pub(crate) fn queue_player_disconnect(&self, player: Arc<Player>) {
501 debug_assert!(
502 player.connection.closed(),
503 "only closed players may enter the disconnect queue"
504 );
505 self.pending_player_disconnects.send(player);
506 }
507
508 pub(crate) fn queue_relocating_player_disconnect(
509 &self,
510 player: Arc<Player>,
511 domain: String,
512 player_data: Arc<PersistentPlayerData>,
513 pending_token: PendingWorldChangeToken,
514 ) {
515 let uuid = player.gameprofile.id;
516 let live_memberships = self
517 .worlds
518 .values()
519 .filter(|world| world.contains_player(&player))
520 .cloned()
521 .collect::<Vec<_>>();
522 if !live_memberships.is_empty() {
523 tracing::error!(
524 player = %player.gameprofile.name,
525 membership_count = live_memberships.len(),
526 "Cleaning live world membership after relocating player admission failed"
527 );
528 for world in live_memberships {
529 world.remove_player_for_world_change(&player);
530 }
531 }
532
533 player.finish_player_transition(pending_token);
534 player.finish_pending_world_change(pending_token);
535 if !self.transition_player_relocation_to_disconnect(&player) {
536 tracing::error!(
537 player = %player.gameprofile.name,
538 "Relocating player lost its exclusive disconnect-save ownership"
539 );
540 return;
541 }
542
543 self.broadcast_player_leave_message(&player);
544 let _ = self.remove_online_player_sync(&player);
545 self.broadcast_to_online(CRemovePlayerInfo { uuids: vec![uuid] });
546 self.pending_player_disconnects
547 .send_prepared(PendingPlayerDisconnect {
548 player,
549 domain,
550 player_data,
551 });
552 }
553
554 pub(super) fn process_player_disconnects(&self) -> Vec<PendingPlayerDisconnect> {
555 self.online_players.iter_players(|_, player| {
557 if player.connection.closed() {
558 self.queue_player_disconnect(Arc::clone(player));
559 }
560 true
561 });
562
563 let mut pending = Vec::new();
564 while let Some(player) = self.pending_player_disconnects.pop() {
565 if let Some(disconnect) = self.process_player_disconnect(player) {
566 pending.push(disconnect);
567 }
568 }
569
570 if !pending.is_empty() {
571 let uuids = pending
574 .iter()
575 .map(|disconnect| disconnect.player.gameprofile.id)
576 .collect();
577 self.broadcast_to_online(CRemovePlayerInfo { uuids });
578 }
579 pending
580 }
581
582 pub(super) fn process_player_disconnect(
583 &self,
584 player: Arc<Player>,
585 ) -> Option<PendingPlayerDisconnect> {
586 let uuid = player.gameprofile.id;
587 if !self.reserve_player_disconnect(&player) {
588 return None;
589 }
590
591 let world = player.get_world();
592 let (player, domain, player_data) = world.detach_player_for_disconnect(Arc::clone(&player));
593
594 self.broadcast_player_leave_message(&player);
596 let player = self.remove_online_player_sync(&player);
597
598 let Some(player) = player else {
599 self.release_player_admission(uuid, PlayerAdmissionState::Disconnecting);
600 return None;
601 };
602
603 Some(PendingPlayerDisconnect {
604 player,
605 domain,
606 player_data: Arc::new(player_data),
607 })
608 }
609
610 async fn save_disconnected_player(&self, pending: PendingPlayerDisconnect) {
611 let PendingPlayerDisconnect {
612 player,
613 domain,
614 player_data,
615 } = pending;
616 let uuid = player.gameprofile.id;
617 let start = Instant::now();
618
619 if let Err(e) = self
620 .player_data_storage
621 .save_domain_data(&domain, uuid, player_data.as_ref())
622 .await
623 {
624 log::error!("Failed to save player domain data for {uuid}: {e}");
625 }
626 if let Err(e) = self
627 .player_data_storage
628 .save_global(
629 uuid,
630 &GlobalPlayerData {
631 last_active_domain: domain,
632 },
633 )
634 .await
635 {
636 log::error!("Failed to save global player data for {uuid}: {e}");
637 }
638
639 player.cleanup();
640 self.release_player_admission(uuid, PlayerAdmissionState::Disconnecting);
641 log::info!("Player {uuid} removed in {:?}", start.elapsed());
642 }
643
644 pub(super) fn start_player_disconnect_saves(self: &Arc<Self>, saves: &mut JoinSet<()>) {
645 let mut pending = self.pending_player_disconnects.drain_prepared();
646 pending.extend(self.process_player_disconnects());
647 for pending in pending {
648 let server = Arc::clone(self);
649 saves.spawn(async move {
650 server.save_disconnected_player(pending).await;
651 });
652 }
653
654 while let Some(result) = saves.try_join_next() {
655 if let Err(error) = result {
656 log::error!("Player disconnect save task failed: {error}");
657 }
658 }
659 }
660
661 pub fn broadcast_to_online<P: ClientPacket>(&self, packet: P) {
663 let Ok(encoded) =
664 EncodedPacket::from_bare(packet, self.config.compression, ConnectionProtocol::Play)
665 else {
666 return;
667 };
668 self.online_players.iter_players(|_, player| {
669 player.connection.send_encoded(encoded.clone());
670 true
671 });
672 }
673
674 pub(super) fn broadcast_to_online_with<P: ClientPacket, F: Fn(&Player) -> P>(&self, packet: F) {
675 self.online_players.iter_players(|_, player| {
676 player.send_packet(packet(player));
677 true
678 });
679 }
680
681 fn sync_tab_list(&self, player: &Arc<Player>) {
686 self.online_players.iter_players(|_, existing_player| {
687 if existing_player.gameprofile.id == player.gameprofile.id {
688 return true;
689 }
690
691 let add_existing = CPlayerInfoUpdate::create_player_initializing(
692 existing_player.gameprofile.id,
693 existing_player.gameprofile.name.clone(),
694 existing_player.gameprofile.properties.clone(),
695 existing_player.game_mode().into(),
696 existing_player.connection.latency(),
697 None,
698 existing_player.shows_hat(),
699 );
700 player.send_packet(add_existing);
701
702 if let Some(session) = existing_player.chat_session()
703 && let Ok(protocol_data) = session.as_data().to_protocol_data()
704 {
705 player.send_packet(CPlayerInfoUpdate::update_chat_session(
706 existing_player.gameprofile.id,
707 protocol_data,
708 ));
709 }
710
711 true
712 });
713
714 let player_info_packet = CPlayerInfoUpdate::create_player_initializing(
715 player.gameprofile.id,
716 player.gameprofile.name.clone(),
717 player.gameprofile.properties.clone(),
718 player.game_mode().into(),
719 player.connection.latency(),
720 None,
721 player.shows_hat(),
722 );
723 self.broadcast_to_online(player_info_packet);
724 }
725
726 pub(super) fn broadcast_player_latency_updates(&self) {
727 let mut latency_entries = Vec::new();
728 self.online_players.iter_players(|uuid, player| {
729 latency_entries.push((*uuid, player.connection.latency()));
730 true
731 });
732
733 if !latency_entries.is_empty() {
734 self.broadcast_to_online(CPlayerInfoUpdate::update_latency(latency_entries));
735 }
736 }
737}