B1: fly-loop-reset runs only /opt/fly/sbin/flysim, a root-owned copy 05-deploy installs from the release tarball after checking it against the tarball's MANIFEST; fixed paths, fly.env parsed as data, env -i. B2: the wrapper pauses fly-watchdog.timer (and waits out a running probe) for the reset, and starts flysim and the timer on every exit. H1: a step that raises or times out is a failed step; the ladder state is saved before the step runs. H2: one --list call computes the build's compatibility once; TimeoutStartSec 25 min. H3: level 2+ resets to the rung below the best or restarts, never lower. Tests run the wrapper against a fake flysim, systemctl and archives.
301 lines
12 KiB
Python
Executable file
301 lines
12 KiB
Python
Executable file
#!/usr/bin/env python3
|
|
"""Out-of-process loop recovery with an escalation ladder; invoked after the watchdog probe.
|
|
|
|
Ladder (infra/docs/loop-recovery.md): restart flysim, then reset to the current rung's
|
|
milestone, then to the rung below. It climbs one step per confirmed trap that outlives
|
|
the previous step, and starts over once the fly reaches a new best rung or has been
|
|
quiet for QUIET seconds. A router model can delay a step, never deny it.
|
|
"""
|
|
import json
|
|
import os
|
|
from pathlib import Path
|
|
import re
|
|
import subprocess
|
|
import sys
|
|
import time
|
|
from urllib.request import Request, urlopen
|
|
|
|
REPORT = Path(os.environ.get("FLY_LOOP_REPORT", "/run/fly/wd/loop.json"))
|
|
STATE = Path(os.environ.get("FLY_LOOP_RECOVERY_STATE", "/var/lib/fly-loop-recover/state.json"))
|
|
HISTORY = Path(os.environ.get("FLY_LOOP_RECOVERY_HISTORY", "/var/lib/fly-loop-recover/history.jsonl"))
|
|
NOTICE = Path(os.environ.get("FLY_RECOVERY_NOTICE", "/run/fly/wd/recovery-notice.json"))
|
|
MILESTONES = Path(os.environ.get("FLY_STATE_DIR", "/srv/fly/state"))
|
|
STATUS_URL = os.environ.get("FLY_STATUS_URL", "http://127.0.0.1:7401/status")
|
|
RESET_BIN = os.environ.get("FLY_LOOP_RESET_BIN", "/opt/fly/bin/fly-loop-reset")
|
|
INTERVAL = int(os.environ.get("WD_LOOP_INTERVAL", "300"))
|
|
COUNTDOWN = int(os.environ.get("FLY_LOOP_COUNTDOWN", "60"))
|
|
SETTLE = 1200 # a report inside this after an action may still hold the old window
|
|
QUIET = 6 * 3600 # this long without a suspected report starts the ladder over
|
|
HOLD = 3 * 3600 # restart-only spacing while the reset budget is spent
|
|
DAY = 24 * 3600
|
|
MAX_RESETS = 2 # milestone resets per DAY
|
|
VETO_LIMIT = 3 # model vetoes of one step before it goes ahead anyway
|
|
MODEL_BUDGET = 90 # seconds across all models in one decision
|
|
VERIFY_TIMEOUT = 240
|
|
LIST_TIMEOUT = 120
|
|
RESTART_TIMEOUT = 180
|
|
RESET_TIMEOUT = 480 # the wrapper waits up to 60 s for the watchdog, then stop, reset, start
|
|
|
|
# docs/design/ladder.md; mirrors RANK_LADDER in flybrain-gb's pokemon_red/mod.rs (a test pins it).
|
|
LADDER = [
|
|
"BOOT", "BEDROOM", "DOWNSTAIRS", "PALLET TOWN", "OAK'S LAB", "GOT A STARTER", "OAK'S PARCEL",
|
|
"POKEDEX", "VIRIDIAN CITY", "VIRIDIAN FOREST", "PEWTER CITY", "BOULDER BADGE", "MT. MOON",
|
|
"CERULEAN CITY", "CASCADE BADGE", "NUGGET BRIDGE", "MET BILL", "VERMILION CITY", "HM CUT",
|
|
"THUNDER BADGE", "ROCK TUNNEL", "LAVENDER TOWN", "CELADON CITY", "SILPH SCOPE", "RAINBOW BADGE",
|
|
"POKE FLUTE", "FUCHSIA CITY", "SOUL BADGE", "SILPH CO. FREED", "MARSH BADGE", "CINNABAR ISLAND",
|
|
"VOLCANO BADGE", "EARTH BADGE", "INDIGO PLATEAU", "BEAT LORELEI", "BEAT BRUNO", "BEAT AGATHA",
|
|
"CHAMPION",
|
|
]
|
|
|
|
PROMPT = (
|
|
"You check a Game Boy Pokemon Red run played by a simulated fly brain through macros "
|
|
"(GO OBJECTIVE, GO WARP, GO ROUTE, FIGHT, ...). A watchdog flagged a suspected loop over "
|
|
"a 10-minute brain window. Stuck means the same few macros repeat with no new places "
|
|
"(places.delta 0) and no rewards; healthy play finds new places, wins battles or earns "
|
|
"rewards. Answer with only the JSON object {\"stuck\": true} or {\"stuck\": false}. "
|
|
"No explanation, no reasoning, no other text."
|
|
)
|
|
|
|
|
|
def load(path, default):
|
|
try:
|
|
return json.loads(path.read_text())
|
|
except (OSError, ValueError):
|
|
return default
|
|
|
|
|
|
def write_json(path, value):
|
|
path.parent.mkdir(parents=True, exist_ok=True)
|
|
tmp = path.with_name(path.name + ".tmp")
|
|
tmp.write_text(json.dumps(value) + "\n")
|
|
os.chmod(tmp, 0o644)
|
|
tmp.replace(path)
|
|
|
|
|
|
def log(event):
|
|
try:
|
|
HISTORY.parent.mkdir(parents=True, exist_ok=True)
|
|
with HISTORY.open("a") as out:
|
|
out.write(json.dumps(event) + "\n")
|
|
except OSError as error:
|
|
print("history not written:", error, file=sys.stderr)
|
|
|
|
|
|
def models():
|
|
names = os.environ.get("FLY_LOOP_MODELS") or os.environ.get("FLY_LOOP_MODEL") or ""
|
|
return [name.strip() for name in names.split(",") if name.strip()]
|
|
|
|
|
|
def parse_stuck(content):
|
|
"""The stuck flag from a free model's reply, which may wrap the JSON in prose or fences."""
|
|
for match in re.finditer(r"\{[^{}]*\}", content or ""):
|
|
try:
|
|
decision = json.loads(match.group(0))
|
|
except ValueError:
|
|
continue
|
|
if type(decision) is dict and type(decision.get("stuck")) is bool:
|
|
return decision["stuck"]
|
|
return None
|
|
|
|
|
|
def verdict(report):
|
|
"""True/False from the first model that answers, None when no router or none answered."""
|
|
url = os.environ.get("FLY_LOOP_ROUTER_URL")
|
|
names = models()
|
|
if not url or not names:
|
|
return None, None
|
|
brief = {key: report.get(key) for key in ("reason", "sequence", "window", "places", "milestone", "dominant", "map")}
|
|
headers = {"Content-Type": "application/json"}
|
|
key = os.environ.get("FLY_LOOP_ROUTER_KEY")
|
|
if key:
|
|
headers["Authorization"] = "Bearer " + key
|
|
deadline = time.monotonic() + MODEL_BUDGET
|
|
for name in names:
|
|
left = deadline - time.monotonic()
|
|
if left < 5:
|
|
break
|
|
payload = {"model": name, "temperature": 0, "max_tokens": 400, "messages": [
|
|
{"role": "system", "content": PROMPT},
|
|
{"role": "user", "content": json.dumps(brief)},
|
|
]}
|
|
request = Request(url.rstrip("/") + "/chat/completions", json.dumps(payload).encode(), headers)
|
|
try:
|
|
with urlopen(request, timeout=min(30, left)) as response:
|
|
answer = json.load(response)
|
|
stuck = parse_stuck(answer["choices"][0]["message"]["content"])
|
|
except (OSError, ValueError, KeyError, IndexError, TypeError) as error:
|
|
print(f"model {name}: {type(error).__name__} {getattr(error, 'code', '')}".rstrip(), file=sys.stderr)
|
|
continue
|
|
if stuck is None:
|
|
print(f"model {name}: no stuck flag in reply", file=sys.stderr)
|
|
continue
|
|
return stuck, name
|
|
return None, None
|
|
|
|
|
|
def rungs():
|
|
found = []
|
|
for path in MILESTONES.glob("milestone-*.checkpoint"):
|
|
number = path.stem.removeprefix("milestone-")
|
|
if number.isdigit():
|
|
found.append(int(number))
|
|
return sorted(found)
|
|
|
|
|
|
def plan(level, rank, best, resets_left):
|
|
"""(action, target rung) for a ladder level; a spent budget or no restorable archive is a restart.
|
|
|
|
Level 1 resets to the current rung's archive; level 2 and beyond to the rung below the best,
|
|
and never lower, so repeated days of a trap cannot walk the run down the ladder.
|
|
"""
|
|
if level == 0 or not resets_left or rank is None:
|
|
return "restart", None
|
|
target = rank if level == 1 else max(best, rank) - 1
|
|
if target > rank or target not in restorable_rungs():
|
|
return "restart", None
|
|
return "reset", target
|
|
|
|
|
|
def restorable_rungs():
|
|
"""The rungs the running build can restore (fly-loop-reset --list), empty on any failure."""
|
|
try:
|
|
done = subprocess.run(["sudo", "-n", RESET_BIN, "--list"], check=False, timeout=LIST_TIMEOUT,
|
|
capture_output=True, text=True)
|
|
except (OSError, subprocess.SubprocessError) as error:
|
|
print("restorable rungs unknown:", type(error).__name__, file=sys.stderr)
|
|
return set()
|
|
if done.returncode != 0:
|
|
return set()
|
|
return {int(line) for line in done.stdout.split() if line.isdigit()}
|
|
|
|
|
|
def status():
|
|
try:
|
|
with urlopen(STATUS_URL, timeout=5) as response:
|
|
return json.load(response)
|
|
except (OSError, ValueError):
|
|
return None
|
|
|
|
|
|
def verify(target, deadline):
|
|
while time.monotonic() < deadline:
|
|
now = status()
|
|
if now and now.get("status") == "running":
|
|
rank = (now.get("milestone") or {}).get("rank")
|
|
if target is None or rank == target:
|
|
return True
|
|
time.sleep(5)
|
|
return False
|
|
|
|
|
|
def act(action, target):
|
|
if action == "restart":
|
|
command, limit = ["sudo", "-n", "systemctl", "restart", "flysim.service"], RESTART_TIMEOUT
|
|
else:
|
|
command, limit = ["sudo", "-n", RESET_BIN, str(target)], RESET_TIMEOUT
|
|
try:
|
|
done = subprocess.run(command, check=False, timeout=limit)
|
|
except (OSError, subprocess.SubprocessError) as error:
|
|
print(f"{action} did not complete: {type(error).__name__}", file=sys.stderr)
|
|
return False
|
|
if done.returncode != 0:
|
|
return False
|
|
return verify(target, time.monotonic() + VERIFY_TIMEOUT)
|
|
|
|
|
|
def notice(base, phase, now):
|
|
try:
|
|
write_json(NOTICE, dict(base, phase=phase, updatedAt=int(now)))
|
|
except OSError as error:
|
|
print("notice not written:", error, file=sys.stderr)
|
|
|
|
|
|
def run(now=None, sleep=time.sleep, clock=None):
|
|
clock = clock or (time.time if now is None else (lambda: now))
|
|
now = int(clock())
|
|
report = load(REPORT, None)
|
|
if type(report) is not dict:
|
|
return "no watchdog report"
|
|
state = load(STATE, {})
|
|
if type(state) is not dict:
|
|
state = {}
|
|
rank = (report.get("milestone") or {}).get("rank")
|
|
level = state.get("level", 0)
|
|
|
|
if type(rank) is int and rank > state.get("bestRank", 0):
|
|
if level:
|
|
log({"at": now, "event": "ladder-reset", "why": "new best rung", "rank": rank})
|
|
state.update(bestRank=rank, level=0)
|
|
level = 0
|
|
last = state.get("lastSuspectedAt")
|
|
if level and last and now - last > QUIET:
|
|
log({"at": now, "event": "ladder-reset", "why": "quiet", "rank": rank})
|
|
state["level"] = level = 0
|
|
state["resets"] = [at for at in state.get("resets", []) if now - at < DAY]
|
|
|
|
at = report.get("at")
|
|
if (report.get("suspected") != 1 or type(at) is not int
|
|
or not 0 <= now - at <= 600 or report.get("action") != "none"):
|
|
state.pop("observedAt", None)
|
|
state["vetoes"] = 0
|
|
write_json(STATE, state)
|
|
return "not a fresh suspected loop"
|
|
|
|
state["lastSuspectedAt"] = now
|
|
acted = state.get("actedAt", 0)
|
|
if at < acted + SETTLE:
|
|
write_json(STATE, state)
|
|
return f"settling after {state.get('lastAction')}"
|
|
observed = state.get("observedAt")
|
|
if not observed or observed < acted + SETTLE:
|
|
state["observedAt"] = at
|
|
write_json(STATE, state)
|
|
return "waiting for second probe"
|
|
if at - observed < INTERVAL - 30:
|
|
write_json(STATE, state)
|
|
return "waiting for next probe"
|
|
|
|
resets_left = MAX_RESETS - len(state["resets"])
|
|
action, target = plan(level, rank, state.get("bestRank", 0), resets_left)
|
|
if action == "restart" and level and now - acted < HOLD:
|
|
write_json(STATE, state)
|
|
return "holding: reset budget spent, next restart after the hold"
|
|
|
|
stuck, model = verdict(report)
|
|
if stuck is False and state.get("vetoes", 0) < VETO_LIMIT:
|
|
state["vetoes"] = state.get("vetoes", 0) + 1
|
|
write_json(STATE, state)
|
|
log({"at": now, "event": "veto", "model": model, "vetoes": state["vetoes"], "rank": rank})
|
|
return f"model {model} says not stuck ({state['vetoes']}/{VETO_LIMIT})"
|
|
why = {True: f"confirmed by {model}", False: "veto limit reached", None: "no model answer, watchdog alone"}[stuck]
|
|
|
|
base = {"v": 1, "id": f"{now}-{action}", "action": action,
|
|
"fromRung": rank, "fromLabel": (report.get("milestone") or {}).get("label"),
|
|
"reason": report.get("reason"), "loop": (report.get("sequence") or [])[:4],
|
|
"stuckSeconds": max(0, now - observed) + 600,
|
|
"announcedAt": now, "executeAt": now + COUNTDOWN}
|
|
if action == "reset":
|
|
base.update(toRung=target, toLabel=LADDER[target] if 0 <= target < len(LADDER) else None)
|
|
notice(base, "countdown", now)
|
|
sleep(COUNTDOWN)
|
|
started = int(clock())
|
|
# Recorded before acting, so a helper killed mid-step has still climbed and spent the reset.
|
|
state.update(level=level + 1, actedAt=started, lastAction=action, vetoes=0)
|
|
state.pop("observedAt", None)
|
|
if action == "reset":
|
|
state["resets"].append(started)
|
|
write_json(STATE, state)
|
|
notice(base, "acting", started)
|
|
ok = act(action, target)
|
|
finished = int(clock())
|
|
notice(base, "done" if ok else "failed", finished)
|
|
state["actedAt"] = finished
|
|
write_json(STATE, state)
|
|
log({"at": finished, "event": action, "target": target, "ok": ok, "why": why, "rank": rank,
|
|
"level": level, "sequence": report.get("sequence"), "map": report.get("map")})
|
|
outcome = "flysim restarted" if action == "restart" else f"reset to rung {target}"
|
|
return f"{outcome} ({why})" if ok else f"{action} FAILED ({why})"
|
|
|
|
|
|
if __name__ == "__main__":
|
|
print(run())
|