""" Orchestration layer for Commonwealth Online server. Wraps the relay server lifecycle, configuration, and admin operations for both CLI and future GUI host applications. """ from __future__ import annotations import json import re import threading from dataclasses import dataclass, asdict from pathlib import Path from typing import Any, Callable from admin_server import AdminServer, DEFAULT_ADMIN_PORT from gns_gameplay_adapter import GnsGameplayAdapter from gns_transport import GnsServerTransport from server_core import FalloutTogetherServer from transport_server import TransportAwareFalloutTogetherServer _IPV4_RE = re.compile(r"^\d{1,3}(?:\.\d{1,3}){3}$") @dataclass class ServerConfig: """Server configuration.""" host: str = "0.0.0.0" port: int = 7777 server_name: str = "Commonwealth Online Server" server_description: str = "" max_players: int = 16 log_verbosity: str = "info" admin_port: int = DEFAULT_ADMIN_PORT bans_path: str | None = None enable_gns_transport: bool = False gns_bridge_path: str | None = None @dataclass class ClientSnapshot: """Snapshot of a connected client.""" player_id: int address: str connected_at: float packets_sent: int packets_received: int label: str @dataclass class ServerStats: """Server statistics snapshot.""" is_running: bool host: str port: str server_name: str server_description: str uptime_seconds: float connected_clients: int clients: list[ClientSnapshot] packets_received: int packets_sent: int transform_packets_received: int transform_packets_broadcast: int world_state_packets_received: int world_state_packets_broadcast: int class ServerService: """ High-level server orchestration facade. Provides lifecycle management, configuration application, and admin operations for the underlying FalloutTogetherServer relay. """ def __init__(self, config: ServerConfig | None = None) -> None: self.config = config or ServerConfig() self._server: FalloutTogetherServer | None = None self._admin: AdminServer | None = None self._gns_adapter: GnsGameplayAdapter | None = None self._log_listeners: list[Callable[..., None]] = [] self._log_lock = threading.RLock() self._serve_thread: threading.Thread | None = None self._running = False self._stop_requested = False self._stop_event = threading.Event() def add_log_listener(self, callback: Callable[..., None]) -> None: """Register a callback for server log messages.""" with self._log_lock: if callback not in self._log_listeners: self._log_listeners.append(callback) def remove_log_listener(self, callback: Callable[..., None]) -> None: """Unregister a log callback.""" with self._log_lock: if callback in self._log_listeners: self._log_listeners.remove(callback) def _dispatch_log(self, message: str, *, level: str = "info") -> None: """Dispatch a log message to all registered listeners.""" with self._log_lock: listeners = list(self._log_listeners) for listener in listeners: try: try: listener(message, level=level) except TypeError: listener(message) except Exception: pass def _resolve_bans_path(self) -> str: if self.config.bans_path: return str(Path(self.config.bans_path)) return str(Path(__file__).resolve().parent / "bans.json") def _create_server(self) -> FalloutTogetherServer: return TransportAwareFalloutTogetherServer( host=self.config.host, port=self.config.port, server_name=self.config.server_name, server_description=self.config.server_description, max_players=self.config.max_players, bans_path=self._resolve_bans_path(), log_verbosity=self.config.log_verbosity, ) def _start_admin(self) -> None: self._admin = AdminServer( handler=self.handle_admin_request, port=self.config.admin_port, log=self._dispatch_log, ) self._admin.start() def _stop_admin(self) -> None: if self._admin is not None: self._admin.stop() self._admin = None def _start_gns(self) -> None: if not self.config.enable_gns_transport or self._server is None: return if self._gns_adapter is not None: return transport = GnsServerTransport( self.config.host, self._server.port, library_path=self.config.gns_bridge_path, ) if transport.local_port != self._server.port: actual_port = transport.local_port transport.close() raise RuntimeError( f"GNS transport bound UDP {actual_port}, expected UDP {self._server.port}." ) adapter = GnsGameplayAdapter(self._server, transport) try: adapter.start() except Exception: transport.close() raise self._gns_adapter = adapter self._dispatch_log( f"GameNetworkingSockets gameplay transport listening on UDP {self._server.port}." ) def _stop_gns(self) -> None: adapter = self._gns_adapter self._gns_adapter = None if adapter is not None: adapter.stop() def start(self) -> None: """Start the server in a background thread.""" if self._running: self._dispatch_log("Server is already running.") return self._stop_requested = False self._stop_event.clear() self._server = self._create_server() self._server.add_log_listener(self._dispatch_log) self._running = True self._start_admin() self._serve_thread = threading.Thread(target=self._serve_forever, daemon=True) self._serve_thread.start() self._dispatch_log("Server started in background thread.") def serve_forever(self) -> None: """Start the server and block until shutdown.""" if self._running: raise RuntimeError("Server is already running.") self._stop_requested = False self._stop_event.clear() self._server = self._create_server() self._server.add_log_listener(self._dispatch_log) self._running = True self._start_admin() self._run_server_loop() def _run_server_loop(self) -> None: try: if self._server is None: return self._server.start() self._start_gns() while not self._stop_requested and self._server.is_running(): self._stop_event.wait(0.1) finally: self._stop_gns() self._stop_admin() self._running = False def _serve_forever(self) -> None: """Internal server loop for background thread.""" self._run_server_loop() def stop(self) -> None: """Stop the server.""" if self._stop_requested and not self._running: return self._stop_requested = True self._stop_event.set() if not self._running and self._server is None: self._dispatch_log("Server is not running.") return self._stop_admin() self._stop_gns() if self._server is not None: self._server.stop() serve_thread = self._serve_thread if serve_thread is not None and serve_thread is not threading.current_thread(): serve_thread.join(timeout=2.0) self._serve_thread = None self._running = False self._dispatch_log("Server stopped.") def is_running(self) -> bool: """Check if the server is running.""" return self._running and (self._server is not None and self._server.is_running()) def get_stats(self) -> ServerStats: """Get current server statistics.""" if not self._server: return ServerStats( is_running=False, host=self.config.host, port=str(self.config.port), server_name=self.config.server_name, server_description=self.config.server_description, uptime_seconds=0.0, connected_clients=0, clients=[], packets_received=0, packets_sent=0, transform_packets_received=0, transform_packets_broadcast=0, world_state_packets_received=0, world_state_packets_broadcast=0, ) core_stats = self._server.get_stats() clients_data = self._server.get_clients() client_snapshots = [ ClientSnapshot( player_id=client["playerId"], address=f"{client['address']}:{client['port']}", connected_at=client["connectedAt"], packets_sent=client["packetsSent"], packets_received=client["packetsReceived"], label=f"{client['address']}:{client['port']}", ) for client in clients_data ] return ServerStats( is_running=core_stats.get("isRunning", False), host=core_stats.get("host", self.config.host), port=str(core_stats.get("port", self.config.port)), server_name=core_stats.get("serverName", self.config.server_name), server_description=core_stats.get( "serverDescription", self.config.server_description ), uptime_seconds=core_stats.get("uptimeSeconds", 0.0), connected_clients=core_stats.get("connectedClients", 0), clients=client_snapshots, packets_received=core_stats.get("packetsReceived", 0), packets_sent=core_stats.get("packetsSent", 0), transform_packets_received=core_stats.get("transformPacketsReceived", 0), transform_packets_broadcast=core_stats.get("transformPacketsBroadcast", 0), world_state_packets_received=core_stats.get("worldStatePacketsReceived", 0), world_state_packets_broadcast=core_stats.get("worldStatePacketsBroadcast", 0), ) def set_server_time(self, hhmm: str) -> tuple[bool, str]: """ Set server time (HHmm format). Returns (success, message). """ if not self._server: return False, "Server is not running." success = self._server.set_server_time(hhmm) if success: return True, f"Server time set to {hhmm}." else: return False, f"Invalid time format. Use HHmm (e.g., 1430 for 14:30)." def set_server_weather(self, fw_console_arg: str) -> tuple[bool, str]: """ Set server weather (form ID or preset name). Returns (success, message). """ if not self._server: return False, "Server is not running." success = self._server.set_server_weather(fw_console_arg) if success: return True, f"Server weather updated to {fw_console_arg}." else: return False, f"Invalid weather ID. Use an 8-digit hex form ID or preset name." def kick_player(self, player_id: int, reason: str = "") -> tuple[bool, str, dict[str, Any] | None]: if not self._server: return False, "Server is not running.", None try: result = self._server.kick_player(int(player_id), reason=reason) return True, f"Kicked player {player_id}.", result except KeyError as error: return False, str(error), None def ban_player(self, player_id: int, reason: str = "") -> tuple[bool, str, dict[str, Any] | None]: if not self._server: return False, "Server is not running.", None try: result = self._server.ban_player(int(player_id), reason=reason) return True, f"Banned player {player_id} (IP {result.get('ip')}).", result except KeyError as error: return False, str(error), None except ValueError as error: return False, str(error), None def ban_ip(self, ip: str, reason: str = "") -> tuple[bool, str, dict[str, Any] | None]: if not self._server: return False, "Server is not running.", None try: result = self._server.ban_ip(ip, reason=reason) return True, f"Banned IP {result.get('ip')}.", result except ValueError as error: return False, str(error), None def unban_ip(self, ip: str) -> tuple[bool, str]: if not self._server: return False, "Server is not running." removed = self._server.unban_ip(ip) if removed: return True, f"Unbanned IP {ip}." return False, f"IP {ip} is not banned." def list_bans(self) -> list[dict[str, Any]]: if not self._server: return [] return self._server.list_bans() def handle_admin_request(self, request: dict[str, Any]) -> dict[str, Any]: """Handle one admin JSON command from the localhost control channel.""" command = str(request.get("cmd") or request.get("command") or "").strip().lower() request_id = request.get("id") def ok(data: dict[str, Any] | None = None, message: str = "") -> dict[str, Any]: response: dict[str, Any] = {"ok": True} if request_id is not None: response["id"] = request_id if message: response["message"] = message if data is not None: response["data"] = data return response def fail(error: str) -> dict[str, Any]: response: dict[str, Any] = {"ok": False, "error": error} if request_id is not None: response["id"] = request_id return response if command in ("ping",): return ok({"pong": True}) if command in ("stats", "status"): stats = self.get_stats() return ok(asdict(stats)) if command in ("clients", "users"): if not self._server: return ok({"total_clients": 0, "clients": []}) clients = [] for client in self._server.get_clients(): address = str(client.get("address", "")) port = client.get("port") endpoint = f"{address}:{port}" if port is not None else address last_transform = client.get("lastTransform") clients.append( { "player_id": client.get("playerId"), "address": endpoint, "label": endpoint, "connected_at": client.get("connectedAt"), "packets_sent": client.get("packetsSent", 0), "packets_received": client.get("packetsReceived", 0), "last_transform": ( dict(last_transform) if isinstance(last_transform, dict) else None ), } ) return ok({"total_clients": len(clients), "clients": clients}) if command == "bans": return ok({"bans": self.list_bans()}) if command == "kick": player_id = request.get("playerId", request.get("player_id")) if player_id is None: return fail("kick requires playerId") reason = str(request.get("reason", "") or "") success, message, result = self.kick_player(int(player_id), reason=reason) return ok(result, message) if success else fail(message) if command == "ban": reason = str(request.get("reason", "") or "") player_id = request.get("playerId", request.get("player_id")) ip = request.get("ip") if player_id is not None: success, message, result = self.ban_player(int(player_id), reason=reason) return ok(result, message) if success else fail(message) if ip: success, message, result = self.ban_ip(str(ip), reason=reason) return ok(result, message) if success else fail(message) return fail("ban requires playerId or ip") if command == "unban": ip = request.get("ip") if not ip: return fail("unban requires ip") success, message = self.unban_ip(str(ip)) return ok({"ip": str(ip)}, message) if success else fail(message) if command == "world_time": hhmm = request.get("hhmm") or request.get("time") if not hhmm: return fail("world_time requires hhmm") success, message = self.set_server_time(str(hhmm)) return ok(message=message) if success else fail(message) if command == "world_weather": weather = request.get("weather") or request.get("fw") if not weather: return fail("world_weather requires weather") success, message = self.set_server_weather(str(weather)) return ok(message=message) if success else fail(message) return fail(f"Unknown admin command: {command or '(empty)'}") def stats_to_json(self, stats: ServerStats) -> str: """Serialize stats to JSON.""" data = asdict(stats) return json.dumps(data, indent=2) def looks_like_ipv4(value: str) -> bool: """Return True when value looks like a dotted IPv4 address.""" if not _IPV4_RE.match(value.strip()): return False parts = value.strip().split(".") return all(0 <= int(part) <= 255 for part in parts)