Restart flysim, then reset to the current rung's milestone, then to the archive below the best rung (never lower), one step per confirmed trap that outlives the previous one. Two resets a day, three-hour restarts once they are spent; the ladder starts over at a new best rung or after six quiet hours. State and history live in the unit's StateDirectory so a reboot does not forget where the ladder stood. A router model list confirms each step; a 'not stuck' answer delays it at most three probes and no answer leaves the watchdog to decide alone. Each step is announced 60 s ahead in /run/fly/wd/recovery-notice.json for the stage's recovery splash. fly-loop-reset is the one new root surface (a sudoers line); 05-deploy now converges config/fly-sudoers so a release can add it.
279 lines
11 KiB
Python
Executable file
279 lines
11 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
|
|
|
|
# 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 archive falls back to restart.
|
|
|
|
Level 1 resets to the current rung's archive, level 2 and beyond to the archive below the
|
|
best rung (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
|
|
ceiling = rank if level == 1 else min(rank, max(best, rank) - 1)
|
|
below = [rung for rung in rungs() if rung <= ceiling]
|
|
if below:
|
|
return "reset", below[-1]
|
|
return "restart", None
|
|
|
|
|
|
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 = ["sudo", "-n", "systemctl", "restart", "flysim.service"]
|
|
else:
|
|
command = ["sudo", "-n", RESET_BIN, str(target)]
|
|
done = subprocess.run(command, check=False, timeout=VERIFY_TIMEOUT)
|
|
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)
|
|
notice(base, "acting", clock())
|
|
ok = act(action, target)
|
|
finished = int(clock())
|
|
notice(base, "done" if ok else "failed", finished)
|
|
|
|
state.update(level=level + 1, actedAt=finished, lastAction=action, vetoes=0)
|
|
state.pop("observedAt", None)
|
|
if action == "reset":
|
|
state["resets"].append(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())
|