import subprocess import threading import os import signal import time import re from pathlib import Path from datetime import datetime from config import config from console_handler import ConsoleHandler from player_tracker import PlayerTracker from mod_manager import ModManager from file_manager import FileManager from health_ping import health_pinger import psutil from functools import wraps from supabase import create_client, Client # Constants for stats monitoring GLOBAL_STATUS_INTERVAL = 5 PROCESS_STATS_INTERVAL = 1 # Supabase Client Initialization try: supabase: Client = create_client(config.SUPABASE_URL, config.SUPABASE_ANON_KEY) except Exception as e: print(f"WARNING: Supabase client initialization failed: {e}") # Set a placeholder to avoid crashes if keys are missing class MockSupabaseStorage: def upload(self, remote_path, data): raise Exception("Supabase not configured.") class MockSupabaseClient: def storage(self): class MockStorage: def from_(self, bucket): return MockSupabaseStorage() return MockStorage() supabase = MockSupabaseClient() class MinecraftServerManager: def __init__(self, socketio=None): self.process = None self.playit_process = None # NEW: Process for the PlayIt.gg agent self.running = False self.starting = False self.stopping = False self.socketio = socketio self.output_thread = None self.process_stats_thread = None self.global_monitor_thread = None self.console_history = [] self.max_history = 1000 self.start_time = None self.players_online = [] self.server_stats = { 'cpu': 0.0, 'memory': 0.0, 'tps': 20.0, 'uptime': 0.0 } self._lock = threading.Lock() # Start the global monitor immediately upon initialization self.start_global_monitor() def get_java_command(self): java_path = os.getenv("JAVA_HOME", "/usr/lib/jvm/temurin-21-jdk-amd64") + "/bin/java" cmd = [ java_path, f"-Xms{config.MIN_RAM}", f"-Xmx{config.MAX_RAM}", ] cmd += config.JAVA_ARGS cmd += [ "-jar", config.SERVER_JAR, "nogui" ] return cmd def log(self, message, msg_type="info"): timestamp = datetime.now().strftime("%H:%M:%S") entry = { 'timestamp': timestamp, 'message': message, 'type': msg_type } with self._lock: self.console_history.append(entry) if len(self.console_history) > self.max_history: self.console_history = self.console_history[-self.max_history:] self._emit_event('console_output', entry) def _emit_event(self, event, data): if self.socketio: self.socketio.emit(event, data) def _emit_status_update(self): self._emit_event('status_update', self.get_status()) def _emit_stats_update(self): self._emit_event('stats_update', self.server_stats) def _emit_player_update(self): self._emit_event('players_update', { 'players': self.players_online, 'count': len(self.players_online) }) # Supabase & Cleanup Logic (Config Save Only) def save_and_cleanup(self): self.log("Starting configuration save (World excluded due to size limit)...", "save") if self.running: self.send_command("save-all") time.sleep(5) try: timestamp = datetime.now().strftime("%Y%m%d_%H%M%S") temp_zip_path = config.BACKUPS_DIR / f"config_save_{timestamp}.zip" # CRITICAL CHANGE: Only save small configuration files save_items = [ str(config.SERVER_DIR / "mods"), str(config.SERVER_DIR / "config"), str(config.SERVER_DIR / "server.properties"), str(config.SERVER_DIR / config.SERVER_JAR), ] # Create zip file. cwd is critical for relative paths in the zip zip_cmd = ["zip", "-r", str(temp_zip_path)] + [str(Path(item).relative_to(config.BASE_DIR)) for item in save_items] subprocess.run(zip_cmd, cwd=str(config.BASE_DIR), check=True) self.log(f"Created local config ZIP: {temp_zip_path.name}", "save") # Upload to Supabase Storage remote_path = f"configs/{temp_zip_path.name}" # Ensure supabase is not the mock client before attempting upload if hasattr(supabase, 'storage'): with open(temp_zip_path, 'rb') as f: supabase.storage().from_(config.SUPABASE_BUCKET_NAME).upload(remote_path, f.read()) self.log(f"Successfully uploaded CONFIG save to Supabase at: {remote_path}", "success") else: self.log("Supabase storage not configured. Skipping upload.", "warning") # Clean up local temporary files self.log("Cleaning up local temporary files...", "cleanup") # Delete temporary uploads subprocess.run(['rm', '-rf', str(config.UPLOADS_DIR / "*")], check=False) # Delete the uploaded local zip file temp_zip_path.unlink(missing_ok=True) self.log("Local temporary files deleted and storage space reclaimed.", "success") except Exception as e: self.log(f"Supabase save/cleanup failed: {str(e)}", "error") return {"success": False, "message": str(e)} return {"success": True, "message": "Configuration saved. World folder retained locally."} # --- NEW: PlayIt.gg Tunnel Method --- def start_playit_tunnel(self): self.log("Starting PlayIt.gg agent...", "internal") try: # The playit command starts the agent. self.playit_process = subprocess.Popen( ["playit"], cwd=str(config.SERVER_DIR), stdout=subprocess.PIPE, stderr=subprocess.STDOUT, # Direct stderr to stdout for consolidation text=True, bufsize=1, start_new_session=True # Ensures the process runs independently ) # Log the message that contains the crucial authorization link self.log("PlayIt.gg agent started. **CHECK CONSOLE FOR AUTHORIZATION LINK!**", "warning") # We don't need a heavy read loop here; the link will appear in the main process logs # if we choose to monitor it separately later. For now, it just starts the agent. except FileNotFoundError: self.log("ERROR: PlayIt.gg agent not found. Did you add it to the Dockerfile?", "error") self.playit_process = None except Exception as e: self.log(f"Failed to start PlayIt.gg agent: {e}", "error") self.playit_process = None def stop_playit_tunnel(self): if self.playit_process and self.playit_process.poll() is None: self.log("Stopping PlayIt.gg agent...", "internal") try: # Send SIGINT to the process group to ensure it cleanly stops os.killpg(os.getpgid(self.playit_process.pid), signal.SIGINT) self.playit_process.wait(timeout=10) self.log("PlayIt.gg agent stopped.", "internal") except Exception as e: self.log(f"Error stopping PlayIt.gg agent, attempting kill: {e}", "warning") self.playit_process.kill() finally: self.playit_process = None # ------------------------------------ # Server Control Methods def start(self): with self._lock: if self.running or self.starting: return {"success": False, "message": "Server is already running or starting"} self.starting = True self._emit_status_update() try: # 1. Ensure EULA is accepted eula_path = config.SERVER_DIR / "eula.txt" with open(eula_path, 'w') as f: f.write("eula=true\n") # 2. Start the Tunnel before the server (if PlayIt.gg is installed) self.start_playit_tunnel() # 3. Start the Minecraft Server cmd = self.get_java_command() self.log(f"Starting server with command: {' '.join(cmd)}") self.process = subprocess.Popen( cmd, cwd=str(config.SERVER_DIR), stdin=subprocess.PIPE, stdout=subprocess.PIPE, stderr=subprocess.STDOUT, text=True, bufsize=1 ) self.running = True self.starting = False self.start_time = datetime.now() self.output_thread = threading.Thread(target=self._read_output, daemon=True) self.output_thread.start() self.process_stats_thread = threading.Thread(target=self._monitor_stats_process, daemon=True) self.process_stats_thread.start() self._emit_status_update() self.log("Server process started successfully") return {"success": True, "message": "Server starting..."} except Exception as e: self.starting = False self._emit_status_update() self.log(f"Failed to start server: {str(e)}", "error") self.stop_playit_tunnel() # Clean up tunnel on failure return {"success": False, "message": str(e)} def stop(self, force=False): with self._lock: if not self.running: return {"success": False, "message": "Server is not running"} self.stopping = True self._emit_status_update() try: if force: if self.process: self.process.kill() else: self.send_command("stop") for _ in range(30): if self.process and self.process.poll() is not None: break time.sleep(1) else: if self.process: self.process.kill() # NEW: Stop the PlayIt tunnel process self.stop_playit_tunnel() self.running = False self.stopping = False self.start_time = None self.players_online = [] self.server_stats = {'cpu': 0.0, 'memory': 0.0, 'tps': 0.0, 'uptime': 0.0} self._emit_status_update() self._emit_player_update() self.log("Server stopped") return {"success": True, "message": "Server stopped"} except Exception as e: self.stopping = False self.log(f"Error stopping server: {str(e)}", "error") self._emit_status_update() return {"success": False, "message": str(e)} def restart(self): self.log("Restarting server...") stop_result = self.stop() if stop_result["success"]: time.sleep(2) return self.start() return stop_result def send_command(self, command): if not self.running or not self.process: return {"success": False, "message": "Server is not running"} try: self.process.stdin.write(command + "\n") self.process.stdin.flush() self.log(f"> {command}", "command") return {"success": True, "message": f"Command sent: {command}"} except Exception as e: return {"success": False, "message": str(e)} # Threaded Methods def start_global_monitor(self): if self.global_monitor_thread is None or not self.global_monitor_thread.is_alive(): self.global_monitor_thread = threading.Thread(target=self._monitor_stats_global, daemon=True) self.global_monitor_thread.start() self.log("Global status monitor started.", "internal") def _monitor_stats_global(self): while True: try: if self.running and self.start_time: with self._lock: self.server_stats['uptime'] = (datetime.now() - self.start_time).total_seconds() self._emit_status_update() self._emit_stats_update() except Exception as e: print(f"Global Stats Monitor Error: {e}") time.sleep(GLOBAL_STATUS_INTERVAL) def _monitor_stats_process(self): while self.running: try: if self.process and self.process.pid: if self.process.poll() is not None: self.running = False self.log("Process stats monitor detected server termination.", "warn") break proc = psutil.Process(self.process.pid) with self._lock: self.server_stats['cpu'] = proc.cpu_percent(interval=None) self.server_stats['memory'] = proc.memory_info().rss / (1024 * 1024) time.sleep(PROCESS_STATS_INTERVAL) except (psutil.NoSuchProcess, psutil.AccessDenied): self.running = False self.log("Stats monitor lost server process.", "warn") break def _read_output(self): try: for line in iter(self.process.stdout.readline, ''): if not line: break line = line.strip() if line: self._process_line(line) self.log(line) except Exception as e: self.log(f"Error reading output: {str(e)}", "error") finally: # Cleanup on server exit self.log("Server process ended") self.running = False self.starting = False self.stopping = False self.process = None self._emit_status_update() # The stop method will call stop_playit_tunnel, but call it here as a safeguard too self.stop_playit_tunnel() def _process_line(self, line): join_match = re.search(r'(\w+) joined the game', line) if join_match: player = join_match.group(1) with self._lock: if player not in self.players_online: self.players_online.append(player) self._emit_player_update() leave_match = re.search(r'(\w+) left the game', line) if leave_match: player = leave_match.group(1) with self._lock: if player in self.players_online: self.players_online.remove(player) self._emit_player_update() if "Done" in line and "For help, type" in line: self._emit_event('server_ready', {'message': 'Server is ready!'}) tps_match = re.search(r'TPS.*?(\d+\.?\d*)', line) if tps_match: with self._lock: self.server_stats['tps'] = float(tps_match.group(1)) # Client Data Methods def get_status(self): return { 'running': self.running, 'starting': self.starting, 'stopping': self.stopping, 'players': self.players_online, 'player_count': len(self.players_online), 'stats': self.server_stats, 'uptime': self.server_stats.get('uptime', 0.0) } def get_console_history(self, limit=100): return self.console_history[-limit:] # Global server manager instance server_manager = None def init_server_manager(socketio=None): global server_manager server_manager = MinecraftServerManager(socketio) return server_manager def get_server_manager(): global server_manager return server_manager