diff --git a/docs/loop-review.md b/docs/loop-review.md index 9ea2b86..de9ddda 100644 --- a/docs/loop-review.md +++ b/docs/loop-review.md @@ -9,8 +9,8 @@ the project." 1. The release watchdog (`infra/bin/fly-watchdog` check 10) flags a suspected loop: few distinct macros, a short sequence repeating, no growth in places explored. It exports `fly_loop_suspected` and writes `/run/fly/wd/loop.json`. It never acts. The separate - `fly-loop-recover.timer` can restart only flysim after two fresh suspected probes; - see `infra/docs/loop-recovery.md`. This does not replace checkpoint-based review. + `fly-loop-recover.timer` unsticks a confirmed trap on its own: flysim restart, then milestone + resets (`infra/docs/loop-recovery.md`). That buys time; it does not replace checkpoint-based review. 2. The coordinator session (Fable) checks that marker on a schedule. On a flag it pulls the live checkpoint read-only (`pct pull`, into `.local/checkpoints/`, never committed), and spawns a review agent with the trap brief: reproduce from the checkpoint with the real diff --git a/infra/05-deploy.sh b/infra/05-deploy.sh index 5d79ce4..29e6620 100644 --- a/infra/05-deploy.sh +++ b/infra/05-deploy.sh @@ -685,10 +685,20 @@ fi # --------------------------------------------------------------------------- log "05-deploy: converging bin/ helpers to /opt/fly/bin" ct_exec "$CTID" -- mkdir -p /opt/fly/bin -for name in fly-watchdog fly-loop-recover fly-recap fly-retention fly-reset-to-milestone flypush flystage-launch flycast-launch wait-for-x wait-for-stage wait-for-health; do +for name in fly-watchdog fly-loop-recover fly-loop-reset fly-recap fly-retention fly-reset-to-milestone flypush flystage-launch flycast-launch wait-for-x wait-for-stage wait-for-health; do converge_file "$CTID" "$INFRA_DIR/bin/$name" "/opt/fly/bin/$name" 0755 root:root >/dev/null done +# The NOPASSWD surface those helpers use (config/fly-sudoers). 02-base installs it at +# provision time; converging it here too means a release that adds a helper's line +# (fly-loop-reset, v0.6.4) does not need a re-provision. Validated before install. +log "05-deploy: converging config/fly-sudoers, validated before install" +TMP_SUDOERS="/tmp/fly-sudoers.$$" +ct_push_file "$CTID" "$INFRA_DIR/config/fly-sudoers" "$TMP_SUDOERS" 0440 +ct_exec "$CTID" -- visudo -c -f "$TMP_SUDOERS" || die "config/fly-sudoers failed visudo -c, refusing to install it" +ct_exec "$CTID" -- install -o root -g root -m 0440 "$TMP_SUDOERS" /etc/sudoers.d/fly-watchdog +ct_exec "$CTID" -- rm -f "$TMP_SUDOERS" + # --------------------------------------------------------------------------- # 5. reload if anything unit-shaped changed # --------------------------------------------------------------------------- diff --git a/infra/bin/fly-loop-recover b/infra/bin/fly-loop-recover index 71de1c0..3ea369f 100755 --- a/infra/bin/fly-loop-recover +++ b/infra/bin/fly-loop-recover @@ -1,67 +1,278 @@ #!/usr/bin/env python3 -"""Conservative, out-of-process loop recovery; invoked after the watchdog probe.""" +"""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", "/run/fly/wd/recovery.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")) -COOLDOWN = 3600 +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") - model = os.environ.get("FLY_LOOP_MODEL") - if not url and not model: - return True - if not url or not model: - return False - payload = {"model": model, "temperature": 0, "messages": [ - {"role": "system", "content": "Classify whether a suspected repeating game macro loop is truly stuck. Return only JSON {\"stuck\":true|false}. No instructions or commands."}, - {"role": "user", "content": json.dumps({key: report.get(key) for key in ("reason", "sequence", "window", "places", "milestone")})}, - ]} + 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 - request = Request(url.rstrip("/") + "/chat/completions", json.dumps(payload).encode(), headers) - with urlopen(request, timeout=20) as response: - answer = json.load(response) - decision = json.loads(answer["choices"][0]["message"]["content"]) - return type(decision) is dict and type(decision.get("stuck")) is bool and decision["stuck"] + 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 run(now=None): - now = int(time.time() if now is None else now) - report = json.loads(REPORT.read_text()) - previous = json.loads(STATE.read_text()) if STATE.exists() else {} +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.write_text(json.dumps({key: previous[key] for key in ("restarted_at",) if key in previous}) + "\n") + state.pop("observedAt", None) + state["vetoes"] = 0 + write_json(STATE, state) return "not a fresh suspected loop" - if "restarted_at" in previous and now - previous["restarted_at"] < COOLDOWN: - return "recovery cooldown" - if previous.get("observed_at") == at: - return "waiting for next probe" - if not previous.get("observed_at") or at - previous["observed_at"] < INTERVAL - 30: - STATE.write_text(json.dumps({"observed_at": at}) + "\n") + + 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" - try: - approved = verdict(report) - except (OSError, ValueError, KeyError, IndexError) as error: - print("router unavailable or invalid:", type(error).__name__, file=sys.stderr) - return "router unavailable" - if not approved: - return "router did not confirm" - subprocess.run(["sudo", "-n", "systemctl", "restart", "flysim.service"], check=True) - STATE.write_text(json.dumps({"observed_at": at, "restarted_at": now}) + "\n") - return "flysim restarted; checkpoint restore keeps the rung" + 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__": diff --git a/infra/bin/fly-loop-reset b/infra/bin/fly-loop-reset new file mode 100755 index 0000000..de16570 --- /dev/null +++ b/infra/bin/fly-loop-reset @@ -0,0 +1,36 @@ +#!/usr/bin/env bash +# infra/bin/fly-loop-reset — the milestone step of fly-loop-recover's ladder, as root. +# +# fly-loop-recover runs as User=fly and reaches this through one NOPASSWD sudoers line +# (config/fly-sudoers). It does what infra/docs/runbook.md "Restart the run from a rung" +# does by hand for a release whose adapter already wrote the archive: stop flysim, promote +# milestone-.checkpoint with fly-reset-to-milestone, start flysim. flysim is started +# again whether or not the reset worked, so the stream never stays down on a failure; +# fly-reset-to-milestone copies both stores aside before it rewrites anything. +# +# Usage: fly-loop-reset +set -euo pipefail + +: "${FLY_STATE_DIR:=/srv/fly/state}" +: "${FLY_SERVICE:=flysim.service}" +: "${FLY_RESET_BIN:=/opt/fly/bin/fly-reset-to-milestone}" + +log() { echo "fly-loop-reset: $*" >&2; } + +RANK="${1:-}" +if [ "$#" -ne 1 ] || ! [[ "$RANK" =~ ^[0-9]{1,2}$ ]]; then + log "usage: fly-loop-reset " + exit 2 +fi +if [ ! -f "${FLY_STATE_DIR}/milestone-${RANK}.checkpoint" ]; then + log "no milestone archive for rung ${RANK}; nothing touched" + exit 1 +fi + +log "stopping ${FLY_SERVICE} to reset to rung ${RANK}" +systemctl stop "$FLY_SERVICE" +status=0 +"$FLY_RESET_BIN" "$RANK" || status=$? +[ "$status" -eq 0 ] || log "fly-reset-to-milestone ${RANK} failed (exit ${status}); starting ${FLY_SERVICE} on the state it left" +systemctl start "$FLY_SERVICE" +exit "$status" diff --git a/infra/config/fly-sudoers b/infra/config/fly-sudoers index f22752d..8963a18 100644 --- a/infra/config/fly-sudoers +++ b/infra/config/fly-sudoers @@ -16,3 +16,7 @@ fly ALL=(root) NOPASSWD: /usr/bin/systemctl restart flycast.service fly ALL=(root) NOPASSWD: /usr/bin/systemctl restart mediamtx.service fly ALL=(root) NOPASSWD: /usr/bin/systemctl restart flypush.service fly ALL=(root) NOPASSWD: /usr/sbin/reboot +# fly-loop-recover.service (also User=fly): its ladder restarts flysim (granted above) and +# resets to a milestone archive through this one root wrapper, which validates the rung, +# stops/resets/starts flysim and nothing else (infra/docs/loop-recovery.md). +fly ALL=(root) NOPASSWD: /opt/fly/bin/fly-loop-reset diff --git a/infra/docs/loop-recovery.md b/infra/docs/loop-recovery.md index 4982482..8f3df98 100644 --- a/infra/docs/loop-recovery.md +++ b/infra/docs/loop-recovery.md @@ -1,29 +1,82 @@ # Automated loop recovery -`fly-watchdog` remains report-only. The separate `fly-loop-recover.timer` checks its -`/run/fly/wd/loop.json` every five minutes. Two fresh suspected reports from -separate watchdog probes, at least ~5 minutes apart, cause **only** a -`flysim.service` restart. The ordinary checkpoint restore keeps the current -rung and learned brain state, but clears session macro ledgers. This does not -press buttons or promote a milestone archive. A successful restart imposes a -one-hour cooldown, including across intervening clear reports. The timer's -journal records every decision; inspect it with -`journalctl -u fly-loop-recover.service`. Disable the timer to stop automatic -recovery: `systemctl disable --now fly-loop-recover.timer`. +`fly-watchdog` remains report-only. The separate `fly-loop-recover.timer` reads its +`/run/fly/wd/loop.json` every five minutes and unsticks a confirmed trap on its own, climbing +a ladder one step per trap that outlives the previous step: -By default confirmation is deterministic, using watchdog's existing signal. -An optional OpenAI-compatible local router can veto a recovery: provision a -root-managed `/etc/fly/loop-recovery.env` readable by the `fly` account, -containing `FLY_LOOP_ROUTER_URL` (base URL ending in `/v1`) and -`FLY_LOOP_MODEL` (an available free-tier model). A private router can additionally -use `FLY_LOOP_ROUTER_KEY`; provision it outside this public checkout and limit -file permissions to `0640 root:fly`. Never put credentials in the unit or the -repository. When either router setting is present, both must be set; malformed or -unavailable responses prevent recovery. The model may only return `{"stuck": true|false}` and cannot -choose commands, buttons, or checkpoint paths. Confirm the model actually -exists and is reachable from the release container before configuring it. +| level | step | cost | +|---|---|---| +| 0 | `systemctl restart flysim.service` | none: the restore keeps the rung and learned state, clears session macro ledgers | +| 1 | reset to the current rung's milestone archive | progress since the rung was first reached | +| 2+ | reset to the archive below the best rung (never lower) | one visible rung | -This is an unstick mechanism, not a macro bug fix. A recurring trap still needs -the checkpoint-based loop review in `docs/loop-review.md`. The live recovery -must be recorded in the host claim log by the operator reviewing the unit -journal; the unit itself has no host claim-log access. +- **Confirmed** means two fresh suspected watchdog reports (at most 10 minutes old, `action: none`) + from separate probes about 5 minutes apart. A clear report in between starts the count again. +- **Outlives** means the reports that confirm it again are at least 20 minutes after the step, so + their 10-brain-minute window lies after it. The helper waits that long after every step. +- **Budget:** at most two milestone resets per 24 hours. With the budget spent the step is a + restart, at most every three hours, until a reset is free again. +- **Starting over:** the ladder returns to level 0 when the fly reaches a new best rung, or after + six hours with no suspected report. + +A milestone step runs `/opt/fly/bin/fly-loop-reset ` through sudo (the one line in +`config/fly-sudoers`): stop flysim, `fly-reset-to-milestone`, start flysim. flysim is started again +even when the reset fails. The reset copies both stores to `/srv/fly/state.reset-` first, as +in the runbook, and needs no deploy when the running release wrote the archive or +`FLY_ACCEPT_ADAPTERS` already names its adapter. After each step the helper waits up to four +minutes for `/status` to report `running` (and the target rank for a reset); otherwise the step is +recorded as failed and the ladder still climbs. + +## On stream + +Every step is announced 60 seconds ahead (`FLY_LOOP_COUNTDOWN`) in +`/run/fly/wd/recovery-notice.json` (`FLY_RECOVERY_NOTICE`), which the stage's recovery splash +reads. It holds the contract below, written atomically, mode 0644: + +```json +{"v": 1, "id": "1790629095-reset", "phase": "countdown|acting|done|failed", + "action": "restart|reset", "fromRung": 12, "fromLabel": "MT. MOON", + "toRung": 11, "toLabel": "BOULDER BADGE", "reason": "unrewarded", + "loop": ["GO OBJECTIVE", "GO WARP"], "stuckSeconds": 900, + "announcedAt": 1790629095, "executeAt": 1790629155, "updatedAt": 1790629160} +``` + +`toRung`/`toLabel` are present for a reset only. Consumers ignore a notice whose `updatedAt` is +more than 15 minutes old. + +## Model confirmation + +An OpenAI-compatible router can be asked before each step. The model may only answer +`{"stuck": true|false}` (the reply may wrap it in prose; the first such object counts). It cannot +choose commands, rungs or buttons. `false` delays the step by one probe, at most three times in a +row; then the step goes ahead. If no model answers (unset, rate-limited, down, malformed), the +watchdog's confirmation stands alone. A model can delay recovery by about 15 minutes, never deny +it. + +Provision a root-managed `/etc/fly/loop-recovery.env`, mode `0640 root:fly`, outside this public +checkout: + +``` +FLY_LOOP_ROUTER_URL= +FLY_LOOP_ROUTER_KEY= +FLY_LOOP_MODELS=,,... +``` + +Models are tried in order within a 90-second budget. Prefer fast free-tier chat models, and put +providers with generous free limits first: free OpenRouter models share a small daily quota and +the router's circuit breaker can close the whole provider for a while. Check each model against +a stuck and a healthy report before listing it. `FLY_LOOP_MODEL` (one model) is still read when +`FLY_LOOP_MODELS` is unset. + +## Operating it + +- Decisions: `journalctl -u fly-loop-recover.service`. Every step, veto and ladder restart is also + appended to `/var/lib/fly-loop-recover/history.jsonl`. +- Ladder state: `/var/lib/fly-loop-recover/state.json` (the unit's `StateDirectory`, so it + survives a reboot). Deleting it starts the ladder over. +- Stop automatic recovery: `systemctl disable --now fly-loop-recover.timer`. +- Record automatic steps you find in the journal in the host claim log when you next claim the + container; the unit has no access to that log. + +This is an unstick mechanism, not a macro bug fix. A recurring trap still needs the +checkpoint-based loop review in `docs/loop-review.md`. diff --git a/infra/tests/lint.sh b/infra/tests/lint.sh index 4da888a..82e9897 100755 --- a/infra/tests/lint.sh +++ b/infra/tests/lint.sh @@ -1485,7 +1485,7 @@ fi echo "--- loop recovery tests ---" if python3 -m unittest discover -s "$REPO_ROOT/infra/tests" -p 'test_loop_recover.py' >/dev/null 2>&1; then - pass "loop recovery: fresh probes, cooldown, and router refusal" + pass "loop recovery: ladder, budget, reboot-safe state, model delay and fallback, splash notice" else fail "loop recovery tests failed; run python3 -m unittest discover -s infra/tests -p test_loop_recover.py -v" fi diff --git a/infra/tests/test_loop_recover.py b/infra/tests/test_loop_recover.py index 7d0c481..881971c 100644 --- a/infra/tests/test_loop_recover.py +++ b/infra/tests/test_loop_recover.py @@ -2,89 +2,205 @@ import importlib.machinery import importlib.util import json from pathlib import Path +import re import tempfile import unittest from unittest.mock import patch -script = Path(__file__).resolve().parents[1] / "bin/fly-loop-recover" +repo = Path(__file__).resolve().parents[2] +script = repo / "infra/bin/fly-loop-recover" spec = importlib.util.spec_from_loader("recover", importlib.machinery.SourceFileLoader("recover", str(script))) recover = importlib.util.module_from_spec(spec) spec.loader.exec_module(recover) +T0 = 1_790_000_000 + class RecoveryTests(unittest.TestCase): def setUp(self): self.temp = tempfile.TemporaryDirectory() self.addCleanup(self.temp.cleanup) root = Path(self.temp.name) - recover.REPORT = root / "loop.json" - recover.STATE = root / "recovery.json" - self.report = {"suspected": 1, "at": 1000, "action": "none", "reason": "unrewarded"} - self.save() + self.root = root + for name, value in (("REPORT", "loop.json"), ("STATE", "var/state.json"), + ("HISTORY", "var/history.jsonl"), ("NOTICE", "run/notice.json")): + patcher = patch.object(recover, name, root / value) + patcher.start() + self.addCleanup(patcher.stop) + milestones = root / "state" + milestones.mkdir() + for rung in (1, 9, 10, 11, 12): + (milestones / f"milestone-{rung}.checkpoint").write_text("x") + patcher = patch.object(recover, "MILESTONES", milestones) + patcher.start() + self.addCleanup(patcher.stop) + self.env = patch.dict("os.environ", {}, clear=False) + self.env.start() + self.addCleanup(self.env.stop) + for key in ("FLY_LOOP_ROUTER_URL", "FLY_LOOP_MODELS", "FLY_LOOP_MODEL", "FLY_LOOP_ROUTER_KEY"): + recover.os.environ.pop(key, None) + self.acts = [] + self.act = patch.object(recover, "act", side_effect=lambda action, target: self.acts.append((action, target)) or True) + self.act.start() + self.addCleanup(self.act.stop) + self.report(suspected=1, at=T0) - def save(self): - recover.REPORT.write_text(json.dumps(self.report)) + def report(self, **fields): + base = {"suspected": 1, "action": "none", "reason": "unrewarded", "sequence": ["GO WARP"], + "milestone": {"rank": 12, "label": "MT. MOON"}, "map": 61} + base.update(fields) + if "rank" in fields: + base["milestone"] = {"rank": base.pop("rank"), "label": "X"} + recover.REPORT.write_text(json.dumps(base)) - def test_two_probes_and_cooldown(self): - with patch.object(recover, "verdict", return_value=True), patch.object(recover.subprocess, "run") as restart: - self.assertIn("second probe", recover.run(1001)) - self.assertIn("next probe", recover.run(1002)) - self.report["at"] = 1300 - self.save() - self.assertIn("restarted", recover.run(1301)) - restart.assert_called_once_with(["sudo", "-n", "systemctl", "restart", "flysim.service"], check=True) - self.assertIn("cooldown", recover.run(1302)) + def tick(self, at, **fields): + """One watchdog probe at `at` and the timer running right after it.""" + self.report(at=at, **fields) + return recover.run(at + 5, sleep=lambda seconds: None) - def test_stale_and_clear_never_restart(self): - with patch.object(recover.subprocess, "run") as restart: - self.assertIn("not a fresh", recover.run(1700)) - self.report.update(at=1700, suspected=0) - self.save() - self.assertIn("not a fresh", recover.run(1700)) - restart.assert_not_called() + def confirm(self, start, **fields): + """Two probes ~5 min apart after `start`; returns the second decision.""" + self.assertIn("second probe", self.tick(start, **fields)) + return self.tick(start + 300, **fields) - def test_router_failure_does_not_act(self): - recover.run(1000) - self.report["at"] = 1300 - self.save() - with patch.object(recover, "verdict", side_effect=ValueError("bad response")), patch.object(recover.subprocess, "run") as restart: - self.assertEqual(recover.run(1300), "router unavailable") - restart.assert_not_called() + def test_ladder_restart_then_current_rung_then_rung_below(self): + self.assertIn("flysim restarted", self.confirm(T0)) + self.assertEqual(self.acts, [("restart", None)]) + self.assertIn("settling", self.tick(T0 + 600)) + self.assertIn("reset to rung 12", self.confirm(T0 + 300 + recover.SETTLE + 10)) + start = T0 + 2 * (300 + recover.SETTLE + 10) + self.assertIn("reset to rung 11", self.confirm(start, rank=12)) + self.assertEqual([a for a in self.acts], [("restart", None), ("reset", 12), ("reset", 11)]) - def test_cleared_probe_preserves_cooldown(self): - with patch.object(recover, "verdict", return_value=True), patch.object(recover.subprocess, "run") as restart: - recover.run(1000) - self.report["at"] = 1300 - self.save() - recover.run(1300) - self.report.update(at=1600, suspected=0) - self.save() - recover.run(1600) - self.report.update(at=1900, suspected=1) - self.save() - recover.run(1900) - self.report["at"] = 2200 - self.save() - self.assertIn("cooldown", recover.run(2200)) - restart.assert_called_once() + def test_spent_reset_budget_holds_to_restarts_three_hours_apart(self): + at = T0 + for _ in range(3): + self.confirm(at) + at += 300 + recover.SETTLE + 10 + self.assertEqual(len(self.acts), 3) + self.assertIn("holding", self.confirm(at)) + # the trap never cleared while holding, so the first probe after the hold acts + self.assertIn("flysim restarted", self.tick(at + recover.HOLD)) + self.assertEqual(self.acts[-1], ("restart", None)) - def test_router_rejection_does_not_restart(self): - recover.run(1000) - self.report["at"] = 1300 - self.save() - with patch.object(recover, "verdict", return_value=False), patch.object(recover.subprocess, "run") as restart: - self.assertIn("did not confirm", recover.run(1300)) - restart.assert_not_called() + def test_deeper_resets_never_go_below_the_rung_under_the_best(self): + state = {"level": 5, "bestRank": 12, "resets": [], "actedAt": 0} + recover.write_json(recover.STATE, state) + self.confirm(T0, rank=11) + self.assertEqual(self.acts, [("reset", 11)]) - def test_partial_router_configuration_cannot_bypass_veto(self): - with patch.dict("os.environ", {"FLY_LOOP_ROUTER_URL": "http://router/v1"}, clear=True): - self.assertFalse(recover.verdict(self.report)) - with patch.dict("os.environ", {"FLY_LOOP_MODEL": "free"}, clear=True): - self.assertFalse(recover.verdict(self.report)) + def test_new_best_rung_starts_the_ladder_over(self): + recover.write_json(recover.STATE, {"level": 2, "bestRank": 12, "resets": [], "actedAt": 0}) + self.confirm(T0, rank=13) + self.assertEqual(self.acts, [("restart", None)]) - def test_unconfigured_router_uses_deterministic_confirmation(self): - with patch.dict("os.environ", {}, clear=True): - self.assertTrue(recover.verdict(self.report)) + def test_quiet_hours_start_the_ladder_over(self): + recover.write_json(recover.STATE, {"level": 2, "bestRank": 12, "actedAt": 0, + "lastSuspectedAt": T0 - recover.QUIET - 1}) + self.confirm(T0) + self.assertEqual(self.acts, [("restart", None)]) + + def test_stale_clear_or_acted_reports_never_act(self): + self.report(at=T0 - 700) + self.assertIn("not a fresh", recover.run(T0)) + self.report(suspected=0, at=T0) + self.assertIn("not a fresh", recover.run(T0 + 1)) + self.report(at=T0, action="restart") + self.assertIn("not a fresh", recover.run(T0 + 2)) + self.assertEqual(self.acts, []) + + def test_a_clear_probe_breaks_the_streak(self): + self.tick(T0) + self.tick(T0 + 300, suspected=0) + self.assertIn("second probe", self.tick(T0 + 600)) + self.assertEqual(self.acts, []) + + def test_state_survives_a_reboot_and_ignores_a_missing_or_corrupt_file(self): + self.confirm(T0) + self.assertTrue(recover.STATE.exists()) + self.assertEqual(json.loads(recover.STATE.read_text())["level"], 1) + recover.STATE.write_text("{nope") + self.assertIn("second probe", self.tick(T0 + 5000)) + + def test_model_vetoes_delay_a_step_but_never_deny_it(self): + with patch.object(recover, "verdict", return_value=(False, "m")): + self.tick(T0) + for n in range(1, recover.VETO_LIMIT + 1): + self.assertIn(f"({n}/{recover.VETO_LIMIT})", self.tick(T0 + 300 * n)) + self.assertIn("veto limit reached", self.tick(T0 + 300 * (recover.VETO_LIMIT + 1))) + self.assertEqual(self.acts, [("restart", None)]) + + def test_no_model_answer_falls_back_to_the_watchdog(self): + with patch.object(recover, "verdict", return_value=(None, None)): + self.assertIn("watchdog alone", self.confirm(T0)) + + def test_models_are_tried_in_order_until_one_answers(self): + recover.os.environ.update(FLY_LOOP_ROUTER_URL="http://router/v1", FLY_LOOP_MODELS="a, b ,c") + replies = iter([OSError("429"), {"choices": [{"message": {"content": "thinking...\n```json\n{\"stuck\": true}\n```"}}]}]) + asked = [] + + class Reply: + def __init__(self, body): + self.body = body + + def __enter__(self): + return self + + def __exit__(self, *exc): + return False + + def read(self): + return json.dumps(self.body).encode() + + def fake(request, timeout): + asked.append(json.loads(request.data)["model"]) + reply = next(replies) + if isinstance(reply, Exception): + raise reply + return Reply(reply) + + with patch.object(recover, "urlopen", side_effect=fake): + self.assertEqual(recover.verdict({"reason": "unrewarded"}), (True, "b")) + self.assertEqual(asked, ["a", "b"]) + + def test_parse_stuck_tolerates_prose_and_rejects_non_booleans(self): + self.assertIs(recover.parse_stuck('{"stuck":false}'), False) + self.assertIs(recover.parse_stuck('Answer: {"stuck": true} done'), True) + self.assertIsNone(recover.parse_stuck('{"stuck": "yes"}')) + self.assertIsNone(recover.parse_stuck("")) + self.assertIsNone(recover.parse_stuck(None)) + + def test_notice_goes_countdown_acting_done_for_the_splash(self): + phases = [] + real = recover.notice + with patch.object(recover, "notice", side_effect=lambda base, phase, now: phases.append(phase) or real(base, phase, now)): + self.confirm(T0 + 300 + recover.SETTLE) + recover.write_json(recover.STATE, dict(json.loads(recover.STATE.read_text()), actedAt=0)) + self.confirm(T0 + 2 * (300 + recover.SETTLE)) + notice = json.loads(recover.NOTICE.read_text()) + self.assertEqual(phases, ["countdown", "acting", "done"] * 2) + self.assertEqual((notice["v"], notice["action"], notice["toRung"], notice["toLabel"], notice["fromRung"]), + (1, "reset", 12, "MT. MOON", 12)) + self.assertEqual(notice["executeAt"] - notice["announcedAt"], recover.COUNTDOWN) + + def test_failed_action_is_reported_and_still_climbs(self): + self.act.stop() + with patch.object(recover, "act", return_value=False): + self.assertIn("FAILED", self.confirm(T0)) + self.act.start() + self.assertEqual(json.loads(recover.NOTICE.read_text())["phase"], "failed") + self.assertEqual(json.loads(recover.STATE.read_text())["level"], 1) + + def test_history_records_every_step(self): + self.confirm(T0) + lines = [json.loads(line) for line in recover.HISTORY.read_text().splitlines()] + self.assertEqual(lines[-1]["event"], "restart") + self.assertTrue(lines[-1]["ok"]) + + def test_ladder_labels_match_the_rust_table(self): + source = (repo / "services/flysim/crates/flybrain-gb/src/pokemon_red/mod.rs").read_text() + table = re.search(r"RANK_LADDER: \[&str; \d+\] = \[(.*?)\];", source, re.S).group(1) + self.assertEqual(recover.LADDER, re.findall(r'"([^"]*)"', table)) if __name__ == "__main__": diff --git a/infra/units/fly-loop-recover.service b/infra/units/fly-loop-recover.service index 429fb2b..a2dff4d 100644 --- a/infra/units/fly-loop-recover.service +++ b/infra/units/fly-loop-recover.service @@ -1,8 +1,12 @@ [Unit] -Description=Check confirmed macro loops for minimal flysim recovery +Description=Recover confirmed macro loops: flysim restart, then milestone resets [Service] Type=oneshot User=fly EnvironmentFile=-/etc/fly/loop-recovery.env +# The ladder's state and history outlive a reboot (/run does not). +StateDirectory=fly-loop-recover +# A step is a 60 s on-stream countdown, up to 90 s of model calls and 4 min to verify. +TimeoutStartSec=12min ExecStart=/opt/fly/bin/fly-loop-recover