feat(operators): agent markdown drive management via Hatch and SSH runbook
This commit is contained in:
Executable
+430
@@ -0,0 +1,430 @@
|
||||
#!/usr/bin/env python3
|
||||
"""agent_md.py — Access and modify Muse agent .md files via Hatch gateway and SSH.
|
||||
|
||||
Enables operators to inspect, audit, diff, and inject operational DRIVE into
|
||||
Muse agents across the fleet (muse, pip, 646, opm, def, dev).
|
||||
"""
|
||||
|
||||
import difflib
|
||||
import json
|
||||
import os
|
||||
import re
|
||||
import sys
|
||||
import time
|
||||
from pathlib import Path
|
||||
|
||||
# Ensure muse_cli can be imported
|
||||
sys.path.insert(0, os.path.expanduser("~/.local/lib/python3.14/site-packages"))
|
||||
try:
|
||||
from muse_cli.gateway import Gateway, load_cookies, AuthError, GatewayError
|
||||
except ImportError:
|
||||
Gateway = None
|
||||
|
||||
NETVM_ROOT = Path("/home/super/Projects/NetVM")
|
||||
SHARED_OPERATORS = NETVM_ROOT / "shared" / "operators"
|
||||
VALID_ACCOUNTS = ["muse", "pip", "646", "opm", "def", "dev"]
|
||||
|
||||
TARGET_MD_FILES = [
|
||||
"SOUL.md",
|
||||
"PROACTIVE_PREFERENCES.md",
|
||||
"HEARTBEAT.md",
|
||||
"AGENTS.md",
|
||||
"MEMORY.md",
|
||||
"USER.md",
|
||||
"TOOLS.md",
|
||||
"IDENTITY.md",
|
||||
]
|
||||
|
||||
# Tunnel / Port inventory
|
||||
TUNNEL_PORTS = {
|
||||
"muse-main": {"port": 2224, "terminal": 7681, "user": "muse"},
|
||||
"muse": {"port": 2225, "terminal": 7682, "user": "hatch"},
|
||||
"646": {"port": 2226, "terminal": 7683, "user": "hatch"},
|
||||
"pip": {"port": 2227, "terminal": 7684, "user": "hatch"},
|
||||
"opm": {"port": 2228, "terminal": 7685, "user": "hatch"},
|
||||
"def": {"port": 2229, "terminal": 7686, "user": "hatch"},
|
||||
"dev": {"port": 2230, "terminal": 7687, "user": "hatch"},
|
||||
}
|
||||
|
||||
|
||||
def get_gateway(account: str) -> "Gateway":
|
||||
"""Obtain an authenticated Gateway connection for an account."""
|
||||
if not Gateway:
|
||||
raise RuntimeError("muse_cli.gateway module is not available")
|
||||
conf_dir = Path.home() / ".config" / "muse-cli" / account
|
||||
cfile = conf_dir / "cookies.txt"
|
||||
if not cfile.exists():
|
||||
raise FileNotFoundError(f"No cookies found for account '{account}' at {cfile}")
|
||||
cookies = load_cookies(str(cfile))
|
||||
if not cookies.strip():
|
||||
raise ValueError(f"Cookies file for '{account}' is empty")
|
||||
return Gateway(cookies)
|
||||
|
||||
|
||||
def list_files(account: str, path: str = "") -> list:
|
||||
"""List files in the agent container filesystem via Hatch."""
|
||||
gw = get_gateway(account)
|
||||
res = gw.call_json("fs.list", body={"path": path})
|
||||
return res.get("entries", [])
|
||||
|
||||
|
||||
def read_md(account: str, filename: str, max_bytes: int = 200000) -> dict:
|
||||
"""Read a markdown file from the agent container via Hatch."""
|
||||
gw = get_gateway(account)
|
||||
offset = 0
|
||||
chunks = []
|
||||
chunk_len = min(65536, max_bytes)
|
||||
|
||||
while True:
|
||||
body = {"path": filename, "offset": offset, "len": chunk_len}
|
||||
res = gw.call_json("fs.read", body=body)
|
||||
text = res.get("text", "")
|
||||
if not text and res.get("data_base64"):
|
||||
import base64
|
||||
text = base64.b64decode(res["data_base64"]).decode("utf-8", "replace")
|
||||
chunks.append(text)
|
||||
offset += len(text.encode("utf-8"))
|
||||
if res.get("eof") or offset >= max_bytes or not text:
|
||||
break
|
||||
|
||||
full_text = "".join(chunks)
|
||||
return {
|
||||
"ok": True,
|
||||
"account": account,
|
||||
"filename": filename,
|
||||
"size": len(full_text.encode("utf-8")),
|
||||
"text": full_text,
|
||||
"eof": True,
|
||||
}
|
||||
|
||||
|
||||
def write_md(account: str, filename: str, text: str, overwrite: bool = True, append: bool = False) -> dict:
|
||||
"""Write content to a file in the agent container via Hatch."""
|
||||
gw = get_gateway(account)
|
||||
body = {
|
||||
"path": filename,
|
||||
"overwrite": overwrite,
|
||||
"append": append,
|
||||
"create_parent": True,
|
||||
"text": text,
|
||||
}
|
||||
res = gw.call_json("fs.write", body=body)
|
||||
return {
|
||||
"ok": True,
|
||||
"account": account,
|
||||
"filename": filename,
|
||||
"bytes_written": res.get("bytes_written", len(text.encode("utf-8"))),
|
||||
"path": res.get("path", f"/{filename}"),
|
||||
}
|
||||
|
||||
|
||||
def audit_agents(accounts: list = None) -> dict:
|
||||
"""Audit markdown files and operational DRIVE across fleet agents."""
|
||||
accounts = accounts or VALID_ACCOUNTS
|
||||
results = {}
|
||||
|
||||
for acct in accounts:
|
||||
acct_res = {
|
||||
"account": acct,
|
||||
"connected": False,
|
||||
"files": {},
|
||||
"drive_status": {},
|
||||
"issues": [],
|
||||
"drive_score": 0,
|
||||
}
|
||||
try:
|
||||
gw = get_gateway(acct)
|
||||
acct_res["connected"] = True
|
||||
acct_res["vm_id"] = gw.vm_id
|
||||
|
||||
entries = gw.call_json("fs.list", body={"path": ""}).get("entries", [])
|
||||
entry_map = {e["name"]: e for e in entries}
|
||||
|
||||
for tf in TARGET_MD_FILES:
|
||||
if tf in entry_map:
|
||||
info = entry_map[tf]
|
||||
acct_res["files"][tf] = {
|
||||
"exists": True,
|
||||
"size": info.get("size", 0),
|
||||
"modified": info.get("modifiedAt", ""),
|
||||
}
|
||||
else:
|
||||
acct_res["files"][tf] = {
|
||||
"exists": False,
|
||||
"size": 0,
|
||||
"modified": None,
|
||||
}
|
||||
|
||||
# Analyze DRIVE indicators
|
||||
# 1. HEARTBEAT.md checklist
|
||||
hb_info = acct_res["files"].get("HEARTBEAT.md", {})
|
||||
if not hb_info.get("exists"):
|
||||
acct_res["drive_status"]["heartbeat"] = "MISSING"
|
||||
acct_res["issues"].append("HEARTBEAT.md missing (no recurring checks)")
|
||||
elif hb_info.get("size", 0) <= 120:
|
||||
acct_res["drive_status"]["heartbeat"] = "EMPTY_CHECKLIST"
|
||||
acct_res["issues"].append("HEARTBEAT.md has empty checklist (background runner idle)")
|
||||
else:
|
||||
acct_res["drive_status"]["heartbeat"] = "ACTIVE"
|
||||
acct_res["drive_score"] += 25
|
||||
|
||||
# 2. PROACTIVE_PREFERENCES.md
|
||||
pro_info = acct_res["files"].get("PROACTIVE_PREFERENCES.md", {})
|
||||
if not pro_info.get("exists"):
|
||||
acct_res["drive_status"]["proactive"] = "MISSING"
|
||||
acct_res["issues"].append("PROACTIVE_PREFERENCES.md missing")
|
||||
elif pro_info.get("size", 0) <= 500:
|
||||
acct_res["drive_status"]["proactive"] = "BLANK_TEMPLATE"
|
||||
acct_res["issues"].append("PROACTIVE_PREFERENCES.md unconfigured (never reaches out)")
|
||||
else:
|
||||
acct_res["drive_status"]["proactive"] = "CONFIGURED"
|
||||
acct_res["drive_score"] += 25
|
||||
|
||||
# 3. SOUL.md
|
||||
soul_info = acct_res["files"].get("SOUL.md", {})
|
||||
if not soul_info.get("exists"):
|
||||
acct_res["drive_status"]["soul"] = "MISSING"
|
||||
acct_res["issues"].append("SOUL.md missing")
|
||||
elif soul_info.get("size", 0) <= 850:
|
||||
acct_res["drive_status"]["soul"] = "PASSIVE_STOCK"
|
||||
acct_res["issues"].append("SOUL.md is passive stock template (no operator drive)")
|
||||
else:
|
||||
acct_res["drive_status"]["soul"] = "OPERATOR_SOUL"
|
||||
acct_res["drive_score"] += 25
|
||||
|
||||
# 4. TOOLS.md & USER.md
|
||||
tools_info = acct_res["files"].get("TOOLS.md", {})
|
||||
user_info = acct_res["files"].get("USER.md", {})
|
||||
if tools_info.get("size", 0) > 400 and user_info.get("size", 0) > 400:
|
||||
acct_res["drive_status"]["context"] = "FULL_CONTEXT"
|
||||
acct_res["drive_score"] += 25
|
||||
else:
|
||||
acct_res["drive_status"]["context"] = "PARTIAL_OR_EMPTY"
|
||||
acct_res["issues"].append("TOOLS.md or USER.md missing operational conventions")
|
||||
|
||||
except Exception as e:
|
||||
acct_res["error"] = str(e)
|
||||
acct_res["issues"].append(f"Connection failed: {e}")
|
||||
|
||||
results[acct] = acct_res
|
||||
|
||||
return results
|
||||
|
||||
|
||||
def diff_md(account: str, filename: str) -> dict:
|
||||
"""Compare an agent's container file against the shared operator template."""
|
||||
local_path = SHARED_OPERATORS / filename
|
||||
if not local_path.exists():
|
||||
raise FileNotFoundError(f"Local template {local_path} not found")
|
||||
|
||||
local_content = local_path.read_text(encoding="utf-8")
|
||||
remote_data = read_md(account, filename)
|
||||
remote_content = remote_data.get("text", "")
|
||||
|
||||
diff = list(difflib.unified_diff(
|
||||
remote_content.splitlines(keepends=True),
|
||||
local_content.splitlines(keepends=True),
|
||||
fromfile=f"{account}:{filename} (remote)",
|
||||
tofile=f"shared/operators/{filename} (local)",
|
||||
))
|
||||
|
||||
return {
|
||||
"ok": True,
|
||||
"account": account,
|
||||
"filename": filename,
|
||||
"identical": len(diff) == 0,
|
||||
"diff": "".join(diff),
|
||||
"remote_size": len(remote_content.encode("utf-8")),
|
||||
"local_size": len(local_content.encode("utf-8")),
|
||||
}
|
||||
|
||||
|
||||
def inject_drive(account: str, force: bool = False) -> dict:
|
||||
"""Inject high-drive operational instructions into the agent's container."""
|
||||
updates = []
|
||||
|
||||
# 1. SOUL.md
|
||||
soul_text = (SHARED_OPERATORS / "SOUL.md").read_text(encoding="utf-8")
|
||||
r_soul = write_md(account, "SOUL.md", soul_text, overwrite=True)
|
||||
updates.append({"file": "SOUL.md", "bytes": r_soul["bytes_written"]})
|
||||
|
||||
# 2. PROACTIVE_PREFERENCES.md
|
||||
pro_text = (SHARED_OPERATORS / "PROACTIVE_PREFERENCES.md").read_text(encoding="utf-8")
|
||||
r_pro = write_md(account, "PROACTIVE_PREFERENCES.md", pro_text, overwrite=True)
|
||||
updates.append({"file": "PROACTIVE_PREFERENCES.md", "bytes": r_pro["bytes_written"]})
|
||||
|
||||
# 3. HEARTBEAT.md
|
||||
hb_text = (SHARED_OPERATORS / "HEARTBEAT.md").read_text(encoding="utf-8")
|
||||
r_hb = write_md(account, "HEARTBEAT.md", hb_text, overwrite=True)
|
||||
updates.append({"file": "HEARTBEAT.md", "bytes": r_hb["bytes_written"]})
|
||||
|
||||
# 4. USER.md
|
||||
user_text = (SHARED_OPERATORS / "USER.md").read_text(encoding="utf-8")
|
||||
r_user = write_md(account, "USER.md", user_text, overwrite=True)
|
||||
updates.append({"file": "USER.md", "bytes": r_user["bytes_written"]})
|
||||
|
||||
# 5. TOOLS.md
|
||||
tools_text = (SHARED_OPERATORS / "TOOLS.md").read_text(encoding="utf-8")
|
||||
r_tools = write_md(account, "TOOLS.md", tools_text, overwrite=True)
|
||||
updates.append({"file": "TOOLS.md", "bytes": r_tools["bytes_written"]})
|
||||
|
||||
# 6. AGENTS.md (preserve existing custom lessons if present)
|
||||
agents_template = (SHARED_OPERATORS / "AGENTS.md").read_text(encoding="utf-8")
|
||||
try:
|
||||
remote_agents = read_md(account, "AGENTS.md").get("text", "")
|
||||
if "## Lessons" in remote_agents and len(remote_agents) > len(agents_template):
|
||||
# Extract custom lessons from remote and merge
|
||||
custom_lessons = remote_agents.split("## Lessons", 1)[1]
|
||||
merged_agents = agents_template.rstrip() + "\n\n## Lessons" + custom_lessons
|
||||
r_agents = write_md(account, "AGENTS.md", merged_agents, overwrite=True)
|
||||
else:
|
||||
r_agents = write_md(account, "AGENTS.md", agents_template, overwrite=True)
|
||||
except Exception:
|
||||
r_agents = write_md(account, "AGENTS.md", agents_template, overwrite=True)
|
||||
updates.append({"file": "AGENTS.md", "bytes": r_agents["bytes_written"]})
|
||||
|
||||
return {
|
||||
"ok": True,
|
||||
"account": account,
|
||||
"action": "inject_drive",
|
||||
"updates": updates,
|
||||
"message": f"Successfully injected high-drive operator files into {account} container",
|
||||
}
|
||||
|
||||
|
||||
def get_ssh_info(account: str = None) -> dict:
|
||||
"""Return SSH connection coordinates and reverse tunnel configuration."""
|
||||
jump_host = "34.139.37.135"
|
||||
if account:
|
||||
entry = TUNNEL_PORTS.get(account, {"port": 2226, "terminal": 7683, "user": "hatch"})
|
||||
port = entry["port"]
|
||||
user = entry["user"]
|
||||
cmd = f"ssh -o StrictHostKeyChecking=no -p {port} {user}@localhost"
|
||||
proxy_cmd = f"ssh -o StrictHostKeyChecking=no -J super@{jump_host} -p {port} {user}@localhost"
|
||||
return {
|
||||
"account": account,
|
||||
"jump_host": jump_host,
|
||||
"port": port,
|
||||
"container_user": user,
|
||||
"terminal_port": entry.get("terminal"),
|
||||
"direct_from_vm": cmd,
|
||||
"jump_command": proxy_cmd,
|
||||
"cat_example": f"cat file.md | {proxy_cmd} 'cat > /home/hatch/file.md'",
|
||||
}
|
||||
return {
|
||||
"jump_host": jump_host,
|
||||
"tunnels": TUNNEL_PORTS,
|
||||
}
|
||||
|
||||
|
||||
def main():
|
||||
import argparse
|
||||
parser = argparse.ArgumentParser(description="Manage Muse agent .md files via Hatch and SSH")
|
||||
sub = parser.add_subparsers(dest="cmd")
|
||||
|
||||
p_audit = sub.add_parser("audit", help="Audit .md files and DRIVE across all agents")
|
||||
p_audit.add_argument("accounts", nargs="*", help="Optional account filter")
|
||||
p_audit.add_argument("--json", action="store_true", help="Output JSON")
|
||||
|
||||
p_list = sub.add_parser("list", help="List container files via Hatch")
|
||||
p_list.add_argument("account", help="Agent account")
|
||||
p_list.add_argument("path", nargs="?", default="", help="Subdirectory path")
|
||||
|
||||
p_read = sub.add_parser("read", help="Read a markdown file via Hatch")
|
||||
p_read.add_argument("account", help="Agent account")
|
||||
p_read.add_argument("filename", help="Filename (e.g. SOUL.md)")
|
||||
|
||||
p_write = sub.add_parser("write", help="Write a markdown file via Hatch")
|
||||
p_write.add_argument("account", help="Agent account")
|
||||
p_write.add_argument("filename", help="Filename (e.g. SOUL.md)")
|
||||
p_write.add_argument("--content", help="Text content to write")
|
||||
p_write.add_argument("--file", help="Local file to copy content from")
|
||||
|
||||
p_diff = sub.add_parser("diff", help="Diff remote file against shared operator template")
|
||||
p_diff.add_argument("account", help="Agent account")
|
||||
p_diff.add_argument("filename", help="Filename (e.g. SOUL.md)")
|
||||
|
||||
p_drive = sub.add_parser("inject-drive", help="Inject high-drive operator files into agent")
|
||||
p_drive.add_argument("account", help="Agent account")
|
||||
p_drive.add_argument("--force", action="store_true", help="Force overwrite")
|
||||
|
||||
p_sync_all = sub.add_parser("sync-all", help="Inject high-drive files across all active agents")
|
||||
|
||||
p_ssh = sub.add_parser("ssh-info", help="Get SSH tunnel dial-in information")
|
||||
p_ssh.add_argument("account", nargs="?", help="Optional agent account")
|
||||
|
||||
args = parser.parse_args()
|
||||
if not args.cmd:
|
||||
parser.print_help()
|
||||
sys.exit(1)
|
||||
|
||||
if args.cmd == "audit":
|
||||
res = audit_agents(args.accounts or None)
|
||||
if args.json:
|
||||
print(json.dumps(res, indent=2))
|
||||
else:
|
||||
print(f"\n{'='*70}\nMUSE AGENT .MD & DRIVE AUDIT REPORT\n{'='*70}")
|
||||
for acct, d in res.items():
|
||||
if not d.get("connected"):
|
||||
print(f"\n[AGENT {acct.upper()}] ✗ Connection failed: {d.get('error')}")
|
||||
continue
|
||||
score = d.get("drive_score", 0)
|
||||
status_color = "HIGH DRIVE" if score >= 75 else ("PARTIAL" if score >= 50 else "LOW DRIVE / STALE")
|
||||
print(f"\n[AGENT {acct.upper()}] DRIVE Score: {score}/100 ({status_color}) VM: {d.get('vm_id', 'unknown')}")
|
||||
for fname, finfo in d.get("files", {}).items():
|
||||
ex = "✓" if finfo.get("exists") else "✗"
|
||||
sz = f"{finfo.get('size', 0):6} bytes"
|
||||
mod = (finfo.get("modified") or "")[:19]
|
||||
print(f" {ex} {fname:24} {sz} {mod}")
|
||||
if d.get("issues"):
|
||||
print(" Issues:")
|
||||
for iss in d["issues"]:
|
||||
print(f" • {iss}")
|
||||
print(f"\n{'='*70}\n")
|
||||
|
||||
elif args.cmd == "list":
|
||||
entries = list_files(args.account, args.path)
|
||||
print(json.dumps(entries, indent=2))
|
||||
|
||||
elif args.cmd == "read":
|
||||
r = read_md(args.account, args.filename)
|
||||
print(r.get("text", ""))
|
||||
|
||||
elif args.cmd == "write":
|
||||
content = args.content
|
||||
if args.file:
|
||||
content = Path(args.file).read_text(encoding="utf-8")
|
||||
if content is None:
|
||||
print("Error: provide --content or --file", file=sys.stderr)
|
||||
sys.exit(2)
|
||||
res = write_md(args.account, args.filename, content)
|
||||
print(json.dumps(res, indent=2))
|
||||
|
||||
elif args.cmd == "diff":
|
||||
res = diff_md(args.account, args.filename)
|
||||
if res["identical"]:
|
||||
print(f"{args.account}:{args.filename} matches local shared/operators/{args.filename} exactly.")
|
||||
else:
|
||||
print(res["diff"])
|
||||
|
||||
elif args.cmd == "inject-drive":
|
||||
res = inject_drive(args.account, force=args.force)
|
||||
print(json.dumps(res, indent=2))
|
||||
|
||||
elif args.cmd == "sync-all":
|
||||
results = {}
|
||||
for acct in VALID_ACCOUNTS:
|
||||
try:
|
||||
results[acct] = inject_drive(acct)
|
||||
print(f"✓ Injected DRIVE into {acct}")
|
||||
except Exception as e:
|
||||
results[acct] = {"ok": False, "error": str(e)}
|
||||
print(f"✗ Failed {acct}: {e}")
|
||||
|
||||
elif args.cmd == "ssh-info":
|
||||
res = get_ssh_info(args.account)
|
||||
print(json.dumps(res, indent=2))
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
main()
|
||||
Reference in New Issue
Block a user