1use super::tick_overload::TickOverloadGuard;
2use super::world_tick_workers::{WorldTickWorkerError, WorldTickWorkers};
3use super::{
4 Arc, CCommandSuggestions, CHUNK_SENDING_TPS, COMMAND_DATA_AUTOSAVE_INTERVAL,
5 COMMAND_REQUESTS_PER_TICK, COMMAND_RESUMPTIONS_PER_TICK, CancellationToken, ChunkPos,
6 ChunkSender, CommandExecutionContext, CommandExecutionOwner, CommandRequest,
7 CommandResultCallback, CommandSender, CommandSource, Duration, EncodedChunk,
8 ExecutionCommandSource, ExecutionStop, GameTickTaskGuard, GlobalPlayerData, Instant, JoinSet,
9 MenuRemovalStatus, NetworkConnection, PendingCommandExecutionQueue, PersistentPlayerData,
10 Player, SEND_PLAYER_INFO_INTERVAL, SLOW_CHUNK_TICK_THRESHOLD, Server, StringReader,
11 SuggestionError, Suggestions, TAB_LIST_UPDATE_INTERVAL, TabListTickStats, ThreadPool, World,
12 command_suggestions_packet, sleep, spawn_blocking,
13};
14use steel_registry::vanilla_custom_stats;
15use steel_utils::threading::{available_worker_threads, worker_threads_for_available};
16use steel_utils::translations;
17
18impl Server {
19 pub(super) fn advance_server_tick(&self) -> (u64, bool) {
20 let (tick_count, runs_normally) = {
21 let mut tick_manager = self.tick_rate_manager.write();
22 tick_manager.tick();
23 let runs_normally = tick_manager.runs_normally();
24 tick_manager.increment_tick_count();
25 (tick_manager.tick_count, runs_normally)
26 };
27 self.server_tick_changed.notify_waiters();
28 (tick_count, runs_normally)
29 }
30
31 pub async fn run(self: Arc<Self>, cancel_token: CancellationToken) {
33 self.packet_processor.open_after_tick();
34 let packet_worker_count =
35 worker_threads_for_available(self.config.packet_workers, available_worker_threads());
36 let mut packet_handles = Vec::with_capacity(packet_worker_count);
37 for worker_id in 0..packet_worker_count {
38 let s = self.clone();
39 let t = cancel_token.clone();
40 packet_handles.push(tokio::spawn(async move {
41 if let Err(error) = spawn_blocking(move || s.packet_processor.run(&s)).await {
42 log::error!("Gameplay packet worker {worker_id} failed: {error}");
43 t.cancel();
44 }
45 }));
46 }
47 let packet_supervisor_cancel = cancel_token.clone();
48 let packet_workers = async move {
49 for handle in packet_handles {
50 if let Err(error) = handle.await {
51 log::error!("Gameplay packet supervisor failed: {error}");
52 packet_supervisor_cancel.cancel();
53 }
54 }
55 };
56 let game_handle = {
57 let s = self.clone();
58 let t = cancel_token.clone();
59 let task_guard = GameTickTaskGuard::new(self.clone(), cancel_token.clone());
60 tokio::spawn(async move {
61 let _task_guard = task_guard;
62 s.run_game_tick(t).await;
63 })
64 };
65 let chunk_send_handle = {
66 let s = self.clone();
67 let t = cancel_token.clone();
68 tokio::spawn(async move { s.run_chunk_sending_tick(t).await })
69 };
70 let ((), game_result, chunk_send_result) =
71 tokio::join!(packet_workers, game_handle, chunk_send_handle);
72 for (task, result) in [
73 ("Game tick", game_result),
74 ("Chunk sending tick", chunk_send_result),
75 ] {
76 if let Err(error) = result {
77 log::error!("{task} task failed: {error}");
78 }
79 }
80 }
81
82 pub async fn save_and_shutdown(&self) {
90 if let Err(error) = self.flush_known_players().await {
91 log::error!("Failed to flush known player cache during shutdown: {error}");
92 }
93
94 let players = self.get_players();
95 for player in &players {
96 let _ = self.reserve_player_disconnect(player);
99 player.disconnect(translations::MULTIPLAYER_DISCONNECT_SERVER_SHUTDOWN.msg());
100 assert_eq!(
101 player.remove_all_menus(),
102 MenuRemovalStatus::Complete,
103 "shutdown menu removal must run after packet processing stops"
104 );
105 }
106
107 for world in self.worlds.values() {
108 world.chunk_map.stop_generation_refill_loop();
109 world.chunk_map.task_tracker.close();
110 world.chunk_map.task_tracker.wait().await;
111 }
112
113 let mut players_to_save = Vec::new();
114 for player in players {
115 let domain = player.get_world().domain().to_owned();
116 player.award_custom_stat(&vanilla_custom_stats::LEAVE_GAME);
117 let data = PersistentPlayerData::from_player(&player);
118 player.store_ender_pearls_with_player();
119 players_to_save.push((player, domain, data));
120 }
121
122 log::info!("Saving world data...");
123 let command_data = self.save_command_data().await;
124 match command_data.scoreboards {
125 Ok(saved) => log::info!("Saved {saved} domain scoreboards"),
126 Err(error) => log::error!("Failed to save domain scoreboards: {error}"),
127 }
128 match command_data.storage {
129 Ok(saved) => log::info!("Saved {saved} domain command storages"),
130 Err(error) => log::error!("Failed to save domain command storage: {error}"),
131 }
132 let mut total_saved = 0;
133 for world in self.worlds.values() {
134 world.cleanup(&mut total_saved).await;
135 }
136 log::info!("Saved {total_saved} chunks");
137
138 log::info!("Saving player data...");
139 let mut saved = 0;
140 for (player, domain, data) in players_to_save {
141 let uuid = player.gameprofile.id;
142 match self
143 .player_data_storage
144 .save_domain_data(&domain, uuid, &data)
145 .await
146 {
147 Ok(()) => {
148 saved += 1;
149 }
150 Err(e) => {
151 log::error!("Failed to save player {uuid} domain data during shutdown: {e}");
152 }
153 }
154 if let Err(e) = self
155 .player_data_storage
156 .save_global(
157 uuid,
158 &GlobalPlayerData {
159 last_active_domain: domain,
160 },
161 )
162 .await
163 {
164 log::error!("Failed to save player {uuid} global data during shutdown: {e}");
165 }
166 }
167 log::info!("Saved {saved} players");
168 }
169
170 #[expect(
172 clippy::too_many_lines,
173 reason = "the ordered tick phases and their shutdown joins remain easier to audit together"
174 )]
175 async fn run_game_tick(self: Arc<Self>, cancel_token: CancellationToken) {
176 let world_tick_workers = match WorldTickWorkers::spawn(self.worlds.values()) {
177 Ok(workers) => workers,
178 Err(error) => {
179 log::error!("Failed to start world tick workers: {error}");
180 cancel_token.cancel();
181 return;
182 }
183 };
184 let mut next_tick_time = Instant::now();
185 let mut overload_guard = TickOverloadGuard::new();
186 let mut next_command_data_autosave = Instant::now() + COMMAND_DATA_AUTOSAVE_INTERVAL;
187 let mut player_info_ticks = 0_u64;
188 let mut pending_command_executions = PendingCommandExecutionQueue::<CommandSource>::new();
189 let mut player_disconnect_saves = JoinSet::new();
190 let mut command_data_autosaves = JoinSet::new();
191
192 loop {
193 if cancel_token.is_cancelled() {
194 break;
195 }
196
197 let (nanoseconds_per_tick, should_sprint_this_tick) = {
198 let mut tick_manager = self.tick_rate_manager.write();
199 let nanoseconds_per_tick = tick_manager.nanoseconds_per_tick;
200 let (should_sprint, sprint_report) = tick_manager.check_should_sprint_this_tick();
201 drop(tick_manager);
202
203 if let Some(report) = sprint_report {
204 self.broadcast_sprint_report(&report);
205 self.broadcast_ticking_state();
206 }
207
208 (nanoseconds_per_tick, should_sprint)
209 };
210
211 if should_sprint_this_tick {
212 next_tick_time = Instant::now();
213 overload_guard.restart_report_gap(next_tick_time);
214 } else {
215 let now = Instant::now();
216 if now < next_tick_time {
217 tokio::select! {
218 () = cancel_token.cancelled() => break,
219 () = sleep(next_tick_time - now) => {}
220 }
221 } else {
222 overload_guard.skip_backlog_if_overloaded(
223 now,
224 &mut next_tick_time,
225 nanoseconds_per_tick,
226 );
227 }
228 next_tick_time += Duration::from_nanos(nanoseconds_per_tick);
229 }
230
231 if cancel_token.is_cancelled() {
232 break;
233 }
234
235 let tick_start = Instant::now();
236 self.packet_processor.close_for_tick().await;
237 self.start_player_disconnect_saves(&mut player_disconnect_saves);
238
239 let (tick_count, runs_normally) = self.advance_server_tick();
240
241 self.tick_pending_command_executions(&mut pending_command_executions);
242 self.tick_command_requests(&mut pending_command_executions);
243 if let Err(error) = self
244 .tick_worlds_game(&world_tick_workers, tick_count, runs_normally)
245 .await
246 {
247 log::error!("World game tick failed: {error}");
248 cancel_token.cancel();
249 break;
250 }
251 player_info_ticks += 1;
252 if player_info_ticks > SEND_PLAYER_INFO_INTERVAL {
253 let _span = tracing::trace_span!("broadcast_latency").entered();
254 self.broadcast_player_latency_updates();
255 player_info_ticks = 0;
256 }
257 self.tick_jobs(tick_count, runs_normally);
258 self.process_player_joins();
259
260 {
261 let server = self.clone();
262 let _ =
263 spawn_blocking(move || server.process_world_changes(tick_count, runs_normally))
264 .await;
265 }
266
267 self.process_domain_switches();
268
269 self.tick_command_data_autosave(
270 &mut next_command_data_autosave,
271 &mut command_data_autosaves,
272 );
273
274 let tab_list_tick_stats = self.record_tick_and_capture_tab_stats(
275 tick_count,
276 tick_start.elapsed().as_nanos() as u64,
277 );
278
279 if let Some(tick_stats) = tab_list_tick_stats {
280 self.broadcast_tab_list(tick_stats);
281 }
282
283 if should_sprint_this_tick {
284 let mut tick_manager = self.tick_rate_manager.write();
285 tick_manager.end_tick_work();
286 }
287
288 self.packet_processor.open_after_tick();
289 if should_sprint_this_tick || Instant::now() >= next_tick_time {
290 self.packet_processor.wait_for_overload_progress().await;
291 }
292 }
293
294 self.jobs.cancel_all();
295 pending_command_executions.cancel_all();
296 self.command_requests.clear();
297 self.packet_processor.stop();
298 self.start_player_disconnect_saves(&mut player_disconnect_saves);
299 self.pending_player_disconnects.clear();
300 while let Some(result) = player_disconnect_saves.join_next().await {
301 if let Err(error) = result {
302 log::error!("Player disconnect save task failed during shutdown: {error}");
303 }
304 }
305 while let Some(result) = command_data_autosaves.join_next().await {
306 if let Err(error) = result {
307 log::error!("Command data autosave task failed during shutdown: {error}");
308 }
309 }
310 }
311
312 fn record_tick_and_capture_tab_stats(
313 &self,
314 tick_count: u64,
315 tick_duration_nanos: u64,
316 ) -> Option<TabListTickStats> {
317 let mut tick_manager = self.tick_rate_manager.write();
318 tick_manager.record_tick_time(tick_duration_nanos);
319 tick_count
320 .is_multiple_of(TAB_LIST_UPDATE_INTERVAL)
321 .then(|| TabListTickStats::capture(&tick_manager))
322 }
323
324 async fn autosave_command_data(&self) {
325 tracing::debug!("Command data autosave started");
326 let results = self.save_command_data().await;
327 match results.scoreboards {
328 Ok(saved) => tracing::debug!(saved, "Domain scoreboard autosave completed"),
329 Err(error) => tracing::error!(%error, "Domain scoreboard autosave failed"),
330 }
331 match results.storage {
332 Ok(saved) => tracing::debug!(saved, "Domain command-storage autosave completed"),
333 Err(error) => tracing::error!(%error, "Domain command-storage autosave failed"),
334 }
335 }
336
337 fn tick_command_data_autosave(
338 self: &Arc<Self>,
339 next_autosave: &mut Instant,
340 saves: &mut JoinSet<()>,
341 ) {
342 while let Some(result) = saves.try_join_next() {
343 if let Err(error) = result {
344 log::error!("Command data autosave task failed: {error}");
345 }
346 }
347 if Instant::now() < *next_autosave {
348 return;
349 }
350 if saves.is_empty() {
351 let server = Arc::clone(self);
352 saves.spawn(async move {
353 server.autosave_command_data().await;
354 });
355 } else {
356 tracing::warn!("Skipping command data autosave while the previous save runs");
357 }
358 *next_autosave = Instant::now() + COMMAND_DATA_AUTOSAVE_INTERVAL;
359 }
360
361 fn tick_pending_command_executions(
362 &self,
363 pending: &mut PendingCommandExecutionQueue<CommandSource>,
364 ) {
365 let stats = pending.tick(COMMAND_RESUMPTIONS_PER_TICK, |owner| owner.is_current(self));
366 if stats.polled == COMMAND_RESUMPTIONS_PER_TICK && stats.pending > 0 {
367 tracing::debug!(
368 polled = stats.polled,
369 finished = stats.finished,
370 pending = stats.pending,
371 "Command resumption tick reached per-tick processing limit"
372 );
373 }
374 }
375
376 fn tick_command_requests(
377 self: &Arc<Self>,
378 pending: &mut PendingCommandExecutionQueue<CommandSource>,
379 ) {
380 let mut handled = 0;
381 for _ in 0..COMMAND_REQUESTS_PER_TICK {
382 let Some(request) = self
383 .command_requests
384 .pop_front_runnable(|owner| !pending.blocks(owner.key()))
385 else {
386 break;
387 };
388 handled += 1;
389
390 match request {
391 CommandRequest::Execute { owner, command } => {
392 if !owner.is_current(self) {
393 continue;
394 }
395 self.execute_command_request(pending, owner, &command);
396 }
397 CommandRequest::Suggestions {
398 owner,
399 transaction_id,
400 input,
401 } => {
402 if !owner.is_current(self) {
403 continue;
404 }
405 let Some(player) = owner.sender().get_player() else {
406 tracing::error!("command suggestion request has a non-player owner");
407 continue;
408 };
409 self.send_command_suggestions(player, transaction_id, &input);
410 }
411 }
412 }
413
414 if handled == COMMAND_REQUESTS_PER_TICK {
415 tracing::debug!(handled, "Command request tick reached its processing limit");
416 }
417 }
418
419 fn execute_command_request(
420 self: &Arc<Self>,
421 pending: &mut PendingCommandExecutionQueue<CommandSource>,
422 owner: CommandExecutionOwner,
423 command: &str,
424 ) {
425 let source = CommandSource::new(owner.sender().clone(), Arc::clone(self));
426 let command = command.strip_prefix('/').unwrap_or(command);
427 let chain = {
428 let dispatcher = self.command_dispatcher.read();
429 let parse = dispatcher.parse(command, source.clone());
430 dispatcher.context_chain(parse)
431 };
432 let chain = match chain {
433 Ok(chain) => chain,
434 Err(error) => {
435 source.handle_error(&error, false);
436 return;
437 }
438 };
439
440 let mut execution = CommandExecutionContext::for_source(&source);
441 execution.queue_initial_command(chain, source, CommandResultCallback::empty());
442 if execution.run() == ExecutionStop::Suspended && !pending.push_suspended(owner, execution)
443 {
444 tracing::error!("suspended command execution could not be retained");
445 }
446 }
447
448 fn send_command_suggestions(
449 self: &Arc<Self>,
450 player: &Arc<Player>,
451 transaction_id: i32,
452 input: &str,
453 ) {
454 let suggestions =
455 self.build_command_suggestions(CommandSender::Player(Arc::clone(player)), input);
456 match suggestions {
457 Ok(suggestions) => {
458 player.send_packet(command_suggestions_packet(transaction_id, &suggestions));
459 }
460 Err(error) => {
461 tracing::warn!(%error, "failed to build command suggestions");
462 player.send_packet(CCommandSuggestions::new(transaction_id, 0, 0, Vec::new()));
463 }
464 }
465 }
466
467 pub(super) fn build_command_suggestions(
468 self: &Arc<Self>,
469 sender: CommandSender,
470 input: &str,
471 ) -> Result<Suggestions, SuggestionError> {
472 let source = CommandSource::new(sender, Arc::clone(self));
473 let mut reader = StringReader::new(input);
474 if reader.peek() == Some('/') {
475 reader.skip();
476 }
477 let dispatcher = self.command_dispatcher.read();
478 let parse = dispatcher.parse_reader(reader, source);
479 dispatcher.completion_suggestions(&parse)
480 }
481
482 async fn run_chunk_sending_tick(self: Arc<Self>, cancel_token: CancellationToken) {
484 let nanos_per_tick = 1_000_000_000 / CHUNK_SENDING_TPS;
485 let mut next_tick_time = Instant::now();
486
487 loop {
488 if cancel_token.is_cancelled() {
489 break;
490 }
491
492 let now = Instant::now();
493 if now < next_tick_time {
494 tokio::select! {
495 () = cancel_token.cancelled() => break,
496 () = sleep(next_tick_time - now) => {}
497 }
498 }
499 next_tick_time += Duration::from_nanos(nanos_per_tick);
500
501 if cancel_token.is_cancelled() {
502 break;
503 }
504
505 let server = self.clone();
506 let _ = spawn_blocking(move || {
507 server.tick_chunk_sending();
508 })
509 .await;
510 }
511 }
512
513 fn tick_chunk_sending(&self) {
518 let tick_start = Instant::now();
519 for world in self.worlds.values() {
520 let mut encode_cache = rustc_hash::FxHashMap::default();
521 world.players.iter_players(|_uuid, player| {
522 Self::send_chunks_for_player(
523 player,
524 world,
525 &mut encode_cache,
526 self.chunk_encoding_pool.as_ref(),
527 );
528 true
529 });
530 }
531
532 let elapsed = tick_start.elapsed();
533 if elapsed >= SLOW_CHUNK_TICK_THRESHOLD {
534 tracing::warn!(?elapsed, "Chunk sending tick slow");
535 }
536 }
537
538 fn send_chunks_for_player(
541 player: &Arc<Player>,
542 world: &Arc<World>,
543 encode_cache: &mut rustc_hash::FxHashMap<ChunkPos, EncodedChunk>,
544 encoding_pool: &ThreadPool,
545 ) {
546 let chunk_pos = *player.last_chunk_pos.lock();
547 let connection = &player.connection;
548
549 let prepared = {
551 let mut sender = player.chunk_sender().lock();
552 sender.prepare_batch(world, chunk_pos, &player.chunk_send_epoch)
553 };
554
555 let Some(batch) = prepared else {
556 return;
557 };
558
559 let compression = connection.compression();
561 let encoded = ChunkSender::encode_batch(&batch, encode_cache, compression, encoding_pool);
562
563 let tracking_view = player.last_tracking_view.lock();
567 let Some(view) = *tracking_view else {
568 return;
569 };
570 if !world.contains_player(player) {
571 return;
572 }
573 let sent_chunks = {
574 let mut sender = player.chunk_sender().lock();
575 sender.commit_batch(&batch, encoded, connection, &player.chunk_send_epoch)
576 };
577
578 if sent_chunks.is_empty() {
579 return;
580 }
581
582 let sent_chunks = player.chunk_sender().lock().sent_chunks_snapshot();
583 world
584 .entity_tracker()
585 .update_player(player, &view, |chunk| sent_chunks.contains(&chunk));
586 }
587
588 #[tracing::instrument(level = "trace", skip(self, workers), name = "tick_worlds")]
589 pub(super) async fn tick_worlds_game(
590 &self,
591 workers: &WorldTickWorkers,
592 tick_count: u64,
593 runs_normally: bool,
594 ) -> Result<(), WorldTickWorkerError> {
595 if runs_normally {
596 self.worlds.advance_domain_game_times();
597 }
598 let all_timings = workers.tick_all(tick_count, runs_normally).await?;
599 for (i, timings) in all_timings.iter().enumerate() {
600 if timings.elapsed < SLOW_CHUNK_TICK_THRESHOLD {
601 continue;
602 }
603 let cm = &timings.chunk_map;
604 let scheduling = &cm.scheduling;
605 tracing::warn!(
606 world = i,
607 elapsed = ?timings.elapsed,
608 tick_count,
609 entity_tick = ?timings.entity_tick,
610 ticket_updates = ?scheduling.ticket_updates,
611 block_entity_unloads = ?scheduling.block_entity_unloads,
612 readiness_demotions = ?scheduling.readiness_demotions,
613 lifecycle_commit = ?scheduling.lifecycle_commit,
614 readiness_reconcile = ?scheduling.readiness_reconcile,
615 post_process_generation = ?scheduling.post_process_generation,
616 post_process_chunk_count = scheduling.post_process_chunk_count,
617 post_process_position_count = scheduling.post_process_position_count,
618 readiness_candidate_count = scheduling.readiness_candidate_count,
619 ticking_snapshot_rebuild = ?scheduling.ticking_snapshot_rebuild,
620 rebuilt_ticking_chunk_count = scheduling.rebuilt_ticking_chunk_count,
621 schedule_generation = ?scheduling.schedule_generation,
622 scheduled_count = scheduling.scheduled_count,
623 run_generation = ?scheduling.run_generation,
624 process_unloads = ?scheduling.process_unloads,
625 broadcast_changes = ?cm.broadcast_changes,
626 collect_tickable = ?cm.collect_tickable,
627 tick_chunks = ?cm.tick_chunks,
628 tick_block_entities = ?cm.tick_block_entities,
629 tickable_count = cm.tickable_count,
630 total_chunks = cm.total_chunks,
631 lookup_cache_holder_hits = cm.lookup_cache.holder_hits,
632 lookup_cache_missing_hits = cm.lookup_cache.missing_hits,
633 lookup_cache_scc_lookups = cm.lookup_cache.scc_lookups,
634 lookup_cache_foreign_map_bypasses = cm.lookup_cache.foreign_map_bypasses,
635 lookup_cache_evictions = cm.lookup_cache.evictions,
636 "Game tick slow"
637 );
638 }
639 Ok(())
640 }
641
642 pub(super) fn tick_jobs(self: &Arc<Self>, tick_count: u64, runs_normally: bool) {
643 let stats = self
644 .jobs
645 .tick(Arc::downgrade(self), tick_count, runs_normally);
646 if stats.polled > 0 && stats.pending > 0 && tick_count.is_multiple_of(100) {
647 tracing::debug!(
648 polled = stats.polled,
649 finished = stats.finished,
650 pending = stats.pending,
651 "Server jobs pending"
652 );
653 }
654 }
655}
656
657#[cfg(test)]
658mod tests {
659 use std::sync::Arc;
660
661 use super::Server;
662 use crate::{
663 player::ResetReason,
664 test_support::{TestPlayerBuilder, fresh_test_world, insert_ready_full_chunk},
665 };
666 use rustc_hash::FxHashMap;
667 use steel_utils::ChunkPos;
668
669 #[test]
670 fn chunk_send_commit_rechecks_live_world_membership() {
671 let world = fresh_test_world("chunk_send_membership_revalidation");
672 let center = ChunkPos::new(0, 0);
673 insert_ready_full_chunk(&world, center);
674 let player = TestPlayerBuilder::new(Arc::clone(&world), "ChunkTester", 1).build();
675 assert!(world.add_player(Arc::clone(&player), ResetReason::InitialJoin));
676 assert!(world.players.remove_player_sync(&player).is_some());
677
678 let encoding_pool = rayon::ThreadPoolBuilder::new().num_threads(1).build();
679 let Ok(encoding_pool) = encoding_pool else {
680 panic!("test chunk encoding pool should initialize");
681 };
682 let mut encode_cache = FxHashMap::default();
683 Server::send_chunks_for_player(&player, &world, &mut encode_cache, &encoding_pool);
684
685 let sender = player.chunk_sender().lock();
686 assert!(sender.pending_chunks.contains(¢er));
687 assert!(!sender.is_chunk_sent(center));
688 assert_eq!(sender.unacknowledged_batch_count_for_test(), 0);
689 drop(sender);
690
691 assert!(world.players.insert(Arc::clone(&player)));
692 world.remove_player_for_world_change(&player);
693 }
694}