DOCS (the ask: after a compaction this session lost track of how the console
works and launched the wrong program, twice).
* NEW context/operator-console.md -- the dedicated topic that was missing.
Leads with the thing I got wrong: btoperator.py is the PySide6 GUI the
operator uses; btconsole.py is the headless relay it spawns. Then ports,
the route table, the roster/seat/identity model, the full round lifecycle
with every launch gate, liveness, mode-specific traps, and log locations.
* NEW docs/OPERATOR_GUIDE.md -- sysop-facing: start the console, set up a
mission, watch pods arrive, launch, run back-to-back rounds, what to press
when LAUNCH looks dead, a troubleshooting table keyed on the exact log
lines, and what to save BEFORE restarting a session (Start Session
truncates operator_relay.log, so restarting to clear a problem destroys the
evidence of it).
* CLAUDE.md: two Quick Lookup rows + a DO-NOT entry naming the two programs,
so the distinction survives the next compaction.
* context/multiplayer.md: a pointer out of the scattered console notes to the
new topic (they were buried across ~8 places in a large file, which is
exactly why they evaporated).
FIXES -- both are regressions in my own previous commit, found by the review
pass, and one would have made a games night WORSE:
* THE REAPER WOULD HAVE KILLED HEALTHY PLAYERS. It was gated only on "no
mission running", which a round RESET satisfies -- so it was armed for the
whole BETWEEN-ROUNDS wait, and that is a period when a pod is legitimately
byte-silent: its seat beacon is write-only for the process lifetime
(L4NET.CPP: "the relay ignores its silence") and its console pad has no egg
yet so it cannot ACK. A real night showed 12-minute and 5-minute gaps; the
180s deadline would have dropped healthy pods and forced their clients to
relaunch. It now runs ONLY in the active staging window, skips pads with no
egg and conns that have not HELLO'd, and last_seen is also stamped from
inbound UDP (a pod streaming updates while its TCP idles was being counted
as silent). Half-open detection is keepalive's job; this is just a backstop.
* UnboundLocalError in the operator UI. My end_sent reset was an `elif` in
the chain that assigns `head`, so that branch left `head` unbound -- a crash
on the first status refresh after End Mission, which is exactly the path the
relay's "StopMission sent" line produces. Moved out of the chain.
Plus one pre-existing wedge with the same symptom as the reported bug, live-
proven in operator_relay.log (~5 minutes of a night lost): _abort_round clears
eggs_released BEFORE the survivors' sockets close, so _maybe_reset_round's own
`if not self.eggs_released: return` skips the template restore forever -- the
roster stays trimmed, the release gates can never be met, and walk-ups get
ROSTER FULL. _rearm_for_new_round's restore is therefore now UNCONDITIONAL (it
was gated on eggs_released, which made Re-arm useless in the one state that most
needs it) and it clears round_hold_until so an abort's settle window is not
inherited. Aborts are ordinary: nothing on the wire distinguishes a straggler's
late FIN from a pod dying mid-load.
scratchpad/test_relay_rearm.py now 19 checks, all passing, including the two new
regression guards (a 15-minute-silent waiting pod is NOT reaped; a pad that was
never sent an egg is NOT reaped) and the mid-pair no-re-arm guard.
Still open, recorded in the new topic's frontmatter: the UDP endpoint map trusts
the sender's self-declared fromHost; the egg-ACK is a fixed-offset parse of one
recv with no reassembly; remote-operator mode can never enable LAUNCH.
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
1927 lines
91 KiB
Python
1927 lines
91 KiB
Python
#!/usr/bin/env python
|
|
"""btconsole.py -- minimal Tesla/WinTesla CONSOLE emulator (the operator station).
|
|
|
|
The pod-side engine (MUNGA_L4/L4NET.CPP) boots `-net <port>` into ConsoleOnly state:
|
|
it listens on TCP <port> and waits for a console to connect and stream the mission
|
|
egg. The real console software is absent from every archive, so this tool speaks
|
|
the console's wire protocol:
|
|
|
|
packet := NetworkPacketHeader + NetworkManager::ReceiveEggFileMessage
|
|
off 0 int32 clientID = 0 (NetworkClient::NetworkManagerClientID)
|
|
off 4 int32 gameID = 0 (NETWORK.cpp:121 -- always 0)
|
|
off 8 int32 fromHost = 1 (FirstLegalHostID == the console)
|
|
off 12 int32 timeStamp = 0 (Time; unused by the egg handler)
|
|
off 16 uint32 messageLength = 1024 (sizeof(ReceiveEggFileMessage))
|
|
off 20 int32 messageID = 3 (NetworkManager::ReceiveEggFileMessageID)
|
|
off 24 uint32 messageFlags = 1 (Receiver::Message::ReliableFlag)
|
|
off 28 int32 sequenceNumber (0..n-1; -1 is the solo-local path only)
|
|
off 32 int32 notationFileLength (total egg bytes)
|
|
off 36 int32 thisMessageLength (bytes used in this chunk, <= 1000)
|
|
off 40 char notationData[1000]
|
|
total 1040 bytes per chunk
|
|
(constants verified from the port build via `btl4.exe` BT_NET_PROBE=1)
|
|
|
|
The pod reassembles chunks sequentially (L4NET.CPP:773-816: seq 0 allocates,
|
|
append until notationFileLength reached), sets NormalState, CreateMission ->
|
|
StartConnecting (the egg's [pilots] list forms the pod mesh), and sends an
|
|
AcknowledgeEggFile back on this socket. The console must STAY CONNECTED (a
|
|
disconnect trips the pod's console-loss path, which also closes the game
|
|
listener -- an engine bug).
|
|
|
|
Usage (legacy dial-out console -- the pod-authentic direction):
|
|
python btconsole.py <egg-file> <host:port> [<host:port> ...]
|
|
Example (two instances on one box, -net 1501 / -net 1601):
|
|
python btconsole.py MP.EGG 127.0.0.1:1501 127.0.0.1:1601
|
|
|
|
RELAY MODE (D1 -- internet play; pods dial OUT to this process):
|
|
python btconsole.py --relay <consolePort> <egg-file> [--bind ADDR] [--udp-drop PCT]
|
|
Example:
|
|
python btconsole.py --relay 1500 MP_RELAY.EGG
|
|
The relay listens on <consolePort>/tcp (console protocol: streams the egg to each
|
|
pod that connects, sends the LAUNCH pair once ALL pods have the egg), on
|
|
<consolePort>+1/tcp (the GAME relay: envelope-framed frame router, see
|
|
RelayEnvelope below), and on <consolePort>+1/udp (the unreliable channel
|
|
forwarder). Pods run with BT_RELAY=<host>:<consolePort> and a unique
|
|
BT_SELF=<their [pilots] entry>.
|
|
|
|
RelayEnvelope (game TCP): { int32 route; uint32 length; } + length payload bytes
|
|
route >= 2 : client->relay: unicast to that hostID
|
|
relay->client: the sender's hostID (frames also self-identify
|
|
via NetworkPacketHeader.fromHost)
|
|
route == -1: broadcast to every registered pod except the sender
|
|
route == -2: HELLO (payload: int32 magic 'BTR1', int32 myHostID)
|
|
route == -3: PEER_UP (payload: int32 hostID) relay->client only
|
|
route == -4: PEER_DOWN (payload: int32 hostID) relay->client only
|
|
UDP envelope: { int32 route; int32 fromHost; uint32 seq; } + frame
|
|
route as above; -2 = HELLO/keepalive (no payload) -> relay answers -5
|
|
(HELLO-ACK). The relay learns each pod's public endpoint from every
|
|
inbound datagram and forwards verbatim; unknown target endpoint falls
|
|
back to wrapping the frame onto the target's game TCP connection.
|
|
"""
|
|
import collections
|
|
import os
|
|
import random
|
|
import selectors
|
|
# player callsign/mech seat requests mutate the session egg via eggmodel;
|
|
# import guarded so a stripped install still relays with roster defaults.
|
|
try:
|
|
import eggmodel
|
|
except ImportError:
|
|
eggmodel = None
|
|
import socket
|
|
import struct
|
|
import sys
|
|
import threading
|
|
import time
|
|
|
|
CHUNK = 1000
|
|
|
|
# LIVENESS (2026-07-26, operator report "I have to reboot the console"). Nothing
|
|
# on the relay ever noticed a pod that vanished WITHOUT closing its TCP -- a
|
|
# sleeping laptop, a dropped Wi-Fi link, a killed process, a NAT idle-timeout on
|
|
# the internet relay. There was no SO_KEEPALIVE and no idle deadline anywhere,
|
|
# so the dead connection stayed in console_conns with acked=True forever, which
|
|
# (a) permanently blocked the round reset (it requires NO conn acked) and (b) held
|
|
# a phantom seat in the launch count. Only restarting the relay cleared it.
|
|
KEEPALIVE_IDLE_SECONDS = 30 # start probing after this much silence
|
|
KEEPALIVE_INTERVAL_SECONDS = 10
|
|
DEAD_PEER_SECONDS = 180.0 # app-level deadline: no bytes at all for this long
|
|
SEND_TIMEOUT_SECONDS = 10.0 # a blocking sendall may not stall the whole relay
|
|
|
|
|
|
def _enable_keepalive(sock):
|
|
"""Turn on TCP keepalive so the OS reaps half-open peers for us. Best
|
|
effort: the tunables are platform-specific and a plain SO_KEEPALIVE is still
|
|
a large improvement over nothing."""
|
|
try:
|
|
sock.setsockopt(socket.SOL_SOCKET, socket.SO_KEEPALIVE, 1)
|
|
except OSError:
|
|
return
|
|
try: # Windows: one ioctl sets both
|
|
sock.ioctl(socket.SIO_KEEPALIVE_VALS,
|
|
(1, KEEPALIVE_IDLE_SECONDS * 1000,
|
|
KEEPALIVE_INTERVAL_SECONDS * 1000))
|
|
except (AttributeError, OSError, ValueError):
|
|
for opt, val in (("TCP_KEEPIDLE", KEEPALIVE_IDLE_SECONDS),
|
|
("TCP_KEEPINTVL", KEEPALIVE_INTERVAL_SECONDS),
|
|
("TCP_KEEPCNT", 4)):
|
|
code = getattr(socket, opt, None)
|
|
if code is not None:
|
|
try:
|
|
sock.setsockopt(socket.IPPROTO_TCP, code, val)
|
|
except OSError:
|
|
pass
|
|
|
|
|
|
def _send_all_guarded(sock, payload, what=""):
|
|
"""Blocking sendall with a TIMEOUT. The relay is single-threaded: every
|
|
pod-facing send used to flip the socket to blocking with no timeout, so one
|
|
wedged peer could freeze the entire relay -- no launch tick, no UDP fan-out,
|
|
no operator control channel -- for as long as the OS kept retransmitting.
|
|
Raises OSError on failure/timeout; callers already handle that by dropping.
|
|
"""
|
|
sock.settimeout(SEND_TIMEOUT_SECONDS)
|
|
try:
|
|
sock.sendall(payload)
|
|
finally:
|
|
try:
|
|
sock.setblocking(False)
|
|
except OSError:
|
|
pass
|
|
|
|
|
|
PKT_FMT_HDR = "<iiii" # clientID, gameID, fromHost, timeStamp
|
|
PKT_FMT_MSG = "<IiIiii" # messageLength, messageID, messageFlags, seq, fileLen, thisLen
|
|
MESSAGE_LENGTH = 1024 # sizeof(ReceiveEggFileMessage)
|
|
EGG_MESSAGE_ID = 3 # NetworkManager::ReceiveEggFileMessageID
|
|
RELIABLE_FLAG = 1
|
|
CONSOLE_HOST_ID = 1 # FirstLegalHostID
|
|
|
|
# The LAUNCH: pods load to WaitingForLaunch and sit until the console sends
|
|
# Application::RunMissionMessage (the operator's launch button). The handler
|
|
# ladder needs it twice: WaitingForLaunch->LaunchingMission (dispatches the
|
|
# player MissionStarting/translocation), then LaunchingMission->RunningMission.
|
|
# Sending it early (state==CreatingMission) hits the handler's Fail() -- hence
|
|
# the settle delay. Constants from BT_NET_PROBE=1.
|
|
APPLICATION_CLIENT_ID = 4 # NetworkClient::ApplicationClientID
|
|
RUN_MISSION_MESSAGE_ID = 5 # Application::RunMissionMessageID
|
|
STOP_MISSION_MESSAGE_ID = 6 # Application::StopMissionMessageID (APP.h:383)
|
|
LAUNCH_SETTLE_SECONDS = 20.0 # egg -> first launch (pods must reach WaitingForLaunch)
|
|
LAUNCH_STEP_SECONDS = 4.0 # first launch -> second (Launching -> Running)
|
|
STOP_GRACE_SECONDS = 2.0 # past the pods' own zero so the ranking shows
|
|
|
|
|
|
def run_mission_packet():
|
|
pkt = struct.pack(PKT_FMT_HDR, APPLICATION_CLIENT_ID, 0, CONSOLE_HOST_ID, 0)
|
|
pkt += struct.pack("<IiI", 12, RUN_MISSION_MESSAGE_ID, RELIABLE_FLAG)
|
|
assert len(pkt) == 16 + 12
|
|
return pkt
|
|
|
|
|
|
def stop_mission_packet():
|
|
# Application::StopMissionMessage {exitCode=NullExitCodeID} (id 6,
|
|
# APPMSG.cpp:21). The 1995 console ENDED the timed mission -- the pods
|
|
# only display the countdown (secondsRemainingInGame, APP.cpp:652) and
|
|
# pop the ranking window for the final 30s (DIRECTOR.cpp:113); the stop
|
|
# itself was always console-side, which is us.
|
|
pkt = struct.pack(PKT_FMT_HDR, APPLICATION_CLIENT_ID, 0, CONSOLE_HOST_ID, 0)
|
|
pkt += struct.pack("<IiIi", 16, STOP_MISSION_MESSAGE_ID, RELIABLE_FLAG, 0)
|
|
assert len(pkt) == 16 + 16
|
|
return pkt
|
|
|
|
|
|
def parse_egg_mission_length(egg_path):
|
|
"""[mission] length= in seconds (0 = untimed)."""
|
|
section = None
|
|
try:
|
|
for raw in open(egg_path, "r", encoding="latin-1"):
|
|
line = raw.strip()
|
|
if line.startswith("[") and line.endswith("]"):
|
|
section = line[1:-1].lower()
|
|
elif section == "mission" and "=" in line:
|
|
key, _, value = line.partition("=")
|
|
if key.strip().lower() == "length":
|
|
return float(value.strip())
|
|
except (OSError, ValueError):
|
|
pass
|
|
return 0.0
|
|
|
|
|
|
def egg_wire_bytes(egg_text_bytes):
|
|
"""Convert raw egg TEXT to the wire form NotationFile::ReadText expects:
|
|
NUL-separated lines (NOTATION.cpp:1043 walks `strchr(buffer,'\\0')+1` per
|
|
line). The real 1996 console pre-processed the egg the same way; sending
|
|
the raw file text makes the whole egg parse as one garbage line ("ERROR:
|
|
no map in egg!")."""
|
|
lines = egg_text_bytes.replace(b"\r\n", b"\n").split(b"\n")
|
|
return b"".join(line + b"\0" for line in lines)
|
|
|
|
|
|
def egg_packets(egg_bytes):
|
|
total = len(egg_bytes)
|
|
seq = 0
|
|
off = 0
|
|
while off < total:
|
|
chunk = egg_bytes[off:off + CHUNK]
|
|
pkt = struct.pack(PKT_FMT_HDR, 0, 0, CONSOLE_HOST_ID, 0)
|
|
pkt += struct.pack(PKT_FMT_MSG, MESSAGE_LENGTH, EGG_MESSAGE_ID, RELIABLE_FLAG,
|
|
seq, total, len(chunk))
|
|
pkt += chunk.ljust(CHUNK, b"\0")
|
|
assert len(pkt) == 16 + MESSAGE_LENGTH
|
|
yield pkt
|
|
off += len(chunk)
|
|
seq += 1
|
|
|
|
|
|
def serve_pod(hostport, egg_bytes):
|
|
host, port = hostport.rsplit(":", 1)
|
|
port = int(port)
|
|
name = f"{host}:{port}"
|
|
# The pod may not be listening yet -- retry like the pods do to each other.
|
|
while True:
|
|
try:
|
|
s = socket.create_connection((host, port), timeout=5)
|
|
break
|
|
except OSError as e:
|
|
print(f"[{name}] waiting for pod ({e})")
|
|
time.sleep(1)
|
|
print(f"[{name}] connected; streaming egg ({len(egg_bytes)} bytes)")
|
|
n = 0
|
|
for pkt in egg_packets(egg_bytes):
|
|
s.sendall(pkt)
|
|
n += 1
|
|
print(f"[{name}] egg sent in {n} chunk(s); launch in {LAUNCH_SETTLE_SECONDS}s")
|
|
s.settimeout(2.0)
|
|
launch_at = time.time() + LAUNCH_SETTLE_SECONDS
|
|
launches_sent = 0
|
|
while True:
|
|
if launches_sent < 2 and time.time() >= launch_at:
|
|
s.sendall(run_mission_packet())
|
|
launches_sent += 1
|
|
launch_at = time.time() + LAUNCH_STEP_SECONDS
|
|
print(f"[{name}] RunMission #{launches_sent} sent")
|
|
try:
|
|
data = s.recv(4096)
|
|
if not data:
|
|
print(f"[{name}] pod closed the console socket", flush=True)
|
|
return
|
|
# Expect the AcknowledgeEggFile (and possibly other console traffic).
|
|
if len(data) >= 24:
|
|
(clid, gid, fh, ts) = struct.unpack_from("<iiii", data, 0)
|
|
(mlen, mid) = struct.unpack_from("<Ii", data, 16)
|
|
print(f"[{name}] pod->console packet: clientID={clid} fromHost={fh} "
|
|
f"msgID={mid} len={mlen} ({len(data)} bytes)")
|
|
else:
|
|
print(f"[{name}] pod->console {len(data)} bytes")
|
|
except socket.timeout:
|
|
continue
|
|
except OSError as e:
|
|
# Previously UNHANDLED: a reset/abort here killed the thread silently
|
|
# (buffered stdout lost) and the console process exited -- leaving the
|
|
# pods holding dead console sockets. Log it and keep the thread alive
|
|
# (idle) so the process + diagnosis survive.
|
|
print(f"[{name}] console socket error: {e!r} -- idling", flush=True)
|
|
while True:
|
|
time.sleep(60)
|
|
|
|
|
|
#############################################################################
|
|
# RELAY MODE (D1)
|
|
#############################################################################
|
|
|
|
ENV_FMT_TCP = "<iI" # route, length
|
|
ENV_FMT_UDP = "<iiI" # route, fromHost, seq
|
|
ENV_TCP_SIZE = 8
|
|
ENV_UDP_SIZE = 12
|
|
ROUTE_BCAST = -1
|
|
ROUTE_HELLO = -2
|
|
ROUTE_PEER_UP = -3
|
|
ROUTE_PEER_DOWN = -4
|
|
ROUTE_UDP_ACK = -5
|
|
ROUTE_SEAT_REQUEST = -6 # client->relay: assign me a free roster seat
|
|
ROUTE_SEAT_ASSIGN = -7 # relay->client: int32 hostID + NUL-terminated tag
|
|
ROUTE_SEAT_FULL = -8 # relay->client: no free seats
|
|
ROUTE_MATCHLOG = -9 # client->relay: matchlog upload (filename NUL + bytes)
|
|
ROUTE_READY = -10 # client->relay: mission load complete (READY light)
|
|
ROUTE_REJOIN = -11 # relay->client: abandon the round, rejoin (peer died mid-load)
|
|
# the 18 certified mech tags (mirror of the FE kVehicles catalog)
|
|
VEHICLE_TAGS = {"blkhawk", "loki", "bhk1", "madcat", "thor", "owens",
|
|
"own1", "thr1", "lok1", "mad1", "avatar", "ava1",
|
|
"sunder", "snd1", "vulture", "vul1", "lok2", "mad2"}
|
|
MATCHLOG_CAP = 8 * 1024 * 1024 # matchlog upload frame cap (client caps at 8MB)
|
|
SEAT_RESERVE_SECONDS = 60.0
|
|
ABORT_SETTLE_SECONDS = 8.0 # issue #33: egg-release hold after a round abort
|
|
|
|
# issue #38: the vehicle table's CAMO colors (a color= outside this set
|
|
# paints the mech silent-gray -- 'Red' is a PATCH color, the camo is
|
|
# 'Crimson'; caught live when a test egg authored color=Red).
|
|
VALID_CAMO_COLORS = {"Black", "Brown", "Crimson", "Green", "Grey", "Tan",
|
|
"White"}
|
|
|
|
# REMOTE OPERATOR CONTROL (2026-07-24): a TCP line protocol so an operator
|
|
# GUI on ANOTHER machine can drive this relay (only the relay host needs
|
|
# port forwarding). Additive: stdin commands keep working unchanged.
|
|
# client -> relay: AUTH <secret> | launch | stop | ping |
|
|
# get mission | set key=value;key=value
|
|
# relay -> client: "OK ..." on auth, then the relay's own timestamped
|
|
# log stream (the GUI already parses those lines).
|
|
# One authenticated operator at a time -- a new AUTH displaces the old.
|
|
# A slow/stalled operator NEVER blocks the relay: writes are buffered per
|
|
# connection and the connection is dropped if the buffer overflows.
|
|
CONTROL_PORT_OFFSET = 7 # control port = console port + 7
|
|
CONTROL_AUTH_TIMEOUT = 10.0 # unauthenticated sockets die after this
|
|
CONTROL_OUTBUF_MAX = 262144 # slow-reader cutoff (bytes)
|
|
CONTROL_SET_KEYS = {"map", "time", "weather", "scenario", "temperature",
|
|
"length"} # whitelist for the 'set' command
|
|
|
|
def control_secret():
|
|
"""Shared secret for control-channel auth, persisted next to the eggs.
|
|
Auto-generated on first use; the operator shares it out-of-band."""
|
|
import secrets as _secrets
|
|
# The relay always runs with cwd = content\ (operator GUI and all rigs
|
|
# set it); the secret lives beside the eggs it guards.
|
|
path = "operator_secret.txt"
|
|
try:
|
|
with open(path, "r", encoding="ascii") as f:
|
|
value = f.read().strip()
|
|
if value:
|
|
return value
|
|
except OSError:
|
|
pass
|
|
value = _secrets.token_hex(16)
|
|
try:
|
|
with open(path, "w", encoding="ascii") as f:
|
|
f.write(value + "\n")
|
|
except OSError:
|
|
pass
|
|
return value
|
|
|
|
|
|
# Module-level sink list consulted by _StampedOut (the wrapper is installed
|
|
# before any Relay exists). Guarded by a lock: the stdin thread prints too.
|
|
# _CONTROL_HISTORY keeps the last N log lines so an operator who connects
|
|
# MID-SESSION gets recent state replayed (their GUI rebuilds the roster from
|
|
# parsed lines -- without replay it would start blank).
|
|
_CONTROL_SINKS = []
|
|
_CONTROL_LOCK = threading.Lock()
|
|
_CONTROL_HISTORY = collections.deque(maxlen=400)
|
|
|
|
|
|
def _control_tee(text):
|
|
"""Best-effort fan-out of log text to authenticated control conns."""
|
|
data = text.encode("utf-8", "replace")
|
|
with _CONTROL_LOCK:
|
|
_CONTROL_HISTORY.append(data)
|
|
for conn in list(_CONTROL_SINKS):
|
|
conn.queue_out(data)
|
|
HELLO_MAGIC = 0x31525442 # 'BTR1' little-endian
|
|
FRAME_CAP = 1600 # NETWORKMANAGER_BUFFER_SIZE (NETWORK.h:89)
|
|
FIRST_GAME_HOST_ID = 2 # FirstLegalHostID(1) == the console; pods 2..N+1
|
|
STATS_PERIOD = 10.0
|
|
|
|
# LAN auto-discovery (BT_RELAY=auto): pods broadcast DISC_PROBE on this UDP
|
|
# port; the relay answers DISC_REPLY + <H consolePort>. The pod combines the
|
|
# replier's source IP with the advertised port. Broadcast never crosses
|
|
# routers, so this is LAN-only by construction; internet players keep using
|
|
# the explicit host:port.
|
|
DISCOVERY_PORT = 15999
|
|
DISC_PROBE = b"BTR1DISC"
|
|
DISC_REPLY = b"BTR1HERE"
|
|
|
|
|
|
def parse_egg_roster(egg_path):
|
|
"""Count the [pilots] pilot= entries -> expected hostIDs 2..N+1 in egg order.
|
|
The relay's ONLY game knowledge; it never parses game frames."""
|
|
entries = []
|
|
in_pilots = False
|
|
for raw in open(egg_path, "rb").read().replace(b"\r\n", b"\n").split(b"\n"):
|
|
line = raw.strip()
|
|
if line.startswith(b"["):
|
|
in_pilots = (line.lower() == b"[pilots]")
|
|
continue
|
|
if in_pilots and line.lower().startswith(b"pilot="):
|
|
entries.append(line.split(b"=", 1)[1].decode("ascii", "replace"))
|
|
# issue #38: an unknown camo color paints the mech SILENT GRAY (the
|
|
# engine tolerates the vehicle-table miss by design) -- warn the
|
|
# operator at load instead of letting a hand-edited egg ship it.
|
|
for raw in open(egg_path, "rb").read().replace(b"\r\n", b"\n").split(b"\n"):
|
|
line = raw.strip()
|
|
if line.lower().startswith(b"color="):
|
|
color = line.split(b"=", 1)[1].decode("ascii", "replace").strip()
|
|
if color and color not in VALID_CAMO_COLORS:
|
|
print(f"[relay] WARNING: egg color={color!r} is not a camo "
|
|
f"color (valid: {sorted(VALID_CAMO_COLORS)}) -- that "
|
|
f"mech will paint GRAY", flush=True)
|
|
return entries
|
|
|
|
|
|
class RelayGameConn:
|
|
"""One accepted game-TCP connection: ACCEPTED -> (HELLO) REGISTERED -> GONE."""
|
|
def __init__(self, sock, addr):
|
|
self.sock = sock
|
|
self.addr = addr
|
|
self.buf = b""
|
|
self.host_id = None # None until HELLO registers it
|
|
self.last_seen = time.time() # liveness deadline (see DEAD_PEER_SECONDS)
|
|
|
|
def name(self):
|
|
hid = self.host_id if self.host_id is not None else "?"
|
|
return f"game[{self.addr[0]}:{self.addr[1]} host={hid}]"
|
|
|
|
|
|
class RelayConsoleConn:
|
|
"""One accepted console connection (legacy raw protocol, relay-terminated)."""
|
|
def __init__(self, sock, addr):
|
|
self.sock = sock
|
|
self.addr = addr
|
|
self.egg_sent = False
|
|
self.acked = False # pod sent AcknowledgeEggFile -- a REAL pod.
|
|
# LIVENESS (2026-07-26): a pod whose machine sleeps, whose Wi-Fi drops,
|
|
# or which is killed without a FIN leaves a HALF-OPEN socket. With no
|
|
# keepalive and no deadline it used to sit here forever holding
|
|
# acked=True, which permanently blocked _maybe_reset_round (it requires
|
|
# NO conn acked) and held a phantom seat in the launch count. Only a
|
|
# console restart cleared it -- the operator's "reboot the console".
|
|
self.last_seen = time.time()
|
|
# Internet hardening: the console port is exposed, and random scanners
|
|
# connect and would otherwise count as pods. Only ACKED connections
|
|
# count toward the launch gate; scanners never speak the protocol.
|
|
|
|
|
|
class RelayControlConn:
|
|
"""One operator control connection. All sends are buffered/non-blocking
|
|
so a stalled remote operator can never stall the relay."""
|
|
|
|
def __init__(self, sock, addr):
|
|
self.sock = sock
|
|
self.addr = addr
|
|
self.buf = b""
|
|
self.outbuf = b""
|
|
self.authed = False
|
|
self.connected_at = time.time()
|
|
self.dead = False
|
|
|
|
def name(self):
|
|
return "%s:%d" % (self.addr[0], self.addr[1])
|
|
|
|
def queue_out(self, data):
|
|
if self.dead or not self.authed:
|
|
return
|
|
self.outbuf += data
|
|
if len(self.outbuf) > CONTROL_OUTBUF_MAX:
|
|
self.dead = True # slow reader: cut it loose
|
|
return
|
|
self.flush_out()
|
|
|
|
def flush_out(self):
|
|
if self.dead or not self.outbuf:
|
|
return
|
|
try:
|
|
n = self.sock.send(self.outbuf)
|
|
self.outbuf = self.outbuf[n:]
|
|
except (BlockingIOError, InterruptedError):
|
|
pass
|
|
except OSError:
|
|
self.dead = True
|
|
|
|
|
|
class Relay:
|
|
def __init__(self, console_port, egg_path, bind_addr, udp_drop_pct,
|
|
manual_launch=False, reserved_tags=None):
|
|
self.console_port = console_port
|
|
self.game_port = console_port + 1
|
|
self.bind_addr = bind_addr
|
|
self.udp_drop_pct = udp_drop_pct
|
|
self.manual_launch = manual_launch
|
|
self._reserved_tags = set(reserved_tags or ())
|
|
self.launch_requested = False # set by the operator (stdin)
|
|
self._launch_blocked_warned = False # warned once about empty seats
|
|
self.stop_requested = False # operator 'stop' (End Mission)
|
|
self.mission_length = parse_egg_mission_length(egg_path)
|
|
self.mission_started_at = None # set when RunMission #2 fires
|
|
self.stop_sent = False
|
|
self.eggs_done_at = None
|
|
self.egg_bytes = egg_wire_bytes(open(egg_path, "rb").read())
|
|
self.egg_path = egg_path
|
|
# WALK-UP SEAT PREFS (2026-07-22): callsign/mech requested in the
|
|
# SEAT_REQUEST payload, keyed by assigned host id. Eggs are HELD
|
|
# until the roster completes so every pod's egg carries every
|
|
# player's callsign (an early pod's egg would otherwise miss late
|
|
# joiners' labels); pods animate in WAITING FOR MISSION ASSIGNMENT
|
|
# meanwhile.
|
|
self.seat_prefs = {}
|
|
self.eggs_released = False
|
|
# SEAT RECLAIM (2026-07-23): tag -> expiry; set when a seated player
|
|
# drops (round end / crash). A rejoin presenting that tag gets the
|
|
# SAME seat back (static ordering, real-pod style); other joiners
|
|
# skip reclaim-held seats until the grace expires.
|
|
self.seat_reclaim = {}
|
|
# ROUND RESET originals (2026-07-22): trim/prefs FINALIZE the egg for
|
|
# one round; when every pod leaves, the round resets to these so the
|
|
# NEXT round's joins hold/rewrite a FRESH egg. (Field report: a
|
|
# rejoin after a crash got the stale finalized egg -- the new mech
|
|
# request was recorded but never applied.)
|
|
self.orig_egg_path = egg_path
|
|
self.orig_egg_bytes = self.egg_bytes
|
|
# PRESENCE BEACONS (2026-07-22, "launch with whoever connects"): the
|
|
# pod KEEPS its seat-request TCP connection open until it exits; a
|
|
# live beacon = a seated player. A beacon dropping before the seat
|
|
# was claimed (game-side HELLO) FREES the seat. The operator's
|
|
# LAUNCH trims the roster to the seats present.
|
|
self.seat_beacons = {}
|
|
self.roster = parse_egg_roster(egg_path)
|
|
self.orig_roster = list(self.roster) # round-reset restore point
|
|
self.expected_ids = set(range(FIRST_GAME_HOST_ID,
|
|
FIRST_GAME_HOST_ID + len(self.roster)))
|
|
# OPERATOR-RESERVED SEATS: roster seats the operator flies LOCALLY (with
|
|
# an explicit BT_SELF). Those instances take ~15s to boot, so a remote
|
|
# join.bat (no BT_SELF -> SEAT_REQUEST) would otherwise be assigned the
|
|
# operator's seat first, and the operator's later HELLO would collide
|
|
# ("already registered" -> dropped -> game never launches). Seat
|
|
# ASSIGNMENT skips these host_ids; the explicit HELLO still claims them.
|
|
self.reserved_host_ids = set()
|
|
for i, tag in enumerate(self.roster):
|
|
if tag in self._reserved_tags:
|
|
self.reserved_host_ids.add(FIRST_GAME_HOST_ID + i)
|
|
if self.reserved_host_ids:
|
|
print(f"[relay] operator-reserved seats (no auto-assign): "
|
|
f"{sorted(self.reserved_host_ids)}", flush=True)
|
|
self.sel = selectors.DefaultSelector()
|
|
self.console_conns = []
|
|
self.game_conns = [] # RelayGameConn (incl. unregistered)
|
|
self.by_host = {} # hostID -> RelayGameConn
|
|
self.seat_reservations = {} # hostID -> expiry (auto-seats)
|
|
self.udp_endpoint = {} # hostID -> (addr, port)
|
|
self.udp_sock = None
|
|
self.launch_at = None
|
|
self.launches_sent = 0
|
|
self.stats = {"tcp_rx": 0, "tcp_tx": 0, "udp_rx": 0, "udp_tx": 0,
|
|
"udp_dropped": 0, "udp_tcp_fallback": 0}
|
|
self.stats_at = time.time() + STATS_PERIOD
|
|
self.ctl_conns = [] # RelayControlConn list
|
|
self.ctl_secret = control_secret()
|
|
|
|
# ---------------- lifecycle ----------------
|
|
|
|
def run(self):
|
|
print(f"[relay] roster: {len(self.roster)} pilot(s) -> hostIDs "
|
|
f"{sorted(self.expected_ids)}: {self.roster}", flush=True)
|
|
con_l = self._listener(self.console_port)
|
|
game_l = self._listener(self.game_port)
|
|
self.udp_sock = socket.socket(socket.AF_INET, socket.SOCK_DGRAM)
|
|
self.udp_sock.bind((self.bind_addr, self.game_port))
|
|
self.udp_sock.setblocking(False)
|
|
self.sel.register(con_l, selectors.EVENT_READ, ("con_listen", None))
|
|
self.sel.register(game_l, selectors.EVENT_READ, ("game_listen", None))
|
|
self.sel.register(self.udp_sock, selectors.EVENT_READ, ("udp", None))
|
|
# Remote-operator control channel (best-effort; the game relay must
|
|
# come up even if this port is taken).
|
|
self.ctl_port = self.console_port + CONTROL_PORT_OFFSET
|
|
try:
|
|
ctl_l = self._listener(self.ctl_port)
|
|
self.sel.register(ctl_l, selectors.EVENT_READ, ("ctl_listen", None))
|
|
print(f"[ctl] operator control on port {self.ctl_port} "
|
|
f"(secret in content\\operator_secret.txt)", flush=True)
|
|
except OSError as e:
|
|
print(f"[ctl] control port {self.ctl_port} unavailable ({e!r}) "
|
|
f"-- remote operators disabled; stdin still works",
|
|
flush=True)
|
|
# LAN auto-discovery responder (best-effort: another relay may own it).
|
|
self.disc_sock = None
|
|
try:
|
|
self.disc_sock = socket.socket(socket.AF_INET, socket.SOCK_DGRAM)
|
|
self.disc_sock.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
|
|
self.disc_sock.bind((self.bind_addr, DISCOVERY_PORT))
|
|
self.disc_sock.setblocking(False)
|
|
self.sel.register(self.disc_sock, selectors.EVENT_READ,
|
|
("disc", None))
|
|
print(f"[relay] LAN discovery answering on udp/{DISCOVERY_PORT} "
|
|
f"(pods may use BT_RELAY=auto)", flush=True)
|
|
except OSError as e:
|
|
print(f"[relay] LAN discovery unavailable ({e!r}) -- "
|
|
f"pods must use the explicit address", flush=True)
|
|
self.disc_sock = None
|
|
print(f"[relay] console {self.bind_addr}:{self.console_port} | game tcp/udp "
|
|
f"{self.bind_addr}:{self.game_port} | udp-drop {self.udp_drop_pct}%",
|
|
flush=True)
|
|
while True:
|
|
for key, _ in self.sel.select(timeout=0.5):
|
|
kind, obj = key.data
|
|
if kind == "con_listen":
|
|
self._accept_console(key.fileobj)
|
|
elif kind == "game_listen":
|
|
self._accept_game(key.fileobj)
|
|
elif kind == "console":
|
|
self._console_read(obj)
|
|
elif kind == "game":
|
|
self._game_read(obj)
|
|
elif kind == "udp":
|
|
self._udp_read()
|
|
elif kind == "disc":
|
|
self._disc_read()
|
|
elif kind == "ctl_listen":
|
|
self._ctl_accept(key.fileobj)
|
|
elif kind == "ctl":
|
|
self._ctl_read(obj)
|
|
self._tick_control()
|
|
self._reap_dead_peers()
|
|
self._tick_launch()
|
|
self._tick_stop()
|
|
self._tick_settle()
|
|
self._tick_stats()
|
|
|
|
def _listener(self, port):
|
|
s = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
|
|
s.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
|
|
s.bind((self.bind_addr, port))
|
|
s.listen(16)
|
|
s.setblocking(False)
|
|
return s
|
|
|
|
# ---------------- console side (legacy protocol, relay-terminated) ----------------
|
|
|
|
def _accept_console(self, listener):
|
|
sock, addr = listener.accept()
|
|
_enable_keepalive(sock)
|
|
sock.setblocking(False)
|
|
conn = RelayConsoleConn(sock, addr)
|
|
self.console_conns.append(conn)
|
|
self.sel.register(sock, selectors.EVENT_READ, ("console", conn))
|
|
# Egg delivery: HELD until every pod's console pad is connected
|
|
# (walk-up seat prefs must be in every pod's egg copy). NOT gated on
|
|
# game-side registration -- a pod HELLOs only after it parses the
|
|
# egg's roster, so that gate would deadlock. Seat requests precede
|
|
# the console pad, so by the Nth pad every pref is already in.
|
|
if self.eggs_released:
|
|
self._send_egg(conn)
|
|
elif len(self.console_conns) >= len(self.roster):
|
|
self._release_eggs()
|
|
else:
|
|
print(f"[relay] console conn {addr[0]}:{addr[1]}: egg HELD "
|
|
f"({len(self.console_conns)}/{len(self.roster)} pods "
|
|
f"present)", flush=True)
|
|
|
|
def _send_egg(self, conn):
|
|
addr = conn.addr # NOT getpeername(): it raises on a dead peer
|
|
try:
|
|
n = 0
|
|
for pkt in egg_packets(self.egg_bytes):
|
|
_send_all_guarded(conn.sock, pkt, "egg chunk")
|
|
n += 1
|
|
conn.egg_sent = True
|
|
print(f"[relay] console conn {addr[0]}:{addr[1]}: egg sent "
|
|
f"({n} chunks); awaiting pod ACK", flush=True)
|
|
except OSError as e:
|
|
print(f"[relay] console conn {addr[0]}:{addr[1]}: egg send failed {e!r}",
|
|
flush=True)
|
|
self._drop_console(conn)
|
|
return
|
|
finally:
|
|
try:
|
|
conn.sock.setblocking(False)
|
|
except OSError:
|
|
pass
|
|
|
|
def _reload_egg_file(self):
|
|
"""BETWEEN-ROUNDS MISSION EDITS (2026-07-23): re-read the session egg
|
|
FILE so operator changes (arena/time/weather/length via the console's
|
|
Apply button) take effect next round without a session restart. The
|
|
roster SHAPE stays fixed mid-session (seat ids are positional): a
|
|
row-count change is refused loudly."""
|
|
try:
|
|
new_roster = parse_egg_roster(self.orig_egg_path)
|
|
if len(new_roster) != len(self.orig_roster):
|
|
print(f"[relay] egg file roster changed "
|
|
f"{len(self.orig_roster)}->{len(new_roster)} seats -- "
|
|
f"IGNORED (restart the session to resize)", flush=True)
|
|
return
|
|
self.orig_egg_bytes = egg_wire_bytes(
|
|
open(self.orig_egg_path, "rb").read())
|
|
self.egg_bytes = self.orig_egg_bytes
|
|
self.egg_path = self.orig_egg_path
|
|
new_length = parse_egg_mission_length(self.orig_egg_path)
|
|
if new_length != self.mission_length:
|
|
print(f"[relay] mission length now {new_length:.0f}s",
|
|
flush=True)
|
|
self.mission_length = new_length
|
|
except OSError as e:
|
|
print(f"[relay] egg reload failed ({e!r}) -- keeping the "
|
|
f"previous mission", flush=True)
|
|
|
|
def _release_eggs(self, present=None):
|
|
if self.eggs_released:
|
|
return
|
|
# issue #33: after a ROUND ABORT, hold the next release until the
|
|
# settle window passes -- rejoiners re-ACK within a second and an
|
|
# instant re-release hands the flapping pod a fresh round to abort.
|
|
if time.time() < getattr(self, "round_hold_until", 0):
|
|
if not getattr(self, "_hold_noted", False):
|
|
self._hold_noted = True
|
|
print(f"[relay] egg release HELD "
|
|
f"({self.round_hold_until - time.time():.0f}s settle "
|
|
f"after the round abort)", flush=True)
|
|
return
|
|
self._hold_noted = False
|
|
self._reload_egg_file() # pick up between-rounds edits
|
|
self.eggs_released = True
|
|
if present is not None and len(present) < len(self.roster):
|
|
self._trim_roster(present)
|
|
if self.seat_prefs:
|
|
self._apply_seat_prefs()
|
|
for conn in list(self.console_conns):
|
|
if not conn.egg_sent:
|
|
self._send_egg(conn)
|
|
|
|
def _trim_roster(self, present):
|
|
"""LAUNCH WITH WHOEVER CONNECTS: shrink the session to the seats
|
|
with live beacons (or claims). Positions shift, so seat ids REMAP --
|
|
a pod re-derives its host id from its TAG's position in the received
|
|
egg, so prefs/beacons re-key by position and every pod stays
|
|
consistent. Runs pre-release only (no pod has HELLO'd yet)."""
|
|
keep_idx = [h - FIRST_GAME_HOST_ID for h in present]
|
|
keep_tags = [self.roster[i] for i in keep_idx if 0 <= i < len(self.roster)]
|
|
old_ids = list(present)
|
|
print(f"[relay] LAUNCH trim: {len(keep_tags)}/{len(self.roster)} "
|
|
f"seats present -- roster shrinks to {keep_tags}", flush=True)
|
|
if eggmodel is not None:
|
|
try:
|
|
doc = eggmodel.EggDoc.load(self.egg_path)
|
|
for i, pl in enumerate(doc.pilots()):
|
|
if i not in keep_idx:
|
|
doc.remove_section(pl["address"])
|
|
doc.replace_section("pilots",
|
|
["pilot=%s" % t for t in keep_tags])
|
|
trimmed = os.path.join(os.path.dirname(self.egg_path) or ".",
|
|
"_relay_trimmed.egg")
|
|
doc.save(trimmed)
|
|
self.egg_path = trimmed # prefs rewrite loads THIS
|
|
self.egg_bytes = egg_wire_bytes(
|
|
doc.emit().encode("latin-1", "replace"))
|
|
except Exception as e:
|
|
print(f"[relay] roster trim FAILED ({e!r}) -- launching the "
|
|
f"FULL roster (empty seats will stall!)", flush=True)
|
|
return
|
|
# re-key everything by the new positions
|
|
remap = {old: FIRST_GAME_HOST_ID + n for n, old in enumerate(old_ids)}
|
|
self.seat_prefs = {remap[h]: v for h, v in self.seat_prefs.items()
|
|
if h in remap}
|
|
new_beacons = {}
|
|
for old, conn in self.seat_beacons.items():
|
|
if old in remap:
|
|
conn.seat_host = remap[old]
|
|
new_beacons[remap[old]] = conn
|
|
self.seat_beacons = new_beacons
|
|
self.roster = keep_tags
|
|
self.expected_ids = set(range(FIRST_GAME_HOST_ID,
|
|
FIRST_GAME_HOST_ID + len(keep_tags)))
|
|
|
|
def _apply_seat_prefs(self):
|
|
# Rewrite vehicle= + the callsign system from the walk-up prefs.
|
|
# eggmodel's rasterizer needs PySide6 (present when launched from the
|
|
# operator app); without it the vehicle still applies and the
|
|
# callsign falls back to the text name= (labels keep the roster
|
|
# default bitmaps).
|
|
if eggmodel is None:
|
|
print("[relay] eggmodel unavailable -- seat prefs NOT applied",
|
|
flush=True)
|
|
return
|
|
try:
|
|
doc = eggmodel.EggDoc.load(self.egg_path)
|
|
pilots = doc.pilots()
|
|
names = [p.get("name") or p["address"] for p in pilots]
|
|
for host_id, (callsign, mech) in sorted(self.seat_prefs.items()):
|
|
i = host_id - FIRST_GAME_HOST_ID
|
|
if not (0 <= i < len(pilots)):
|
|
continue
|
|
if mech:
|
|
doc.set_kv(pilots[i]["address"], "vehicle", mech)
|
|
if callsign:
|
|
names[i] = callsign
|
|
try:
|
|
doc.set_callsigns(eggmodel.make_callsigns(names))
|
|
except Exception as e:
|
|
print(f"[relay] callsign rasterizer unavailable ({e!r}) -- "
|
|
f"vehicles applied, callsign bitmaps unchanged",
|
|
flush=True)
|
|
self.egg_bytes = egg_wire_bytes(
|
|
doc.emit().encode("latin-1", "replace"))
|
|
for host_id, (callsign, mech) in sorted(self.seat_prefs.items()):
|
|
i = host_id - FIRST_GAME_HOST_ID
|
|
if 0 <= i < len(pilots):
|
|
print(f"[relay] player {host_id - FIRST_GAME_HOST_ID + 1} "
|
|
f"egg: vehicle={mech or pilots[i].get('vehicle')} "
|
|
f"callsign={callsign or names[i]!r}", flush=True)
|
|
except Exception as e:
|
|
print(f"[relay] seat-pref egg rewrite FAILED ({e!r}) -- "
|
|
f"serving the original egg", flush=True)
|
|
|
|
def _maybe_reset_round(self):
|
|
# All pods gone from a FINALIZED round -> restore the template so the
|
|
# next round's joins are honored (held egg, fresh prefs/trim). Live
|
|
# beacons (players already waiting for the next round) are re-keyed
|
|
# by their TAG's position in the restored roster.
|
|
if not self.eggs_released:
|
|
return
|
|
if self.by_host or any(c.acked for c in self.console_conns):
|
|
return
|
|
self.eggs_released = False
|
|
self.egg_path = self.orig_egg_path
|
|
self._reload_egg_file()
|
|
self.egg_bytes = self.orig_egg_bytes
|
|
self.roster = list(self.orig_roster)
|
|
self.expected_ids = set(range(FIRST_GAME_HOST_ID,
|
|
FIRST_GAME_HOST_ID + len(self.roster)))
|
|
self.launches_sent = 0
|
|
self.launch_at = None
|
|
self.launch_requested = False
|
|
self.eggs_done_at = None
|
|
self._launch_blocked_warned = False
|
|
self.stop_sent = False
|
|
self.stop_requested = False
|
|
self.mission_started_at = None
|
|
old_beacons = self.seat_beacons
|
|
old_prefs = self.seat_prefs
|
|
self.seat_beacons = {}
|
|
self.seat_prefs = {}
|
|
self.seat_reservations = {}
|
|
for old_id, conn in old_beacons.items():
|
|
tag = getattr(conn, "seat_tag", None)
|
|
if tag in self.roster:
|
|
new_id = FIRST_GAME_HOST_ID + self.roster.index(tag)
|
|
conn.seat_host = new_id
|
|
self.seat_beacons[new_id] = conn
|
|
if old_id in old_prefs:
|
|
self.seat_prefs[new_id] = old_prefs[old_id]
|
|
print(f"[relay] round RESET -- roster/egg restored "
|
|
f"({len(self.seat_beacons)} player(s) already waiting)",
|
|
flush=True)
|
|
|
|
def _reap_dead_peers(self):
|
|
"""Age out peers that stopped speaking entirely.
|
|
|
|
TCP keepalive (set on accept) handles most half-open sockets, but it is
|
|
best-effort and platform-dependent, and a peer can also be alive at the
|
|
TCP level while its game process is gone. A silent connection is what
|
|
used to pin the relay: an acked console conn blocks the round reset
|
|
forever, and a registered game conn holds its by_host seat, so BOTH
|
|
launch re-arm paths stay unreachable and the operator has to restart.
|
|
|
|
Deliberately generous (DEAD_PEER_SECONDS): a pod legitimately sends
|
|
nothing while sitting at the front end between rounds, so this must only
|
|
catch the genuinely gone. Never applied while a mission is running --
|
|
an in-mission pod that goes quiet is the netcode's problem, not ours.
|
|
"""
|
|
if self.launches_sent >= 2:
|
|
return # mission running: leave peers alone
|
|
if not self.eggs_released:
|
|
#
|
|
# BETWEEN ROUNDS / BEFORE THE EGG GOES OUT, SILENCE IS CORRECT and
|
|
# can last a long time (a real night showed 12- and 5-minute gaps
|
|
# while the operator set up the next mission). A waiting pod has
|
|
# nothing to send: its seat beacon is deliberately write-only for the
|
|
# process lifetime -- L4NET.CPP: "the relay ignores its silence" --
|
|
# and its console pad has no egg yet, so it cannot ACK. Reaping here
|
|
# would drop HEALTHY players and force their client to relaunch,
|
|
# which is worse than the ghost this reaper exists to remove.
|
|
# (Caught in review 2026-07-26, before it ever ran on a real night.)
|
|
#
|
|
return
|
|
#
|
|
# So the reaper only runs in the ACTIVE STAGING WINDOW: the egg is out and
|
|
# we are waiting for ACK/READY, which is exactly when a pod that has gone
|
|
# quiet is a pod that is gone. A truly half-open socket is caught sooner
|
|
# and more cheaply by TCP keepalive (set on accept, and the pod sets it on
|
|
# its own beacon too) -- this deadline is only the backstop for platforms
|
|
# where the keepalive tunables do not take.
|
|
#
|
|
now = time.time()
|
|
for conn in list(self.console_conns):
|
|
if not conn.egg_sent:
|
|
continue # nothing was asked of it yet
|
|
if now - getattr(conn, "last_seen", now) > DEAD_PEER_SECONDS:
|
|
print(f"[relay] console conn {conn.addr[0]}:{conn.addr[1]} "
|
|
f"REAPED -- silent for "
|
|
f"{now - conn.last_seen:.0f}s while staging "
|
|
f"(acked={conn.acked}); it can no longer hold the round "
|
|
f"open", flush=True)
|
|
self._drop_console(conn)
|
|
for conn in list(self.game_conns):
|
|
if conn.host_id is None:
|
|
continue # beacon / not yet HELLO'd: silent by design
|
|
if now - getattr(conn, "last_seen", now) > DEAD_PEER_SECONDS:
|
|
self._drop_game(conn, "silent for %.0fs while staging -- reaped"
|
|
% (now - conn.last_seen))
|
|
|
|
def _pods_ready(self):
|
|
"""REAL pods = console connections that ACKED the egg (scanners never
|
|
speak the protocol, so they can't hold a roster slot)."""
|
|
return sum(1 for c in self.console_conns if c.acked)
|
|
|
|
def _check_launch_gate(self):
|
|
# BACK-TO-BACK MISSIONS (2026-07-18): a previous mission finished
|
|
# (launches_sent==2) and now all seats have RE-ACKED -- i.e. every pod
|
|
# exited, re-ran join.bat, and reconnected fresh. Reset the launch
|
|
# state so the relay (and the operator's LAUNCH button) re-arm for a
|
|
# new round -- no Stop/Start needed. The binary needs nothing: a
|
|
# rejoined pod is a brand-new process at WaitingForLaunch. (Only ever
|
|
# reached on a pod ACK, which never happens mid-mission.)
|
|
if self.launches_sent >= 2 and self._pods_ready() >= len(self.roster):
|
|
print("[relay] all seats rejoined after the last mission -- "
|
|
"re-arming for a NEW mission", flush=True)
|
|
self.launches_sent = 0
|
|
self.stop_sent = False
|
|
self.stop_requested = False
|
|
self.mission_started_at = None
|
|
self.launch_at = None
|
|
self.launch_requested = False
|
|
self._launch_blocked_warned = False
|
|
self.eggs_done_at = None
|
|
# fall through -> the normal gate re-prints WAITING FOR OPERATOR
|
|
|
|
# GLOBAL launch: arm the timer when the LAST pod has ACKED its egg
|
|
# (auto mode), or wait for the operator's 'launch' command (manual
|
|
# mode -- the operator app's Launch button writes it to our stdin).
|
|
if self._pods_ready() >= len(self.roster) and self.launches_sent == 0 \
|
|
and self.eggs_done_at is None:
|
|
self.eggs_done_at = time.time()
|
|
if self.manual_launch:
|
|
print("[relay] all pods ACKED the egg; WAITING FOR OPERATOR "
|
|
"LAUNCH", flush=True)
|
|
else:
|
|
self.launch_at = self.eggs_done_at + LAUNCH_SETTLE_SECONDS
|
|
print(f"[relay] all pods ACKED the egg; LAUNCH pair in "
|
|
f"{LAUNCH_SETTLE_SECONDS}s", flush=True)
|
|
|
|
def _console_read(self, conn):
|
|
try:
|
|
data = conn.sock.recv(4096)
|
|
except OSError:
|
|
data = b""
|
|
if data:
|
|
conn.last_seen = time.time()
|
|
if not data:
|
|
print(f"[relay] console conn {conn.addr[0]}:{conn.addr[1]} closed",
|
|
flush=True)
|
|
self._drop_console(conn)
|
|
return
|
|
if len(data) >= 24:
|
|
(clid, _gid, fh, _ts) = struct.unpack_from("<iiii", data, 0)
|
|
(mlen, mid) = struct.unpack_from("<Ii", data, 16)
|
|
print(f"[relay] pod->console clientID={clid} fromHost={fh} msgID={mid} "
|
|
f"len={mlen}", flush=True)
|
|
# AcknowledgeEggFile (clientID=0 NetworkManager, msgID=4): this
|
|
# connection is a REAL pod -- count it toward the launch gate.
|
|
if clid == 0 and mid == 4 and not conn.acked:
|
|
conn.acked = True
|
|
print(f"[relay] pod ACK from {conn.addr[0]}:{conn.addr[1]} "
|
|
f"({self._pods_ready()}/{len(self.roster)} ready)",
|
|
flush=True)
|
|
self._check_launch_gate()
|
|
|
|
def _drop_console(self, conn):
|
|
try:
|
|
self.sel.unregister(conn.sock)
|
|
except (KeyError, ValueError):
|
|
pass
|
|
try:
|
|
conn.sock.close()
|
|
except OSError:
|
|
pass
|
|
if conn in self.console_conns:
|
|
self.console_conns.remove(conn)
|
|
self._maybe_reset_round()
|
|
|
|
def _log_launch_readiness(self):
|
|
# PRE-LAUNCH READINESS (2026-07-18): make the "launching short" footgun
|
|
# VISIBLE. A roster seat with no registered pod means the mission will
|
|
# STALL -- every pod waits on that empty seat's connection and the
|
|
# All-connections-completed gate never fires. Log filled vs empty
|
|
# seats so the operator sees it BEFORE the confusion.
|
|
filled, empty = [], []
|
|
for i, tag in enumerate(self.roster):
|
|
host_id = FIRST_GAME_HOST_ID + i
|
|
(filled if host_id in self.by_host else empty).append(
|
|
f"seat{i + 1}({tag})")
|
|
print(f"[relay] LAUNCH readiness: {len(filled)}/{len(self.roster)} "
|
|
f"seats filled -- {', '.join(filled) or 'NONE'}", flush=True)
|
|
if empty:
|
|
print(f"[relay] *** WARNING: {len(empty)} EMPTY seat(s): "
|
|
f"{', '.join(empty)} -- the mission will STALL (every pod "
|
|
f"waits on the empty seat). Stop, reduce the roster to the "
|
|
f"players present, and re-launch.", flush=True)
|
|
|
|
def _rearm_for_new_round(self, why):
|
|
"""OPERATOR-DRIVEN RE-ARM (fixes the "LAUNCH does nothing after a round"
|
|
wedge, operator report 2026-07-25).
|
|
|
|
The two automatic re-arm paths are both ALL-OR-NOTHING:
|
|
`_check_launch_gate` needs EVERY seat of the last round to re-ACK, and
|
|
`_maybe_reset_round` needs EVERY pod to be gone. A games night lives
|
|
between those: most pods rejoin, one player closes their window or
|
|
crashes. In that state `launches_sent` stayed latched at 2, the manual
|
|
LAUNCH branch (guarded on `launches_sent == 0`) ignored every press, and
|
|
because the "not all seats filled" diagnostic sits inside that same
|
|
guard, NOTHING was printed -- the operator saw a dead button and had to
|
|
restart the session. (The operator confirmed the mirror image: with a
|
|
stable group, where everyone does rejoin, relaunching worked fine.)
|
|
|
|
So: an explicit operator LAUNCH is now itself sufficient authority to
|
|
start a new round. Reset only the ROUND state -- seats, beacons, prefs
|
|
and connections are left alone, so whoever is here stays here.
|
|
"""
|
|
print(f"[relay] RE-ARM for a new round ({why}); "
|
|
f"previous round state cleared", flush=True)
|
|
self.launches_sent = 0
|
|
self.launch_at = None
|
|
self.stop_sent = False
|
|
# A stop the operator asked for during the LAST round must not carry
|
|
# into this one -- see _tick_stop, where a stale latch StopMissioned
|
|
# the next mission in the tick it launched.
|
|
self.stop_requested = False
|
|
self.mission_started_at = None
|
|
self.eggs_done_at = None
|
|
self._launch_blocked_warned = False
|
|
# Re-open the round for joins/edits: restore the untrimmed roster and the
|
|
# template egg, so round 2+ does not serve last round's trimmed egg (new
|
|
# joiners got ROSTER FULL) and between-round mission edits are picked up.
|
|
#
|
|
# UNCONDITIONAL, and that matters. An earlier cut gated this on
|
|
# `eggs_released`, which made it useless in the one state that most needs
|
|
# it: _abort_round clears eggs_released BEFORE the survivors' sockets
|
|
# close, so _maybe_reset_round's own `if not self.eggs_released: return`
|
|
# then skips the template restore forever. The roster stays trimmed to
|
|
# last round's attendance, the release gates (which all compare against
|
|
# len(self.roster)) can never be met, walk-ups get ROSTER FULL, and the
|
|
# relay sits dead -- live-proven on 2026-07-25, nearly 5 minutes of a
|
|
# games night lost. Re-arm has to be able to dig out of exactly that.
|
|
self.eggs_released = False
|
|
self.egg_path = self.orig_egg_path
|
|
self.egg_bytes = self.orig_egg_bytes
|
|
self.roster = list(self.orig_roster)
|
|
self.expected_ids = set(range(FIRST_GAME_HOST_ID,
|
|
FIRST_GAME_HOST_ID + len(self.roster)))
|
|
self.round_hold_until = 0 # do not inherit an abort's settle window
|
|
self._reload_egg_file()
|
|
for c in self.console_conns:
|
|
c.acked = False # last round's ACK must not count
|
|
c.egg_sent = False # every pod needs THIS round's egg
|
|
|
|
def _tick_launch(self):
|
|
# RE-ARM ON DEMAND: an operator LAUNCH while a previous round is still
|
|
# latched (launches_sent >= 1 -- 2 after a normal round, 1 if a round
|
|
# was aborted between the RunMission pair) is a deliberate new-round
|
|
# request, not a no-op. This must run BEFORE the gate below, whose
|
|
# `launches_sent == 0` guard is exactly what used to swallow the press.
|
|
# `launch_at is not None` means a launch SEQUENCE IS IN FLIGHT: the
|
|
# RunMission pair is timed, and `launch_requested` deliberately stays set
|
|
# across it (it is only cleared once #2 has gone out). Re-arming on that
|
|
# state resets launches_sent between #1 and #2, so the pair never
|
|
# completes and the relay re-releases eggs every tick -- a re-arm storm
|
|
# that kicks every pod into an identity-resync loop. Caught on the rig,
|
|
# 2026-07-26. So: only re-arm when nothing is in flight.
|
|
if (self.manual_launch and self.launch_requested
|
|
and self.launch_at is None and self.launches_sent >= 1):
|
|
self._rearm_for_new_round("operator LAUNCH after a finished round")
|
|
# fall through: the request stands and is served by the gate below.
|
|
|
|
# Manual mode: the operator's 'launch' arms the timer, still honouring
|
|
# the settle window (pods must reach WaitingForLaunch or the handler
|
|
# Fail()s -- the same reason the auto timer waits).
|
|
if (self.manual_launch and self.launch_requested
|
|
and self.launch_at is None and self.launches_sent == 0):
|
|
if not self.eggs_released:
|
|
# LAUNCH WITH WHOEVER CONNECTS: the operator fired before the
|
|
# roster filled -- shrink the session to the seated players
|
|
# and release the eggs; the launch pair fires once they ACK.
|
|
present = sorted(set(self.seat_beacons) | set(self.by_host))
|
|
if not present:
|
|
if not self._launch_blocked_warned:
|
|
self._launch_blocked_warned = True
|
|
print("[relay] LAUNCH pressed but NO players are "
|
|
"seated yet -- waiting for joins", flush=True)
|
|
return
|
|
self._launch_blocked_warned = False
|
|
self._release_eggs(present)
|
|
return # gate proceeds as their ACKs land
|
|
if self.eggs_done_at is not None:
|
|
self.launch_at = max(time.time(),
|
|
self.eggs_done_at + LAUNCH_SETTLE_SECONDS)
|
|
wait = max(0.0, self.launch_at - time.time())
|
|
print(f"[relay] operator LAUNCH received; firing in "
|
|
f"{wait:.0f}s", flush=True)
|
|
self._log_launch_readiness()
|
|
elif not self._launch_blocked_warned:
|
|
# Operator pressed LAUNCH but the gate is NOT ready (not every
|
|
# roster seat has ACKed its egg) -- previously this did NOTHING
|
|
# visible ("I pressed launch and nothing happened"). Say WHY.
|
|
self._launch_blocked_warned = True
|
|
print("[relay] LAUNCH pressed but NOT all seats have ACKed "
|
|
"their egg yet -- holding:", flush=True)
|
|
self._log_launch_readiness()
|
|
print("[relay] if a listed seat is never coming back, press "
|
|
"RE-ARM to re-open the round and launch with whoever is "
|
|
"here (no session restart needed)", flush=True)
|
|
if self.launch_at is None or self.launches_sent >= 2:
|
|
return
|
|
if time.time() < self.launch_at:
|
|
return
|
|
# ALL-READY GATE (2026-07-22, live finding round 3): firing the
|
|
# mission while a pod is still LOADING drowns that pod's load in the
|
|
# running mission's update streams (same event queue) -- the operator,
|
|
# who always joins last, loaded 5+ minutes or crashed while testers
|
|
# took ~20s. Hold the RunMission pair until every REGISTERED pod has
|
|
# sent READY (route -10): everyone enters together, 1995-style.
|
|
# Escape hatches: a pod that never readies is announced every 10s
|
|
# (the operator sees who), and a 3-minute cap force-fires so one
|
|
# wedged pod cannot hold the night hostage.
|
|
if self.launches_sent == 0:
|
|
not_ready = [h for h, c in sorted(self.by_host.items())
|
|
if not getattr(c, "ready", False)]
|
|
if not_ready and time.time() < self.launch_at + 180.0:
|
|
now = time.time()
|
|
if now - getattr(self, "_hold_notice_at", 0) >= 10.0:
|
|
self._hold_notice_at = now
|
|
names = ", ".join(
|
|
"PLAYER %d" % (h - FIRST_GAME_HOST_ID + 1)
|
|
for h in not_ready)
|
|
print(f"[relay] launch HELD -- still loading: {names} "
|
|
f"({len(self.by_host) - len(not_ready)}/"
|
|
f"{len(self.by_host)} ready)", flush=True)
|
|
return
|
|
if not_ready:
|
|
print(f"[relay] launch FORCED after 180s hold -- "
|
|
f"{len(not_ready)} pod(s) never readied", flush=True)
|
|
pkt = run_mission_packet()
|
|
for conn in list(self.console_conns):
|
|
try:
|
|
_send_all_guarded(conn.sock, pkt, "launch")
|
|
except OSError as e:
|
|
print(f"[relay] launch send failed to {conn.addr}: {e!r} "
|
|
f"-- dropping that pod", flush=True)
|
|
self._drop_console(conn)
|
|
self.launches_sent += 1
|
|
self.launch_at = time.time() + LAUNCH_STEP_SECONDS
|
|
print(f"[relay] RunMission #{self.launches_sent} sent to "
|
|
f"{len(self.console_conns)} pod(s)", flush=True)
|
|
if self.launches_sent == 2:
|
|
self.mission_started_at = time.time()
|
|
# CONSUME the launch: clear the timer + the request so the NEXT
|
|
# LAUNCH press is a fresh, deliberate new-mission request (else the
|
|
# back-to-back block would re-fire while this mission is running).
|
|
self.launch_at = None
|
|
self.launch_requested = False
|
|
if self.mission_length > 0:
|
|
print(f"[relay] mission clock: {self.mission_length:.0f}s "
|
|
"(StopMission at expiry)", flush=True)
|
|
|
|
def _tick_stop(self):
|
|
# THE MISSION CLOCK (console-side, as in 1995): send StopMission when
|
|
# the egg's [mission] length expires, or on the operator's 'stop'
|
|
# command (the End Mission button).
|
|
if self.launches_sent < 2:
|
|
# NO MISSION IS RUNNING. A stop asked for now (End Mission pressed
|
|
# between rounds, which is a natural operator reflex) used to sit
|
|
# latched and then StopMission the NEXT mission in the very tick it
|
|
# launched -- indistinguishable, from the operator's seat, from
|
|
# "launch did nothing". Consume and announce it instead.
|
|
if self.stop_requested:
|
|
self.stop_requested = False
|
|
print("[relay] END MISSION ignored -- no mission is running "
|
|
"(the latch was cleared, so it cannot kill the next "
|
|
"launch)", flush=True)
|
|
return
|
|
if self.stop_sent:
|
|
return
|
|
expired = (self.mission_length > 0
|
|
and self.mission_started_at is not None
|
|
and time.time() >= self.mission_started_at
|
|
+ self.mission_length + STOP_GRACE_SECONDS)
|
|
if not (expired or self.stop_requested):
|
|
return
|
|
pkt = stop_mission_packet()
|
|
for conn in list(self.console_conns):
|
|
try:
|
|
_send_all_guarded(conn.sock, pkt, "stop")
|
|
except OSError as e:
|
|
print(f"[relay] stop send failed to {conn.addr}: {e!r}",
|
|
flush=True)
|
|
self.stop_sent = True
|
|
why = "operator END MISSION" if self.stop_requested else "time expired"
|
|
print(f"[relay] StopMission sent to {len(self.console_conns)} pod(s) "
|
|
f"({why})", flush=True)
|
|
|
|
# ---------------- game TCP side ----------------
|
|
|
|
def _accept_game(self, listener):
|
|
sock, addr = listener.accept()
|
|
sock.setsockopt(socket.IPPROTO_TCP, socket.TCP_NODELAY, 1)
|
|
_enable_keepalive(sock)
|
|
sock.setblocking(False)
|
|
conn = RelayGameConn(sock, addr)
|
|
self.game_conns.append(conn)
|
|
self.sel.register(sock, selectors.EVENT_READ, ("game", conn))
|
|
print(f"[relay] {conn.name()} accepted (awaiting HELLO)", flush=True)
|
|
|
|
def _game_read(self, conn):
|
|
try:
|
|
data = conn.sock.recv(8192)
|
|
except OSError:
|
|
data = b""
|
|
if data:
|
|
conn.last_seen = time.time()
|
|
if not data:
|
|
self._drop_game(conn, "closed")
|
|
return
|
|
conn.buf += data
|
|
while len(conn.buf) >= ENV_TCP_SIZE:
|
|
route, length = struct.unpack_from(ENV_FMT_TCP, conn.buf, 0)
|
|
# game frames are small; the matchlog upload (route -9) is a whole
|
|
# file -- allow it its own generous cap (matches the client's 8MB)
|
|
cap = MATCHLOG_CAP if route == ROUTE_MATCHLOG else FRAME_CAP
|
|
if length > cap:
|
|
self._drop_game(conn, f"oversize frame ({length})")
|
|
return
|
|
if len(conn.buf) < ENV_TCP_SIZE + length:
|
|
break # partial; wait for more
|
|
payload = conn.buf[ENV_TCP_SIZE:ENV_TCP_SIZE + length]
|
|
conn.buf = conn.buf[ENV_TCP_SIZE + length:]
|
|
self.stats["tcp_rx"] += 1
|
|
self._handle_game_frame(conn, route, payload)
|
|
if conn.sock.fileno() < 0: # dropped mid-loop
|
|
return
|
|
|
|
def _handle_game_frame(self, conn, route, payload):
|
|
if conn.host_id is None and route == ROUTE_MATCHLOG:
|
|
# MATCHLOG AUTO-UPLOAD: a pod's post-mission throwaway dial hands
|
|
# us its match forensic log (payload: filename NUL + file bytes).
|
|
# Saved under matchlogs/ prefixed with the sender address so all
|
|
# peers' files land side by side for tools/matchcheck.py.
|
|
nul = payload.find(b"\0")
|
|
if nul <= 0 or nul > 150:
|
|
self._drop_game(conn, "malformed matchlog upload")
|
|
return
|
|
name = os.path.basename(payload[:nul].decode("ascii", "replace"))
|
|
if not (name.startswith("matchlog_") and name.endswith(".txt")):
|
|
self._drop_game(conn, f"suspicious matchlog name {name!r}")
|
|
return
|
|
data = payload[nul + 1:]
|
|
os.makedirs("matchlogs", exist_ok=True)
|
|
peer = conn.sock.getpeername()[0].replace(":", "_")
|
|
out_path = os.path.join("matchlogs", f"{peer}_{name}")
|
|
with open(out_path, "wb") as f:
|
|
f.write(data)
|
|
print(f"[relay] matchlog saved: {out_path} ({len(data)} bytes)",
|
|
flush=True)
|
|
return
|
|
if conn.host_id is None and route == ROUTE_SEAT_REQUEST:
|
|
# SERVER-ASSIGNED SEATS: hand out the lowest roster seat that is
|
|
# neither claimed (by_host) nor reserved; reserve it briefly so
|
|
# two simultaneous joiners can't race onto the same seat. The
|
|
# requesting connection closes after the reply; the pod re-dials
|
|
# and claims the seat with a normal HELLO.
|
|
now = time.time()
|
|
self.seat_reservations = {h: t for h, t in
|
|
self.seat_reservations.items()
|
|
if t > now or h in self.seat_beacons}
|
|
self.seat_reclaim = {t: e for t, e in self.seat_reclaim.items()
|
|
if e > now}
|
|
claim_tag = ""
|
|
if payload:
|
|
fields = payload.split(b"\0")
|
|
if len(fields) > 2:
|
|
claim_tag = fields[2].decode("ascii", "replace").strip()
|
|
assigned = None
|
|
if claim_tag and claim_tag in self.roster:
|
|
host_id = FIRST_GAME_HOST_ID + self.roster.index(claim_tag)
|
|
if host_id not in self.by_host and host_id not in self.seat_beacons:
|
|
assigned = (host_id, claim_tag)
|
|
self.seat_reclaim.pop(claim_tag, None)
|
|
print(f"[relay] seat RECLAIMED by returning player "
|
|
f"(tag '{claim_tag}')", flush=True)
|
|
# STATIC SEATS (issue #34, night-2 capture): a manually restarted
|
|
# client carries NO claim and used to mint a NEW seat while its
|
|
# old one sat reserved with the same callsign -- duplicate roster
|
|
# entries, and exactly what pods must never do. Fallback: match
|
|
# the requester's (IP + callsign) against remembered seat
|
|
# identities and RECLAIM that seat -- displacing a zombie beacon
|
|
# if one is still squatting on it. (Pods are fixed machines:
|
|
# same box + same name = same seat, always.)
|
|
if assigned is None and payload:
|
|
req_ip = conn.sock.getpeername()[0]
|
|
req_callsign = (payload.split(b"\0")[0]
|
|
.decode("ascii", "replace").strip()[:15])
|
|
if not hasattr(self, "seat_identity"):
|
|
self.seat_identity = {}
|
|
if req_callsign:
|
|
for i, tag in enumerate(self.roster):
|
|
if self.seat_identity.get(tag) != (req_ip, req_callsign):
|
|
continue
|
|
host_id = FIRST_GAME_HOST_ID + i
|
|
if host_id in self.by_host:
|
|
break # actively PLAYING: not ours
|
|
old = self.seat_beacons.get(host_id)
|
|
if old is not None: # zombie duplicate: displace it
|
|
print(f"[relay] displacing stale beacon on seat "
|
|
f"'{tag}' (same player rejoining)",
|
|
flush=True)
|
|
old.seat_host = None
|
|
self.seat_beacons.pop(host_id, None)
|
|
assigned = (host_id, tag)
|
|
self.seat_reclaim.pop(tag, None)
|
|
print(f"[relay] seat RECLAIMED by identity "
|
|
f"({req_callsign}@{req_ip} -> '{tag}')",
|
|
flush=True)
|
|
break
|
|
if assigned is None:
|
|
for i, tag in enumerate(self.roster):
|
|
host_id = FIRST_GAME_HOST_ID + i
|
|
if (host_id in self.by_host or host_id in self.seat_reservations
|
|
or host_id in self.reserved_host_ids
|
|
or self.seat_reclaim.get(tag, 0) > now):
|
|
continue # claimed / racing / held for return
|
|
assigned = (host_id, tag)
|
|
break
|
|
if assigned is None:
|
|
self._send_raw(conn, struct.pack(ENV_FMT_TCP,
|
|
ROUTE_SEAT_FULL, 0))
|
|
print(f"[relay] {conn.name()} seat request: ROSTER FULL",
|
|
flush=True)
|
|
return
|
|
host_id, tag = assigned
|
|
self.seat_reservations[host_id] = now + SEAT_RESERVE_SECONDS
|
|
# WALK-UP PREFS: optional payload {callsign NUL mech NUL}
|
|
pref_note = ""
|
|
if payload:
|
|
parts = payload.split(b"\0")
|
|
callsign = parts[0].decode("ascii", "replace").strip()[:15]
|
|
mech = (parts[1].decode("ascii", "replace").strip().lower()
|
|
if len(parts) > 1 else "")
|
|
callsign = "".join(c for c in callsign
|
|
if c.isalnum() or c in " -_.")
|
|
if mech and mech not in VEHICLE_TAGS:
|
|
print(f"[relay] {conn.name()} requested unknown mech "
|
|
f"{mech!r} -- keeping the roster default", flush=True)
|
|
mech = ""
|
|
if callsign or mech:
|
|
self.seat_prefs[host_id] = (callsign, mech)
|
|
pref_note = f" callsign={callsign!r} mech={mech or '(default)'}"
|
|
payload_out = struct.pack("<i", host_id) + tag.encode() + b"\0"
|
|
self._send_raw(conn, struct.pack(ENV_FMT_TCP, ROUTE_SEAT_ASSIGN,
|
|
len(payload_out)) + payload_out)
|
|
conn.seat_host = host_id # this conn IS the beacon now
|
|
conn.seat_tag = tag # tag survives a round reset
|
|
self.seat_beacons[host_id] = conn
|
|
# issue #34: remember who sat here (IP + callsign) so the same
|
|
# player's claim-less restart reclaims THIS seat, never a new one
|
|
if payload:
|
|
_cs = (payload.split(b"\0")[0]
|
|
.decode("ascii", "replace").strip()[:15])
|
|
if _cs:
|
|
if not hasattr(self, "seat_identity"):
|
|
self.seat_identity = {}
|
|
self.seat_identity[tag] = (conn.sock.getpeername()[0], _cs)
|
|
print(f"[relay] {conn.name()} seat request -> assigned host "
|
|
f"{host_id} '{tag}' (reserved {SEAT_RESERVE_SECONDS:.0f}s)"
|
|
f"{pref_note}", flush=True)
|
|
cs, mech = self.seat_prefs.get(host_id, ("", ""))
|
|
print(f"[relay] PLAYER {host_id - FIRST_GAME_HOST_ID + 1} SEATED "
|
|
f"(host {host_id}) tag='{tag}' "
|
|
f"callsign='{cs}' mech='{mech}' "
|
|
f"({len(self.seat_beacons)} seated)", flush=True)
|
|
return
|
|
if conn.host_id is None:
|
|
# First frame MUST be a valid HELLO.
|
|
if route != ROUTE_HELLO or len(payload) < 8:
|
|
self._drop_game(conn, f"first frame not HELLO (route={route})")
|
|
return
|
|
magic, host_id = struct.unpack_from("<iI", payload, 0)
|
|
host_id = int(host_id)
|
|
if magic != HELLO_MAGIC:
|
|
self._drop_game(conn, f"bad HELLO magic 0x{magic:x}")
|
|
return
|
|
if host_id not in self.expected_ids:
|
|
# issue #33: a plain drop bounced the client straight into a
|
|
# relaunch loop (stale identity -> reject -> relaunch -> same
|
|
# stale identity, ~1/s). Tell it to REJOIN instead (the
|
|
# rejoin path re-requests a seat = identity resync), rate-
|
|
# limited per IP so a truly broken client can't spin fast.
|
|
ip = conn.sock.getpeername()[0] if conn.sock else "?"
|
|
if not hasattr(self, "hello_reject_at"):
|
|
self.hello_reject_at = {}
|
|
now = time.time()
|
|
if now - self.hello_reject_at.get(ip, 0) > 5.0:
|
|
self.hello_reject_at[ip] = now
|
|
try:
|
|
self._send_raw(conn, struct.pack(ENV_FMT_TCP,
|
|
ROUTE_REJOIN, 0))
|
|
except OSError:
|
|
pass
|
|
self._drop_game(conn, f"HELLO hostID {host_id} not in roster "
|
|
f"{sorted(self.expected_ids)} -- "
|
|
f"REJOIN sent (identity resync)")
|
|
return
|
|
if host_id in self.by_host:
|
|
self._drop_game(conn, f"HELLO hostID {host_id} already registered")
|
|
return
|
|
conn.host_id = host_id
|
|
self.by_host[host_id] = conn
|
|
self.seat_reservations.pop(host_id, None) # claim clears reserve
|
|
print(f"[relay] {conn.name()} REGISTERED "
|
|
f"({len(self.by_host)}/{len(self.roster)})", flush=True)
|
|
if len(self.by_host) >= len(self.roster):
|
|
self._release_eggs() # fallback; normally already fired
|
|
# PEER_UP exchange: newcomer learns everyone; everyone learns newcomer.
|
|
for other_id in sorted(self.by_host):
|
|
if other_id != host_id:
|
|
self._send_control(conn, ROUTE_PEER_UP, other_id)
|
|
for other in list(self.by_host.values()): # snapshot: a failed
|
|
# send drops that peer INSIDE the loop (mutates by_host)
|
|
if other is not conn:
|
|
self._send_control(other, ROUTE_PEER_UP, host_id)
|
|
return
|
|
if route == ROUTE_READY:
|
|
conn.ready = True
|
|
ready_count = sum(1 for c in self.by_host.values()
|
|
if getattr(c, "ready", False))
|
|
i = conn.host_id - FIRST_GAME_HOST_ID
|
|
tag = self.roster[i] if 0 <= i < len(self.roster) else "?"
|
|
print(f"[relay] PLAYER {conn.host_id - FIRST_GAME_HOST_ID + 1} "
|
|
f"READY (host {conn.host_id}) tag='{tag}' "
|
|
f"({ready_count}/{len(self.by_host)} loaded)", flush=True)
|
|
return
|
|
if route == ROUTE_BCAST:
|
|
env = struct.pack(ENV_FMT_TCP, conn.host_id, len(payload))
|
|
for other in list(self.by_host.values()):
|
|
if other is not conn:
|
|
self._send_raw(other, env + payload)
|
|
elif route >= FIRST_GAME_HOST_ID:
|
|
target = self.by_host.get(route)
|
|
if target is None:
|
|
return # peer gone; drop silently
|
|
env = struct.pack(ENV_FMT_TCP, conn.host_id, len(payload))
|
|
self._send_raw(target, env + payload)
|
|
else:
|
|
print(f"[relay] {conn.name()} bad route {route} -- ignored", flush=True)
|
|
|
|
def _send_control(self, conn, route, host_id):
|
|
self._send_raw(conn, struct.pack(ENV_FMT_TCP, route, 4)
|
|
+ struct.pack("<i", host_id))
|
|
|
|
def _send_raw(self, conn, data):
|
|
try:
|
|
_send_all_guarded(conn.sock, data)
|
|
self.stats["tcp_tx"] += 1
|
|
except OSError as e:
|
|
self._drop_game(conn, f"send failed {e!r}")
|
|
|
|
def _abort_round(self, why):
|
|
"""ROUND ABORT (connection audit 2026-07-23): a registered pod died
|
|
AFTER egg release but BEFORE launch -- the engine's all-connections
|
|
gate would leave every survivor waiting on it forever. Bounce all
|
|
remaining pods back to their seats (REJOIN) and reset the round; the
|
|
reclaim window puts everyone back in the same seats in seconds.
|
|
|
|
STORM HARDENING (issue #33, night-2 capture): the first abort's
|
|
rejoiners reconnected + re-ACKed within a second, RE-releasing eggs
|
|
while a flapping pod was still bouncing -- its next drop aborted the
|
|
NEW round, cascading (3 aborts in 8 s, every client relaunch-churning
|
|
into the 'Could not create a Direct3D device' 20 s deadline). Two
|
|
dampers: stale console-conn ACK flags are cleared (they counted
|
|
toward the next release), and egg release is HELD for a settle
|
|
window after any abort so the round re-forms only once the flapping
|
|
stops."""
|
|
if not self.eggs_released or self.launches_sent >= 2:
|
|
return
|
|
print(f"[relay] ROUND ABORT ({why}) -- sending every pod back to "
|
|
f"the join wait for a fresh round "
|
|
f"(egg release held {ABORT_SETTLE_SECONDS:.0f}s)", flush=True)
|
|
# reset FIRST so the cascade of resulting drops can't re-trigger
|
|
self.eggs_released = False
|
|
self.launch_at = None
|
|
self.launch_requested = False
|
|
self.eggs_done_at = None
|
|
self._launch_blocked_warned = False
|
|
self.round_hold_until = time.time() + ABORT_SETTLE_SECONDS
|
|
for c in self.console_conns:
|
|
c.acked = False # an aborted round's ACK must not count
|
|
c.egg_sent = False # toward re-releasing the next one
|
|
survivors = list(self.by_host.values())
|
|
for conn in survivors:
|
|
try:
|
|
self._send_raw(conn, struct.pack(ENV_FMT_TCP, ROUTE_REJOIN, 0))
|
|
except OSError:
|
|
pass
|
|
# the full state restore happens in _maybe_reset_round as they drop
|
|
|
|
def _drop_game(self, conn, why):
|
|
print(f"[relay] {conn.name()} dropped: {why}", flush=True)
|
|
try:
|
|
self.sel.unregister(conn.sock)
|
|
except (KeyError, ValueError):
|
|
pass
|
|
try:
|
|
conn.sock.close()
|
|
except OSError:
|
|
pass
|
|
if conn in self.game_conns:
|
|
self.game_conns.remove(conn)
|
|
seat = getattr(conn, "seat_host", None)
|
|
if seat is not None and self.seat_beacons.get(seat) is conn:
|
|
del self.seat_beacons[seat]
|
|
tag_for_seat = getattr(conn, "seat_tag", None)
|
|
if tag_for_seat: # hold for the returning player
|
|
self.seat_reclaim[tag_for_seat] = time.time() + 90.0
|
|
if seat not in self.by_host: # never claimed: player left
|
|
self.seat_reservations.pop(seat, None)
|
|
self.seat_prefs.pop(seat, None)
|
|
i = seat - FIRST_GAME_HOST_ID
|
|
tag = self.roster[i] if 0 <= i < len(self.roster) else "?"
|
|
print(f"[relay] PLAYER {seat - FIRST_GAME_HOST_ID + 1} LEFT "
|
|
f"(host {seat}) tag='{tag}' "
|
|
f"({len(self.seat_beacons)} seated)", flush=True)
|
|
was_registered = (conn.host_id is not None
|
|
and self.by_host.get(conn.host_id) is conn)
|
|
if conn.host_id is not None and self.by_host.get(conn.host_id) is conn:
|
|
del self.by_host[conn.host_id]
|
|
self.udp_endpoint.pop(conn.host_id, None)
|
|
for other in list(self.by_host.values()): # snapshot: a failed
|
|
# send drops that peer INSIDE the loop (mutates by_host)
|
|
self._send_control(other, ROUTE_PEER_DOWN, conn.host_id)
|
|
if was_registered:
|
|
self._abort_round(f"host {conn.host_id} dropped mid-load")
|
|
self._maybe_reset_round()
|
|
|
|
# ---------------- game UDP side ----------------
|
|
|
|
def _udp_read(self):
|
|
while True:
|
|
try:
|
|
data, endpoint = self.udp_sock.recvfrom(4096)
|
|
except (BlockingIOError, InterruptedError):
|
|
return
|
|
except OSError:
|
|
return
|
|
if len(data) < ENV_UDP_SIZE:
|
|
continue
|
|
route, from_host, _seq = struct.unpack_from(ENV_FMT_UDP, data, 0)
|
|
if from_host not in self.by_host:
|
|
continue # not TCP-registered: drop
|
|
self.stats["udp_rx"] += 1
|
|
# NAT-rebind tolerant: every datagram refreshes the endpoint map.
|
|
self.udp_endpoint[from_host] = endpoint
|
|
# ... and counts as LIVENESS. Without this, a pod streaming updates
|
|
# over UDP while its TCP sits idle looks "silent" to the reaper.
|
|
live = self.by_host.get(from_host)
|
|
if live is not None:
|
|
live.last_seen = time.time()
|
|
if route == ROUTE_HELLO:
|
|
ack = struct.pack(ENV_FMT_UDP, ROUTE_UDP_ACK, CONSOLE_HOST_ID, 0)
|
|
try:
|
|
self.udp_sock.sendto(ack, endpoint)
|
|
except OSError:
|
|
pass
|
|
continue
|
|
if route == ROUTE_BCAST:
|
|
targets = [h for h in self.by_host if h != from_host]
|
|
elif route >= FIRST_GAME_HOST_ID and route in self.by_host:
|
|
targets = [route]
|
|
else:
|
|
continue
|
|
for host_id in targets:
|
|
if self.udp_drop_pct and random.random() * 100.0 < self.udp_drop_pct:
|
|
self.stats["udp_dropped"] += 1
|
|
continue
|
|
ep = self.udp_endpoint.get(host_id)
|
|
if ep is not None:
|
|
try:
|
|
self.udp_sock.sendto(data, ep) # forward VERBATIM
|
|
self.stats["udp_tx"] += 1
|
|
continue
|
|
except OSError:
|
|
pass
|
|
# Target's UDP endpoint unknown/broken: wrap the frame onto its
|
|
# game TCP connection (asymmetric-UDP-blockage fallback).
|
|
frame = data[ENV_UDP_SIZE:]
|
|
target = self.by_host.get(host_id)
|
|
if target is not None and frame:
|
|
env = struct.pack(ENV_FMT_TCP, from_host, len(frame))
|
|
self._send_raw(target, env + frame)
|
|
self.stats["udp_tcp_fallback"] += 1
|
|
|
|
# ---------------- LAN auto-discovery ----------------
|
|
|
|
def _disc_read(self):
|
|
while True:
|
|
try:
|
|
data, endpoint = self.disc_sock.recvfrom(64)
|
|
except (BlockingIOError, InterruptedError, OSError):
|
|
return
|
|
if data == DISC_PROBE:
|
|
try:
|
|
self.disc_sock.sendto(
|
|
DISC_REPLY + struct.pack("<H", self.console_port),
|
|
endpoint)
|
|
print(f"[relay] discovery probe from "
|
|
f"{endpoint[0]}:{endpoint[1]} -> answered "
|
|
f"(console port {self.console_port})", flush=True)
|
|
except OSError:
|
|
pass
|
|
|
|
# ---------------- stats ----------------
|
|
|
|
# ---------------- remote operator control ----------------
|
|
|
|
def _ctl_accept(self, listener):
|
|
try:
|
|
sock, addr = listener.accept()
|
|
sock.setblocking(False)
|
|
conn = RelayControlConn(sock, addr)
|
|
self.ctl_conns.append(conn)
|
|
self.sel.register(sock, selectors.EVENT_READ, ("ctl", conn))
|
|
print(f"[ctl] operator connection from {conn.name()} "
|
|
f"(awaiting AUTH)", flush=True)
|
|
except OSError:
|
|
pass
|
|
|
|
def _ctl_drop(self, conn, why):
|
|
try:
|
|
self.sel.unregister(conn.sock)
|
|
except (KeyError, ValueError):
|
|
pass
|
|
try:
|
|
conn.sock.close()
|
|
except OSError:
|
|
pass
|
|
with _CONTROL_LOCK:
|
|
if conn in _CONTROL_SINKS:
|
|
_CONTROL_SINKS.remove(conn)
|
|
if conn in self.ctl_conns:
|
|
self.ctl_conns.remove(conn)
|
|
if conn.authed or why != "auth timeout":
|
|
print(f"[ctl] {conn.name()} dropped: {why}", flush=True)
|
|
|
|
def _ctl_read(self, conn):
|
|
try:
|
|
data = conn.sock.recv(4096)
|
|
except (BlockingIOError, InterruptedError):
|
|
return
|
|
except OSError:
|
|
self._ctl_drop(conn, "closed")
|
|
return
|
|
if not data:
|
|
self._ctl_drop(conn, "closed")
|
|
return
|
|
conn.buf += data
|
|
if len(conn.buf) > 8192:
|
|
self._ctl_drop(conn, "oversize command")
|
|
return
|
|
while b"\n" in conn.buf:
|
|
line, conn.buf = conn.buf.split(b"\n", 1)
|
|
try:
|
|
self._ctl_command(conn, line.decode("utf-8", "replace").strip())
|
|
except Exception as e: # a control bug must NEVER
|
|
print(f"[ctl] command error: {e!r}", flush=True) # kill games
|
|
|
|
def _ctl_command(self, conn, line):
|
|
if not line:
|
|
return
|
|
if not conn.authed:
|
|
parts = line.split(None, 1)
|
|
if len(parts) == 2 and parts[0] == "AUTH" \
|
|
and parts[1].strip() == self.ctl_secret:
|
|
# one operator at a time: the newest displaces the rest
|
|
for other in list(self.ctl_conns):
|
|
if other is not conn and other.authed:
|
|
other.queue_out(b"[ctl] displaced by "
|
|
+ conn.name().encode() + b"\n")
|
|
other.flush_out()
|
|
self._ctl_drop(other, "displaced by new operator")
|
|
conn.authed = True
|
|
try:
|
|
conn.sock.send(("OK bt-relay-control 1 console=%d\n"
|
|
% self.console_port).encode())
|
|
except OSError:
|
|
conn.dead = True
|
|
# replay recent history (bounded), THEN subscribe live --
|
|
# under the lock so no line is lost or duplicated between
|
|
# the replay and the live stream
|
|
with _CONTROL_LOCK:
|
|
conn.queue_out(b"[ctl] --- session history replay ---\n")
|
|
for past in _CONTROL_HISTORY:
|
|
conn.queue_out(past)
|
|
conn.queue_out(b"[ctl] --- live ---\n")
|
|
_CONTROL_SINKS.append(conn)
|
|
print(f"[ctl] operator AUTHENTICATED: {conn.name()} "
|
|
f"(driving this relay)", flush=True)
|
|
else:
|
|
self._ctl_drop(conn, "bad auth")
|
|
return
|
|
cmd = line.split(None, 1)[0].lower()
|
|
if cmd == "launch":
|
|
print(f"[ctl] LAUNCH from {conn.name()} "
|
|
f"(launches_sent={self.launches_sent} "
|
|
f"eggs_released={self.eggs_released} "
|
|
f"pods_acked={self._pods_ready()}/{len(self.roster)})",
|
|
flush=True)
|
|
self.launch_requested = True
|
|
elif cmd == "stop":
|
|
print(f"[ctl] STOP from {conn.name()}", flush=True)
|
|
self.stop_requested = True
|
|
elif cmd in ("rearm", "newround"):
|
|
# Explicit escape hatch: recover the launcher without restarting the
|
|
# session. Before this existed, killing the relay was the only way
|
|
# out of a latched round (operator report 2026-07-25).
|
|
print(f"[ctl] RE-ARM from {conn.name()}", flush=True)
|
|
self._rearm_for_new_round("operator re-arm command")
|
|
elif cmd == "ping":
|
|
conn.queue_out(b"[ctl] pong\n")
|
|
elif cmd == "get":
|
|
self._ctl_get_mission(conn)
|
|
elif cmd == "set":
|
|
self._ctl_set_mission(conn, line[3:].strip())
|
|
else:
|
|
conn.queue_out(("[ctl] unknown command: %s\n" % cmd).encode())
|
|
|
|
def _ctl_get_mission(self, conn):
|
|
try:
|
|
doc = eggmodel.EggDoc.load(self.orig_egg_path)
|
|
m = doc.mission()
|
|
print("[ctl] mission map=%s time=%s weather=%s scenario=%s "
|
|
"temperature=%s length=%s seats=%d"
|
|
% (m.get("map", ""), m.get("time", ""),
|
|
m.get("weather", ""), m.get("scenario", ""),
|
|
m.get("temperature", ""), m.get("length", ""),
|
|
len(self.roster)), flush=True)
|
|
except Exception as e:
|
|
print(f"[ctl] get mission failed: {e!r}", flush=True)
|
|
|
|
def _ctl_set_mission(self, conn, spec):
|
|
if eggmodel is None:
|
|
print("[ctl] set rejected: eggmodel unavailable", flush=True)
|
|
return
|
|
kw = {}
|
|
for pair in spec.split(";"):
|
|
pair = pair.strip()
|
|
if not pair or "=" not in pair:
|
|
continue
|
|
k, v = pair.split("=", 1)
|
|
k = k.strip().lower()
|
|
v = v.strip()
|
|
if k in CONTROL_SET_KEYS and v:
|
|
kw[k] = v
|
|
if not kw:
|
|
print("[ctl] set: nothing applied (allowed keys: %s)"
|
|
% ",".join(sorted(CONTROL_SET_KEYS)), flush=True)
|
|
return
|
|
try:
|
|
doc = eggmodel.EggDoc.load(self.orig_egg_path)
|
|
doc.set_mission(**kw)
|
|
doc.save(self.orig_egg_path)
|
|
print("[ctl] mission settings applied by %s: %s -- takes effect "
|
|
"next round" % (conn.name(),
|
|
" ".join("%s=%s" % kv for kv in
|
|
sorted(kw.items()))), flush=True)
|
|
except Exception as e:
|
|
print(f"[ctl] set failed: {e!r}", flush=True)
|
|
|
|
def _tick_control(self):
|
|
now = time.time()
|
|
for conn in list(self.ctl_conns):
|
|
if conn.dead:
|
|
self._ctl_drop(conn, "send buffer overflow/socket error")
|
|
elif (not conn.authed
|
|
and now - conn.connected_at > CONTROL_AUTH_TIMEOUT):
|
|
self._ctl_drop(conn, "auth timeout")
|
|
else:
|
|
conn.flush_out()
|
|
|
|
def _tick_settle(self):
|
|
# issue #33: the post-abort settle hold blocks egg release while the
|
|
# rejoin flapping calms; once it expires, re-run the release check
|
|
# (the ACKs that arrived DURING the hold produced no later event to
|
|
# trigger it).
|
|
hold = getattr(self, "round_hold_until", 0)
|
|
if hold and time.time() >= hold and not self.eggs_released:
|
|
self.round_hold_until = 0
|
|
waiting = [c for c in self.console_conns if not c.egg_sent]
|
|
if waiting and len(self.console_conns) >= len(self.roster):
|
|
print("[relay] abort settle window over -- releasing eggs to "
|
|
f"{len(waiting)} waiting pod(s)", flush=True)
|
|
self._release_eggs()
|
|
|
|
def _tick_stats(self):
|
|
if time.time() < self.stats_at:
|
|
return
|
|
self.stats_at = time.time() + STATS_PERIOD
|
|
s = self.stats
|
|
print(f"[relay-stats] tcp rx/tx {s['tcp_rx']}/{s['tcp_tx']} | "
|
|
f"udp rx/tx {s['udp_rx']}/{s['udp_tx']} "
|
|
f"(dropped {s['udp_dropped']}, tcp-fallback {s['udp_tcp_fallback']}) | "
|
|
f"registered {sorted(self.by_host)} udp-known {sorted(self.udp_endpoint)}",
|
|
flush=True)
|
|
|
|
|
|
def relay_main(argv):
|
|
"""argv: --relay <consolePort> <egg-file> [--bind ADDR] [--udp-drop PCT]
|
|
[--manual-launch]"""
|
|
args = argv[:]
|
|
bind_addr = "0.0.0.0"
|
|
udp_drop = 0.0
|
|
manual = False
|
|
if "--bind" in args:
|
|
i = args.index("--bind")
|
|
bind_addr = args[i + 1]
|
|
del args[i:i + 2]
|
|
if "--udp-drop" in args:
|
|
i = args.index("--udp-drop")
|
|
udp_drop = float(args[i + 1])
|
|
del args[i:i + 2]
|
|
if "--manual-launch" in args:
|
|
manual = True
|
|
args.remove("--manual-launch")
|
|
reserved_tags = ()
|
|
if "--reserve" in args:
|
|
i = args.index("--reserve")
|
|
reserved_tags = tuple(t for t in args[i + 1].split(",") if t)
|
|
del args[i:i + 2]
|
|
if len(args) != 2:
|
|
print(__doc__)
|
|
return 2
|
|
console_port = int(args[0])
|
|
egg_path = args[1]
|
|
relay = Relay(console_port, egg_path, bind_addr, udp_drop, manual,
|
|
reserved_tags=reserved_tags)
|
|
if not relay.roster:
|
|
print(f"[relay] ERROR: no [pilots] entries in {egg_path}")
|
|
return 2
|
|
# Operator command channel: a daemon thread reads stdin. 'launch' fires
|
|
# the mission (manual mode only); 'stop' ends it early (End Mission).
|
|
def stdin_reader():
|
|
if sys.stdin is None: # pythonw / broken shim pipeline
|
|
print("[relay] no stdin -- operator commands unavailable",
|
|
flush=True)
|
|
return
|
|
for line in sys.stdin:
|
|
command = line.strip().lower()
|
|
if command == "launch":
|
|
# ALWAYS record the press at RECEIPT. Every wedge in the launch
|
|
# state machine used to be silent, so the operator could not
|
|
# tell "the command never arrived" from "the relay ignored it".
|
|
print(f"[relay] operator LAUNCH command received "
|
|
f"(launches_sent={relay.launches_sent} "
|
|
f"eggs_released={relay.eggs_released})", flush=True)
|
|
relay.launch_requested = True
|
|
elif command == "stop":
|
|
print("[relay] operator STOP command received", flush=True)
|
|
relay.stop_requested = True
|
|
elif command in ("rearm", "newround"):
|
|
print("[relay] operator RE-ARM command received", flush=True)
|
|
relay._rearm_for_new_round("operator re-arm command")
|
|
elif command:
|
|
print(f"[relay] unknown operator command: {command!r}",
|
|
flush=True)
|
|
threading.Thread(target=stdin_reader, daemon=True).start()
|
|
try:
|
|
relay.run()
|
|
except KeyboardInterrupt:
|
|
print("[relay] shutting down")
|
|
return 0
|
|
|
|
|
|
class _StampedOut:
|
|
"""Line-timestamping stdout wrapper: every log line gets a wall-clock
|
|
prefix so post-mortems can measure delays instead of guessing (operator
|
|
request 2026-07-22). The GUI's monitor regexes all use .search(), so
|
|
the prefix is transparent to them."""
|
|
|
|
def __init__(self, inner):
|
|
self._inner = inner
|
|
self._at_line_start = True
|
|
|
|
def write(self, text):
|
|
out = []
|
|
for ch in text:
|
|
if self._at_line_start and ch != "\n":
|
|
out.append(time.strftime("%H:%M:%S "))
|
|
self._at_line_start = False
|
|
out.append(ch)
|
|
if ch == "\n":
|
|
self._at_line_start = True
|
|
joined = "".join(out)
|
|
self._inner.write(joined)
|
|
try:
|
|
_control_tee(joined)
|
|
except Exception:
|
|
pass # the tee must never kill logs
|
|
|
|
def flush(self):
|
|
self._inner.flush()
|
|
|
|
def __getattr__(self, name):
|
|
return getattr(self._inner, name)
|
|
|
|
|
|
def main():
|
|
sys.stdout = _StampedOut(sys.stdout)
|
|
if len(sys.argv) >= 2 and sys.argv[1] == "--relay":
|
|
return relay_main(sys.argv[2:])
|
|
if len(sys.argv) < 3:
|
|
print(__doc__)
|
|
return 2
|
|
egg_bytes = egg_wire_bytes(open(sys.argv[1], "rb").read())
|
|
pods = sys.argv[2:]
|
|
threads = [threading.Thread(target=serve_pod, args=(p, egg_bytes), daemon=True)
|
|
for p in pods]
|
|
for t in threads:
|
|
t.start()
|
|
try:
|
|
while any(t.is_alive() for t in threads):
|
|
time.sleep(0.5)
|
|
except KeyboardInterrupt:
|
|
print("console shutting down")
|
|
return 0
|
|
|
|
|
|
if __name__ == "__main__":
|
|
sys.exit(main())
|