Files
box/bin/agent-drive-watchdog.py
T

201 lines
7.8 KiB
Python
Raw Normal View History

#!/usr/bin/env python3
"""
agent-drive-watchdog.py — Automated Drive Watchdog & Healing Daemon for Muse Agents.
Monitors operational DRIVE scores across the agent fleet (muse, pip, 646, opm, def, dev)
via Hatch WebSocket RPC. If any agent's DRIVE score drops below 100 or critical drive
files (HEARTBEAT.md, PROACTIVE_PREFERENCES.md, SOUL.md) are degraded or missing:
1. Detects degraded state and missing checklist/preferences.
2. Selectively auto-heals core drive files using canonical shared templates.
3. Preserves MEMORY.md and agent-generated workspace files.
4. Records state & healing history to /tmp/agent-drive-watchdog.json.
5. Emits structured telemetry to stdout/journal.
Can be run:
- Once: python3 bin/agent-drive-watchdog.py --once
- Continuous loop: python3 bin/agent-drive-watchdog.py --interval 600
- Via systemd timer: agent-drive-watchdog.timer (every 10m)
"""
import argparse
import datetime
import json
import logging
import os
import sys
import time
from pathlib import Path
# Add NetVM bin to path
BASE_DIR = Path(__file__).resolve().parent.parent
sys.path.insert(0, str(BASE_DIR / "bin"))
from agent_md import audit_agents, write_md, read_md, SHARED_OPERATORS, VALID_ACCOUNTS
STATE_FILE = Path("/tmp/agent-drive-watchdog.json")
LOG_FILE = Path("/tmp/agent-drive-watchdog.log")
logging.basicConfig(
level=logging.INFO,
format="%(asctime)s [%(levelname)s] %(message)s",
handlers=[
logging.StreamHandler(sys.stdout),
logging.FileHandler(LOG_FILE, mode="a", encoding="utf-8")
]
)
def load_state() -> dict:
if STATE_FILE.exists():
try:
return json.loads(STATE_FILE.read_text(encoding="utf-8"))
except Exception as e:
logging.warning(f"Failed to read existing state file: {e}")
return {
"last_run": None,
"history": [],
"agent_status": {},
}
def save_state(state: dict):
try:
# Keep last 50 history entries
if len(state.get("history", [])) > 50:
state["history"] = state["history"][-50:]
STATE_FILE.write_text(json.dumps(state, indent=2), encoding="utf-8")
except Exception as e:
logging.error(f"Failed to save state file: {e}")
def heal_agent_drive(account: str, audit_data: dict) -> dict:
"""Selectively auto-heal core drive files for an agent."""
healed = []
issues = audit_data.get("issues", [])
# Check what specifically needs healing
need_soul = any("SOUL.md" in iss for iss in issues) or audit_data.get("files", {}).get("SOUL.md", {}).get("size", 0) < 1000
need_pro = any("PROACTIVE_PREFERENCES.md" in iss for iss in issues) or audit_data.get("files", {}).get("PROACTIVE_PREFERENCES.md", {}).get("size", 0) < 800
need_hb = any("HEARTBEAT.md" in iss for iss in issues) or audit_data.get("files", {}).get("HEARTBEAT.md", {}).get("size", 0) < 300
need_context = any("TOOLS.md or USER.md" in iss for iss in issues)
logging.info(f"[{account.upper()}] Auto-healing drive... need_soul={need_soul}, need_pro={need_pro}, need_hb={need_hb}, need_context={need_context}")
if need_soul:
soul_text = (SHARED_OPERATORS / "SOUL.md").read_text(encoding="utf-8")
res = write_md(account, "SOUL.md", soul_text, overwrite=True)
healed.append({"file": "SOUL.md", "bytes": res.get("bytes_written")})
if need_pro:
pro_text = (SHARED_OPERATORS / "PROACTIVE_PREFERENCES.md").read_text(encoding="utf-8")
res = write_md(account, "PROACTIVE_PREFERENCES.md", pro_text, overwrite=True)
healed.append({"file": "PROACTIVE_PREFERENCES.md", "bytes": res.get("bytes_written")})
if need_hb:
hb_text = (SHARED_OPERATORS / "HEARTBEAT.md").read_text(encoding="utf-8")
res = write_md(account, "HEARTBEAT.md", hb_text, overwrite=True)
healed.append({"file": "HEARTBEAT.md", "bytes": res.get("bytes_written")})
if need_context:
tools_text = (SHARED_OPERATORS / "TOOLS.md").read_text(encoding="utf-8")
res_t = write_md(account, "TOOLS.md", tools_text, overwrite=True)
healed.append({"file": "TOOLS.md", "bytes": res_t.get("bytes_written")})
user_text = (SHARED_OPERATORS / "USER.md").read_text(encoding="utf-8")
res_u = write_md(account, "USER.md", user_text, overwrite=True)
healed.append({"file": "USER.md", "bytes": res_u.get("bytes_written")})
return {"ok": True, "account": account, "healed": healed}
def run_cycle(auto_heal: bool = True) -> dict:
"""Run an audit and healing cycle across all fleet accounts."""
now_iso = datetime.datetime.now(datetime.timezone.utc).isoformat()
logging.info("Starting fleet drive watchdog audit cycle...")
state = load_state()
state["last_run"] = now_iso
try:
audit_results = audit_agents()
except Exception as e:
logging.error(f"Audit failed during cycle: {e}")
return {"ok": False, "error": str(e)}
cycle_report = {
"timestamp": now_iso,
"total_agents": len(audit_results),
"high_drive": 0,
"degraded": 0,
"healed_agents": [],
}
for account, a_data in audit_results.items():
score = a_data.get("drive_score", 0)
issues = a_data.get("issues", [])
status = "HIGH_DRIVE" if score == 100 else "DEGRADED"
if status == "HIGH_DRIVE":
cycle_report["high_drive"] += 1
logging.info(f"Agent {account.upper():6}: DRIVE score 100/100 (HIGH DRIVE)")
else:
cycle_report["degraded"] += 1
logging.warning(f"Agent {account.upper():6}: DRIVE score {score}/100 ({status}) - Issues: {', '.join(issues)}")
if auto_heal:
try:
heal_res = heal_agent_drive(account, a_data)
healed_files = [h["file"] for h in heal_res.get("healed", [])]
logging.info(f"Agent {account.upper():6}: Successfully healed files: {', '.join(healed_files)}")
cycle_report["healed_agents"].append({
"account": account,
"prior_score": score,
"healed_files": healed_files,
})
except Exception as e:
logging.error(f"Agent {account.upper():6}: Healing failed: {e}")
state["agent_status"][account] = {
"score": score,
"status": status,
"issues": issues,
"last_checked": now_iso,
}
state["history"].append(cycle_report)
save_state(state)
logging.info(f"Watchdog cycle complete. High drive: {cycle_report['high_drive']}/{cycle_report['total_agents']}. Degraded: {cycle_report['degraded']}. Healed: {len(cycle_report['healed_agents'])}.")
return cycle_report
def main():
parser = argparse.ArgumentParser(description="NetVM Automated Agent Drive Watchdog & Healing Daemon")
parser.add_argument("--once", action="store_true", help="Run a single audit/healing pass and exit")
parser.add_argument("--no-heal", action="store_true", help="Audit only; do not auto-heal degraded agents")
parser.add_argument("--interval", type=int, default=600, help="Loop interval in seconds (default: 600s / 10m)")
parser.add_argument("--status", action="store_true", help="Print recent watchdog status and exit")
args = parser.parse_args()
if args.status:
state = load_state()
print(json.dumps(state, indent=2))
return
if args.once:
run_cycle(auto_heal=not args.no_heal)
return
logging.info(f"Starting NetVM Agent Drive Watchdog daemon (interval: {args.interval}s)...")
while True:
try:
run_cycle(auto_heal=not args.no_heal)
except Exception as e:
logging.error(f"Unexpected error in watchdog loop: {e}", exc_info=True)
time.sleep(args.interval)
if __name__ == "__main__":
main()