back to scripts

agent-manager.py

python 705 lines secrets redacted

The L2 watchdog — performance review, gap analysis and job-description enforcement across the whole agent fleet.

Note Live script from my home-lab server. Tokens, IDs, phone numbers and other secrets have been replaced with placeholders like <WHATSAPP_GROUP_ID> — everything else is the real, running code.
#!/usr/bin/env python3
"""Agent Manager — L2 watchdog with PDR, gap analysis, and job-description enforcement."""
import datetime
import hashlib
import json
import os
import re
import subprocess
import time
import urllib.request

REPORT_DIR = "/home/lanky/reports"
STATE_DIR = "/var/lib/lanky-agent-manager"
PROMPTS_DIR = "/home/codex/agents/prompts"
LATEST = os.path.join(REPORT_DIR, "agent-manager-latest.txt")
PDR_DIR = os.path.join(REPORT_DIR, "pdr")
HISTORY_FILE = os.path.join(STATE_DIR, "run-history.json")
STATE_FILE = os.path.join(STATE_DIR, "last-alert.sha256")
TOOL_INVENTORY_FILE = os.path.join(STATE_DIR, "tool-inventory.json")
TRAINING_PLAN_FILE = os.path.join(STATE_DIR, "training-plan.json")
GROUP_ID = os.environ.get("LANKY_ALERT_GROUP_ID", "<WHATSAPP_GROUP_ID>")

# Agent roster: name → max acceptable report age in seconds
AGENTS = {
    "media":   3  * 60 * 60,
    "docker":  2  * 60 * 60,
    "dns":     8  * 60 * 60,
    "cyber":   2  * 60 * 60,
    "backup":  13 * 60 * 60,
    "network": 1  * 60 * 60,
    "updates": 25 * 60 * 60,
    "nas":     4  * 60 * 60,
    "ssl":     25 * 60 * 60,
}

# IT-best-practice role roster — agents that SHOULD exist
# Each entry: (slug, description, priority)
REQUIRED_ROLES = [
    ("cyber",   "Security monitoring, intrusion detection, firewall audit",       "high"),
    ("media",   "Media stack (Jellyfin/Sonarr/Radarr/qBittorrent) health",        "high"),
    ("docker",  "Container orchestration, Home Assistant stack, disk/GPU",         "high"),
    ("dns",     "Pi-hole DNS/ad-blocking availability",                            "medium"),
    ("backup",  "Backup recency and integrity verification",                       "high"),
    ("network", "Internet, Starlink dish, Google Wifi mesh availability",          "high"),
    ("updates", "APT patch status, Docker image age, security advisory tracking",  "medium"),
    ("nas",     "QNAP NAS capacity, availability, firmware currency",              "medium"),

    ("ssl",     "TLS/SSL certificate expiry monitoring across all services",       "medium"),
    ("ups",     "UPS/power monitoring and graceful shutdown triggering",           "low"),
    ("siem",    "Log aggregation, anomaly detection, SIEM-style alerting",         "low"),
]

JD_REQUIRED_SECTIONS = [
    "## Role",
    "## Tools",
    "## Alert",
    "## Output Format",
]


def run(command: list, timeout: int = 420) -> subprocess.CompletedProcess:
    try:
        return subprocess.run(
            command, text=True,
            stdout=subprocess.PIPE, stderr=subprocess.STDOUT,
            timeout=timeout, check=False,
        )
    except subprocess.TimeoutExpired as exc:
        return subprocess.CompletedProcess(command, 124, f"Timed out: {exc}")


def send(message: str) -> bool:
    result = run(
        ["sudo", "-u", "lanky", "-H", "openclaw", "message", "send",
         "--channel", "whatsapp", "--account", "default",
         "--target", GROUP_ID, "--message", message],
        timeout=120,
    )
    return result.returncode == 0


def valid_report(path: str) -> tuple:
    try:
        text = open(path, encoding="utf-8", errors="replace").read()
    except OSError as exc:
        return False, f"report unavailable: {exc}"
    required = ("Status", "Actions taken or recommended", "Risks / escalation", "Continuous improvement", "Next check")
    missing = [h for h in required if h.lower() not in text.lower()]
    if len(text) < 120:
        return False, "report is too short"
    if missing:
        return False, "missing headings: " + ", ".join(missing)
    if "FAILED:" in text:
        return False, "report contains a failed remediation/check"
    has_evidence = any(m in text for m in ("FIXED:", "FIXED/", "CHECKED:", "MONITORING:", "SKIPPED:", "NO ACTION:"))
    if "Actions taken or recommended" in text and not has_evidence:
        return False, "report has no remediation or monitoring evidence"
    if "Continuous improvement" in text and not all(
        marker in text
        for marker in ("Missed or weak signal:", "Improvement made:", "Improvement needed:")
    ):
        return False, "continuous improvement section is incomplete"
    return True, "valid"


def score_report(path: str) -> dict:
    """Return a quality score dict for PDR purposes."""
    score = {"freshness": 0, "completeness": 0, "actions": 0, "total": 0, "detail": ""}
    try:
        age = time.time() - os.path.getmtime(path)
        score["freshness"] = 2 if age < 3600 else 1 if age < 14400 else 0
    except OSError:
        score["detail"] = "report file missing"
        return score

    valid, detail = valid_report(path)
    score["completeness"] = 2 if valid else 1 if "too short" not in detail else 0
    score["detail"] = detail

    try:
        text = open(path, encoding="utf-8", errors="replace").read()
        fixed = len(re.findall(r"FIXED", text))
        checked = len(re.findall(r"CHECKED:", text))
        critical = len(re.findall(r"CRITICAL", text))
        score["actions"] = 2 if (fixed + checked) > 0 else 1
        score["critical_count"] = critical
    except OSError:
        pass

    score["total"] = score["freshness"] + score["completeness"] + score["actions"]
    return score


def validate_jd(agent: str) -> tuple:
    """Validate that the agent's job description exists and has required sections."""
    path = os.path.join(PROMPTS_DIR, f"{agent}.md")
    if not os.path.isfile(path):
        return False, f"Job description missing: {path}"
    try:
        text = open(path, encoding="utf-8").read()
    except OSError as exc:
        return False, f"JD unreadable: {exc}"
    missing = [s for s in JD_REQUIRED_SECTIONS if s.lower() not in text.lower()]
    if missing:
        return False, "JD missing sections: " + ", ".join(missing)
    return True, "JD valid"


def load_history() -> dict:
    try:
        return json.load(open(HISTORY_FILE, encoding="utf-8"))
    except (OSError, ValueError):
        return {}


def save_history(history: dict) -> None:
    os.makedirs(STATE_DIR, exist_ok=True)
    with open(HISTORY_FILE, "w", encoding="utf-8") as f:
        json.dump(history, f, indent=2)



def load_json_file(path: str, default):
    try:
        return json.load(open(path, encoding="utf-8"))
    except (OSError, ValueError):
        return default


def save_json_file(path: str, value) -> None:
    os.makedirs(os.path.dirname(path), exist_ok=True)
    tmp = path + ".tmp"
    with open(tmp, "w", encoding="utf-8") as handle:
        json.dump(value, handle, indent=2, ensure_ascii=True)
        handle.write("\n")
    os.replace(tmp, path)


def ignored_capability_name(name: str) -> bool:
    lower = name.lower()
    ignored_bits = (".bak", "bak-", "backup-", ".old", ".tmp", ".disabled", "~")
    return lower.startswith("__") or any(bit in lower for bit in ignored_bits)


def infer_training_owner(name: str, path: str) -> str:
    haystack = (name + " " + path).lower()
    rules = [
        ("media", ("media", "jelly", "sonarr", "radarr", "prowlarr", "torrent", "arr-")),
        ("cyber", ("cyber", "cve", "av", "clam", "security", "hardening")),
        ("docker", ("docker", "container", "watchdog", "gpu", "home-assistant", "ha-", "govee")),
        ("nas", ("nas", "duplicate", "junk", "mount", "storage")),
        ("backup", ("backup",)),
        ("dns", ("dns", "pihole", "pi-hole")),
        ("network", ("network", "ssl", "cert", "telegram", "google-oauth")),
        ("updates", ("update", "apt", "trivy")),
    ]
    for owner, needles in rules:
        if any(needle in haystack for needle in needles):
            return owner
    return "manager"


def scan_tool_capabilities() -> dict:
    roots = ["/home/lanky/scripts", "/home/codex/agents/bin"]
    scripts = []
    for root in roots:
        try:
            names = sorted(os.listdir(root))
        except OSError:
            continue
        for name in names:
            if ignored_capability_name(name):
                continue
            path = os.path.join(root, name)
            try:
                st = os.stat(path)
            except OSError:
                continue
            if not os.path.isfile(path) or not os.access(path, os.X_OK):
                continue
            scripts.append({
                "name": name,
                "path": path,
                "owner_hint": infer_training_owner(name, path),
                "mtime": int(st.st_mtime),
                "size": int(st.st_size),
            })

    prompts = []
    try:
        prompt_names = sorted(os.listdir(PROMPTS_DIR))
    except OSError:
        prompt_names = []
    for name in prompt_names:
        if ignored_capability_name(name):
            continue
        path = os.path.join(PROMPTS_DIR, name)
        if not os.path.isfile(path):
            continue
        try:
            st = os.stat(path)
        except OSError:
            continue
        prompts.append({
            "name": name,
            "path": path,
            "owner_hint": name.rsplit(".", 1)[0],
            "mtime": int(st.st_mtime),
            "size": int(st.st_size),
        })

    unit_result = run([
        "bash", "-lc",
        "systemctl list-unit-files --no-pager 'lanky-*' | awk 'NR>1 && $1 ~ /^lanky-/ {print $1}' | sort"
    ], timeout=30)
    units = [line.strip() for line in unit_result.stdout.splitlines() if line.strip()]
    return {
        "generated": datetime.datetime.now(datetime.timezone.utc).isoformat(),
        "scripts": scripts,
        "prompts": prompts,
        "systemd_units": units,
    }


def capability_key(kind: str, item) -> str:
    if isinstance(item, str):
        return f"{kind}:{item}"
    return f"{kind}:{item.get('path') or item.get('name')}"


def update_training_plan(inventory: dict) -> dict:
    previous = load_json_file(TOOL_INVENTORY_FILE, {})
    plan = load_json_file(TRAINING_PLAN_FILE, {"items": []})
    if not isinstance(plan, dict):
        plan = {"items": []}
    plan.setdefault("items", [])

    previous_keys = set()
    for kind in ("scripts", "prompts", "systemd_units"):
        for item in previous.get(kind, []):
            previous_keys.add(capability_key(kind, item))

    current_items = []
    current_keys = set()
    for kind in ("scripts", "prompts", "systemd_units"):
        for item in inventory.get(kind, []):
            key = capability_key(kind, item)
            current_keys.add(key)
            current_items.append((kind, key, item))

    existing = {entry.get("key") for entry in plan.get("items", [])}
    new_entries = []
    for kind, key, item in current_items:
        if key in previous_keys or key in existing:
            continue
        name = item if isinstance(item, str) else item.get("name", key)
        path = item if isinstance(item, str) else item.get("path", "")
        owner = "manager" if isinstance(item, str) else item.get("owner_hint", "manager")
        suggested = owner if owner in AGENTS else "manager"
        entry = {
            "key": key,
            "status": "pending_manager_review",
            "detected": inventory["generated"],
            "kind": kind,
            "name": name,
            "path": path,
            "suggested_agent": suggested,
            "training_question": f"Should {suggested} learn or use {name}?",
            "safe_first_step": "Review purpose, permissions, and read-only/report-only use before adding to any agent workflow.",
        }
        plan["items"].append(entry)
        new_entries.append(entry)

    removed = sorted(previous_keys - current_keys)
    plan["last_scan"] = inventory["generated"]
    plan["pending_count"] = sum(1 for entry in plan.get("items", []) if entry.get("status") == "pending_manager_review")
    save_json_file(TOOL_INVENTORY_FILE, inventory)
    save_json_file(TRAINING_PLAN_FILE, plan)
    return {
        "new_entries": new_entries,
        "removed": removed,
        "pending": [entry for entry in plan.get("items", []) if entry.get("status") == "pending_manager_review"],
        "inventory_counts": {
            "scripts": len(inventory.get("scripts", [])),
            "prompts": len(inventory.get("prompts", [])),
            "systemd_units": len(inventory.get("systemd_units", [])),
        },
    }


def collect_learning_requests(agents, report_dir: str) -> list:
    """Collect agent requests to learn new skills/rules/tools before scope changes."""
    pattern = re.compile(
        r"(?im)^\s*(?:[-*]\s*)?MANAGER LEARNING REQUEST:\s*(.+?)\s*$"
    )
    requests = []
    for agent in agents:
        path = os.path.join(report_dir, f"agent-{agent}-latest.md")
        try:
            report_text = open(path, encoding="utf-8", errors="replace").read()
        except OSError:
            continue
        for raw in pattern.findall(report_text):
            parts = [part.strip() for part in raw.split("|")]
            while len(parts) < 4:
                parts.append("")
            topic, reason, first_step, approval = parts[:4]
            requests.append({
                "agent": agent,
                "topic": topic,
                "reason": reason,
                "first_step": first_step,
                "approval": approval,
                "raw": raw.strip(),
            })
    return requests


def run_gap_analysis(lines: list, problems: list) -> list:
    """Compare REQUIRED_ROLES against deployed AGENTS. Return list of gap findings."""
    gaps = []
    deployed = set(AGENTS.keys())
    for slug, description, priority in REQUIRED_ROLES:
        if slug not in deployed:
            gap_msg = f"GAP [{priority.upper()}]: No '{slug}' agent — {description}"
            gaps.append(gap_msg)
            if priority == "high":
                problems.append(gap_msg)
    lines.append(f"gap-analysis: {len(gaps)} role gap(s) identified")
    return gaps


def run_pdr(agent: str, history: dict, lines: list, problems: list) -> str:
    """Generate a PDR entry for the agent. Returns PDR summary string."""
    path = os.path.join(REPORT_DIR, f"agent-{agent}-latest.md")
    score = score_report(path)
    jd_ok, jd_detail = validate_jd(agent)

    # Record this run in history
    now_iso = datetime.datetime.now(datetime.timezone.utc).isoformat()
    if agent not in history:
        history[agent] = {"runs": []}
    history[agent]["runs"].append({
        "ts": now_iso,
        "score": score["total"],
        "detail": score["detail"],
        "jd_ok": jd_ok,
    })
    # Keep last 20 runs
    history[agent]["runs"] = history[agent]["runs"][-20:]

    recent_runs = history[agent]["runs"][-5:]
    avg_score = sum(r["score"] for r in recent_runs) / len(recent_runs) if recent_runs else 0
    consecutive_fails = 0
    for r in reversed(recent_runs):
        if r["score"] < 3:
            consecutive_fails += 1
        else:
            break

    rating = "EXCEEDS" if avg_score >= 5.5 else "MEETS" if avg_score >= 4 else "NEEDS IMPROVEMENT" if avg_score >= 2.5 else "UNSATISFACTORY"
    pdr = f"{agent}: rating={rating} avg_score={avg_score:.1f}/6 jd={jd_detail} last_report={score['detail']}"

    if not jd_ok:
        problems.append(f"PDR: {agent} has no valid job description — {jd_detail}")
    if consecutive_fails >= 3:
        problems.append(f"PDR: {agent} has failed last {consecutive_fails} consecutive checks — consider L3 review")
    if rating == "UNSATISFACTORY":
        problems.append(f"PDR: {agent} is UNSATISFACTORY (avg score {avg_score:.1f}/6)")

    return pdr



OLLAMA_URL = os.environ.get("OLLAMA_URL", "http://localhost:11434")
OLLAMA_MODEL = os.environ.get("OLLAMA_MODEL", "llama3.1:8b")
MANAGER_PROMPT_FILE = os.path.join(PROMPTS_DIR, "manager.md")


def ollama_generate(prompt, max_tokens=2000):
    import json as _j, urllib.request as _u
    payload = _j.dumps({"model": OLLAMA_MODEL, "prompt": prompt, "stream": False,
        "keep_alive": 0,
        "options": {"num_predict": max_tokens, "temperature": 0.1}}).encode()
    req = _u.Request(OLLAMA_URL + "/api/generate", data=payload,
        headers={"Content-Type": "application/json"})
    try:
        with _u.urlopen(req, timeout=360) as r:
            return _j.load(r).get("response", "")
    except Exception as exc:
        return "LLM call failed: " + str(exc)


def llm_manager_review(python_report, agents, report_dir):
    nl = chr(10)
    try:
        mgr_prompt = open(MANAGER_PROMPT_FILE, encoding="utf-8").read()
    except OSError:
        mgr_prompt = "You are the Manager agent overseeing a home server agent fleet."
    parts = [mgr_prompt, nl + nl + "--- PYTHON AUDIT REPORT ---" + nl, python_report, nl + "--- ALL AGENT REPORTS ---"]
    for agent in agents:
        path = os.path.join(report_dir, "agent-" + agent + "-latest.md")
        try:
            text = open(path, encoding="utf-8", errors="replace").read()
        except OSError:
            text = "(report not found)"
        parts.append(nl + "--- " + agent.upper() + " ---" + nl + text)
    parts.append(nl + "--- YOUR TASK ---" + nl
        + "Review as Manager. Be concise. Do not repeat the Python audit." + nl
        + "1. TOOL_REQUEST blocks: decide GRANTED/DENIED/ESCALATED." + nl
        + "2. Flag suspicious reports (OK/no evidence, CHECK_FAILED)." + nl
        + "3. Decide MANAGER LEARNING REQUESTS as APPROVED, DENIED, or ESCALATED. If no learning requests are listed, say LEARNING_DECISIONS: none. Otherwise approve only safe in-scope learning or report-only detector work; escalate scope, policy, destructive, credential, firewall, public exposure, or new-tool installs to L3." + nl
        + "4. Review TOOL CAPABILITY SCAN and TRAINING_PLAN candidates. If TRAINING_PLAN lines are present, do not say none: summarize pending count/themes and decide whether candidates are manager-review, approved report-only learning, denied, or escalated to L3." + nl
        + "5. Note JOB_DESCRIPTION_DRIFT if any." + nl
        + "6. Output: FLEET_STATUS, CRITICAL_ESCALATIONS, TOOL_DECISIONS, LEARNING_DECISIONS, TRAINING_PLAN, COVERAGE_GAPS, SUMMARY." + nl)
    return ollama_generate("".join(parts), max_tokens=2000)

def main() -> int:
    os.makedirs(STATE_DIR, exist_ok=True)
    os.makedirs(REPORT_DIR, exist_ok=True)
    os.makedirs(PDR_DIR, exist_ok=True)
    now = time.time()
    lines = []
    problems = []
    history = load_history()

    lines.append("Agent Manager audit")
    lines.append("Generated: " + time.strftime("%Y-%m-%dT%H:%M:%S%z"))
    lines.append("")

    # ── 1. Freshness + execution check for each agent ─────────────────────────
    lines.append("=== Agent Freshness & Execution ===")
    for agent, max_age in AGENTS.items():
        unit = f"lanky-agent@{agent}.service"
        path = os.path.join(REPORT_DIR, f"agent-{agent}-latest.md")
        try:
            age = now - os.path.getmtime(path)
        except OSError:
            age = max_age + 1

        result = run(["systemctl", "show", "-p", "Result", "--value", unit], timeout=30)
        unit_result = result.stdout.strip() or "unknown"
        needs_run = age > max_age or unit_result not in ("success", "unknown")
        action = "none"

        if needs_run:
            rerun = run(["systemctl", "start", unit])
            action = "reran agent"
            if rerun.returncode:
                problems.append(f"{agent}: failed to start: {rerun.stdout.strip()}")
            try:
                age = time.time() - os.path.getmtime(path)
            except OSError:
                age = max_age + 1

        valid, detail = valid_report(path)
        if not valid and not needs_run:
            rerun = run(["systemctl", "start", unit])
            action = "reran invalid agent"
            if rerun.returncode:
                problems.append(f"{agent}: failed to rerun: {rerun.stdout.strip()}")
            valid, detail = valid_report(path)
        if not valid:
            problems.append(f"{agent}: {detail}")

        lines.append(f"  {agent}: age={int(age)}s result={unit_result} report={detail} action={action}")

    # ── 2. PDR for all agents ──────────────────────────────────────────────────
    lines.append("")
    lines.append("=== Performance & Development Review (PDR) ===")
    for agent in AGENTS:
        pdr_line = run_pdr(agent, history, lines, problems)
        lines.append("  " + pdr_line)
    save_history(history)

    # ── 3. Job description audit ───────────────────────────────────────────────
    lines.append("")
    lines.append("=== Job Description Audit ===")
    for agent in AGENTS:
        jd_ok, jd_detail = validate_jd(agent)
        lines.append(f"  {agent}: {jd_detail}")
        if not jd_ok:
            problems.append(f"JD missing or incomplete for {agent}")

    # ── 4. Gap analysis ────────────────────────────────────────────────────────
    lines.append("")
    lines.append("=== IT Role Gap Analysis ===")
    gaps = run_gap_analysis(lines, problems)
    for g in gaps:
        lines.append("  " + g)
    if not gaps:
        lines.append("  All expected IT roles have a deployed agent.")

    # -- 4b. Learning requests -------------------------------------------------
    lines.append("")
    lines.append("=== Manager Learning Requests ===")
    learning_requests = collect_learning_requests(AGENTS, REPORT_DIR)
    if learning_requests:
        for item in learning_requests:
            lines.append(
                "  {agent}: topic={topic} | reason={reason} | first_step={first_step} | approval={approval}".format(**item)
            )
    else:
        lines.append("  No agents requested new learning or scope changes this run.")

    # -- 4c. Tool capability scan and training plan -----------------------------
    lines.append("")
    lines.append("=== Tool Capability Scan ===")
    capability_scan = update_training_plan(scan_tool_capabilities())
    counts = capability_scan["inventory_counts"]
    lines.append(
        f"  Inventory: scripts={counts['scripts']} prompts={counts['prompts']} systemd_units={counts['systemd_units']}"
    )
    if capability_scan["new_entries"]:
        lines.append("  New capabilities detected:")
        for item in capability_scan["new_entries"][:25]:
            lines.append(
                f"  - {item['kind']} {item['name']} -> suggested_agent={item['suggested_agent']} | {item['safe_first_step']}"
            )
        if len(capability_scan["new_entries"]) > 25:
            lines.append(f"  - ... {len(capability_scan['new_entries']) - 25} more new capabilities")
    else:
        lines.append("  No new capabilities detected since the previous manager inventory.")
    if capability_scan["removed"]:
        lines.append(f"  Removed capabilities since last scan: {len(capability_scan['removed'])}")
    pending = capability_scan["pending"]
    lines.append(f"  Training plan pending review: {len(pending)} item(s)")
    for item in pending[:25]:
        lines.append(
            f"  TRAINING_PLAN: {item['suggested_agent']} | {item['name']} | {item['training_question']} | {item['safe_first_step']}"
        )
    if len(pending) > 25:
        lines.append(f"  TRAINING_PLAN: ... {len(pending) - 25} more pending item(s); see {TRAINING_PLAN_FILE}")

    # ── 5. Cyber watch service ─────────────────────────────────────────────────
    lines.append("")
    lines.append("=== Supporting Services ===")
    cyber_watch = run(
        ["systemctl", "show", "-p", "Result", "--value", "lanky-cyber-watch.service"],
        timeout=30,
    ).stdout.strip() or "unknown"
    if cyber_watch not in ("success", "unknown"):
        problems.append(f"cyber-watch: last service result was {cyber_watch}")
    lines.append(f"  cyber-watch: result={cyber_watch}")

    # ── 6. Agent handoffs ──────────────────────────────────────────────────────
    for hf_path in ["/var/lib/lanky-agent-handoff/active.json",
                    "/var/lib/lanky-cyber-handoff/active.json"]:
        try:
            handoffs = json.load(open(hf_path, encoding="utf-8"))
        except (OSError, ValueError):
            handoffs = {}
        if isinstance(handoffs, list):
            handoffs = {h["id"]: h for h in handoffs if isinstance(h, dict)}
        for item in handoffs.values():
            status = item.get("status", "unknown")
            lines.append(
                f"  agent-handoff {item.get('id')}: source={item.get('source')} "
                f"target={item.get('target')} status={status}"
            )
            if status in {"dispatch_failed", "not_acknowledged", "chain_limit", "unknown"}:
                problems.append(
                    f"Agent handoff {item.get('id')} from {item.get('source')} "
                    f"to {item.get('target')} is {status}"
                )

    # ── 7. L3 escalation queue ─────────────────────────────────────────────────
    l3_sync = run(["/home/lanky/scripts/l3-escalation.py", "sync"], timeout=150)
    if l3_sync.returncode:
        problems.append("L3 escalation queue or WhatsApp notification failed")
    l3_queue = "/var/lib/lanky-l3-escalations/active.json"
    try:
        l3_items = json.load(open(l3_queue, encoding="utf-8"))
    except (OSError, ValueError):
        l3_items = {}
    for item in l3_items.values():
        status = item.get("status", "unknown")
        lines.append(f"  l3-change {item.get('id')}: status={status} owner={item.get('owner')}")
        if status in {"approved", "decision_received"}:
            problems.append(
                f"L3 change {item.get('id')} has authorization but still needs implementation"
            )

    # ── 8. Service desk jobs ────────────────────────────────────────────────────
    service_desk_jobs = os.path.join(REPORT_DIR, "service-desk-jobs.json")
    try:
        jobs = json.load(open(service_desk_jobs, encoding="utf-8"))
    except (OSError, ValueError):
        jobs = {}
    current_time = datetime.datetime.now(datetime.timezone.utc)
    for job in jobs.values():
        job_id = job.get("id")
        status = job.get("status", "unknown")
        lines.append(f"  service-desk job {job_id}: status={status} owner={job.get('agent')}")
        if status in {"queued", "running"}:
            try:
                eta = datetime.datetime.fromisoformat(job["eta"])
            except (KeyError, ValueError):
                problems.append(f"Service Desk job {job_id} has no valid ETA")
            else:
                if current_time > eta:
                    problems.append(f"Service Desk job {job_id} is overdue: {job.get('request')}")
        if status in {"completed", "failed"} and job.get("notification_status") == "failed":
            problems.append(
                f"Service Desk job {job_id} completed but WhatsApp notification failed"
            )

    # ── 9. Home Assistant integration status ────────────────────────────────────
    lines.append("")
    lines.append("=== Home Assistant API Status ===")
    ha_token_file = "/home/codex/agents/ha-token"
    ha_url = "http://localhost:8123/api/"
    if os.path.isfile(ha_token_file):
        try:
            token = open(ha_token_file).read().strip()
            import urllib.request as _ur
            req = _ur.Request(ha_url)
            req.add_header("Authorization", "Bearer " + token)
            resp = _ur.urlopen(req, timeout=5).read().decode()
            lines.append("  HA REST API: connected, token valid.")
        except Exception as e:
            lines.append(f"  HA REST API: token file exists but API failed — {str(e)[:60]}")
            problems.append("Home Assistant API token is configured but API is unreachable")
    else:
        lines.append("  HA REST API: no token configured at /home/codex/agents/ha-token (optional).")

    # ── Final report ────────────────────────────────────────────────────────────
    lines.append("")
    if problems:
        lines.append("=== Problems ===")
        for p in problems:
            lines.append("- " + p)
    else:
        lines.append("All agents passed freshness, PDR, JD, gap, and remediation checks.")

    report = "\n".join(lines) + "\n"
    open(LATEST, "w", encoding="utf-8").write(report)

    finding = "\n".join(problems)
    digest = hashlib.sha256(finding.encode()).hexdigest() if finding else ""
    try:
        previous = open(STATE_FILE, encoding="ascii").read().strip()
    except OSError:
        previous = ""

    if finding and digest != previous:
        if send("AGENT MANAGER ALERT: " + str(len(problems)) + " issue(s) found.\n\n" + finding):
            open(STATE_FILE, "w", encoding="ascii").write(digest + "\n")
    elif not finding and previous:
        if send("AGENT MANAGER RESOLVED: all area agents are fresh and passing their checks."):
            open(STATE_FILE, "w", encoding="ascii").write("")


    # LLM Manager Review
    llm_out = llm_manager_review(report, AGENTS, REPORT_DIR)
    if llm_out and not llm_out.startswith("LLM call failed"):
        report += chr(10) + "=== LLM Manager Review ===" + chr(10) + llm_out + chr(10)
        open(LATEST, "w", encoding="utf-8").write(report)

    print(report, end="")
    return 1 if problems else 0


if __name__ == "__main__":
    raise SystemExit(main())

back to scripts