MC / server_manager.py
smodusermc's picture
Update server_manager.py
2802712 verified
Raw History Blame Contribute Delete
16.4 kB
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