flybrain/infra/bin/fly-loop-recover
acamilo f205d95e5e loop recovery: review-ladder fixes, root never runs a fly-writable binary
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.
2026-09-28 21:34:32 +00:00

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())