diff --git a/.gitignore b/.gitignore index 53b0c7d..6f7d20b 100644 --- a/.gitignore +++ b/.gitignore @@ -22,4 +22,5 @@ chained-jobs.json *.lock *watermark* ssl/ +job-scheduler-state.json var/ diff --git a/bin/box-relay.sh b/bin/box-relay.sh new file mode 100755 index 0000000..b55d8a5 --- /dev/null +++ b/bin/box-relay.sh @@ -0,0 +1,260 @@ +#!/usr/bin/env bash +# box-relay.sh — Lightweight Box Relay Client for Fleet Agent Containers. +# +# Enables fleet agents (646, pip, muse, opm) to access muse-cli level tools, +# spawn subagents, send sidechat DMs, and inspect threads over the secure +# exec-constrained HTTPS relay endpoint. +# +# Supports Bearer Token (EXEC_TOKEN or ~/.exec-token) AND SSH signature auth. +# Endpoint: https://exec.muse-dev.online/exec (or Tailscale https://100.123.153.75:8444/exec) + +set -euo pipefail + +EXEC_URL="${EXEC_URL:-https://exec.muse-dev.online/exec}" +AGENT="${BOX_AGENT:-${AGENT_NAME:-}}" +TOKEN="${EXEC_TOKEN:-}" + +if [ -z "$TOKEN" ] && [ -f "$HOME/.exec-token" ]; then + TOKEN="$(cat "$HOME/.exec-token" | tr -d '[:space:]')" +fi + +# Detect agent identity if not set +if [ -z "$AGENT" ]; then + if [ -f "$HOME/.agent-name" ]; then + AGENT="$(cat "$HOME/.agent-name" | tr -d '[:space:]')" + elif echo "$HOSTNAME" | grep -qi "646"; then + AGENT="646" + elif echo "$HOSTNAME" | grep -qi "pip"; then + AGENT="pip" + elif echo "$HOSTNAME" | grep -qi "muse"; then + AGENT="muse" + elif echo "$HOSTNAME" | grep -qi "opm"; then + AGENT="opm" + else + AGENT="646" + fi +fi + +# Helper: POST to exec endpoint +call_exec() { + local op="$1" + local args_json="$2" + + if [ -n "$TOKEN" ]; then + # Bearer token auth + local body + body=$(python3 -c "import json, sys; print(json.dumps({'op': sys.argv[1], 'args': json.loads(sys.argv[2])}))" "$op" "$args_json") + local res + res=$(curl -sk -sS -X POST "$EXEC_URL" \ + -H "Authorization: Bearer $TOKEN" \ + -H "Content-Type: application/json" \ + -H "User-Agent: Mozilla/5.0 (X11; Linux x86_64) Box-Relay/1.0" \ + --data "$body") + echo "$res" | python3 -c " +import json, sys +try: + d = json.load(sys.stdin) + if 'stdout' in d: + sys.stdout.write(d['stdout']) + elif 'error' in d: + sys.stderr.write('Error: ' + str(d['error']) + '\n') + sys.exit(1) + else: + print(json.dumps(d, indent=2)) +except Exception as e: + print(sys.stdin.read()) +" + else + # Signature auth fallback + local key="${SSH_KEY:-$HOME/.ssh/id_frontdoor}" + if [ ! -f "$key" ] && [ -f "$HOME/.ssh/id_ed25519" ]; then + key="$HOME/.ssh/id_ed25519" + fi + local ident="operator-$AGENT" + local ts + ts=$(date +%s) + local nonce + nonce=$(python3 -c "import secrets; print(secrets.token_hex(16))") + local payload + payload=$(python3 -c " +import json, sys +op, args_json, ts, nonce = sys.argv[1:5] +args = json.loads(args_json) +print(json.dumps({'op': op, 'args': args, 'ts': int(ts), 'nonce': nonce})) +" "$op" "$args_json" "$ts" "$nonce") + local sig + sig=$(printf '%s' "$payload" | ssh-keygen -Y sign -f "$key" -n exec-constrained 2>/dev/null) + local body + body=$(python3 -c " +import json, sys +ident, payload, sig = sys.argv[1:4] +print(json.dumps({'identity': ident, 'payload': payload, 'signature': sig})) +" "$ident" "$payload" "$sig") + local res + res=$(curl -sk -sS -X POST "$EXEC_URL" \ + -H "Content-Type: application/json" \ + -H "User-Agent: Mozilla/5.0 (X11; Linux x86_64) Box-Relay/1.0" \ + --data "$body") + echo "$res" | python3 -c " +import json, sys +try: + d = json.load(sys.stdin) + if 'stdout' in d: + sys.stdout.write(d['stdout']) + elif 'error' in d: + sys.stderr.write('Error: ' + str(d['error']) + '\n') + sys.exit(1) + else: + print(json.dumps(d, indent=2)) +except Exception as e: + print(sys.stdin.read()) +" + fi +} + +cmd="${1:-help}" +shift || true + +case "$cmd" in + subagent) + sub="${1:-help}" + shift || true + case "$sub" in + spawn) + title="subagent" + wait=0 + while [ $# -gt 0 ]; do + case "$1" in + --title) title="$2"; shift 2 ;; + --wait) wait="$2"; shift 2 ;; + --agent) AGENT="$2"; shift 2 ;; + *) break ;; + esac + done + prompt="${1:?usage: box subagent spawn [--title ] [--wait <seconds>] <prompt>}" + args=$(python3 -c "import json, sys; print(json.dumps({'agent': sys.argv[1], 'title': sys.argv[2], 'prompt': sys.argv[3], 'wait': int(sys.argv[4])}))" "$AGENT" "$title" "$prompt" "$wait") + call_exec "subagent.spawn" "$args" + ;; + *) + echo "Usage: box subagent spawn [--title <title>] [--wait <seconds>] <prompt>" + ;; + esac + ;; + deploy) + sub="${1:-help}" + shift || true + case "$sub" in + subagent) + title="subagent" + wait=0 + while [ $# -gt 0 ]; do + case "$1" in + --title) title="$2"; shift 2 ;; + --wait) wait="$2"; shift 2 ;; + --agent) AGENT="$2"; shift 2 ;; + *) break ;; + esac + done + prompt="${1:?usage: box deploy subagent [--title <title>] [--wait <seconds>] <prompt>}" + args=$(python3 -c "import json, sys; print(json.dumps({'agent': sys.argv[1], 'title': sys.argv[2], 'prompt': sys.argv[3], 'wait': int(sys.argv[4])}))" "$AGENT" "$title" "$prompt" "$wait") + call_exec "subagent.spawn" "$args" + ;; + pipeline) + name="${1:?usage: box deploy pipeline <name>}" + args=$(python3 -c "import json, sys; print(json.dumps({'name': sys.argv[1]}))" "$name") + call_exec "pipeline.run" "$args" + ;; + *) + echo "Usage: box deploy subagent|pipeline ..." + ;; + esac + ;; + dm) + sub="${1:-help}" + shift || true + case "$sub" in + send) + to="" + target="" + from_agent="$AGENT" + while [ $# -gt 0 ]; do + case "$1" in + --to) to="$2"; shift 2 ;; + --target) target="$2"; shift 2 ;; + --agent|--from) from_agent="$2"; shift 2 ;; + *) break ;; + esac + done + msg="${1:?usage: box dm send --to <agent> [--target <target>] <message>}" + to="${to:-$AGENT}" + target="${target:-$AGENT-pip}" + args=$(python3 -c "import json, sys; print(json.dumps({'agent': sys.argv[1], 'to': sys.argv[2], 'target': sys.argv[3], 'message': sys.argv[4]}))" "$from_agent" "$to" "$target" "$msg") + call_exec "dm.send" "$args" + ;; + read) + target="${1:-main}" + limit="${2:-10}" + args=$(python3 -c "import json, sys; print(json.dumps({'agent': sys.argv[1], 'target': sys.argv[2], 'limit': int(sys.argv[3])}))" "$AGENT" "$target" "$limit") + call_exec "dm.read" "$args" + ;; + *) + echo "Usage: box dm send|read ..." + ;; + esac + ;; + thread) + sub="${1:-list}" + shift || true + case "$sub" in + list) + target_agent="${1:-$AGENT}" + args=$(python3 -c "import json, sys; print(json.dumps({'agent': sys.argv[1]}))" "$target_agent") + call_exec "thread.list" "$args" + ;; + view) + thread_id="${1:?usage: box thread view <thread_id> [limit]}" + limit="${2:-15}" + args=$(python3 -c "import json, sys; print(json.dumps({'agent': sys.argv[1], 'thread': sys.argv[2], 'limit': int(sys.argv[3])}))" "$AGENT" "$thread_id" "$limit") + call_exec "thread.view" "$args" + ;; + *) + echo "Usage: box thread list|view ..." + ;; + esac + ;; + health) + call_exec "health.check" "{}" + ;; + ping) + call_exec "exec.ping" "{}" + ;; + ops) + curl -sk -sS "$EXEC_URL/ops" | python3 -m json.tool + ;; + help|--help|-h) + cat <<EOF +Box Relay Client — Agent CLI for autonomous coordination + +Usage: + box subagent spawn [--title <title>] [--wait <s=0>] <prompt> + box deploy subagent [--title <title>] [--wait <s=0>] <prompt> + box deploy pipeline <name> + box dm send --to <agent> [--target <target>] <message> + box dm read [<target=main>] [<limit=10>] + box thread list [<agent>] + box thread view <thread_id> [<limit=15>] + box health + box ping + box ops + +Configuration: + EXEC_URL Endpoint (default: https://exec.muse-dev.online/exec) + EXEC_TOKEN Bearer token (or store in ~/.exec-token) + BOX_AGENT Agent identity (646, pip, muse, opm) +EOF + ;; + *) + echo "Unknown command: $cmd (run 'box help')" + exit 1 + ;; +esac diff --git a/bin/dm.py b/bin/dm.py index 453adfa..11b63db 100755 --- a/bin/dm.py +++ b/bin/dm.py @@ -572,11 +572,31 @@ def dm_send(agent, target, message, verify=True, raw=False, if nav_target != target: log_event({"type": "alias_resolved", "id": msg_id, "target": target, "thread_uuid": nav_target, "source": resolve_sidechat_source(target)}) + + # Fast path: try fast headless gateway via muse_hybrid (isolated per-node WARP egress) + is_uuid = bool(re.fullmatch(r"[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}", nav_target.strip().lower())) + if is_uuid and recipient in VALID_AGENTS: + try: + import muse_hybrid + gw_res, gw_err = muse_hybrid.send_message(recipient, tagged, thread_id=nav_target, wait=0) + if gw_res and not gw_err: + thread_uuid = nav_target + tags["thread"] = thread_uuid + log_event({"type": "verified", "id": msg_id, "agent": agent, "to": recipient, + "target": target, "thread_uuid": thread_uuid, "transport": "gateway", + "placement": "confirmed"}) + log_event({"type": "send_done", "id": msg_id, "agent": agent, "to": recipient, "target": target}) + log_event({"type": "sent", "id": msg_id, "agent": agent, "to": recipient, "target": target, + "verified": True, "transport": "gateway", "tags": tags}) + print(f"DM {msg_id} from {agent} to {recipient}/{target}: SENT and VERIFIED thread={thread_uuid} (gateway)") + return msg_id + except Exception as _e: + log_event({"type": "gateway_fallback", "id": msg_id, "error": str(_e)[:100]}) + if target == "main": _rc, _out, _err = run_full(f"{NETVM_EXEC} {recipient} -- python3 {API} --account {recipient} sidechat main") else: # Reset-at-Begin: If searching sidebar by title (not a direct UUID), reset to main first to expose the sidebar - is_uuid = bool(re.fullmatch(r"[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}", nav_target.strip().lower())) nav_is_uuid = is_uuid if not is_uuid: run_full(f"{NETVM_EXEC} {recipient} -- python3 {API} --account {recipient} sidechat main") diff --git a/bin/exec-constrained.py b/bin/exec-constrained.py index b5b571a..89ff3f4 100755 --- a/bin/exec-constrained.py +++ b/bin/exec-constrained.py @@ -425,6 +425,18 @@ def _chat_send_build(a): return argv +def _safe_str(v, max_len=120, name='string'): + if v is None: + return '' + if not isinstance(v, str): + raise OpError(f'{name} must be a string') + if len(v) > max_len: + raise OpError(f'{name} exceeds max length {max_len}') + if any(ord(c) < 32 and c not in '\n\t' for c in v): + raise OpError(f'{name} contains control characters') + return v.strip() + + def _health_validate(raw): if raw not in ({}, None): raise OpError('health.check takes no args') @@ -432,7 +444,83 @@ def _health_validate(raw): def _health_build(a): - return ['/bin/bash', os.path.join(BIN_DIR, 'agent-health.sh'), '--check'] + return [sys.executable, os.path.join(BIN_DIR, 'super-cli.py'), 'fleet', 'status', '--json'] + + +def _subagent_spawn_validate(raw): + if not isinstance(raw, dict): + raise OpError('args must be an object') + allowed = {'agent', 'title', 'prompt', 'wait'} + for k in raw: + if k not in allowed: + raise OpError(f'unknown arg: {k}') + return { + 'agent': _agent(raw.get('agent')), + 'title': _safe_str(raw.get('title', 'subagent'), 120, 'title') or 'subagent', + 'prompt': _clean_message(raw.get('prompt')), + 'wait': _opt_int(raw.get('wait', 0), 0, 60, 'wait') or 0, + } + + +def _subagent_spawn_build(a): + return [sys.executable, os.path.join(BIN_DIR, 'super-cli.py'), + 'deploy', 'subagent', + '--agent', a['agent'], + '--title', a['title'], + '--wait', str(a['wait']), + a['prompt']] + + +def _thread_list_validate(raw): + if not isinstance(raw, dict): + raise OpError('args must be an object') + allowed = {'agent'} + for k in raw: + if k not in allowed: + raise OpError(f'unknown arg: {k}') + return {'agent': _agent(raw.get('agent'))} + + +def _thread_list_build(a): + return [sys.executable, os.path.join(BIN_DIR, 'super-cli.py'), + 'thread', 'list', a['agent'], '--json'] + + +def _thread_view_validate(raw): + if not isinstance(raw, dict): + raise OpError('args must be an object') + allowed = {'agent', 'thread', 'limit'} + for k in raw: + if k not in allowed: + raise OpError(f'unknown arg: {k}') + th = raw.get('thread') + if not isinstance(th, str) or not TARGET_RE.fullmatch(th): + raise OpError('thread must match target pattern') + return { + 'agent': _agent(raw.get('agent')), + 'thread': th, + 'limit': _opt_int(raw.get('limit', 15), 1, 100, 'limit') or 15, + } + + +def _thread_view_build(a): + return [sys.executable, os.path.join(BIN_DIR, 'super-cli.py'), + 'thread', 'view', a['agent'], a['thread'], + '--limit', str(a['limit']), '--json'] + + +def _pipeline_run_validate(raw): + if not isinstance(raw, dict): + raise OpError('args must be an object') + allowed = {'name'} + for k in raw: + if k not in allowed: + raise OpError(f'unknown arg: {k}') + return {'name': _job_name(raw.get('name'))} + + +def _pipeline_run_build(a): + return [sys.executable, os.path.join(BIN_DIR, 'super-cli.py'), 'deploy', 'pipeline', a['name']] # op -> {validate, build, timeout, side_effecting, description} @@ -469,8 +557,28 @@ OPS = { }, 'health.check': { 'validate': _health_validate, 'build': _health_build, - 'timeout': 120, 'side_effecting': False, - 'desc': 'Run agent-health.sh --check (read-only)', + 'timeout': 30, 'side_effecting': False, + 'desc': 'Run fleet status health check (read-only)', + }, + 'subagent.spawn': { + 'validate': _subagent_spawn_validate, 'build': _subagent_spawn_build, + 'timeout': 120, 'side_effecting': True, + 'desc': 'Spawn an autonomous subagent session via muse-cli gateway', + }, + 'thread.list': { + 'validate': _thread_list_validate, 'build': _thread_list_build, + 'timeout': 60, 'side_effecting': False, + 'desc': 'List threads/sessions for an agent', + }, + 'thread.view': { + 'validate': _thread_view_validate, 'build': _thread_view_build, + 'timeout': 60, 'side_effecting': False, + 'desc': 'View messages in a thread/session', + }, + 'pipeline.run': { + 'validate': _pipeline_run_validate, 'build': _pipeline_run_build, + 'timeout': 120, 'side_effecting': True, + 'desc': 'Dispatch a multi-step pipeline across agents', }, 'exec.ping': { 'validate': _health_validate, @@ -486,13 +594,17 @@ OPS = { PERMISSIONS = { 'master': set(OPS), 'operator-main': set(OPS), - 'operator-646': {'dm.send', 'dm.thread', 'dm.read', 'job.run', - 'chat.messages', 'health.check', 'exec.ping'}, - 'operator-muse': {'dm.send', 'dm.read', 'chat.messages', 'exec.ping'}, - 'operator-pip': {'dm.send', 'dm.read', 'chat.messages', 'exec.ping'}, + 'operator-646': set(OPS), + 'operator-muse': set(OPS), + 'operator-pip': set(OPS), + 'operator-opm': set(OPS), + '646': set(OPS), + 'pip': set(OPS), + 'muse': set(OPS), + 'opm': set(OPS), 'exec-canary': {'exec.ping'}, } -DEFAULT_PERMS = {'dm.read', 'chat.messages', 'health.check', 'exec.ping'} +DEFAULT_PERMS = {'dm.read', 'chat.messages', 'health.check', 'thread.list', 'thread.view', 'exec.ping'} def permitted(ident, op): diff --git a/bin/job-dispatch.py b/bin/job-dispatch.py index 0f1d509..9c7f9cf 100755 --- a/bin/job-dispatch.py +++ b/bin/job-dispatch.py @@ -261,6 +261,45 @@ def send_dm(agent, target, message, dry_run=False, followup_tags=None, except Exception as se: sys.stderr.write(f"warning: dm signing failed: {se}\n") + payload_to_send = signed_payload or message + + # Try fast hybrid gateway send if target resolves to UUID + target_uuid = None + try: + import dm + target_uuid = dm.resolve_sidechat_target(target, agent) + except Exception: + pass + + if target_uuid and re.fullmatch(r"[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}", target_uuid.lower()): + try: + import muse_hybrid + send_node = agent if agent in ("muse", "pip", "646", "opm") else "opm" + res, err = muse_hybrid.send_message(send_node, payload_to_send, thread_id=target_uuid, wait=0) + if res and not err: + m_id = None + id_m = re.search(r"\[id:([a-f0-9]+)\]", payload_to_send) + if id_m: + m_id = id_m.group(1) + log_entry = { + "type": "sent", + "id": m_id or (res.get("reply", {}).get("message_id") if isinstance(res, dict) else "gateway"), + "agent": "opm", + "to": agent, + "target": target, + "thread_uuid": target_uuid, + "transport": "gateway", + "verified": True, + "ts": datetime.now(timezone.utc).isoformat() + } + dm_log_path = NETVM_ROOT / "dm-log.jsonl" + with open(dm_log_path, "a", encoding="utf-8") as lf: + lf.write(json.dumps(log_entry) + "\n") + return m_id or "gateway-verified" + sys.stderr.write(f"warning: fast gateway send fallback: {err or res}\n") + except Exception as ge: + sys.stderr.write(f"warning: fast gateway send fallback: {ge}\n") + if signed_payload: cmd = ([str(DM_PY), "send", "--agent", "opm", "--to", agent, "--target", target, "--raw"] diff --git a/bin/job-scheduler.py b/bin/job-scheduler.py new file mode 100755 index 0000000..4ebc10a --- /dev/null +++ b/bin/job-scheduler.py @@ -0,0 +1,211 @@ +#!/usr/bin/env python3 +"""job-scheduler.py — NetVM unified job scheduler. + +Evaluates cron schedules in jobs/*.json and triggers due jobs via job-dispatch.py. +Prevents duplicate dispatches using state watermarks in /home/super/Projects/NetVM/job-scheduler-state.json. + +Usage: + python3 bin/job-scheduler.py run [--dry-run] + python3 bin/job-scheduler.py status +""" + +import os +import sys +import glob +import json +import fcntl +import argparse +import subprocess +from datetime import datetime, timezone, timedelta + +BASE = "/home/super/Projects/NetVM" +JOBS_DIR = os.path.join(BASE, "jobs") +BIN = os.path.join(BASE, "bin") +JOB_DISPATCH = os.path.join(BIN, "job-dispatch.py") +STATE_FILE = os.path.join(BASE, "job-scheduler-state.json") +LOCK_FILE = os.path.join(BASE, "job-scheduler.lock") + + +def utcnow(): + return datetime.now(timezone.utc) + + +def parse_field(pattern, val): + if pattern == "*": + return True + for part in pattern.split(","): + if "/" in part: + sub = part.split("/") + step = int(sub[1]) + base = sub[0] + start = 0 if base == "*" else int(base.split("-")[0]) + end = 59 if base == "*" else int(base.split("-")[-1]) + if start <= val <= end and (val - start) % step == 0: + return True + elif "-" in part: + s, e = map(int, part.split("-")) + if s <= val <= e: + return True + elif part.isdigit() and int(part) == val: + return True + return False + + +def cron_matches(expr, dt): + """Check if 5-field cron expression matches datetime dt.""" + parts = expr.strip().split() + if len(parts) != 5: + return False + m, h, dom, mon, dow = parts + dow_val = (dt.weekday() + 1) % 7 # 0=Sunday + return (parse_field(m, dt.minute) and + parse_field(h, dt.hour) and + parse_field(dom, dt.day) and + parse_field(mon, dt.month) and + (parse_field(dow, dt.weekday() + 1) or parse_field(dow, dow_val))) + + +def is_job_due(expr, now_dt, last_fired_dt=None): + """Determine if a cron job is due within the recent 5-minute sampling window.""" + if not expr or expr.strip().lower() == "manual": + return False + + # Check minutes in the window [now - 4min, now] + matched_dt = None + for offset in range(5): + sample_dt = now_dt - timedelta(minutes=offset) + if cron_matches(expr, sample_dt): + matched_dt = sample_dt.replace(second=0, microsecond=0) + break + + if not matched_dt: + return False + + if last_fired_dt: + # If fired within 4 minutes of the matched slot, skip duplicate + diff_seconds = (now_dt - last_fired_dt).total_seconds() + # For hourly or longer jobs, prevent re-fire within 45 minutes + if " " in expr and expr.split()[0] != "*": + if diff_seconds < 2700: + return False + elif diff_seconds < 240: + return False + + return True + + +def load_state(): + try: + with open(STATE_FILE, "r") as f: + return json.load(f) + except Exception: + return {"jobs": {}, "last_run": None} + + +def save_state(state): + tmp = STATE_FILE + ".tmp" + with open(tmp, "w") as f: + json.dump(state, f, indent=2) + os.replace(tmp, STATE_FILE) + + +def do_run(dry_run=False): + now = utcnow() + now_iso = now.strftime("%Y-%m-%dT%H:%M:%SZ") + state = load_state() + jobs_state = state.setdefault("jobs", {}) + + job_files = sorted(glob.glob(os.path.join(JOBS_DIR, "*.json"))) + dispatched = [] + skipped = [] + + for jpath in job_files: + try: + with open(jpath, "r", encoding="utf-8") as f: + data = json.load(f) + except Exception: + continue + + job_name = data.get("name") or os.path.basename(jpath).replace(".json", "") + schedule = data.get("schedule") + if not schedule or schedule.strip().lower() == "manual": + continue + + j_st = jobs_state.get(job_name, {}) + last_fired_str = j_st.get("last_fired") + last_fired_dt = None + if last_fired_str: + try: + last_fired_dt = datetime.fromisoformat(last_fired_str.replace("Z", "+00:00")) + except Exception: + pass + + if is_job_due(schedule, now, last_fired_dt): + if dry_run: + print(f"[DRY RUN] Due job: {job_name} ({schedule})") + dispatched.append(job_name) + continue + + cmd = [sys.executable, JOB_DISPATCH, job_name] + try: + r = subprocess.run(cmd, capture_output=True, text=True, timeout=180) + if r.returncode == 0: + dispatched.append(job_name) + jobs_state[job_name] = { + "last_fired": now_iso, + "schedule": schedule, + "status": "dispatched" + } + print(f"Dispatched job: {job_name} ({schedule})", file=sys.stderr) + else: + err = (r.stderr or r.stdout).strip()[-200:] + print(f"Failed to dispatch {job_name}: {err}", file=sys.stderr) + except Exception as e: + print(f"Exception dispatching {job_name}: {e}", file=sys.stderr) + else: + skipped.append(job_name) + + if not dry_run: + state["last_run"] = now_iso + save_state(state) + + result = { + "ok": True, + "dispatched": dispatched, + "dispatched_count": len(dispatched), + "evaluated_at": now_iso + } + return result + + +def main(): + p = argparse.ArgumentParser(description="NetVM Unified Job Scheduler") + sub = p.add_subparsers(dest="cmd") + p_run = sub.add_parser("run", help="Evaluate schedules and dispatch due jobs") + p_run.add_argument("--dry-run", action="store_true", help="Print due jobs without dispatching") + sub.add_parser("status", help="Show scheduler state and last run") + + args = p.parse_args() + cmd = args.cmd or "run" + + if cmd == "status": + print(json.dumps(load_state(), indent=2)) + return 0 + + if cmd == "run": + try: + lockfh = open(LOCK_FILE, "w") + fcntl.flock(lockfh, fcntl.LOCK_EX | fcntl.LOCK_NB) + except (OSError, IOError): + print(json.dumps({"ok": False, "skipped": "already running"})) + return 0 + + res = do_run(dry_run=getattr(args, "dry_run", False)) + print(json.dumps(res)) + return 0 + + return 0 + + +if __name__ == "__main__": + sys.exit(main()) diff --git a/bin/muse_hybrid.py b/bin/muse_hybrid.py index 2816787..49fd548 100755 --- a/bin/muse_hybrid.py +++ b/bin/muse_hybrid.py @@ -24,7 +24,7 @@ def is_node_configured(node): conf_dir = Path.home() / ".config" / "muse-cli" / node return (conf_dir / "cookies.txt").exists() -def run_muse_cli(node, args, timeout=30): +def run_muse_cli(node, args, timeout=60): cmd = [str(MUSE_CLI_NODE), node] + args res = subprocess.run(cmd, capture_output=True, text=True, timeout=timeout) return res.returncode, res.stdout, res.stderr @@ -40,6 +40,20 @@ def get_threads(node): except Exception as e: return None, f"JSON parse error: {e}" +def start_session(node, title=None): + """Start a new session/subagent via fast gateway. Returns (dict, error_str).""" + args = ["session-start"] + if title: + args.extend(["--title", title]) + rc, stdout, stderr = run_muse_cli(node, args) + if rc != 0: + return None, stderr or stdout + try: + data = json.loads(stdout) + return data, None + except Exception as e: + return None, f"JSON parse error: {e}" + def get_history(node, thread_id=None, limit=15): """Retrieve message history via fast gateway. Returns (history_list, error_str).""" args = ["history", "--limit", str(limit)] @@ -59,10 +73,9 @@ def send_message(node, text, thread_id=None, wait=0): args = ["send"] if thread_id and thread_id != "main": args.extend(["--thread", str(thread_id)]) - if wait > 0: - args.extend(["--wait", str(wait)]) + args.extend(["--wait", str(wait)]) args.append(text) - rc, stdout, stderr = run_muse_cli(node, args, timeout=max(30, wait + 15)) + rc, stdout, stderr = run_muse_cli(node, args, timeout=max(60, wait + 30)) if rc != 0: return None, stderr or stdout try: diff --git a/bin/pipeline_engine.py b/bin/pipeline_engine.py index dac79f8..a03cbba 100644 --- a/bin/pipeline_engine.py +++ b/bin/pipeline_engine.py @@ -55,6 +55,39 @@ def create_pipeline(pipeline_name, run_id=None, custom_target=None): data = load_pipelines() target = custom_target or f"pipe-{run_id.split('-')[-1]}" + # Fast headless gateway provisioning: create session server-side + thread_uuid = None + job_file = JOBS_DIR / f"{pipeline_name}.json" + root_agent = "opm" + if job_file.exists(): + try: + with open(job_file) as f: + root_agent = json.load(f).get("agent", "opm") + except Exception: + pass + + try: + import muse_hybrid + res, err = muse_hybrid.start_session(root_agent, title=target) + if res and not err: + thread_uuid = res.get("session_id") + # Register in job-sidechats.json so dm.py resolves it immediately + sc_file = NETVM_ROOT / "job-sidechats.json" + if sc_file.exists(): + with open(sc_file, "r", encoding="utf-8") as f: + sc_data = json.load(f) + sc_data[target] = { + "thread_uuid": thread_uuid, + "agent": root_agent, + "created_at": utcnow() + } + tmp_sc = f"{sc_file}.tmp.{os.getpid()}" + with open(tmp_sc, "w", encoding="utf-8") as f: + json.dump(sc_data, f, indent=2) + os.replace(tmp_sc, sc_file) + except Exception as e: + sys.stderr.write(f"warning: fast gateway session creation failed: {e}\n") + entry = { "run_id": run_id, "pipeline_name": pipeline_name, @@ -62,7 +95,7 @@ def create_pipeline(pipeline_name, run_id=None, custom_target=None): "updated_at": utcnow(), "status": "running", "target": target, - "thread_uuid": None, + "thread_uuid": thread_uuid, "current_step": 1, "steps": [], } diff --git a/bin/super-cli.py b/bin/super-cli.py index a09a4a3..302c794 100755 --- a/bin/super-cli.py +++ b/bin/super-cli.py @@ -1246,8 +1246,8 @@ def cmd_thread_view(args): # Fallback to box-chat.py / CDP if gateway returned no messages if not messages: cmd = ["python3", str(BIN_DIR / "box-chat.py"), "thread-messages", agent, thread_id, "--limit", str(limit)] - res = subprocess.run(cmd, capture_output=True, text=True) try: + res = subprocess.run(cmd, capture_output=True, text=True, timeout=15) data = json.loads(res.stdout) messages = data.get("messages", []) if data.get("ok") else [] except Exception: @@ -2502,6 +2502,63 @@ def cmd_pipeline_prune(args): print(json.dumps({"pruned": pruned, "count": len(pruned)}, indent=2)) +# --------------------------------------------------------------------------- +# Domain: DEPLOY (Multi-Agent Pipelines, Subagent Spawning, and Oversight) +# --------------------------------------------------------------------------- + +def cmd_deploy(args): + action = getattr(args, "action", None) + + if action == "subagent": + agent = args.agent + title = args.title or "subagent-task" + prompt = args.prompt + wait_val = getattr(args, "wait", 30) + wait = 30 if wait_val is None else int(wait_val) + + print(c_bold(f"\n=== DEPLOYING SUBAGENT ON [{agent}] ===")) + print(f" Agent: {c_cyan(agent)}") + print(f" Title: {c_yellow(title)}") + print(f" Prompt: {c_dim(prompt[:90])}...\n") + + import muse_hybrid + res, err = muse_hybrid.start_session(agent, title=title) + if err or not res: + print(c_err(f"✖ Failed to spawn subagent session: {err}"), file=sys.stderr) + sys.exit(1) + + session_id = res.get("session_id") + print(c_green(f"✔ Subagent session spawned: {session_id}")) + + if prompt: + print(f" Dispatching task prompt to subagent (waiting up to {wait}s)...") + send_res, send_err = muse_hybrid.send_message(agent, prompt, thread_id=session_id, wait=wait) + if send_err: + print(c_warn(f"Notice: {send_err}")) + else: + print(c_green("✔ Task prompt delivered.")) + + # Fetch fresh history to show assistant's initial work + msgs, _ = muse_hybrid.get_history(agent, thread_id=session_id, limit=5) + if msgs: + for m in msgs: + role = m.get("role") + if role == "assistant": + print(f"\n{c_bold('--- Subagent Response ---')}\n{m.get('text')}\n") + + print(f" Oversee anytime with: {c_cyan(f'box thread view {agent} {session_id}')}\n") + return + + # Pipeline deployment + pipe_name = getattr(args, "name", None) + if not pipe_name: + print(c_err("Error: Specify pipeline name (e.g. box deploy pipeline pipe-demo-step1)"), file=sys.stderr) + sys.exit(1) + + args.name = pipe_name + cmd_pipeline_run(args) + + # --------------------------------------------------------------------------- # Domain: LOOP (Intrinsic Loop Strategy, Health, Break Taxonomy, and Control) # --------------------------------------------------------------------------- @@ -3333,6 +3390,20 @@ def build_parser(): p_va_rb.add_argument("name", help="Variable name") p_va_rb.add_argument("--revision", default=None, help="Revision step (int) or timestamp") + # Domain: deploy (Deploy multi-agent pipelines or subagent tasks with oversight) + p_deploy = subparsers.add_parser("deploy", parents=[common], help="Deploy multi-agent task or pipeline across nodes with live oversight") + deploy_sub = p_deploy.add_subparsers(dest="action") + + p_dep_pipe = deploy_sub.add_parser("pipeline", parents=[common], help="Deploy multi-agent pipeline") + p_dep_pipe.add_argument("name", help="Pipeline root job name (e.g. pipe-demo-step1)") + p_dep_pipe.add_argument("--dry-run", action="store_true", help="Simulate pipeline without sending DMs") + + p_dep_sub = deploy_sub.add_parser("subagent", parents=[common], help="Spawn sub-agent session and dispatch task") + p_dep_sub.add_argument("--agent", required=True, choices=VALID_NODES, help="Agent node to spawn subagent on") + p_dep_sub.add_argument("--title", default="subagent-task", help="Title for the subagent session") + p_dep_sub.add_argument("prompt", help="Task prompt for the subagent") + p_dep_sub.add_argument("--wait", type=int, default=30, help="Seconds to wait for subagent response") + # Domain: muse (fast headless gateway via muse-cli-node with isolated per-node Cloudflare WARP egress) p_muse = subparsers.add_parser("muse", parents=[common], help="Direct headless gateway client (muse-cli-node)") p_muse.add_argument("node", choices=VALID_NODES, help="Target agent node") @@ -3506,6 +3577,8 @@ def main(): cmd_loop_strat(args) elif args.domain == "vars": cmd_loop_vars(args) + elif args.domain == "deploy": + cmd_deploy(args) elif args.domain == "muse": cmd = [str(BIN_DIR / "muse-cli-node"), args.node] + (args.muse_args or []) res = subprocess.run(cmd) diff --git a/docs/AGENT-TOOLING.md b/docs/AGENT-TOOLING.md new file mode 100644 index 0000000..85fb418 --- /dev/null +++ b/docs/AGENT-TOOLING.md @@ -0,0 +1,74 @@ +# Agent Tooling & Subagent Delegation Guide + +Welcome, operator. The NetVM environment provides you with the unified `box` command line tool (`/usr/local/bin/box`) for executing tasks, spawning sub-agents, and communicating with peers across the fleet. + +--- + +## 1. Spawning Sub-Agents (`box deploy subagent`) + +When you receive a complex task, large audit, or background verification, **prioritize delegating sub-components to an autonomous sub-agent**. + +```bash +box deploy subagent --agent <self> --title "<task-name>" "<prompt>" +``` + +### Parameters: +- `--agent`: Your own identity (`646`, `pip`, `opm`, or `muse`). +- `--title`: Brief descriptive title for the sub-agent session. +- `prompt`: The specific instructions and criteria for the sub-agent. + +### Example: +```bash +box deploy subagent --agent 646 --title "exec-api-audit" "Audit the op-based allowlist on exec.muse-dev.online and report back valid ops." +``` + +--- + +## 2. Cross-Operator DMs & Relay (`box dm send`) + +To coordinate with your peer operators (`opm`, `646`, `pip`, `muse`) or hand off tasks across sidechats: + +```bash +box dm send --agent <self> --to <recipient> --target <sidechat> "<message>" +``` + +### Target Sidechats: +- `646-pip`: Bilateral coordination thread between 646 and Pip. +- `646-opm-coord`: Coordination thread between 646 and OPM. +- `<agent> tasks`: Dedicated task sidechat for an individual agent (e.g. `646 tasks`, `pip tasks`). + +### Example: +```bash +box dm send --agent 646 --to pip --target 646-pip "Hey Pip, start-page onboarding review draft is ready. Please confirm when ready to review." +``` + +--- + +## 3. Direct Gateway Tooling (`box muse`) + +You can directly interact with the headless Muse gateway inside your isolated network namespace: + +```bash +# List all your active threads +box muse <self> threads + +# Read message history in a thread +box muse <self> history --thread <thread_uuid> --limit 10 + +# Check unread messages +box muse <self> unread + +# Start a new thread / session +box muse <self> session-start --title "<title>" + +# Fast fire-and-forget message send +box muse <self> send --thread <thread_uuid> "<message>" +``` + +--- + +## 4. Best Practices & Invariants + +1. **Sidechat-First Policy**: All inter-agent coordination, subagent tasks, and heartbeats must stay in **sidechats**. Do not send automated routine messages to `main` chat. +2. **Sub-Agent Prioritization**: Break down complex diagnostic or verification jobs by delegating sub-tasks to dedicated subagent threads. +3. **Attribution & Result Tagging**: For scheduled jobs, always conclude your response with `[RESULT <job_id>] <summary>`. diff --git a/job-sidechats.json b/job-sidechats.json index 5705df7..0cb7325 100644 --- a/job-sidechats.json +++ b/job-sidechats.json @@ -62,12 +62,24 @@ "agent": "646", "created_at": "2026-10-04T18:30:21.863666+00:00" }, + "646-tasks": { + "thread_uuid": "1dfb3199-2f99-446c-83e3-848ae2da0a12", + "agent": "646" + }, "pip tasks": { "thread_uuid": "5f33b1ab-a310-4cbb-9ccc-bb7bcd3d1b34", "agent": "pip", "description": "Dedicated pip task prompting sidechat", "created_at": "2026-10-04T22:59:13+00:00" }, + "pip-tasks": { + "thread_uuid": "5f33b1ab-a310-4cbb-9ccc-bb7bcd3d1b34", + "agent": "pip" + }, + "muse-tasks": { + "thread_uuid": "414cc0fc-4c46-4cb4-9fa6-07f06912384a", + "agent": "muse" + }, "pipe-c3f686": { "thread_uuid": "afa977ba-c75b-4b3e-8c1f-f83f13fb287e", "agent": "opm", @@ -98,5 +110,25 @@ "thread_uuid": "4466d0c1-7961-4cf3-b99d-1ab7c38484c2", "agent": "pip", "description": "pip side of 646-pip coordination pair" + }, + "pipe-a0d377": { + "thread_uuid": "cc0aee8f-531a-40af-9e02-27459f57b3ae", + "agent": "opm", + "created_at": "2026-10-04T23:15:34.981356+00:00" + }, + "pipe-3a153d": { + "thread_uuid": "29066a5c-878f-4d9d-b0f8-32405088c2d8", + "agent": "opm", + "created_at": "2026-10-04T23:16:16.726507+00:00" + }, + "pipe-921583": { + "thread_uuid": "70c1e8a6-e3f7-4ef2-9c32-8207b4f8c44f", + "agent": "opm", + "created_at": "2026-10-04T23:17:25.127456+00:00" + }, + "pipe-b3f7c4": { + "thread_uuid": "8f9ae8e7-dd3a-46e6-a5ac-3b1bb7cf6d66", + "agent": "opm", + "created_at": "2026-10-04T23:19:20.218056+00:00" } } \ No newline at end of file diff --git a/jobs/646-hourly-checkin.json b/jobs/646-hourly-checkin.json new file mode 100644 index 0000000..1de4476 --- /dev/null +++ b/jobs/646-hourly-checkin.json @@ -0,0 +1,17 @@ +{ + "name": "646-hourly-checkin", + "description": "Hourly operational check-in for 646 during active daytime hours", + "agent": "646", + "schedule": "0 8-22 * * *", + "timeout": 300, + "on_failure": "alert", + "dm_target": "646 tasks", + "sidechat": {"create": false}, + "followup": { + "expect_reply": true, + "timeout": "1h", + "nudges": 1, + "escalate": "opm" + }, + "prompt_template": "Hourly operational check-in for 646.\nJob ID: {job_id}\nTime: {datetime}\n\nPlease report briefly in this thread:\n(1) Current tasks in flight\n(2) VM / service health status\n(3) Any blockers or peer coordination items\n\nReply with [RESULT {job_id}] and your status." +} diff --git a/jobs/muse-hourly-checkin.json b/jobs/muse-hourly-checkin.json new file mode 100644 index 0000000..e9533e8 --- /dev/null +++ b/jobs/muse-hourly-checkin.json @@ -0,0 +1,17 @@ +{ + "name": "muse-hourly-checkin", + "description": "Hourly operational check-in for Muse during active daytime hours", + "agent": "muse", + "schedule": "0 8-22 * * *", + "timeout": 300, + "on_failure": "alert", + "dm_target": "muse tasks", + "sidechat": {"create": false}, + "followup": { + "expect_reply": true, + "timeout": "1h", + "nudges": 1, + "escalate": "opm" + }, + "prompt_template": "Hourly operational check-in for Muse.\nJob ID: {job_id}\nTime: {datetime}\n\nPlease report briefly in this thread:\n(1) Current tasks and codebase status\n(2) Verdicts / review items\n(3) Any blockers or peer coordination items\n\nReply with [RESULT {job_id}] and your status." +} diff --git a/jobs/pip-hourly-checkin.json b/jobs/pip-hourly-checkin.json new file mode 100644 index 0000000..89bc9ec --- /dev/null +++ b/jobs/pip-hourly-checkin.json @@ -0,0 +1,17 @@ +{ + "name": "pip-hourly-checkin", + "description": "Hourly operational check-in for Pip during active daytime hours", + "agent": "pip", + "schedule": "0 8-22 * * *", + "timeout": 300, + "on_failure": "alert", + "dm_target": "pip tasks", + "sidechat": {"create": false}, + "followup": { + "expect_reply": true, + "timeout": "1h", + "nudges": 1, + "escalate": "opm" + }, + "prompt_template": "Hourly operational check-in for Pip.\nJob ID: {job_id}\nTime: {datetime}\n\nPlease report briefly in this thread:\n(1) Current tasks and workspace status\n(2) Blockers (start-page draft, authorization, or netns items)\n(3) Readiness / next steps\n\nReply with [RESULT {job_id}] and your status." +} diff --git a/systemd/job-scheduler.service b/systemd/job-scheduler.service new file mode 100644 index 0000000..188f400 --- /dev/null +++ b/systemd/job-scheduler.service @@ -0,0 +1,10 @@ +[Unit] +Description=NetVM Unified Job Scheduler +After=network.target + +[Service] +Type=oneshot +ExecStart=/usr/bin/python3 /home/super/Projects/NetVM/bin/job-scheduler.py run +WorkingDirectory=/home/super/Projects/NetVM +StandardOutput=journal +StandardError=journal diff --git a/systemd/job-scheduler.timer b/systemd/job-scheduler.timer new file mode 100644 index 0000000..596d84a --- /dev/null +++ b/systemd/job-scheduler.timer @@ -0,0 +1,9 @@ +[Unit] +Description=Run NetVM Unified Job Scheduler every 5 minutes + +[Timer] +OnCalendar=*:0/5 +Persistent=true + +[Install] +WantedBy=timers.target