feat(kpi): add autonomous worker auto-spawn engine and watchdog reconciliation
This commit is contained in:
+148
-1
@@ -407,6 +407,128 @@ def get_live_advisory_block(node: str) -> str:
|
||||
return "\n".join(lines)
|
||||
|
||||
|
||||
NODE_SIDECHATS = {
|
||||
"646": "646 tasks",
|
||||
"opm": "heartbeat",
|
||||
"pip": "646-pip-coord",
|
||||
"dev": "dev-coord",
|
||||
"def": "def-coord",
|
||||
"muse": "646-muse-coord",
|
||||
}
|
||||
|
||||
|
||||
def find_pending_work_for_node(node: str) -> Optional[Dict[str, Any]]:
|
||||
"""Find assigned pending job or swarm slot for an agent node."""
|
||||
# 1. Look for node-specific auto-work jobs
|
||||
if JOBS_DIR.exists():
|
||||
candidates = sorted(list(JOBS_DIR.glob(f"auto-work-{node}-*.json")) + list(JOBS_DIR.glob(f"{node}-*.json")))
|
||||
for c in candidates:
|
||||
try:
|
||||
with open(c, "r", encoding="utf-8") as f:
|
||||
data = json.load(f)
|
||||
job_agent = data.get("agent")
|
||||
if job_agent and job_agent != node:
|
||||
continue
|
||||
job_name = c.stem
|
||||
return {
|
||||
"type": "job",
|
||||
"name": job_name,
|
||||
"path": str(c),
|
||||
"cmd": f"{sys.executable} {BIN_DIR}/job-dispatch.py {job_name}",
|
||||
}
|
||||
except Exception:
|
||||
continue
|
||||
|
||||
# 2. Check pending swarm slots
|
||||
try:
|
||||
from swarm_worker.poller import find_pending_slots
|
||||
slots = find_pending_slots()
|
||||
if slots:
|
||||
slot = slots[0]
|
||||
sw_id = slot.get("swarm_id", "swarm")
|
||||
idx = slot.get("slot_index", 0)
|
||||
return {
|
||||
"type": "swarm",
|
||||
"name": f"swarm-{sw_id}-s{idx}",
|
||||
"path": None,
|
||||
"cmd": f"{sys.executable} {BIN_DIR}/swarm_worker/daemon.py",
|
||||
}
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
return None
|
||||
|
||||
|
||||
def auto_spawn_workers(nodes: Optional[List[str]] = None, dry_run: bool = False) -> List[Dict[str, Any]]:
|
||||
"""Reconcile idle agents and auto-spawn background tmux workers to execute pending work."""
|
||||
target_nodes = nodes or VALID_NODES
|
||||
results = []
|
||||
|
||||
for node in target_nodes:
|
||||
# Check active tmux workers for this node
|
||||
active_tmux = len(get_agent_tmux_workers(node))
|
||||
if active_tmux > 0:
|
||||
results.append({
|
||||
"node": node,
|
||||
"action": "skip",
|
||||
"reason": f"Active tmux worker already running ({active_tmux})",
|
||||
})
|
||||
continue
|
||||
|
||||
# Check work availability
|
||||
work = find_pending_work_for_node(node)
|
||||
if not work:
|
||||
results.append({
|
||||
"node": node,
|
||||
"action": "idle",
|
||||
"reason": "No pending jobs or swarm slots",
|
||||
})
|
||||
continue
|
||||
|
||||
session_label = f"worker-{work['name'][:18]}"
|
||||
cmd_to_run = f"{work['cmd']} > /tmp/tmux-{node}-{session_label}.log 2>&1"
|
||||
|
||||
if dry_run:
|
||||
results.append({
|
||||
"node": node,
|
||||
"action": "would_spawn",
|
||||
"session": session_label,
|
||||
"work_type": work["type"],
|
||||
"work_name": work["name"],
|
||||
"command": work["cmd"],
|
||||
})
|
||||
continue
|
||||
|
||||
# Execute spawn
|
||||
spawn_res = spawn_tmux_worker(node, session_label, cmd_to_run)
|
||||
if spawn_res.get("ok"):
|
||||
# Send sidechat notification
|
||||
try:
|
||||
from invite_handler import send_loopback_notice
|
||||
chat = NODE_SIDECHATS.get(node, "646 tasks")
|
||||
msg = f"[BOX-AUTO-WORKER] Spawned background tmux worker '{session_label}' executing {work['type']} ({work['name']}). Logs at /tmp/tmux-{node}-{session_label}.log"
|
||||
send_loopback_notice(recipient=node, target=chat, message=msg)
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
results.append({
|
||||
"node": node,
|
||||
"action": "spawned",
|
||||
"session": session_label,
|
||||
"work_type": work["type"],
|
||||
"work_name": work["name"],
|
||||
"command": work["cmd"],
|
||||
})
|
||||
else:
|
||||
results.append({
|
||||
"node": node,
|
||||
"action": "error",
|
||||
"error": spawn_res.get("error", "Unknown spawn error"),
|
||||
})
|
||||
|
||||
return results
|
||||
|
||||
|
||||
def spawn_tmux_worker(node: str, session: str, command: str) -> Dict[str, Any]:
|
||||
"""Spawn an autonomous tmux worker session on the agent's netns or shared socket."""
|
||||
# Ensure session name is prefixed
|
||||
@@ -484,7 +606,10 @@ def main():
|
||||
p_spawn.add_argument("node", choices=VALID_NODES, help="Agent node")
|
||||
p_spawn.add_argument("session", help="Session label")
|
||||
p_spawn.add_argument("worker_command", help="Command to execute inside worker")
|
||||
p_spawn.add_argument("--json", action="store_true")
|
||||
p_autospawn = subparsers.add_parser("auto-spawn", help="Auto-spawn background tmux workers for idle nodes with pending work")
|
||||
p_autospawn.add_argument("--node", choices=VALID_NODES, default=None, help="Filter by node")
|
||||
p_autospawn.add_argument("--dry-run", action="store_true", help="Report what would be spawned without executing")
|
||||
p_autospawn.add_argument("--json", action="store_true")
|
||||
|
||||
args = parser.parse_args()
|
||||
|
||||
@@ -555,6 +680,28 @@ def main():
|
||||
print(f"✘ Failed to spawn worker: {res.get('error')}", file=sys.stderr)
|
||||
sys.exit(1)
|
||||
|
||||
elif args.command == "auto-spawn":
|
||||
nodes = [args.node] if getattr(args, "node", None) else None
|
||||
results = auto_spawn_workers(nodes=nodes, dry_run=args.dry_run)
|
||||
if args.json:
|
||||
print(json.dumps(results, indent=2))
|
||||
return
|
||||
print(f"\n=== AUTO-SPAWN WORKER RECONCILIATION {'(DRY-RUN)' if args.dry_run else ''} ===")
|
||||
for r in results:
|
||||
n = r.get("node")
|
||||
act = r.get("action")
|
||||
if act == "spawned":
|
||||
print(f" ✔ @{n:<5} : SPAWNED session '{r.get('session')}' ({r.get('work_type')}: {r.get('work_name')})")
|
||||
elif act == "would_spawn":
|
||||
print(f" ? @{n:<5} : WOULD SPAWN session '{r.get('session')}' ({r.get('work_type')}: {r.get('work_name')})")
|
||||
elif act == "skip":
|
||||
print(f" - @{n:<5} : SKIP ({r.get('reason')})")
|
||||
elif act == "idle":
|
||||
print(f" - @{n:<5} : IDLE ({r.get('reason')})")
|
||||
elif act == "error":
|
||||
print(f" ✘ @{n:<5} : ERROR ({r.get('error')})")
|
||||
print()
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
main()
|
||||
|
||||
Reference in New Issue
Block a user