#!/usr/bin/env python3 """Resume Claude Code background sessions that stalled on a recoverable error. Claude Code blocks and exits when it hits a usage limit, an API outage, or a network failure. The process is gone, so nothing scheduled inside the session can revive it. This runs outside every session and resumes the ones whose cause has cleared. """ import argparse import fcntl import json import os import re import shutil import subprocess import sys from datetime import datetime, timedelta from pathlib import Path from zoneinfo import ZoneInfo # Where `claude` installs itself, and what a scheduler's minimal PATH omits. CLAUDE_FALLBACK_DIRS = ("~/.local/bin", "/usr/local/bin", "/opt/homebrew/bin") MEMINFO = Path("/proc/meminfo") DEFAULT_PROMPT = ( "Continue where you left off. Your previous run was interrupted by a " "transient failure (usage limit, API outage, or network), which has cleared." ) DEFAULT_PREFIX = "X " # Every quota Claude Code can report as "You've hit your · ", including the ones that cost money rather than time. QUOTA_BLOCK = re.compile( r"(session|weekly|opus|sonnet|fable \d+|fast|spend|rate.?|usage( credit)?) ?limit" r"|out of (usage|extra usage)" r"|usage credits", re.IGNORECASE, ) TRANSIENT_BLOCK = re.compile( r"api (error|unavailable)" r"|overloaded" r"|connection (error|reset|closed|failure)" r"|network error" r"|fetch failed" r"|socket hang up" r"|etimedout|econnreset|econnrefused|enotfound", re.IGNORECASE, ) # A session parked on a human decision must never be auto-resumed, and an # entitlement someone else controls never clears by waiting. Vetoes even when # the same detail also names a limit. NEEDS_HUMAN = re.compile( r"input needed|waiting for|permission|approval" r"|awaiting (your )?(direction|decision|input|answer|go.?ahead)" r"|needs? (your )?(direction|decision|input)" r"|seat type does ?n.t include" r"|disabled by your admin" r"|disabled for your org" r"|limit is set to \$0", re.IGNORECASE, ) # Money, not time: cleared by raising the cap, by the month rolling over, or by # the subscription window putting the session back on included quota. None of # those announce themselves, so poll rather than wait out a window. SPEND_BLOCK = re.compile( r"spend limit|usage credit|out of (usage|extra usage)|usage credits", re.IGNORECASE ) WEEKLY_LIMIT = re.compile(r"weekly limit", re.IGNORECASE) SESSION_WINDOW = timedelta(hours=5) SPEND_WINDOW = timedelta(hours=1) WEEKLY_WINDOW = timedelta(days=7) RESETS_AT = re.compile( r"resets\s+(?:(?Pmon|tue|wed|thu|fri|sat|sun)[a-z]*\s+)?" r"(?P\d{1,2})(?::(?P\d{2}))?\s*(?Pam|pm)?" r"(?:\s*\((?P[A-Za-z_]+/[A-Za-z_]+)\))?", re.IGNORECASE, ) WEEKDAYS = {"mon": 0, "tue": 1, "wed": 2, "thu": 3, "fri": 4, "sat": 5, "sun": 6} def log(message): print(f"{datetime.now().astimezone().isoformat(timespec='seconds')} {message}", flush=True) def memory_strain(min_available_mb, max_swap_used_pct, meminfo=MEMINFO): """Why the host cannot take another session, or None when it can. Also None where the host does not report memory this way: a guard that cannot read the numbers must not be the thing that stops recovery. """ try: lines = meminfo.read_text().splitlines() except OSError: return None fields = {} for line in lines: key, _, rest = line.partition(":") amount = rest.split() if amount and amount[0].isdigit(): fields[key] = int(amount[0]) if "MemAvailable" not in fields: return None available_mb = fields["MemAvailable"] // 1024 if available_mb < min_available_mb: return f"{available_mb} MB available, need {min_available_mb}" swap_total = fields.get("SwapTotal", 0) if swap_total: used_pct = round(100 * (swap_total - fields.get("SwapFree", 0)) / swap_total) if used_pct > max_swap_used_pct: return f"swap {used_pct}% used, max {max_swap_used_pct}%" return None def resolve_claude(name): """Absolute path to the Claude Code binary, or None if it cannot be found.""" found = shutil.which(name) if found: return found for directory in CLAUDE_FALLBACK_DIRS: candidate = Path(directory).expanduser() / name if candidate.is_file() and os.access(candidate, os.X_OK): return str(candidate) return None def single_instance(path): """Exclusive lock for this run, or None when another run already holds it. Overlapping runs both read the ledger before either writes it, so both resume the same session — two agents racing on one git worktree. """ path.parent.mkdir(parents=True, exist_ok=True) handle = path.open("w") try: fcntl.flock(handle, fcntl.LOCK_EX | fcntl.LOCK_NB) except OSError: handle.close() return None return handle def active_sessions(claude_bin): """Sessions Claude Code still considers active, or None if they cannot be read.""" try: done = subprocess.run( [claude_bin, "agents", "--json"], capture_output=True, text=True, timeout=120, check=True, ) except (OSError, subprocess.SubprocessError) as err: log(f"error: could not list sessions: {err}") return None try: return json.loads(done.stdout) except json.JSONDecodeError as err: log(f"error: unparseable session list: {err}") return None def pid_alive(pid): if not isinstance(pid, int): return False try: os.kill(pid, 0) except ProcessLookupError: return False except PermissionError: return True return True def parse_reset(detail, blocked_at): """The reset moment named in `detail`, or None. `blocked_at` anchors day-less times.""" found = RESETS_AT.search(detail) if not found: return None hour = int(found.group("hour")) minute = int(found.group("minute") or 0) meridiem = (found.group("meridiem") or "").lower() if meridiem == "pm" and hour != 12: hour += 12 elif meridiem == "am" and hour == 12: hour = 0 if hour > 23 or minute > 59: return None zone = None if found.group("tz"): try: zone = ZoneInfo(found.group("tz")) except Exception: zone = None anchor = blocked_at.astimezone(zone) if zone else blocked_at target = anchor.replace(hour=hour, minute=minute, second=0, microsecond=0) day = (found.group("dow") or "").lower()[:3] if day in WEEKDAYS: target += timedelta(days=(WEEKDAYS[day] - target.weekday()) % 7) if target < anchor: target += timedelta(days=7) elif target < anchor: target += timedelta(days=1) return target def resumable(detail): """True when `detail` names a cause that clears without anyone acting.""" if NEEDS_HUMAN.search(detail): return False return bool(QUOTA_BLOCK.search(detail) or TRANSIENT_BLOCK.search(detail)) def reset_moment(detail, blocked_at): """When the block should have lifted, or None if it can be retried at once. A stated time is clamped to the limit window: state.json is sometimes rewritten after the reset already happened, which would otherwise push the parsed time-of-day a full day into the future. """ if not QUOTA_BLOCK.search(detail): return None if WEEKLY_LIMIT.search(detail): window = WEEKLY_WINDOW elif SPEND_BLOCK.search(detail): window = SPEND_WINDOW else: window = SESSION_WINDOW stated = parse_reset(detail, blocked_at) bound = blocked_at + window return min(stated, bound) if stated else bound def load_ledger(path): try: ledger = json.loads(path.read_text()) except (OSError, json.JSONDecodeError): ledger = {} ledger.setdefault("attempts", {}) ledger.setdefault("resumes", []) return ledger def save_ledger(path, ledger): path.parent.mkdir(parents=True, exist_ok=True) path.write_text(json.dumps(ledger, indent=2, sort_keys=True)) def candidates(sessions, jobs_dir, now): """Stalled background sessions, oldest block first.""" found = [] for session in sessions: job_id, session_id = session.get("id"), session.get("sessionId") if session.get("kind") != "background" or not job_id or not session_id: continue if pid_alive(session.get("pid")): continue state_file = jobs_dir / job_id / "state.json" try: state = json.loads(state_file.read_text()) except (OSError, json.JSONDecodeError): continue if state.get("state") != "blocked": continue detail = str(state.get("detail") or "") if not resumable(detail): continue blocked_at = datetime.fromtimestamp(state_file.stat().st_mtime).astimezone() reset_at = reset_moment(detail, blocked_at) if reset_at and reset_at > now: log(f"{job_id}: waiting for reset at {reset_at.isoformat(timespec='minutes')}") continue found.append({ "cwd": session.get("cwd"), "detail": detail, "id": job_id, "blocked_at": blocked_at, "name": session.get("name"), "session_id": session_id, "state_file": state_file, }) return sorted(found, key=lambda job: job["blocked_at"]) def mark_superseded(state_file, prefix): """Prefix the husk's name so it reads apart from its successor, or say why not. `claude agents` takes the name from state.json and offers no command to set it, so the file is the only seam. """ try: state = json.loads(state_file.read_text()) except (OSError, json.JSONDecodeError) as err: return str(err) name = state.get("name") if not isinstance(name, str) or name.startswith(prefix): return None state["name"] = prefix + name flags = state.get("respawnFlags") if isinstance(flags, list) and "--name" in flags: at = flags.index("--name") + 1 if at < len(flags): flags[at] = state["name"] staged = state_file.with_name(state_file.name + ".new") try: staged.write_text(json.dumps(state, indent=2)) staged.replace(state_file) except OSError as err: return str(err) return None def resume(job, claude_bin, prompt, prefix, dry_run): cwd = job["cwd"] if not cwd or not Path(cwd).is_dir(): log(f"{job['id']}: skipped, working directory is gone ({cwd})") return False if dry_run: log(f"{job['id']}: would resume '{job['name']}' in {cwd} — {job['detail']}") return False command = [claude_bin, "--bg", "--resume", job["session_id"]] if job["name"]: command += ["--name", job["name"]] try: done = subprocess.run( command + [prompt], capture_output=True, cwd=cwd, text=True, timeout=180, check=True, ) except subprocess.TimeoutExpired: # The launch may have registered before we gave up, and a second try # would put two agents on one conversation. log(f"{job['id']}: resume timed out, counting the attempt") return True except (OSError, subprocess.SubprocessError) as err: log(f"{job['id']}: resume failed: {err}") return False log(f"{job['id']}: resumed '{job['name']}' — {job['detail']}") # The husk stays 'blocked' and active forever otherwise, and would be # resumed again into a duplicate session working the same conversation. try: subprocess.run( [claude_bin, "stop", job["id"]], capture_output=True, text=True, timeout=60, check=True, ) except (OSError, subprocess.SubprocessError) as err: log(f"{job['id']}: could not retire stale entry: {err}") problem = mark_superseded(job["state_file"], prefix) if problem: log(f"{job['id']}: could not rename stale entry: {problem}") return True def main(): parser = argparse.ArgumentParser( description=__doc__, formatter_class=argparse.RawDescriptionHelpFormatter ) parser.add_argument("--claude-bin", default=os.environ.get("CLAUDE_BIN", "claude")) parser.add_argument( "--claude-home", type=Path, default=Path(os.environ.get("CLAUDE_HOME", Path.home() / ".claude")), ) parser.add_argument("--dry-run", action="store_true", help="report candidates, resume nothing") parser.add_argument( "--max-attempts", type=int, default=int(os.environ.get("WATCHDOG_MAX_ATTEMPTS", "1")), help="resumes of any one session, ever (default: 1; its successor is " "eligible on its own if it stalls again)", ) parser.add_argument( "--max-per-day", type=int, default=int(os.environ.get("WATCHDOG_MAX_PER_DAY", "50")), help="runaway backstop across all sessions per 24h (default: 50)", ) parser.add_argument( "--max-per-run", type=int, default=int(os.environ.get("WATCHDOG_MAX_PER_RUN", "0")), help="resumes per invocation, 0 for every eligible session (default: 0)", ) parser.add_argument( "--max-swap-used-pct", type=int, default=int(os.environ.get("WATCHDOG_MAX_SWAP_USED_PCT", "50")), help="swap fill above which the host is already trading, so no resume (default: 50)", ) parser.add_argument( "--min-available-mb", type=int, default=int(os.environ.get("WATCHDOG_MIN_AVAILABLE_MB", "2048")), help="memory a resume needs the host to have spare, in MB (default: 2048)", ) parser.add_argument("--prompt", default=os.environ.get("WATCHDOG_PROMPT", DEFAULT_PROMPT)) parser.add_argument( "--superseded-prefix", default=os.environ.get("WATCHDOG_SUPERSEDED_PREFIX", DEFAULT_PREFIX), help="prepended to the name of a husk once its successor runs " f"(default: {DEFAULT_PREFIX!r})", ) args = parser.parse_args() claude_bin = resolve_claude(args.claude_bin) if not claude_bin: log(f"error: '{args.claude_bin}' is not in PATH or " f"{', '.join(CLAUDE_FALLBACK_DIRS)} — resuming nothing") return 1 state_dir = args.claude_home / "session-watchdog" lock = None if not args.dry_run: lock = single_instance(state_dir / ".lock") # released when main returns if lock is None: log("another run is still going, skipping") return 0 now = datetime.now().astimezone() ledger_path = state_dir / "ledger.json" ledger = load_ledger(ledger_path) recent = [at for at in ledger["resumes"] if now.timestamp() - at < 86400] budget = args.max_per_day - len(recent) if args.max_per_run > 0: budget = min(budget, args.max_per_run) if budget <= 0: log(f"circuit breaker: {len(recent)} resumes in the last 24h (max {args.max_per_day})") return 0 strain = memory_strain(args.min_available_mb, args.max_swap_used_pct) if strain: log(f"holding off: {strain}") return 0 sessions = active_sessions(claude_bin) if sessions is None: return 1 resumed = 0 for job in candidates(sessions, args.claude_home / "jobs", now): if resumed >= budget: break tried = ledger["attempts"].get(job["session_id"], 0) if tried >= args.max_attempts: continue # A session launched seconds ago has not grown into its memory yet, so # charge every resume of this run for the headroom it is about to take. strain = memory_strain(args.min_available_mb * (resumed + 1), args.max_swap_used_pct) if strain: log(f"holding off: {strain}") break if not resume(job, claude_bin, args.prompt, args.superseded_prefix, args.dry_run): continue ledger["attempts"][job["session_id"]] = tried + 1 ledger["resumes"] = recent + [now.timestamp()] recent = ledger["resumes"] save_ledger(ledger_path, ledger) resumed += 1 if not resumed: log("nothing to resume") return 0 if __name__ == "__main__": sys.exit(main())