Skip to main content

steel_core/server/
mod.rs

1//! This module contains the `Server` struct, which is the main entry point for the server.
2mod broadcasting;
3/// Tick-polled server jobs.
4pub mod jobs;
5mod packet_processor;
6mod pregen;
7/// The registry cache for the server.
8pub mod registry_cache;
9mod run_loop;
10mod service_keys;
11mod tick_overload;
12/// The tick rate manager for the server.
13pub mod tick_rate_manager;
14mod world_tick_workers;
15/// Domain-aware loaded world map.
16pub mod worlds;
17
18use crate::bootstrap::init_globals;
19use crate::chunk::{
20    chunk_request::{ChunkRequest, ChunkRequestHandle, ChunkRequestState, ChunkTicketKind},
21    status::ChunkStatus,
22};
23use crate::command::brigadier::{StringReader, SuggestionError, Suggestions};
24use crate::command::execution::{
25    CommandExecutionContext, CommandResultCallback, CommandSource, ExecutionCommandSource,
26    ExecutionStop,
27};
28use crate::command::sender::{CommandExecutionOwner, CommandSender};
29use crate::command::storage::DomainCommandStorage;
30use crate::command::{
31    COMMAND_REQUESTS_PER_TICK, COMMAND_RESUMPTIONS_PER_TICK, CommandCompletion, CommandDispatcher,
32    CommandQueueFull, CommandRegistry, CommandRequest, CommandRequestQueue,
33    PendingCommandExecutionQueue, client_permission_event, command_suggestions_packet,
34    command_tree_packet, create_registered_dispatcher,
35};
36use crate::config::{ResolvedWorldConfig, RuntimeConfig, WorldsConfig, validate_login_security};
37use crate::entity::{
38    Entity, EntityBase, PendingWorldChangeToken, RemovalReason, SharedEntity, change_entity_world,
39};
40
41use crate::chunk_saver::{ChunkStorage, PersistentEntity, registry::WorldStorageRegistry};
42use crate::level_data::{GameTimeSource, LevelDataManager, RespawnData, WorldGenerationSettings};
43use crate::permission::{
44    OP_GROUP, PermissionGroupManager, PermissionGroupManagerError, PermissionGroupUpdateError,
45    PermissionGroupsConfig, PermissionMetadataExpression, PermissionRuleExpression, PermissionSet,
46    PermissionSubjectIndex, PermissionSubjectState,
47};
48use crate::player::chunk_sender::{ChunkSender, EncodedChunk};
49use crate::player::connection::NetworkConnection;
50use crate::player::connection::ScheduledPlayPacket;
51use crate::player::player_data::{
52    PersistentEnderPearl, PersistentPlayerData, PersistentRootVehicle,
53};
54use crate::player::player_data_storage::{GlobalPlayerData, PlayerDataStorage};
55use crate::player::player_inventory::MenuRemovalStatus;
56use crate::player::{
57    DomainResidenceToken, GameProfile, KnownPlayer, KnownPlayerNameLookup, KnownPlayers, Player,
58    ProfileLookupError, ResetReason, is_valid_player_name, lookup_online_profile, offline_uuid,
59};
60use crate::portal::{
61    PortalKind, TeleportPostTransition, TeleportTransition, WorldChangeRequest, end_gateway,
62    end_portal, nether_portal,
63};
64use crate::scoreboard::DomainScoreboards;
65use crate::server::jobs::{FnServerJob, ServerJobContext, ServerJobQueue};
66use crate::server::packet_processor::PacketProcessor;
67pub(crate) use crate::server::packet_processor::PlayerPacketTransition;
68use crate::server::registry_cache::RegistryCache;
69use crate::server::service_keys::ServiceKeyStore;
70use crate::server::worlds::WorldMap;
71use crate::world::player_spawn_finder::{PlayerSpawnSearch, PlayerSpawnSearchPoll};
72use crate::world::{PlayerMap, World, WorldConfig};
73use crate::worldgen::WorldGeneratorRegistry;
74use crate::worldgen::registry::GeneratorOutput;
75use crossbeam::queue::SegQueue;
76use glam::DVec3;
77use rayon::{ThreadPool, ThreadPoolBuilder};
78use rustc_hash::FxHashMap;
79use std::sync::atomic::AtomicI32;
80use std::{
81    collections::BTreeSet,
82    io, mem,
83    sync::{Arc, mpsc},
84    time::{Duration, Instant},
85};
86use steel_crypto::{key_store::KeyStore, signature::ProfileKeyValidator};
87use steel_protocol::packet_traits::{ClientPacket, EncodedPacket};
88use steel_protocol::packets::game::{
89    CCommandSuggestions, CEntityEvent, CLogin, CPlayerInfoUpdate, CRemovePlayerInfo,
90    CSetDefaultSpawnPosition, CSystemChat, CTabList, CTickingState, CTickingStep,
91    CommonPlayerSpawnInfo, RelativeMovement,
92};
93use steel_protocol::utils::ConnectionProtocol;
94use steel_registry::vanilla_game_rules::{
95    ALLOW_ENTERING_NETHER_USING_PORTALS, IMMEDIATE_RESPAWN, LIMITED_CRAFTING, REDUCED_DEBUG_INFO,
96};
97use steel_registry::{
98    RegistryEntry, dimension_type::DimensionTypeRef, vanilla_dimension_types, vanilla_entities,
99};
100use steel_utils::{
101    BlockPos, ChunkPos, Identifier,
102    locks::{AsyncMutex, SyncMutex, SyncRwLock},
103    text::DisplayResolutor,
104    threading::{DEBUG_STACK_SIZE, available_worker_threads},
105    translations,
106};
107use text_components::{Modifier, TextComponent, format::Color};
108use tick_rate_manager::{SprintReport, TickRateManager};
109use tokio::{
110    runtime::Runtime,
111    sync::Notify,
112    task::{JoinSet, spawn_blocking},
113    time::sleep,
114};
115use tokio_util::sync::CancellationToken;
116use uuid::Uuid;
117
118/// Interval in ticks between tab list updates (20 ticks = 1 second).
119const TAB_LIST_UPDATE_INTERVAL: u64 = 20;
120/// Interval in ticks between player info broadcasts (600 ticks = 30 seconds).
121/// Matches vanilla `PlayerList.SEND_PLAYER_INFO_INTERVAL`.
122const SEND_PLAYER_INFO_INTERVAL: u64 = 600;
123/// Wall-clock interval between saves of command-owned persistent server data.
124/// Matches vanilla's intended five-minute autosave cadence.
125const COMMAND_DATA_AUTOSAVE_INTERVAL: Duration = Duration::from_secs(300);
126
127#[derive(Clone, Copy)]
128struct TabListTickStats {
129    tps: f32,
130    recent_mspt: f32,
131    average_mspt: f32,
132    p95_mspt: f32,
133}
134
135impl TabListTickStats {
136    fn capture(tick_manager: &TickRateManager) -> Self {
137        Self {
138            tps: tick_manager.get_tps(),
139            recent_mspt: tick_manager.get_smoothed_mspt(),
140            average_mspt: tick_manager.get_average_mspt(),
141            p95_mspt: tick_manager.get_p95(),
142        }
143    }
144}
145
146/// Results from saving every command-owned persistent data set.
147pub struct CommandDataSaveResults {
148    /// Number of dirty domain scoreboards written, or the save error.
149    pub scoreboards: io::Result<usize>,
150    /// Number of dirty domain command-storage values written, or the save error.
151    pub storage: io::Result<usize>,
152}
153
154mod known_players;
155
156use known_players::KnownPlayerCacheState;
157
158/// Tick rate for the chunk sending loop.
159const CHUNK_SENDING_TPS: u64 = 20;
160
161/// Work duration at which background chunk work is considered slow.
162const SLOW_CHUNK_TICK_THRESHOLD: Duration = Duration::from_millis(50);
163
164fn configured_chunk_generation_threads(configured_threads: Option<usize>) -> Option<usize> {
165    cap_positive_thread_count(configured_threads, available_worker_threads())
166}
167
168fn configured_chunk_encoding_threads(configured_threads: Option<usize>) -> Option<usize> {
169    cap_positive_thread_count(configured_threads, available_worker_threads())
170}
171
172fn cap_positive_thread_count(
173    configured_threads: Option<usize>,
174    available_threads: usize,
175) -> Option<usize> {
176    let configured_threads = configured_threads.filter(|&threads| threads > 0)?;
177    Some(configured_threads.min(available_threads.max(1)))
178}
179
180#[cfg(test)]
181mod tests;
182
183#[derive(Clone, Copy)]
184struct PreparedSpawn {
185    position: DVec3,
186    rotation: (f32, f32),
187}
188
189fn apply_default_spawn(player: &Arc<Player>, world: &Arc<World>, spawn: PreparedSpawn) {
190    player.base().set_position_local(spawn.position);
191    player.set_rotation(spawn.rotation);
192    player.restore_game_modes(world.default_gamemode, None);
193    player
194        .abilities
195        .lock()
196        .update_for_game_mode(world.default_gamemode);
197}
198
199fn is_allowed_to_enter_portal(source_world: &World, target_world: &World) -> bool {
200    is_allowed_to_enter_portal_target(
201        is_nether_dimension_type(target_world),
202        source_world.get_game_rule(&ALLOW_ENTERING_NETHER_USING_PORTALS),
203    )
204}
205
206const fn is_allowed_to_enter_portal_target(
207    target_is_nether: bool,
208    allow_entering_nether_using_portals: bool,
209) -> bool {
210    if !target_is_nether {
211        return true;
212    }
213
214    allow_entering_nether_using_portals
215}
216
217fn can_teleport_between_worlds(
218    entity: &dyn Entity,
219    source_world: &World,
220    target_world: &World,
221    projectile_owner_seen_credits: impl Fn(&uuid::Uuid) -> Option<bool>,
222) -> bool {
223    if is_end_return_transition(source_world.dimension_type, target_world.dimension_type) {
224        return can_entity_return_from_end_to_overworld(entity, projectile_owner_seen_credits);
225    }
226
227    true
228}
229
230fn is_end_return_transition(
231    source_dimension_type: DimensionTypeRef,
232    target_dimension_type: DimensionTypeRef,
233) -> bool {
234    source_dimension_type == &vanilla_dimension_types::THE_END
235        && target_dimension_type == &vanilla_dimension_types::OVERWORLD
236}
237
238fn is_nether_dimension_type(world: &World) -> bool {
239    world.dimension_type == &vanilla_dimension_types::THE_NETHER
240}
241
242fn can_entity_return_from_end_to_overworld(
243    entity: &dyn Entity,
244    projectile_owner_seen_credits: impl Fn(&uuid::Uuid) -> Option<bool>,
245) -> bool {
246    if entity.entity_type() == &vanilla_entities::ENDER_PEARL
247        && entity
248            .projectile_owner_uuid()
249            .and_then(|uuid| projectile_owner_seen_credits(&uuid))
250            == Some(false)
251    {
252        return false;
253    }
254
255    direct_passengers_allow_end_return(entity)
256}
257
258fn direct_passengers_allow_end_return(entity: &dyn Entity) -> bool {
259    for passenger in entity.passengers() {
260        if passenger
261            .as_player()
262            .is_some_and(|player| !player.has_seen_credits())
263        {
264            return false;
265        }
266    }
267
268    true
269}
270
271fn local_respawn_data_for_world(world: &World) -> RespawnData {
272    let level_data = world.level_data.read();
273    let data = level_data.data();
274    RespawnData::of(world.key.clone(), data.spawn_pos(), data.spawn.angle, 0.0)
275}
276
277fn generation_settings_for_world(
278    world_entry: &ResolvedWorldConfig,
279    generator_output: &GeneratorOutput,
280) -> WorldGenerationSettings {
281    WorldGenerationSettings::from_generator_config(
282        world_entry.generator_config.generator().clone(),
283        &generator_output.config,
284        generator_output.dimension_type.key.clone(),
285        generator_output.dimension_type.min_y,
286        generator_output.dimension_type.height,
287    )
288}
289
290fn world_config_registries() -> Result<(WorldGeneratorRegistry, WorldStorageRegistry), String> {
291    let generator_registry = WorldGeneratorRegistry::new_with_builtins()
292        .map_err(|e| format!("failed to initialize world generator registry: {e}"))?;
293    let storage_registry = WorldStorageRegistry::new_with_builtins()
294        .map_err(|e| format!("failed to initialize world storage registry: {e}"))?;
295    Ok((generator_registry, storage_registry))
296}
297
298struct DomainPlayerState {
299    world: Arc<World>,
300    data: DomainPlayerData,
301    spawn_chunk_request: ChunkRequestHandle,
302}
303
304struct UnpreparedDomainPlayerState {
305    world: Arc<World>,
306    explicit_target: bool,
307    data: UnpreparedDomainPlayerData,
308}
309
310enum UnpreparedDomainPlayerData {
311    SavedRestored { data: Box<PersistentPlayerData> },
312    SavedWithoutLocation { data: Box<PersistentPlayerData> },
313    FirstVisit,
314}
315
316enum DomainPlayerData {
317    SavedRestored {
318        data: Box<PersistentPlayerData>,
319    },
320    SavedWithoutLocation {
321        data: Box<PersistentPlayerData>,
322        spawn: PreparedSpawn,
323    },
324    FirstVisit {
325        spawn: PreparedSpawn,
326    },
327}
328
329struct DomainSwitchRequest {
330    player: Arc<Player>,
331    target_domain: String,
332    target_world: Option<Arc<World>>,
333    pending_token: PendingWorldChangeToken,
334}
335
336/// Failure while atomically editing one player's persisted permission state.
337#[derive(Debug, thiserror::Error)]
338pub enum PlayerPermissionUpdateError<E> {
339    /// The caller rejected the proposed edit.
340    #[error("{0}")]
341    Edit(E),
342    /// The edit assigns a group that is not configured.
343    #[error("unknown permission group '{0}'")]
344    UnknownGroup(String),
345    /// The permission snapshot could not be persisted.
346    #[error("failed to update player permissions: {0}")]
347    Storage(io::Error),
348}
349
350impl<E> From<io::Error> for PlayerPermissionUpdateError<E> {
351    fn from(value: io::Error) -> Self {
352        Self::Storage(value)
353    }
354}
355
356mod permissions;
357
358#[cfg(test)]
359use permissions::validate_player_permission_group_update;
360
361mod player_admission;
362mod player_lifecycle;
363
364pub use player_admission::{DuplicatePlayerWaitError, PlayerJoinReservation};
365use player_admission::{PlayerAdmissionState, PlayerDisconnectQueue, PlayerJoinQueue};
366
367mod world_changes;
368
369use jobs::domain_switch::DomainSwitchJob;
370use jobs::teleport::{
371    EndGatewayTeleportJob, EndPortalTeleportJob, EnderPearlRestoreJob, NetherPortalTeleportJob,
372    RootVehicleRestoreJob, WorldSpawnTeleportJob, clear_pending_world_change,
373    portal_entity_still_valid,
374};
375
376/// The main server struct.
377pub struct Server {
378    /// Runtime configuration (view distance, compression, etc.).
379    pub config: Arc<RuntimeConfig>,
380    /// Runtime permission groups and their persistence boundary.
381    pub permission_groups: PermissionGroupManager,
382    /// The cancellation token for graceful shutdown.
383    pub cancel_token: CancellationToken,
384    /// The key store for the server.
385    pub key_store: KeyStore,
386    /// The registry cache for the server.
387    pub registry_cache: RegistryCache,
388    /// A list of all the worlds on the server.
389    pub worlds: WorldMap,
390    /// Players currently connected to the server, independent of world membership.
391    online_players: PlayerMap,
392    /// UUIDs reserved by a join or disconnect/save lifecycle transition.
393    player_admissions: SyncMutex<FxHashMap<Uuid, PlayerAdmissionState>>,
394    /// Wakes verified logins waiting for an older session with the same UUID to leave.
395    player_admission_changed: Notify,
396    /// Wakes connection lifecycle work waiting for a specific server tick.
397    server_tick_changed: Notify,
398    /// The tick rate manager for the server.
399    pub tick_rate_manager: SyncRwLock<TickRateManager>,
400    /// The number of minutes required for a player to be idle for them to be kicked (timed out) from the server.
401    /// If this is equal to 0, no kicking will happen.
402    pub player_idle_timeout: AtomicI32,
403    /// Command scoreboards isolated by Steel domain.
404    pub scoreboards: DomainScoreboards,
405    /// Command NBT storage isolated by Steel domain.
406    pub(crate) command_storage: DomainCommandStorage,
407    /// Saves and dispatches commands to appropriate handlers.
408    command_dispatcher: SyncRwLock<CommandDispatcher>,
409    /// Steel-owned permission keys exposed for command autocomplete.
410    command_permission_keys: Vec<String>,
411    /// Command work submitted from connection and console tasks.
412    command_requests: CommandRequestQueue,
413    /// Decoded serverbound play packets handled during the inter-tick phase.
414    packet_processor: PacketProcessor,
415    /// Dedicated worker pool for CPU-heavy chunk persistence and packet encoding.
416    chunk_encoding_pool: Arc<ThreadPool>,
417    /// Jobs resumed from a known point in the server game tick.
418    pub jobs: ServerJobQueue,
419    /// Player data storage for saving/loading player state.
420    pub player_data_storage: PlayerDataStorage,
421    /// Persisted permission state indexed by player UUID.
422    player_permission_states: SyncRwLock<PermissionSubjectIndex>,
423    /// Serializes persistence and cache publication for player permission edits.
424    player_permission_updates: AsyncMutex<()>,
425    /// Player identities and coalesced persistence state.
426    known_players: SyncMutex<KnownPlayerCacheState>,
427    /// Wakes shutdown when the single known-player save worker becomes idle.
428    known_player_save_idle: Notify,
429    /// HTTP client used by online-mode name-to-profile lookups.
430    profile_lookup_client: reqwest::Client,
431    /// Cached Mojang service keys used to validate player-key certificates.
432    service_keys: Arc<ServiceKeyStore>,
433    /// Player joins prepared by async I/O and finalized at the game tick safe point.
434    pending_player_joins: PlayerJoinQueue,
435    /// Disconnected players waiting to be detached at the next game tick safe point.
436    pending_player_disconnects: PlayerDisconnectQueue,
437    /// Queued world changes to process after the tick.
438    pub pending_world_changes: SyncMutex<Vec<(SharedEntity, WorldChangeRequest)>>,
439    /// Queued domain switches to process after world ticks.
440    pending_domain_switches: SyncMutex<Vec<DomainSwitchRequest>>,
441}
442
443struct GameTickTaskGuard {
444    server: Arc<Server>,
445    cancel_token: CancellationToken,
446}
447
448impl GameTickTaskGuard {
449    const fn new(server: Arc<Server>, cancel_token: CancellationToken) -> Self {
450        Self {
451            server,
452            cancel_token,
453        }
454    }
455}
456
457impl Drop for GameTickTaskGuard {
458    fn drop(&mut self) {
459        self.server.packet_processor.stop();
460        self.cancel_token.cancel();
461    }
462}
463
464impl Server {
465    /// Returns the current server tick number.
466    pub fn current_tick(&self) -> u64 {
467        self.tick_rate_manager.read().tick_count
468    }
469
470    /// Waits until the server reaches `target_tick`.
471    pub async fn wait_until_tick(&self, target_tick: u64) {
472        loop {
473            let tick_changed = self.server_tick_changed.notified();
474            tokio::pin!(tick_changed);
475            tick_changed.as_mut().enable();
476
477            if self.current_tick() >= target_tick {
478                return;
479            }
480
481            tick_changed.await;
482        }
483    }
484
485    pub(crate) fn permission_rule_suggestions(&self) -> Vec<String> {
486        let mut suggestions = self
487            .command_permission_keys
488            .iter()
489            .cloned()
490            .collect::<BTreeSet<_>>();
491        let config = self.permission_groups.config_snapshot();
492        for group in config.groups.values() {
493            suggestions.extend(group.allow.iter().cloned());
494            suggestions.extend(group.deny.iter().cloned());
495        }
496        for (_, state) in self.player_permission_states.read().entries() {
497            suggestions.extend(state.overrides().entries().iter().map(|entry| {
498                PermissionRuleExpression::new(entry.key().clone(), entry.context().clone())
499                    .to_string()
500            }));
501        }
502        suggestions.into_iter().collect()
503    }
504
505    pub(crate) fn permission_metadata_suggestions(&self) -> Vec<String> {
506        let mut suggestions = BTreeSet::new();
507        let config = self.permission_groups.config_snapshot();
508        for group in config.groups.values() {
509            suggestions.extend(group.metadata.iter().map(|rule| rule.key.clone()));
510        }
511        for (_, state) in self.player_permission_states.read().entries() {
512            suggestions.extend(state.metadata_overrides().entries().iter().map(|entry| {
513                PermissionMetadataExpression::new(entry.key().clone(), entry.context().clone())
514                    .to_string()
515            }));
516        }
517        suggestions.into_iter().collect()
518    }
519
520    /// Creates a new server with only Steel's built-in commands.
521    pub async fn new(
522        chunk_runtime: Arc<Runtime>,
523        cancel_token: CancellationToken,
524        config: RuntimeConfig,
525        worlds_config: WorldsConfig,
526        permission_groups: PermissionGroupManager,
527    ) -> Result<Self, String> {
528        Self::new_with_commands(
529            chunk_runtime,
530            cancel_token,
531            config,
532            worlds_config,
533            permission_groups,
534            CommandRegistry::new(),
535        )
536        .await
537    }
538
539    /// Creates a new server and atomically merges startup command extensions after built-ins.
540    #[expect(
541        clippy::too_many_lines,
542        reason = "server initialization is a single cohesive flow"
543    )]
544    pub async fn new_with_commands(
545        chunk_runtime: Arc<Runtime>,
546        cancel_token: CancellationToken,
547        config: RuntimeConfig,
548        worlds_config: WorldsConfig,
549        permission_groups: PermissionGroupManager,
550        command_registry: CommandRegistry,
551    ) -> Result<Self, String> {
552        validate_login_security(config.online_mode, config.encryption).map_err(str::to_owned)?;
553        let config = Arc::new(config);
554        init_globals();
555        log::info!(
556            "SteelMC is not affiliated with Mojang or Microsoft. Use is subject to the Minecraft EULA: https://aka.ms/MinecraftEULA"
557        );
558
559        // Authlib starts this fetch alongside server initialization and waits on first use.
560        // It runs whatever the login mode, because `handle_chat_session_update` reads these
561        // keys with no online-mode gate, as vanilla does.
562        let service_keys = Arc::new(
563            ServiceKeyStore::new(config.services_server.as_deref())
564                .map_err(|error| format!("failed to configure Minecraft services keys: {error}"))?,
565        );
566        let service_keys_ready = service_keys.start(cancel_token.clone());
567
568        let registry_cache = RegistryCache::new(config.compression);
569
570        let (generator_registry, storage_registry) = world_config_registries()?;
571        let resolved_worlds = worlds_config
572            .validate_and_resolve(&generator_registry, &storage_registry)
573            .map_err(|e| format!("failed to validate worlds.toml: {e}"))?;
574
575        let mut world_storage = storage_registry.resolve_worlds(&resolved_worlds)?;
576
577        let generation_pool: Arc<ThreadPool> = Arc::new({
578            let mut builder = ThreadPoolBuilder::new().thread_name(|i| format!("rayon-gen-{i}"));
579            if let Some(chunk_generation_threads) =
580                configured_chunk_generation_threads(config.chunk_generation_threads)
581            {
582                builder = builder.num_threads(chunk_generation_threads);
583            }
584            // Debug builds have deep call chains in density functions that overflow the default 2 MB stack
585            if cfg!(debug_assertions) {
586                builder = builder.stack_size(DEBUG_STACK_SIZE);
587            }
588            builder
589                .build()
590                .map_err(|e| format!("failed to create generation thread pool: {e}"))?
591        });
592        let chunk_encoding_pool = Arc::new({
593            let mut builder =
594                ThreadPoolBuilder::new().thread_name(|i| format!("rayon-chunk-encode-{i}"));
595            if let Some(chunk_encoding_threads) =
596                configured_chunk_encoding_threads(config.chunk_encoding_threads)
597            {
598                builder = builder.num_threads(chunk_encoding_threads);
599            }
600            builder
601                .build()
602                .map_err(|e| format!("failed to create chunk encoding thread pool: {e}"))?
603        });
604
605        let player_data_storage = PlayerDataStorage::from_selection(
606            resolved_worlds.save_path.clone(),
607            &resolved_worlds.player_storage,
608        )
609        .await
610        .map_err(|e| format!("failed to create player data storage: {e}"))?;
611        let player_permission_states = player_data_storage
612            .load_permission_subjects()
613            .await
614            .map_err(|error| format!("failed to load player permissions: {error}"))?;
615        let known_players = player_data_storage
616            .load_known_players()
617            .await
618            .map_err(|error| format!("failed to load known players: {error}"))?;
619        let mut worlds = WorldMap::new(
620            resolved_worlds.default_domain.clone(),
621            &resolved_worlds.domains,
622            &resolved_worlds.worlds,
623        );
624
625        let mut construct_world = async |world_entry: &ResolvedWorldConfig,
626                                         game_time_source: GameTimeSource|
627               -> Result<Arc<World>, String> {
628            let storage_output = world_storage
629                .remove(&world_entry.key)
630                .ok_or_else(|| format!("world {} has no resolved storage", world_entry.key))?;
631            let world_seed = LevelDataManager::load_seed_or_default(
632                storage_output.level_data_path.as_deref(),
633                world_entry.seed,
634            )
635            .await
636            .map_err(|e| {
637                format!(
638                    "failed to load level data seed for {}: {e}",
639                    world_entry.key
640                )
641            })?;
642            let generator_output = generator_registry
643                .create(
644                    storage_output.level_data_path.as_deref(),
645                    &world_entry.generator_config,
646                    world_seed,
647                    generation_pool.clone(),
648                )
649                .map_err(|e| format!("failed to create generator for {}: {e}", world_entry.key))?;
650            let generation_settings = generation_settings_for_world(world_entry, &generator_output);
651            let world = World::new_with_config_and_encoding_pool(
652                chunk_runtime.clone(),
653                world_entry.key.clone(),
654                generator_output.dimension_type,
655                world_seed,
656                WorldConfig {
657                    game_time_source,
658                    storage: storage_output.storage,
659                    level_data_path: storage_output
660                        .level_data_path
661                        .map(|path| path.to_string_lossy().into_owned()),
662                    generator: Arc::new(generator_output.generator),
663                    generation_settings,
664                    view_distance: config.view_distance,
665                    simulation_distance: config.simulation_distance,
666                    max_chained_neighbor_updates: config.max_chained_neighbor_updates,
667                    compression: config.compression,
668                    is_flat: generator_output.is_flat,
669                    sea_level: generator_output.sea_level,
670                    default_gamemode: world_entry.default_gamemode,
671                    difficulty: world_entry.difficulty,
672                },
673                generation_pool.clone(),
674                Arc::clone(&chunk_encoding_pool),
675            )
676            .await
677            .map_err(|e| format!("failed to create world {}: {e}", world_entry.key))?;
678            world
679                .initialize_spawn_if_needed()
680                .await
681                .map_err(|e| format!("failed to initialize spawn for {}: {e}", world_entry.key))?;
682            Ok(world)
683        };
684        for domain in &resolved_worlds.domains {
685            let primary_config = resolved_worlds
686                .worlds
687                .iter()
688                .find(|world| world.key == domain.default_world && world.domain == domain.name)
689                .ok_or_else(|| {
690                    format!(
691                        "domain {} has no configured primary {}",
692                        domain.name, domain.default_world
693                    )
694                })?;
695            let primary = construct_world(primary_config, GameTimeSource::Primary).await?;
696            let clock = Arc::clone(&primary.game_time);
697            worlds.insert(primary_config.key.clone(), primary);
698            for world_entry in resolved_worlds
699                .worlds
700                .iter()
701                .filter(|world| world.domain == domain.name && world.key != domain.default_world)
702            {
703                let world =
704                    construct_world(world_entry, GameTimeSource::Derived(Arc::clone(&clock)))
705                        .await?;
706                worlds.insert(world_entry.key.clone(), world);
707            }
708        }
709        worlds.validate_game_times()?;
710
711        let scoreboards = DomainScoreboards::load(&worlds)
712            .await
713            .map_err(|error| format!("failed to load domain scoreboards: {error}"))?;
714        let command_storage = DomainCommandStorage::load(&worlds)
715            .await
716            .map_err(|error| format!("failed to load domain command storage: {error}"))?;
717        let registered_commands = create_registered_dispatcher(command_registry)
718            .map_err(|error| format!("failed to register commands: {error}"))?;
719        let command_permission_keys = registered_commands
720            .permissions
721            .into_iter()
722            .map(|permission| permission.as_str().to_owned())
723            .collect();
724
725        // Steel finishes the initial attempt before opening its listener, except offline,
726        // where `enforces_secure_chat` needs online mode so nothing acts on the result.
727        if config.online_mode && service_keys_ready.await.is_err() {
728            log::error!("Minecraft services key fetch task stopped before its initial attempt");
729        }
730
731        Ok(Server {
732            config,
733            permission_groups,
734            cancel_token,
735            key_store: KeyStore::create(),
736            worlds,
737            online_players: PlayerMap::new(),
738            player_admissions: SyncMutex::new(FxHashMap::default()),
739            player_admission_changed: Notify::new(),
740            server_tick_changed: Notify::new(),
741            registry_cache,
742            tick_rate_manager: SyncRwLock::new(TickRateManager::new()),
743            player_idle_timeout: AtomicI32::new(0),
744            scoreboards,
745            command_storage,
746            command_dispatcher: SyncRwLock::new(registered_commands.dispatcher),
747            command_permission_keys,
748            command_requests: CommandRequestQueue::new(),
749            packet_processor: PacketProcessor::new(),
750            chunk_encoding_pool,
751            jobs: ServerJobQueue::new(),
752            player_data_storage,
753            player_permission_states: SyncRwLock::new(player_permission_states),
754            player_permission_updates: AsyncMutex::new(()),
755            known_players: SyncMutex::new(KnownPlayerCacheState::new(known_players)),
756            known_player_save_idle: Notify::new(),
757            profile_lookup_client: reqwest::Client::new(),
758            service_keys,
759            pending_player_joins: PlayerJoinQueue::new(),
760            pending_player_disconnects: PlayerDisconnectQueue::new(),
761            pending_world_changes: SyncMutex::new(vec![]),
762            pending_domain_switches: SyncMutex::new(vec![]),
763        })
764    }
765
766    /// Returns the current player-certificate validator, if service keys are available.
767    pub fn profile_key_signature_validator(&self) -> Option<Arc<ProfileKeyValidator>> {
768        self.service_keys.profile_key_validator()
769    }
770
771    /// Returns whether secure chat can currently be enforced.
772    #[must_use]
773    pub fn enforces_secure_chat(&self) -> bool {
774        self.config.enforce_secure_chat
775            && self.config.online_mode
776            && self.profile_key_signature_validator().is_some()
777    }
778
779    /// Saves all dirty domain command storage through domain default worlds.
780    pub async fn save_command_storage(&self) -> io::Result<usize> {
781        self.command_storage.save(&self.worlds).await
782    }
783
784    /// Saves all command-owned persistent data while allowing each data set to fail independently.
785    pub async fn save_command_data(&self) -> CommandDataSaveResults {
786        CommandDataSaveResults {
787            scoreboards: self.scoreboards.save(&self.worlds).await,
788            storage: self.save_command_storage().await,
789        }
790    }
791
792    /// Queues a command for execution at the start of the next game tick.
793    pub fn submit_command(
794        &self,
795        sender: CommandSender,
796        command: String,
797    ) -> Result<(), CommandQueueFull> {
798        self.command_requests.submit(CommandRequest::Execute {
799            owner: CommandExecutionOwner::capture(sender, self),
800            command,
801        })
802    }
803
804    pub(crate) fn submit_command_suggestions(
805        &self,
806        player: Arc<Player>,
807        transaction_id: i32,
808        input: String,
809    ) -> Result<(), CommandQueueFull> {
810        self.command_requests.submit(CommandRequest::Suggestions {
811            owner: CommandExecutionOwner::capture(CommandSender::Player(player), self),
812            transaction_id,
813            input,
814        })
815    }
816
817    /// Schedules a decoded play packet for the inter-tick packet phase.
818    pub(crate) fn schedule_play_packet(
819        &self,
820        player: Arc<Player>,
821        packet: ScheduledPlayPacket,
822        payload_bytes: usize,
823    ) {
824        self.packet_processor
825            .schedule(player, packet, payload_bytes);
826    }
827
828    /// Pauses later packets while `player` is replaced by a new incarnation.
829    pub(crate) fn begin_player_packet_transition(
830        &self,
831        player: &Arc<Player>,
832    ) -> Option<PlayerPacketTransition> {
833        if player.connection.closed() || !player.session.is_current_player(player) {
834            return None;
835        }
836        self.packet_processor.pause_player_session(&player.session)
837    }
838
839    /// Resumes packets retained by an exact player-replacement transition.
840    pub(crate) fn finish_player_packet_transition(
841        &self,
842        transition: PlayerPacketTransition,
843    ) -> bool {
844        self.packet_processor.resume_player_session(transition)
845    }
846
847    /// Discards all pending packet work for a closed player session.
848    pub(crate) fn discard_player_packets(&self, player: &Player) {
849        self.packet_processor
850            .discard_player_session(&player.session);
851    }
852
853    /// Returns Brigadier completions visible to a command sender.
854    pub fn command_completions(
855        self: &Arc<Self>,
856        sender: CommandSender,
857        input: &str,
858    ) -> Vec<CommandCompletion> {
859        if !CommandExecutionOwner::capture(sender.clone(), self).is_current(self) {
860            return Vec::new();
861        }
862        match self.build_command_suggestions(sender, input) {
863            Ok(suggestions) => {
864                let range = suggestions.range();
865                suggestions
866                    .list()
867                    .iter()
868                    .map(|suggestion| {
869                        CommandCompletion::new(
870                            range.start(),
871                            range.len(),
872                            suggestion.text().to_owned(),
873                        )
874                    })
875                    .collect()
876            }
877            Err(error) => {
878                tracing::warn!(%error, "failed to build command suggestions");
879                Vec::new()
880            }
881        }
882    }
883}