Merge feat/loop-recover-ladder: automatic loop recovery climbs restart, rung reset, rung below
This commit is contained in:
commit
3922b82762
9 changed files with 908 additions and 129 deletions
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -272,6 +272,21 @@ if [ -n "$RELEASE_TARBALL" ]; then
|
|||
cpu_pin chown -R fly:fly "$release_path"
|
||||
fi
|
||||
|
||||
# A root-owned flysim for fly-loop-reset, which runs as root and so must never execute the
|
||||
# fly-owned copy above (the fly account could replace it). Taken from the tarball on the host
|
||||
# and checked against the tarball's own MANIFEST on every deploy, so a release installed by
|
||||
# an earlier deploy is covered too.
|
||||
log "05-deploy: installing a root-owned flysim at /opt/fly/sbin/flysim for fly-loop-reset"
|
||||
root_flysim="$(mktemp)"
|
||||
tar -xzOf "$RELEASE_TARBALL" ./flysim > "$root_flysim"
|
||||
want_sha="$(tar -xzOf "$RELEASE_TARBALL" ./MANIFEST | awk '$2 == "flysim" || $2 == "./flysim" {print $1; exit}')"
|
||||
[ -n "$want_sha" ] && [ "$(sha256sum "$root_flysim" | cut -d' ' -f1)" = "$want_sha" ] \
|
||||
|| { rm -f "$root_flysim"; die "the tarball's flysim does not match its MANIFEST; not installing /opt/fly/sbin/flysim"; }
|
||||
ct_exec "$CTID" -- install -d -o root -g root -m 0755 /opt/fly/sbin
|
||||
ct_push_file "$CTID" "$root_flysim" /opt/fly/sbin/flysim.new 0755
|
||||
ct_exec "$CTID" -- sh -c 'chown root:root /opt/fly/sbin/flysim.new && mv -f /opt/fly/sbin/flysim.new /opt/fly/sbin/flysim'
|
||||
rm -f "$root_flysim"
|
||||
|
||||
log "05-deploy: overlaying infra-owned stage/serve.{mjs,sh} (docs/design/infra.md's 'own tiny static server' clarification)"
|
||||
converge_file "$CTID" "$INFRA_DIR/config/serve.mjs" "${release_path}/stage/serve.mjs" 0644 fly:fly >/dev/null
|
||||
converge_file "$CTID" "$INFRA_DIR/config/serve.sh" "${release_path}/stage/serve.sh" 0755 fly:fly >/dev/null
|
||||
|
|
@ -685,10 +700,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
|
||||
# ---------------------------------------------------------------------------
|
||||
|
|
|
|||
|
|
@ -1,67 +1,300 @@
|
|||
#!/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
|
||||
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")
|
||||
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
|
||||
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)
|
||||
with urlopen(request, timeout=20) as response:
|
||||
try:
|
||||
with urlopen(request, timeout=min(30, left)) 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"]
|
||||
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 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.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)
|
||||
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__":
|
||||
|
|
|
|||
172
infra/bin/fly-loop-reset
Executable file
172
infra/bin/fly-loop-reset
Executable file
|
|
@ -0,0 +1,172 @@
|
|||
#!/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: stop flysim, promote milestone-<N>.checkpoint with fly-reset-to-milestone,
|
||||
# start flysim.
|
||||
#
|
||||
# Root never executes anything the fly account can write. The release tree under
|
||||
# /opt/fly/releases is fly-owned, so the flysim run here is the root-owned copy 05-deploy
|
||||
# installs at /opt/fly/sbin/flysim straight from the release tarball; without it every rung is
|
||||
# refused and the ladder falls back to restarts. Paths are fixed (test overrides are honoured
|
||||
# only when not root), and /etc/fly/fly.env is read as KEY=VALUE data, never sourced.
|
||||
#
|
||||
# fly-reset-to-milestone does not check that the build can restore the archive, and a flysim
|
||||
# that refuses every checkpoint refuses to start: a black stream. So a rung is only used when
|
||||
# its archive's compatibility string equals this build's, or differs only in an adapter id that
|
||||
# FLY_ACCEPT_ADAPTERS names (05-deploy.sh's adapter_migration_accepted).
|
||||
#
|
||||
# fly-watchdog restarts a flysim whose /healthz fails, which during a reset would start it on a
|
||||
# half-rewritten store; its timer is stopped for the reset and started again on every exit path.
|
||||
#
|
||||
# Usage: fly-loop-reset <rung> reset to that rung
|
||||
# fly-loop-reset --list print the rungs this build can restore, one per line
|
||||
set -euo pipefail
|
||||
|
||||
STATE_DIR=/srv/fly/state
|
||||
ENV_FILE=/etc/fly/fly.env
|
||||
FLYSIM=/opt/fly/sbin/flysim
|
||||
RESET_BIN=/opt/fly/bin/fly-reset-to-milestone
|
||||
if [ "$(id -u)" -ne 0 ]; then
|
||||
STATE_DIR="${FLY_LOOP_RESET_TEST_STATE_DIR:-$STATE_DIR}"
|
||||
ENV_FILE="${FLY_LOOP_RESET_TEST_ENV_FILE:-$ENV_FILE}"
|
||||
FLYSIM="${FLY_LOOP_RESET_TEST_FLYSIM:-$FLYSIM}"
|
||||
RESET_BIN="${FLY_LOOP_RESET_TEST_RESET_BIN:-$RESET_BIN}"
|
||||
fi
|
||||
SERVICE=flysim.service
|
||||
WATCHDOG_TIMER=fly-watchdog.timer
|
||||
WATCHDOG_SERVICE=fly-watchdog.service
|
||||
readonly STATE_DIR ENV_FILE FLYSIM RESET_BIN SERVICE WATCHDOG_TIMER WATCHDOG_SERVICE
|
||||
|
||||
log() { echo "fly-loop-reset: $*" >&2; }
|
||||
usage() { log "usage: fly-loop-reset <rung> | --list"; exit 2; }
|
||||
|
||||
LIST=0
|
||||
RANK=""
|
||||
[ "$#" -eq 1 ] || usage
|
||||
if [ "$1" = --list ]; then
|
||||
LIST=1
|
||||
elif [[ "$1" =~ ^[0-9]{1,2}$ ]]; then
|
||||
RANK="$1"
|
||||
else
|
||||
usage
|
||||
fi
|
||||
|
||||
# FLY_* lines of fly.env as `env` arguments. Values may be quoted the systemd way; nothing
|
||||
# is evaluated.
|
||||
env_args=()
|
||||
accepted=""
|
||||
if [ -r "$ENV_FILE" ]; then
|
||||
while IFS= read -r line || [ -n "$line" ]; do
|
||||
[[ "$line" =~ ^(FLY_[A-Z0-9_]*)=(.*)$ ]] || continue
|
||||
key="${BASH_REMATCH[1]}"
|
||||
value="${BASH_REMATCH[2]}"
|
||||
if [[ "$value" =~ ^\"(.*)\"$ ]] || [[ "$value" =~ ^\'(.*)\'$ ]]; then
|
||||
value="${BASH_REMATCH[1]}"
|
||||
fi
|
||||
env_args+=("$key=$value")
|
||||
if [ "$key" = FLY_ACCEPT_ADAPTERS ]; then
|
||||
accepted="$value"
|
||||
fi
|
||||
done < "$ENV_FILE"
|
||||
fi
|
||||
clean_env() { env -i PATH=/usr/sbin:/usr/bin:/sbin:/bin HOME=/root "${env_args[@]}" "$@"; }
|
||||
|
||||
# The same rule as 05-deploy.sh: exactly one '/'-separated field differs, it is the adapter
|
||||
# (index 1), and the archive's adapter id is listed.
|
||||
migration_accepted() {
|
||||
local old="$1" new="$2" i differing=0 index=-1 entry
|
||||
local -a old_parts new_parts
|
||||
IFS='/' read -r -a old_parts <<< "$old"
|
||||
IFS='/' read -r -a new_parts <<< "$new"
|
||||
[ "${#old_parts[@]}" -eq "${#new_parts[@]}" ] || return 1
|
||||
for ((i = 0; i < ${#old_parts[@]}; i++)); do
|
||||
if [ "${old_parts[$i]}" != "${new_parts[$i]}" ]; then
|
||||
differing=$((differing + 1))
|
||||
index=$i
|
||||
fi
|
||||
done
|
||||
[ "$differing" -eq 1 ] && [ "$index" -eq 1 ] || return 1
|
||||
for entry in ${accepted//,/ }; do
|
||||
[ "$entry" = "${old_parts[1]}" ] && return 0
|
||||
done
|
||||
return 1
|
||||
}
|
||||
|
||||
# The root copy must be the flysim that runs: a manual rollback of /opt/fly/current without a
|
||||
# deploy would otherwise approve archives the running build refuses.
|
||||
LIVE_FLYSIM=/opt/fly/current/flysim
|
||||
[ "$(id -u)" -eq 0 ] || LIVE_FLYSIM="${FLY_LOOP_RESET_TEST_LIVE_FLYSIM:-$FLYSIM}"
|
||||
same_build() {
|
||||
[ -f "$LIVE_FLYSIM" ] \
|
||||
&& [ "$(sha256sum < "$FLYSIM" | cut -d' ' -f1)" = "$(sha256sum < "$LIVE_FLYSIM" | cut -d' ' -f1)" ]
|
||||
}
|
||||
|
||||
build_compat=""
|
||||
if [ -x "$FLYSIM" ] && [ "$(stat -c %u "$FLYSIM")" = "$(id -u)" ] && same_build; then
|
||||
build_compat="$(clean_env timeout 90 "$FLYSIM" --print-compatibility 2>/dev/null | tail -n1 || true)"
|
||||
fi
|
||||
if [ -z "$build_compat" ]; then
|
||||
log "no compatibility string from $FLYSIM (missing, not owned by $(id -un), not the running build, or failed); no rung is restorable"
|
||||
[ "$LIST" -eq 1 ] && exit 0
|
||||
exit 3
|
||||
fi
|
||||
|
||||
restorable() {
|
||||
local archive="${STATE_DIR}/milestone-$1.checkpoint" compat
|
||||
[ -f "$archive" ] || return 1
|
||||
compat="$(head -c 262144 "$archive" 2>/dev/null | grep -a -o -m1 '"compatibility"[[:space:]]*:[[:space:]]*"[^"]*"' | head -n1 | cut -d'"' -f4 || true)"
|
||||
[ -n "$compat" ] || return 1
|
||||
[ "$compat" = "$build_compat" ] || migration_accepted "$compat" "$build_compat"
|
||||
}
|
||||
|
||||
if [ "$LIST" -eq 1 ]; then
|
||||
for archive in "${STATE_DIR}"/milestone-*.checkpoint; do
|
||||
[ -e "$archive" ] || continue
|
||||
rung="${archive##*/milestone-}"
|
||||
rung="${rung%.checkpoint}"
|
||||
if [[ "$rung" =~ ^[0-9]{1,2}$ ]] && restorable "$rung"; then
|
||||
echo "$rung"
|
||||
fi
|
||||
done | sort -n
|
||||
exit 0
|
||||
fi
|
||||
|
||||
if ! restorable "$RANK"; then
|
||||
log "rung ${RANK}'s archive is missing or not restorable by this build (FLY_ACCEPT_ADAPTERS='${accepted}'); nothing touched"
|
||||
exit 3
|
||||
fi
|
||||
|
||||
# shellcheck disable=SC2317 # invoked by the EXIT trap
|
||||
finish() {
|
||||
local status=$?
|
||||
systemctl start "$SERVICE" || log "starting ${SERVICE} failed"
|
||||
systemctl start "$WATCHDOG_TIMER" || log "starting ${WATCHDOG_TIMER} failed"
|
||||
exit "$status"
|
||||
}
|
||||
trap finish EXIT
|
||||
trap 'exit 143' TERM INT HUP
|
||||
|
||||
log "pausing ${WATCHDOG_TIMER} and stopping ${SERVICE} to reset to rung ${RANK}"
|
||||
systemctl stop "$WATCHDOG_TIMER"
|
||||
# A oneshot probe mid-run is "activating", which `is-active` does not count; wait on ActiveState.
|
||||
watchdog_idle() {
|
||||
case "$(systemctl show -p ActiveState --value "$WATCHDOG_SERVICE")" in
|
||||
inactive|failed) return 0 ;;
|
||||
*) return 1 ;;
|
||||
esac
|
||||
}
|
||||
for _ in $(seq 1 60); do
|
||||
watchdog_idle && break
|
||||
sleep 1
|
||||
done
|
||||
if ! watchdog_idle; then
|
||||
log "${WATCHDOG_SERVICE} still running after 60 s; not resetting"
|
||||
exit 4
|
||||
fi
|
||||
systemctl stop "$SERVICE"
|
||||
status=0
|
||||
clean_env FLY_BIN="$FLYSIM" "$RESET_BIN" "$RANK" || status=$?
|
||||
[ "$status" -eq 0 ] || log "fly-reset-to-milestone ${RANK} failed (exit ${status}); starting ${SERVICE} on the state it left"
|
||||
exit "$status"
|
||||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -1,29 +1,104 @@
|
|||
# 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. The same three-hour spacing applies
|
||||
when a reset level finds no restorable rung and restarts instead.
|
||||
- **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 only uses an archive the running build can restore. `fly-loop-reset --list`
|
||||
compares each archive's compatibility string with `flysim --print-compatibility` and accepts an
|
||||
adapter-only difference that `FLY_ACCEPT_ADAPTERS` in `/etc/fly/fly.env` names (the rule
|
||||
`05-deploy.sh` applies). A flysim that refuses every checkpoint refuses to start, so a step whose
|
||||
rung is not restorable is a restart instead, never a lower rung. Keep `FLY_ACCEPT_ADAPTERS` in the
|
||||
release env file so a deploy does not drop it.
|
||||
|
||||
The step runs `/opt/fly/bin/fly-loop-reset <rung>` through sudo (the one line in
|
||||
`config/fly-sudoers`). As root it:
|
||||
|
||||
- runs only the root-owned `/opt/fly/sbin/flysim`, and only while it is byte-identical to the
|
||||
flysim `/opt/fly/current` points to, that `05-deploy.sh` installs from the release
|
||||
tarball after checking it against the tarball's MANIFEST (never the fly-owned release tree).
|
||||
Without that copy (an infra-only deploy, or a manual rollback) no rung is restorable and the
|
||||
ladder only restarts; after a release deploy, check as `fly` that
|
||||
`sudo -n /opt/fly/bin/fly-loop-reset --list` names the current rung;
|
||||
- reads `fly.env` as `KEY=VALUE` data, never sources it, and ignores the caller's environment;
|
||||
- stops `fly-watchdog.timer` and waits for a running probe to finish, so the watchdog cannot
|
||||
start flysim on a half-rewritten store; stops flysim; runs `fly-reset-to-milestone`; and on every
|
||||
exit path starts flysim and the watchdog timer again.
|
||||
|
||||
The reset copies both stores to `/srv/fly/state.reset-<UTC>` first, as in the runbook, and
|
||||
clears milestone archives above the rung. 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. The step is recorded before it runs, so a helper killed
|
||||
mid-step has still climbed and spent the reset. Do not stop `fly-loop-recover.service` during a
|
||||
step: that kills the reset too, and the wrapper's exit trap starts flysim on whatever state is
|
||||
left.
|
||||
|
||||
## 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=<base URL ending in /v1>
|
||||
FLY_LOOP_ROUTER_KEY=<router key>
|
||||
FLY_LOOP_MODELS=<model>,<fallback model>,...
|
||||
```
|
||||
|
||||
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`.
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -2,14 +2,20 @@ 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)
|
||||
RESTORABLE_RUNGS = recover.restorable_rungs
|
||||
WRAPPER = repo / "infra/bin/fly-loop-reset"
|
||||
|
||||
T0 = 1_790_000_000
|
||||
|
||||
|
||||
class RecoveryTests(unittest.TestCase):
|
||||
|
|
@ -17,74 +23,333 @@ class RecoveryTests(unittest.TestCase):
|
|||
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.unrestorable = set()
|
||||
restorable = patch.object(recover, "restorable_rungs", side_effect=lambda: {1, 9, 10, 11, 12} - self.unrestorable)
|
||||
restorable.start()
|
||||
self.addCleanup(restorable.stop)
|
||||
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_an_unrestorable_rung_below_the_best_means_a_restart_never_a_lower_rung(self):
|
||||
self.unrestorable = {11}
|
||||
recover.write_json(recover.STATE, {"level": 2, "bestRank": 12, "resets": [], "actedAt": 0})
|
||||
self.confirm(T0)
|
||||
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_a_fly_already_below_the_rung_under_its_best_is_restarted_not_reset(self):
|
||||
recover.write_json(recover.STATE, {"level": 2, "bestRank": 12, "resets": [], "actedAt": 0})
|
||||
self.confirm(T0, rank=10)
|
||||
self.assertEqual(self.acts, [("restart", None)])
|
||||
|
||||
def test_no_restorable_archive_means_a_restart(self):
|
||||
self.unrestorable = {1, 9, 10, 11, 12}
|
||||
recover.write_json(recover.STATE, {"level": 1, "bestRank": 12, "resets": [], "actedAt": 0})
|
||||
self.confirm(T0)
|
||||
self.assertEqual(self.acts, [("restart", None)])
|
||||
|
||||
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_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_a_step_that_raises_or_times_out_is_a_failed_step_with_state_saved(self):
|
||||
self.act.stop()
|
||||
with patch.object(recover.subprocess, "run", side_effect=recover.subprocess.TimeoutExpired("sudo", 1)):
|
||||
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_state_is_saved_before_the_step_runs(self):
|
||||
seen = []
|
||||
self.act.stop()
|
||||
with patch.object(recover, "act", side_effect=lambda a, t: seen.append(json.loads(recover.STATE.read_text())) or True):
|
||||
self.confirm(T0)
|
||||
self.act.start()
|
||||
self.assertEqual(seen[0]["level"], 1)
|
||||
self.assertNotIn("observedAt", seen[0])
|
||||
|
||||
def test_restorable_rungs_parses_the_list_and_is_empty_on_any_failure(self):
|
||||
ok = recover.subprocess.CompletedProcess([], 0, stdout="11\n12\n", stderr="")
|
||||
with patch.object(recover.subprocess, "run", return_value=ok):
|
||||
self.assertEqual(RESTORABLE_RUNGS(), {11, 12})
|
||||
with patch.object(recover.subprocess, "run", side_effect=OSError("no sudo")):
|
||||
self.assertEqual(RESTORABLE_RUNGS(), set())
|
||||
with patch.object(recover.subprocess, "run", return_value=recover.subprocess.CompletedProcess([], 1, "12", "")):
|
||||
self.assertEqual(RESTORABLE_RUNGS(), set())
|
||||
|
||||
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))
|
||||
|
||||
|
||||
BUILD = "lif-1ms-f64-v2/pokered-unique8-v7/abc/fly-kc-mbon-rstdp-v2"
|
||||
|
||||
|
||||
@unittest.skipIf(recover.os.geteuid() == 0, "the wrapper ignores test overrides as root")
|
||||
class WrapperTests(unittest.TestCase):
|
||||
"""infra/bin/fly-loop-reset against a fake flysim, systemctl, reset tool and archives."""
|
||||
|
||||
def setUp(self):
|
||||
self.temp = tempfile.TemporaryDirectory()
|
||||
self.addCleanup(self.temp.cleanup)
|
||||
root = Path(self.temp.name)
|
||||
self.root = root
|
||||
self.calls = root / "calls.log"
|
||||
(root / "state").mkdir()
|
||||
(root / "bin").mkdir()
|
||||
# The watchdog probe is mid-run ("activating") for the first two looks.
|
||||
self.script(root / "bin/systemctl", f'''echo "systemctl $*" >> {self.calls}
|
||||
if [ "$1" = show ]; then
|
||||
n=$(grep -c "^systemctl show" {self.calls})
|
||||
if [ "$n" -le 2 ]; then echo activating; else echo inactive; fi
|
||||
fi''')
|
||||
self.script(root / "flysim", f'[ "$1" = --print-compatibility ] && echo "{BUILD}"')
|
||||
self.script(root / "reset", f'echo "reset $* FLY_BIN=$FLY_BIN" >> {self.calls}')
|
||||
self.env_file = root / "fly.env"
|
||||
self.env_file.write_text('FLY_MACRO_MODE=macros\nGAME_TITLE=Pokemon (Red) $(touch pwned)\n')
|
||||
self.archive(12, BUILD)
|
||||
self.archive(11, BUILD.replace("-v7", "-v6"))
|
||||
self.archive(1, BUILD.replace("-v7", "-v5"))
|
||||
|
||||
def script(self, path, body):
|
||||
path.write_text("#!/bin/sh\n" + body + "\n")
|
||||
path.chmod(0o755)
|
||||
|
||||
def archive(self, rung, compat):
|
||||
(self.root / f"state/milestone-{rung}.checkpoint").write_bytes(
|
||||
b"FLYSIM01" + json.dumps({"generation": 1, "compatibility": compat}).encode() + b"\x00" * 64)
|
||||
|
||||
def run_wrapper(self, *args):
|
||||
env = {"PATH": f"{self.root}/bin:/usr/bin:/bin", "FLY_LOOP_RESET_TEST_STATE_DIR": str(self.root / "state"),
|
||||
"FLY_LOOP_RESET_TEST_ENV_FILE": str(self.env_file), "FLY_LOOP_RESET_TEST_FLYSIM": str(self.root / "flysim"),
|
||||
"FLY_LOOP_RESET_TEST_RESET_BIN": str(self.root / "reset")}
|
||||
env.update(getattr(self, "extra_env", {}))
|
||||
return recover.subprocess.run([str(WRAPPER), *args], env=env, capture_output=True, text=True, timeout=60)
|
||||
|
||||
def log(self):
|
||||
return self.calls.read_text() if self.calls.exists() else ""
|
||||
|
||||
def test_list_names_only_rungs_this_build_restores(self):
|
||||
self.assertEqual(self.run_wrapper("--list").stdout.split(), ["12"])
|
||||
self.env_file.write_text("FLY_ACCEPT_ADAPTERS=pokered-unique8-v6\n")
|
||||
self.assertEqual(self.run_wrapper("--list").stdout.split(), ["11", "12"])
|
||||
self.assertFalse((self.root / "pwned").exists())
|
||||
|
||||
def test_reset_pauses_the_watchdog_and_always_starts_flysim_again(self):
|
||||
done = self.run_wrapper("12")
|
||||
self.assertEqual(done.returncode, 0, done.stderr)
|
||||
self.assertEqual([line.split()[:3] for line in self.log().splitlines() if not line.startswith("systemctl is-active")], [
|
||||
["systemctl", "stop", "fly-watchdog.timer"],
|
||||
["systemctl", "show", "-p"], ["systemctl", "show", "-p"], ["systemctl", "show", "-p"], ["systemctl", "show", "-p"],
|
||||
["systemctl", "stop", "flysim.service"],
|
||||
["reset", "12", f"FLY_BIN={self.root}/flysim"],
|
||||
["systemctl", "start", "flysim.service"], ["systemctl", "start", "fly-watchdog.timer"]])
|
||||
|
||||
def test_a_failed_reset_still_starts_flysim_and_the_watchdog(self):
|
||||
self.script(self.root / "reset", f'echo "reset $*" >> {self.calls}; exit 7')
|
||||
self.assertEqual(self.run_wrapper("12").returncode, 7)
|
||||
self.assertIn("systemctl start flysim.service", self.log())
|
||||
self.assertIn("systemctl start fly-watchdog.timer", self.log())
|
||||
|
||||
def test_an_unrestorable_or_missing_rung_touches_nothing(self):
|
||||
for rung in ("11", "5"):
|
||||
self.assertEqual(self.run_wrapper(rung).returncode, 3)
|
||||
self.assertEqual(self.log(), "")
|
||||
|
||||
def test_no_build_compatibility_means_nothing_is_restorable(self):
|
||||
(self.root / "flysim").unlink()
|
||||
self.assertEqual(self.run_wrapper("--list").stdout, "")
|
||||
self.assertEqual(self.run_wrapper("12").returncode, 3)
|
||||
self.assertEqual(self.log(), "")
|
||||
|
||||
def test_a_root_copy_that_is_not_the_running_build_restores_nothing(self):
|
||||
other = self.root / "live-flysim"
|
||||
self.script(other, "echo other build")
|
||||
self.extra_env = {"FLY_LOOP_RESET_TEST_LIVE_FLYSIM": str(other)}
|
||||
self.assertEqual(self.run_wrapper("--list").stdout, "")
|
||||
self.assertEqual(self.run_wrapper("12").returncode, 3)
|
||||
|
||||
def test_a_watchdog_probe_that_never_finishes_blocks_the_reset(self):
|
||||
self.script(self.root / "bin/systemctl", f'echo "systemctl $*" >> {self.calls}; [ "$1" != show ] || echo activating')
|
||||
self.script(self.root / "bin/sleep", "exit 0")
|
||||
self.assertEqual(self.run_wrapper("12").returncode, 4)
|
||||
self.assertNotIn("systemctl stop flysim.service", self.log())
|
||||
self.assertIn("systemctl start fly-watchdog.timer", self.log())
|
||||
|
||||
def test_bad_arguments_are_refused(self):
|
||||
for args in ((), ("--check", "12"), ("12", "13"), ("../12",), ("123",), ("-1",)):
|
||||
self.assertEqual(self.run_wrapper(*args).returncode, 2, args)
|
||||
self.assertEqual(self.log(), "")
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
|
|
|
|||
|
|
@ -1,8 +1,13 @@
|
|||
[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: the rung list (<=2 min), up to 90 s of model calls, a 60 s on-stream countdown,
|
||||
# the reset (<=8 min) and 4 min to verify. systemd must never kill the wrapper mid-reset.
|
||||
TimeoutStartSec=25min
|
||||
ExecStart=/opt/fly/bin/fly-loop-recover
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue