feat(operators): amendment workflow and automated drive watchdog daemon
This commit is contained in:
Executable
+200
@@ -0,0 +1,200 @@
|
||||
#!/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()
|
||||
Reference in New Issue
Block a user