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 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 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(¤t, 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(¤t, 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(¤t, 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 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 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 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 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}