Split the monolithic test server into a reusable server core and a thin terminal launcher, and add a PySide6 developer GUI. Added server_core.py (FalloutTogetherServer) with lifecycle control, thread-safe client snapshots, stats and log listener support; added client_session.py for per-client state; added dev_server_app.py GUI and requirements.txt. Updated server.py to use the new server core, and revised server/README.md and docs/dev-log.md to document the new structure, usage, and validation steps. Protocol behavior (welcome/transform/disconnect handling and newline-separated JSON) remains unchanged.
382 lines
13 KiB
Python
382 lines
13 KiB
Python
from __future__ import annotations
|
|
|
|
import json
|
|
import socket
|
|
import threading
|
|
import time
|
|
from collections.abc import Callable
|
|
from typing import Any
|
|
|
|
from client_session import ClientSession
|
|
|
|
|
|
HOST = "127.0.0.1"
|
|
PORT = 7777
|
|
ACCEPT_TIMEOUT_SECONDS = 0.5
|
|
|
|
|
|
class FalloutTogetherServer:
|
|
def __init__(self, host: str = HOST, port: int = PORT) -> None:
|
|
self.host = host
|
|
self.port = port
|
|
|
|
self._lock = threading.RLock()
|
|
self._log_lock = threading.RLock()
|
|
self._clients: dict[socket.socket, ClientSession] = {}
|
|
self._log_listeners: list[Callable[[str], None]] = []
|
|
self._server_socket: socket.socket | None = None
|
|
self._accept_thread: threading.Thread | None = None
|
|
self._running = False
|
|
self._started_at: float | None = None
|
|
self._next_player_id = 1
|
|
|
|
self._stats: dict[str, int] = {
|
|
"clientsConnected": 0,
|
|
"clientsDisconnected": 0,
|
|
"packetsReceived": 0,
|
|
"packetsSent": 0,
|
|
"packetsBroadcast": 0,
|
|
"transformPacketsReceived": 0,
|
|
"transformPacketsBroadcast": 0,
|
|
"disconnectPacketsBroadcast": 0,
|
|
}
|
|
|
|
def start(self) -> None:
|
|
with self._lock:
|
|
if self._running:
|
|
return
|
|
|
|
self._prepare_server_socket()
|
|
self._accept_thread = threading.Thread(target=self._accept_loop, daemon=True)
|
|
self._accept_thread.start()
|
|
|
|
def serve_forever(self) -> None:
|
|
with self._lock:
|
|
if self._running:
|
|
raise RuntimeError("Server is already running.")
|
|
|
|
self._prepare_server_socket()
|
|
|
|
self._accept_loop()
|
|
|
|
def stop(self) -> None:
|
|
clients: list[ClientSession]
|
|
server_socket: socket.socket | None
|
|
|
|
with self._lock:
|
|
self._running = False
|
|
server_socket = self._server_socket
|
|
self._server_socket = None
|
|
clients = list(self._clients.values())
|
|
|
|
if server_socket is not None:
|
|
try:
|
|
server_socket.close()
|
|
except OSError:
|
|
pass
|
|
|
|
for client in clients:
|
|
self._close_client_socket(client)
|
|
|
|
accept_thread = self._accept_thread
|
|
if accept_thread is not None and accept_thread is not threading.current_thread():
|
|
accept_thread.join(timeout=1.0)
|
|
|
|
def is_running(self) -> bool:
|
|
with self._lock:
|
|
return self._running
|
|
|
|
def get_clients(self) -> list[dict[str, Any]]:
|
|
with self._lock:
|
|
clients = list(self._clients.values())
|
|
|
|
return [client.to_snapshot() for client in clients]
|
|
|
|
def get_stats(self) -> dict[str, Any]:
|
|
with self._lock:
|
|
connected_clients = len(self._clients)
|
|
started_at = self._started_at
|
|
stats = dict(self._stats)
|
|
stats.update(
|
|
{
|
|
"host": self.host,
|
|
"port": self.port,
|
|
"isRunning": self._running,
|
|
"startedAt": started_at,
|
|
"uptimeSeconds": time.time() - started_at if started_at is not None else 0.0,
|
|
"connectedClients": connected_clients,
|
|
"nextPlayerId": self._next_player_id,
|
|
}
|
|
)
|
|
|
|
return stats
|
|
|
|
def add_log_listener(self, callback: Callable[[str], None]) -> None:
|
|
with self._log_lock:
|
|
if callback not in self._log_listeners:
|
|
self._log_listeners.append(callback)
|
|
|
|
def remove_log_listener(self, callback: Callable[[str], None]) -> None:
|
|
with self._log_lock:
|
|
if callback in self._log_listeners:
|
|
self._log_listeners.remove(callback)
|
|
|
|
def _prepare_server_socket(self) -> None:
|
|
server_socket = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
|
|
try:
|
|
if hasattr(socket, "SO_EXCLUSIVEADDRUSE"):
|
|
server_socket.setsockopt(socket.SOL_SOCKET, socket.SO_EXCLUSIVEADDRUSE, 1)
|
|
server_socket.bind((self.host, self.port))
|
|
server_socket.listen()
|
|
server_socket.settimeout(ACCEPT_TIMEOUT_SECONDS)
|
|
except OSError:
|
|
server_socket.close()
|
|
raise
|
|
|
|
self._server_socket = server_socket
|
|
self._running = True
|
|
self._started_at = time.time()
|
|
|
|
self._log(f"Fallout 4 Together local test server listening on {self.host}:{self.port}")
|
|
self._log("Waiting for newline-separated JSON transform packets...")
|
|
|
|
def _accept_loop(self) -> None:
|
|
try:
|
|
while self.is_running():
|
|
with self._lock:
|
|
server_socket = self._server_socket
|
|
|
|
if server_socket is None:
|
|
break
|
|
|
|
try:
|
|
connection, address = server_socket.accept()
|
|
except socket.timeout:
|
|
# Wake periodically so Ctrl+C is handled promptly in Windows terminals.
|
|
continue
|
|
except OSError:
|
|
if self.is_running():
|
|
self._log("Server socket closed unexpectedly.")
|
|
break
|
|
|
|
thread = threading.Thread(target=self._handle_client, args=(connection, address), daemon=True)
|
|
thread.start()
|
|
finally:
|
|
with self._lock:
|
|
server_socket = self._server_socket
|
|
self._running = False
|
|
self._server_socket = None
|
|
|
|
if server_socket is not None:
|
|
try:
|
|
server_socket.close()
|
|
except OSError:
|
|
pass
|
|
|
|
def _assign_client(self, connection: socket.socket, address: tuple[str, int]) -> ClientSession:
|
|
with self._lock:
|
|
# The server owns player IDs so early clients do not need to coordinate
|
|
# identity with each other or know anything about remote players yet.
|
|
client = ClientSession(
|
|
connection=connection,
|
|
address=address,
|
|
player_id=self._next_player_id,
|
|
connected_at=time.time(),
|
|
)
|
|
self._next_player_id += 1
|
|
self._clients[connection] = client
|
|
self._stats["clientsConnected"] += 1
|
|
|
|
return client
|
|
|
|
def _remove_client(self, client: ClientSession) -> bool:
|
|
with self._lock:
|
|
removed = self._clients.pop(client.connection, None) is not None
|
|
if removed:
|
|
self._stats["clientsDisconnected"] += 1
|
|
return removed
|
|
|
|
def _handle_client(self, connection: socket.socket, address: tuple[str, int]) -> None:
|
|
client = self._assign_client(connection, address)
|
|
self._log(f"Client connected: {client.label} (player {client.player_id})")
|
|
|
|
with connection:
|
|
buffer = ""
|
|
|
|
try:
|
|
self._send_packet(
|
|
client,
|
|
{
|
|
"type": "welcome",
|
|
"playerId": client.player_id,
|
|
"serverTime": time.time(),
|
|
},
|
|
)
|
|
|
|
while self.is_running():
|
|
chunk = connection.recv(4096)
|
|
if not chunk:
|
|
break
|
|
|
|
buffer += chunk.decode("utf-8", errors="replace")
|
|
while "\n" in buffer:
|
|
line, buffer = buffer.split("\n", 1)
|
|
self._handle_line(client, line.strip())
|
|
except ConnectionResetError:
|
|
self._log(f"Client disconnected unexpectedly: {client.label} (player {client.player_id})")
|
|
except OSError as error:
|
|
if self.is_running():
|
|
self._log(f"Client connection error: {client.label} (player {client.player_id}): {error}")
|
|
finally:
|
|
self._disconnect_client(client)
|
|
|
|
def _handle_line(self, client: ClientSession, line: str) -> None:
|
|
if not line:
|
|
return
|
|
|
|
received_at = time.time()
|
|
client.record_received(received_at)
|
|
|
|
with self._lock:
|
|
self._stats["packetsReceived"] += 1
|
|
|
|
try:
|
|
packet = json.loads(line)
|
|
except json.JSONDecodeError as error:
|
|
self._log(f"Invalid JSON from {client.label}: {error}: {line}")
|
|
return
|
|
|
|
if not isinstance(packet, dict):
|
|
self._log(f"Invalid packet from {client.label}: expected JSON object: {packet}")
|
|
return
|
|
|
|
if packet.get("type") == "transform":
|
|
packet["playerId"] = client.player_id
|
|
packet["serverTime"] = time.time()
|
|
client.record_transform(packet)
|
|
with self._lock:
|
|
self._stats["transformPacketsReceived"] += 1
|
|
|
|
self._print_packet(client, packet)
|
|
|
|
if packet.get("type") == "transform":
|
|
self._broadcast_transform(client, packet)
|
|
|
|
def _send_packet(self, client: ClientSession, packet: dict[str, Any], *, broadcast: bool = False) -> None:
|
|
encoded = json.dumps(packet, separators=(",", ":")).encode("utf-8") + b"\n"
|
|
client.connection.sendall(encoded)
|
|
client.record_sent(broadcast=broadcast)
|
|
|
|
with self._lock:
|
|
self._stats["packetsSent"] += 1
|
|
if broadcast:
|
|
self._stats["packetsBroadcast"] += 1
|
|
|
|
def _disconnect_client(self, client: ClientSession) -> None:
|
|
if not self._remove_client(client):
|
|
return
|
|
|
|
self._close_client_socket(client)
|
|
self._log(f"Client disconnected: {client.label} (player {client.player_id})")
|
|
self._broadcast_disconnect(client)
|
|
|
|
def _close_client_socket(self, client: ClientSession) -> None:
|
|
try:
|
|
client.connection.close()
|
|
except OSError:
|
|
pass
|
|
|
|
def _broadcast_transform(self, sender: ClientSession, packet: dict[str, Any]) -> None:
|
|
with self._lock:
|
|
recipients = [client for client in self._clients.values() if client.connection != sender.connection]
|
|
|
|
failed_recipients: list[ClientSession] = []
|
|
successful_sends = 0
|
|
for recipient in recipients:
|
|
try:
|
|
# Transforms go to every other client only. The sender already knows
|
|
# its own movement; echoing it back would create duplicate local state.
|
|
self._send_packet(recipient, packet, broadcast=True)
|
|
successful_sends += 1
|
|
except OSError:
|
|
failed_recipients.append(recipient)
|
|
|
|
with self._lock:
|
|
self._stats["transformPacketsBroadcast"] += successful_sends
|
|
|
|
for recipient in failed_recipients:
|
|
self._disconnect_client(recipient)
|
|
|
|
self._log(f"Broadcast transform from player {sender.player_id} to {successful_sends} other client(s)")
|
|
|
|
def _broadcast_disconnect(self, disconnected_client: ClientSession) -> None:
|
|
packet = {
|
|
"type": "disconnect",
|
|
"playerId": disconnected_client.player_id,
|
|
"serverTime": time.time(),
|
|
}
|
|
|
|
with self._lock:
|
|
recipients = list(self._clients.values())
|
|
|
|
failed_recipients: list[ClientSession] = []
|
|
successful_sends = 0
|
|
for recipient in recipients:
|
|
try:
|
|
self._send_packet(recipient, packet, broadcast=True)
|
|
successful_sends += 1
|
|
except OSError:
|
|
failed_recipients.append(recipient)
|
|
|
|
with self._lock:
|
|
self._stats["disconnectPacketsBroadcast"] += successful_sends
|
|
|
|
for recipient in failed_recipients:
|
|
self._disconnect_client(recipient)
|
|
|
|
self._log(
|
|
f"Broadcast disconnect for player {disconnected_client.player_id} "
|
|
f"to {successful_sends} other client(s)"
|
|
)
|
|
|
|
def _print_packet(self, client: ClientSession, packet: dict[str, Any]) -> None:
|
|
if packet.get("type") != "transform":
|
|
self._log(f"Packet from {client.label}: {packet}")
|
|
return
|
|
|
|
try:
|
|
x = float(packet["x"])
|
|
y = float(packet["y"])
|
|
z = float(packet["z"])
|
|
angle_z = float(packet["angleZ"])
|
|
except (KeyError, TypeError, ValueError):
|
|
self._log(f"Malformed transform packet from {client.label}: {packet}")
|
|
return
|
|
|
|
movement_type = packet.get("movementType", "normal")
|
|
extra_fields = []
|
|
for field_name in ("playerId", "cellId", "worldspaceId", "clientTime", "serverTime"):
|
|
if field_name in packet:
|
|
extra_fields.append(f"{field_name}={packet[field_name]}")
|
|
|
|
extra_details = ""
|
|
if extra_fields:
|
|
extra_details = ", " + ", ".join(extra_fields)
|
|
|
|
self._log(
|
|
"Transform from "
|
|
f"{client.label}: x={x:.2f}, y={y:.2f}, z={z:.2f}, angleZ={angle_z:.2f}, "
|
|
f"movementType={movement_type}{extra_details}"
|
|
)
|
|
|
|
def _log(self, message: str) -> None:
|
|
with self._log_lock:
|
|
listeners = list(self._log_listeners)
|
|
|
|
for listener in listeners:
|
|
try:
|
|
listener(message)
|
|
except Exception:
|
|
# Future UI listeners must not be able to break server networking.
|
|
pass
|