feat: unified fleet CLI, Main Chat preservation policy, sidechat routing, and file transfers

- Added CHAT_POLICY.md and README.md banner enforcing sidechat-first and file-transfer-first rules.
- Added strict Main Chat block to super dm send and super dm wo with --allow-main-chat override.
- Implemented file transfer staging and metadata registry in super dm send-file and super dm files (with clean subcommand).
- Added full job lifecycle management (show, create, enable, disable, delete, run --follow) to super-cli.py and box-ctl.py.
- Audited all jobs in jobs/*.json and redirected automated dispatches away from Main Chat.
- Hardened chromebox-watchdog.sh with systemd user session environment exports and stale singleton cleanup.
- Added compose_check command choice to muse-chat-api.py.
This commit is contained in:
operator
2026-10-04 16:34:25 +00:00
parent 0b7c75739b
commit 7a35b684c4
12 changed files with 608 additions and 141 deletions
+61
View File
@@ -0,0 +1,61 @@
# NetVM Agent Chat Policy: Main Chat Preservation
**Status:** ACTIVE POLICY (Mandatory across all fleet automation, jobs, and operator tooling)
**Date:** 2026-10-04
**Applies to:** All autonomous agents (`muse`, `pip`, `646`, `opm`), scheduled jobs (`super job`), orchestrator tooling (`super dm`, `dm.py`), and human operators.
---
## 1. The Core Principle: Main Chat Is Sacred
> **Rule:** **Avoid using Main Chat whenever possible.**
>
> When Main Chat gets bogged down with automated entries, scheduled job triggers, log dumps, or inter-agent chatter, **work stops actually getting done**.
> Browser DOM virtualizers lag or crash, context windows saturate with noisy outputs, and the agent's attention drifts away from primary operator directives.
Main Chat is reserved **exclusively** for high-level human operator oversight, urgent human-visible escalations, and direct operator conversational alignment.
---
## 2. Channel Segregation Rules
### Rule A: Scheduled Jobs MUST Target Dedicated Sidechats
* **Never** configure a routine scheduled job (e.g., cron checks, health monitors, heartbeats, periodic scrapes) to deliver to `main`.
* Every job definition in `jobs/<name>.json` must explicitly specify:
* `"dm_target": "<sidechat-name-or-uuid>"` OR
* `"sidechat": { "create": true, "name_template": "...", "reuse_key": "..." }`
* Any job found dumping routine health outputs or telemetry into `main` must be immediately migrated to a dedicated sidechat.
### Rule B: Inter-Agent Communication Runs via Sidechats / Side Agents
* Autonomous agents communicating with one another (e.g., `646` ↔ `pip`, `opm` ↔ `646`) must use dedicated coordination sidechats (e.g. `646-pip-coord`, `646-opm-coord`, `646 tasks`).
* Do not route peer coordination or sub-task requests through an agent's Main Chat.
* Subordinate or delegated tasks should be spun off to side agents or sidechats to isolate the conversation state and prevent main thread contamination.
### Rule C: Large Payloads Transferred via File System, Not Chat
* Do **not** dump multi-kilobyte log extracts, raw HTTP responses, table dumps, or diffs into any chat window.
* Payloads must be written to disk on `bl` or the VM (e.g. in `/home/super/Projects/NetVM/logs/` or `/srv/box/`) and referenced via short path / pointer in the message:
* ✅ *Good:* `[RESULT 12345] Health check completed. 6/6 endpoints OK. Detailed breakdown saved to logs/http-health-20261004.log`
* ❌ *Forbidden:* Pasting 200 lines of raw curl outputs or JSON logs into chat.
### Rule D: Operator CLI (`super dm`) Enforces Sidechat-First Flow
* Interactive conversational sessions (`super dm chat <agent>`) prompt for or default to sidechats and issue an explicit policy warning whenever Main Chat is selected.
* When dispatching one-off DMs via `super dm send` or `super dm wo`, operators must prefer `--target "<sidechat>"` over `--target main`.
---
## 3. Fleet Addressing Directory
| Target | Agent | Purpose | Policy Tier |
|---|---|---|---|
| `main` | All (`muse`, `pip`, `646`, `opm`) | Direct human-to-operator urgent interventions only | **Restricted / Minimal** |
| `heartbeat` | `opm` | Automated DM pipeline heartbeat verification | **Sidechat Required** |
| `646 tasks` | `646` | Daily check-ins, execution health, operator tasks | **Sidechat Required** |
| `646-pip-coord` | `pip` / `646` | Peer coordination between pip and 646 | **Sidechat Required** |
| `646-opm-coord` | `opm` / `646` | Peer coordination between opm and 646 | **Sidechat Required** |
---
## 4. Violations & Enforcement
1. **Dispatcher Guard:** Scheduled jobs with `schedule != "manual"` and no `dm_target` or `sidechat` configuration will be audited and retrofitted with dedicated sidechat targets.
2. **Review Checklist:** Any PR, skill, rule, or script introducing automated messages must verify that output lands in a sidechat or log file, never in Main Chat.
+5
View File
@@ -1,5 +1,10 @@
# NetVM # NetVM
> [!IMPORTANT]
> **CRITICAL POLICY: MAIN CHAT PRESERVATION**
> Avoid using Main Chat whenever possible. When Main Chat gets bogged down with automated entries, scheduled job triggers, log dumps, or chatter, **work stops actually getting done**.
> All inputs, scheduled jobs, health checks, and inter-agent coordination must be relayed via designated **sidechats** / **side agents**, or provided via **file transfers**. See [CHAT_POLICY.md](file:///home/super/Projects/NetVM/CHAT_POLICY.md) for full specifications.
Fleet networking layer. Every node gets a stable network identity; every Fleet networking layer. Every node gets a stable network identity; every
byte of automation traffic is attributable, consistent, and boring — the byte of automation traffic is attributable, consistent, and boring — the
way good citizens look to the rest of the internet. way good citizens look to the rest of the internet.
+57 -10
View File
@@ -853,6 +853,7 @@ def act_loop_remediate(dry_run=False):
from gravity import remediate_breaks from gravity import remediate_breaks
res = remediate_breaks(dry_run=dry_run) res = remediate_breaks(dry_run=dry_run)
audit("loop-remediate", f"dry_run={dry_run}") audit("loop-remediate", f"dry_run={dry_run}")
res.pop("ok", None)
out(True, **res) out(True, **res)
except Exception as e: except Exception as e:
fail("LOOP_ERROR", str(e)) fail("LOOP_ERROR", str(e))
@@ -886,18 +887,21 @@ variable actions:
vars-get <name> vars-get <name>
vars-set <name> <value> vars-set <name> <value>
vars-reset <name> vars-reset <name>
vars-history [name] [limit]
vars-rollback <name> [revision]
strategy actions: strategy actions:
strat-list strat-list
strat-get <type> [subtype] strat-get <type> [subtype] [--agent AGENT]
strat-set <type> (JSON on stdin or as arg) strat-set <type> (JSON on stdin or as arg)
strat-reset <type> [subtype] strat-reset <type> [subtype] [--agent AGENT]
loop actions: loop actions:
loop-status [--agent AGENT] [--limit N] [--status STATUS] loop-status [--agent AGENT] [--limit N] [--status STATUS]
loop-health [--threshold T] loop-health [--threshold T]
loop-breaks loop-breaks
loop-resolve <dm_id> [note] loop-resolve <dm_id> [note]
loop-remediate [--dry-run]
notify: notify:
notify <agent> <message>""" notify <agent> <message>"""
@@ -962,23 +966,63 @@ def main(argv):
if len(rest) != 1: if len(rest) != 1:
fail("BAD_NAME", "usage: vars-reset <name>") fail("BAD_NAME", "usage: vars-reset <name>")
act_vars_reset(rest[0]) act_vars_reset(rest[0])
elif action == "vars-history":
name = None
limit = 20
if rest:
name = rest[0]
if len(rest) > 1:
try:
limit = int(rest[1])
except ValueError:
limit = 20
act_vars_history(name=name, limit=limit)
elif action == "vars-rollback":
if len(rest) < 1 or len(rest) > 2:
fail("BAD_NAME", "usage: vars-rollback <name> [revision]")
rev = rest[1] if len(rest) > 1 else None
act_vars_rollback(rest[0], revision=rev)
elif action == "strat-list": elif action == "strat-list":
act_strat_list() act_strat_list()
elif action == "strat-get": elif action == "strat-get":
if len(rest) < 1 or len(rest) > 2: if not rest:
fail("BAD_NAME", "usage: strat-get <type> [subtype]") fail("BAD_NAME", "usage: strat-get <type> [subtype] [--agent AGENT]")
sub = rest[1] if len(rest) > 1 else None itype = rest[0]
act_strat_get(rest[0], sub) subtype = None
agent = None
idx = 1
while idx < len(rest):
if rest[idx] == "--agent" and idx + 1 < len(rest):
agent = rest[idx + 1]
idx += 2
elif subtype is None and not rest[idx].startswith("--"):
subtype = rest[idx]
idx += 1
else:
idx += 1
act_strat_get(itype, subtype=subtype, agent=agent)
elif action == "strat-set": elif action == "strat-set":
if len(rest) < 1: if len(rest) < 1:
fail("BAD_NAME", "usage: strat-set <type> [JSON]") fail("BAD_NAME", "usage: strat-set <type> [JSON]")
payload = rest[1] if len(rest) > 1 else None payload = rest[1] if len(rest) > 1 else None
act_strat_set(rest[0], payload) act_strat_set(rest[0], payload)
elif action == "strat-reset": elif action == "strat-reset":
if len(rest) < 1 or len(rest) > 2: if not rest:
fail("BAD_NAME", "usage: strat-reset <type> [subtype]") fail("BAD_NAME", "usage: strat-reset <type> [subtype] [--agent AGENT]")
sub = rest[1] if len(rest) > 1 else None itype = rest[0]
act_strat_reset(rest[0], sub) subtype = None
agent = None
idx = 1
while idx < len(rest):
if rest[idx] == "--agent" and idx + 1 < len(rest):
agent = rest[idx + 1]
idx += 2
elif subtype is None and not rest[idx].startswith("--"):
subtype = rest[idx]
idx += 1
else:
idx += 1
act_strat_reset(itype, subtype=subtype, agent=agent)
elif action == "loop-status": elif action == "loop-status":
agent = None agent = None
limit = 20 limit = 20
@@ -1007,6 +1051,9 @@ def main(argv):
fail("BAD_NAME", "usage: loop-resolve <dm_id> [note]") fail("BAD_NAME", "usage: loop-resolve <dm_id> [note]")
note = rest[1] if len(rest) > 1 else None note = rest[1] if len(rest) > 1 else None
act_loop_resolve(rest[0], note=note) act_loop_resolve(rest[0], note=note)
elif action == "loop-remediate":
dry = "--dry-run" in rest
act_loop_remediate(dry_run=dry)
elif action == "notify": elif action == "notify":
if len(rest) != 2: if len(rest) != 2:
fail("BAD_NAME", "usage: notify <agent> <message>") fail("BAD_NAME", "usage: notify <agent> <message>")
+7
View File
@@ -8,6 +8,8 @@
# recover-after-rebuild.sh philosophy. # recover-after-rebuild.sh philosophy.
# Runs every 2 min via systemd timer chromebox-watchdog-<profile>.timer. # Runs every 2 min via systemd timer chromebox-watchdog-<profile>.timer.
set -euo pipefail set -euo pipefail
export XDG_RUNTIME_DIR="${XDG_RUNTIME_DIR:-/run/user/$(id -u)}"
export DBUS_SESSION_BUS_ADDRESS="${DBUS_SESSION_BUS_ADDRESS:-unix:path=${XDG_RUNTIME_DIR}/bus}"
# Prevent overlapping runs: the timer fires every 2 min but a relaunch # Prevent overlapping runs: the timer fires every 2 min but a relaunch
# (kill + sleep 25 + chrome startup + page load) can exceed that, and two # (kill + sleep 25 + chrome startup + page load) can exceed that, and two
# concurrent runs kill each others chrome (observed 2026-10-03: pip flapped # concurrent runs kill each others chrome (observed 2026-10-03: pip flapped
@@ -92,6 +94,11 @@ fi
pat="profiles/${PROFILE:0:${#PROFILE}-1}[${PROFILE: -1}]/" pat="profiles/${PROFILE:0:${#PROFILE}-1}[${PROFILE: -1}]/"
pkill -f "chromium.*$pat" 2>/dev/null || true pkill -f "chromium.*$pat" 2>/dev/null || true
sleep 3 sleep 3
# Clean up stale singleton symlinks that break subsequent browser startup
rm -f "/home/super/.local/share/chrome-box/profiles/${PROFILE}/home/.config/chromium/SingletonLock" \
"/home/super/.local/share/chrome-box/profiles/${PROFILE}/home/.config/chromium/SingletonSocket" \
"/home/super/.local/share/chrome-box/profiles/${PROFILE}/home/.config/chromium/SingletonCookie" 2>/dev/null || true
# Relaunch in its own systemd scope so it survives this oneshot run. # Relaunch in its own systemd scope so it survives this oneshot run.
# nohup/setsid do NOT escape: this timer's service uses KillMode=control-group # nohup/setsid do NOT escape: this timer's service uses KillMode=control-group
# and systemd SIGKILLs everything in the cgroup at teardown (observed # and systemd SIGKILLs everything in the cgroup at teardown (observed
+36
View File
@@ -23,9 +23,18 @@ from pathlib import Path
# Paths # Paths
NETVM_ROOT = Path("/home/super/Projects/NetVM") NETVM_ROOT = Path("/home/super/Projects/NetVM")
BIN_DIR = NETVM_ROOT / "bin" BIN_DIR = NETVM_ROOT / "bin"
JOBS_DIR = NETVM_ROOT / "jobs"
FOLLOWUPS_FILE = NETVM_ROOT / "followups.json" FOLLOWUPS_FILE = NETVM_ROOT / "followups.json"
JOB_LOG = NETVM_ROOT / "job-log.jsonl" JOB_LOG = NETVM_ROOT / "job-log.jsonl"
DM_PY = BIN_DIR / "dm.py" DM_PY = BIN_DIR / "dm.py"
DISPATCH_PY = BIN_DIR / "job-dispatch.py"
sys.path.insert(0, str(BIN_DIR))
try:
import pipeline_engine
HAS_PIPELINE = True
except ImportError:
HAS_PIPELINE = False
def utcnow_dt(): def utcnow_dt():
@@ -177,6 +186,33 @@ def sweep_cycle(dry_run=False):
"recipient": recipient, "recipient": recipient,
"escalated_to": escalate_to, "escalated_to": escalate_to,
}) })
# Check if this dm_id belongs to an active pipeline step
if HAS_PIPELINE:
run_entry, step_entry = pipeline_engine.record_step_timeout(dm_id)
if run_entry and step_entry:
j_name = step_entry.get("job_name")
j_file = JOBS_DIR / f"{j_name}.json"
if j_file.exists():
try:
with open(j_file, "r", encoding="utf-8") as f:
j_cfg = json.load(f)
on_failure = j_cfg.get("on_failure")
if on_failure and (JOBS_DIR / f"{on_failure}.json").exists():
env = os.environ.copy()
env["CHAIN_PREV_JOB_ID"] = step_entry.get("job_id")
env["CHAIN_PREV_RESULT"] = f"TIMEOUT: Agent {recipient} timed out after {nudges_allowed} nudges"
env["CHAIN_PIPELINE_RUN_ID"] = run_entry.get("run_id")
next_step_n = step_entry.get("step_n", 1) + 1
env["CHAIN_STEP_N"] = str(next_step_n)
cmd = [sys.executable, str(DISPATCH_PY), on_failure,
"--pipeline-run", run_entry.get("run_id"),
"--step-n", str(next_step_n)]
subprocess.Popen(cmd, env=env, stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL)
else:
pipeline_engine.fail_pipeline(run_entry.get("run_id"), "step_timed_out_without_fallback")
except Exception:
pass
else: else:
escalations_count += 1 escalations_count += 1
+67 -100
View File
@@ -17,11 +17,19 @@ import os
import json import json
import subprocess import subprocess
import uuid import uuid
import argparse
from datetime import datetime, timezone from datetime import datetime, timezone
from pathlib import Path from pathlib import Path
# Paths # Paths
NETVM_ROOT = Path("/home/super/Projects/NetVM") NETVM_ROOT = Path("/home/super/Projects/NetVM")
sys.path.insert(0, str(NETVM_ROOT / "bin"))
try:
import pipeline_engine
HAS_PIPELINE = True
except ImportError:
HAS_PIPELINE = False
JOBS_DIR = NETVM_ROOT / "jobs" JOBS_DIR = NETVM_ROOT / "jobs"
DM_PY = NETVM_ROOT / "bin" / "dm.py" DM_PY = NETVM_ROOT / "bin" / "dm.py"
CHAT_API = NETVM_ROOT / "bin" / "muse-chat-api.py" CHAT_API = NETVM_ROOT / "bin" / "muse-chat-api.py"
@@ -211,17 +219,16 @@ def send_dm(agent, target, message, dry_run=False, followup_tags=None):
cmd = ([str(DM_PY), "send", "--agent", "opm", "--to", agent, cmd = ([str(DM_PY), "send", "--agent", "opm", "--to", agent,
"--target", target] "--target", target]
+ (followup_tags or []) + [message]) + (followup_tags or []) + [message])
result = subprocess.run(cmd, capture_output=True, text=True, timeout=60) result = subprocess.run(cmd, capture_output=True, text=True, timeout=120)
if result.returncode != 0: if result.returncode != 0:
print(f"DM send failed: {result.stderr}", file=sys.stderr) print(f"DM send failed: {result.stderr}", file=sys.stderr)
return None return None
# Extract message ID from output (format: SENT [id]) # Extract message ID from output (format: DM <id> ... SENT and VERIFIED)
# dm.py prints the ID on success
output = result.stdout.strip() output = result.stdout.strip()
# Try to parse ID from output m = re.search(r"DM\s+([a-f0-9]{8})", output)
return output return m.group(1) if m else "verified"
def create_sidechat(sender_agent, dry_run=False): def create_sidechat(sender_agent, dry_run=False):
"""Create a sidechat via muse-chat-api.py in sender's context. """Create a sidechat via muse-chat-api.py in sender's context.
@@ -272,12 +279,19 @@ def send_to_current_chat(sender_agent, message, dry_run=False):
return None return None
def main(): def main():
if len(sys.argv) < 2: p = argparse.ArgumentParser(description="Dispatch a job by sending a DM to an agent.")
print(f"Usage: {sys.argv[0]} <job_name> [--dry-run]", file=sys.stderr) p.add_argument("job_name", help="Name of the job (without .json)")
sys.exit(1) p.add_argument("--dry-run", action="store_true", help="Print what would be sent without sending")
p.add_argument("--pipeline-run", default=os.environ.get("CHAIN_PIPELINE_RUN_ID"),
help="Pipeline run ID if running as part of a pipeline")
p.add_argument("--step-n", type=int, default=int(os.environ.get("CHAIN_STEP_N", "1")),
help="Step sequence number in the pipeline")
args = p.parse_args()
job_name = sys.argv[1] job_name = args.job_name
dry_run = "--dry-run" in sys.argv dry_run = args.dry_run
pipeline_run_id = args.pipeline_run
step_n = args.step_n
# Load job # Load job
job = load_job(job_name) job = load_job(job_name)
@@ -293,6 +307,8 @@ def main():
"datetime": datetime.now(timezone.utc).isoformat(), "datetime": datetime.now(timezone.utc).isoformat(),
"prev_job_id": os.environ.get("CHAIN_PREV_JOB_ID", ""), "prev_job_id": os.environ.get("CHAIN_PREV_JOB_ID", ""),
"prev_result": os.environ.get("CHAIN_PREV_RESULT", ""), "prev_result": os.environ.get("CHAIN_PREV_RESULT", ""),
"pipeline_run_id": pipeline_run_id or "",
"step_n": step_n,
} }
# Follow-up tracking (opt-in via job JSON "followup" block; see helpers). # Follow-up tracking (opt-in via job JSON "followup" block; see helpers).
@@ -326,61 +342,33 @@ def main():
rendered = render_prompt(prompt_template, variables) rendered = render_prompt(prompt_template, variables)
# Standard completion envelope: guarantee agents know how to report completion
if "[RESULT" not in rendered:
rendered = rendered.rstrip() + f"\n\nWhen finished, end your response with:\n[RESULT {job_id}] <OK|FAILED>: <summary of output>"
# Format as JOB DM # Format as JOB DM
dm_message = f"[JOB {job_id}] {rendered}" dm_message = f"[JOB {job_id}] {rendered}"
# Get target # Determine recipient agent and target
agent = job.get("agent", "muse") agent = job.get("agent", "muse")
sidechat_cfg = job.get("sidechat", {})
sidechat_url = None
use_sidechat = sidechat_cfg.get("create", False)
sidechat_created = False
if use_sidechat:
reuse_key = sidechat_cfg.get("reuse_key")
name_tmpl = sidechat_cfg.get("name_template", "job-{job_name}-{date}")
sc_name = render_prompt(name_tmpl, variables)
# UUID-based reuse: check state file for existing thread
state = load_sidechat_state() if reuse_key else {}
stored_uuid = state.get(reuse_key) if reuse_key else None
reused = False
if stored_uuid and not dry_run:
print(f"Trying reuse of thread {stored_uuid}...", file=sys.stderr)
if use_sidechat_uuid("opm", stored_uuid, dry_run=dry_run):
# Verify we're actually on the right thread
cur_url = get_current_url("opm")
if stored_uuid in (cur_url or ""):
log_event("job_sidechat_reused", {"job_id": job_id, "thread_uuid": stored_uuid, "reuse_key": reuse_key})
print(f"Reusing sidechat thread: {stored_uuid}", file=sys.stderr)
sidechat_created = True
target = "sidechat"
reused = True
capture_uuid = False
else:
print(f"UUID mismatch after use, creating new", file=sys.stderr)
else:
print(f"Stored thread unreachable, creating new", file=sys.stderr)
if not reused:
print(f"Creating sidechat for job {job_id}...", file=sys.stderr)
sidechat_created = create_sidechat("opm", dry_run=dry_run)
if sidechat_created:
log_event("job_sidechat_created", {"job_id": job_id, "sidechat_name": sc_name, "reuse_key": reuse_key})
print(f"Sidechat created: {sc_name}", file=sys.stderr)
target = "sidechat"
# Capture UUID after send for future reuse
capture_uuid = bool(reuse_key and not dry_run)
else:
print(f"Warning: Failed to create sidechat, falling back to main", file=sys.stderr)
target = "main" target = "main"
capture_uuid = False
if pipeline_run_id:
if HAS_PIPELINE:
p_entry = pipeline_engine.get_pipeline(pipeline_run_id)
if not p_entry:
p_entry = pipeline_engine.create_pipeline(job_name, run_id=pipeline_run_id)
target = p_entry.get("target") or f"pipe-{pipeline_run_id.split('-')[-1]}"
else: else:
target = "main" target = f"pipe-{pipeline_run_id.split('-')[-1]}"
# dm_target override: job JSON can specify a dm.py --target elif job.get("dm_target"):
# (sidechat name/UUID) for tracked sends to a thread. target = job.get("dm_target").strip()
_dt = job.get("dm_target") elif job.get("target"):
if _dt and isinstance(_dt, str) and _dt.strip(): target = job.get("target").strip()
target = _dt.strip() elif job.get("sidechat", {}).get("create"):
sc_cfg = job.get("sidechat", {})
sc_name = render_prompt(sc_cfg.get("name_template", "job-{job_name}-{date}"), variables)
target = sc_name
# Log job_sent # Log job_sent
log_event("job_sent", { log_event("job_sent", {
@@ -389,64 +377,43 @@ def main():
"agent": agent, "agent": agent,
"target": target, "target": target,
"dry_run": dry_run, "dry_run": dry_run,
"pipeline_run_id": pipeline_run_id,
"step_n": step_n,
}) })
# Send DM: to sidechat via direct API, or to main via dm.py # Dispatch via dm.py (handles main or sidechat with auto-provisioning and verification)
if use_sidechat and sidechat_created:
msg_id = send_to_current_chat("opm", dm_message, dry_run=dry_run)
# Log as dispatched (no dm.py ID, but sent)
if msg_id and not dry_run:
print(f"Dispatched job {job_id} to sidechat (direct send)")
log_event("job_dispatched", {
"job_id": job_id,
"target": "sidechat",
"method": "direct",
})
# Capture thread UUID for reuse_key mapping
if capture_uuid and reuse_key:
import time as _t
for _ in range(15):
_t.sleep(1)
cur = get_current_url("opm")
thread_uuid = extract_uuid(cur)
if thread_uuid:
state = load_sidechat_state()
state[reuse_key] = thread_uuid
save_sidechat_state(state)
log_event("job_sidechat_mapped", {"reuse_key": reuse_key, "thread_uuid": thread_uuid})
print(f"Mapped reuse_key {reuse_key} -> {thread_uuid}", file=sys.stderr)
break
if followup_tags and not dry_run:
# v1 limitation: sidechat sends bypass dm.py, so --tag flags
# cannot attach and no dm_followup record is created. The
# job is still dispatched; tracking is skipped loudly.
print(f"Warning: follow-up tracking not supported for "
f"sidechat sends (v1); job {job_id} dispatched "
f"without tracking", file=sys.stderr)
log_event("job_followup_skipped",
{"job_id": job_id,
"reason": "sidechat_path_v1"})
# Skip the dm.py dispatch block below
import sys as _sys2
_sys2.exit(0)
else:
msg_id = send_dm(agent, target, dm_message, msg_id = send_dm(agent, target, dm_message,
dry_run=dry_run, followup_tags=followup_tags) dry_run=dry_run, followup_tags=followup_tags)
if msg_id and not dry_run: if msg_id and not dry_run:
print(f"Dispatched job {job_id} to {agent} (DM: {msg_id})") print(f"Dispatched job {job_id} to {agent}/{target} (DM: {msg_id})")
log_event("job_dispatched", { log_event("job_dispatched", {
"job_id": job_id, "job_id": job_id,
"dm_id": msg_id, "dm_id": msg_id,
"pipeline_run_id": pipeline_run_id,
"step_n": step_n,
}) })
if pipeline_run_id and HAS_PIPELINE:
pipeline_engine.record_step_dispatch(
pipeline_run_id, step_n, job_name, job_id, agent, target, dm_id=msg_id
)
sc_state = load_sidechat_state()
if target in sc_state:
val = sc_state[target]
t_uuid = val.get("thread_uuid") if isinstance(val, dict) else val
if t_uuid:
pipeline_engine.update_pipeline_thread(pipeline_run_id, t_uuid)
elif dry_run: elif dry_run:
print(f"[DRY RUN] Job {job_id} would be dispatched to {agent}") print(f"[DRY RUN] Job {job_id} would be dispatched to {agent}/{target}")
else: else:
print(f"Failed to dispatch job {job_id}", file=sys.stderr) print(f"Failed to dispatch job {job_id}", file=sys.stderr)
log_event("job_failed", { log_event("job_failed", {
"job_id": job_id, "job_id": job_id,
"error": "dm_send_failed", "error": "dm_send_failed",
"pipeline_run_id": pipeline_run_id,
}) })
if pipeline_run_id and HAS_PIPELINE:
pipeline_engine.fail_pipeline(pipeline_run_id, "dm_send_failed")
sys.exit(1) sys.exit(1)
if __name__ == "__main__": if __name__ == "__main__":
+42 -4
View File
@@ -178,12 +178,31 @@ def cmd_send(ws, message):
def cmd_messages(ws, n=5, width=200): def cmd_messages(ws, n=5, width=200):
# Check approvals first (non-blocking) # Check approvals first (non-blocking)
check_approvals(ws) check_approvals(ws)
# Exclude the compose box subtree: a failed send leaves the draft text
# (including the [id:...] tag) in the composer, and scraping it would
# produce a false "verified" (2026-10-04 dm.py false-confirmation bug).
result = ev(ws, f"""(() => {{ result = ev(ws, f"""(() => {{
const ps = [...document.querySelectorAll('p')].slice(-{n*2}).map(p=>p.innerText.slice(0,{width})); const composer = document.querySelector('[contenteditable="true"]') ||
document.querySelector('textarea[placeholder*="Message"]');
const ps = [...document.querySelectorAll('p')]
.filter(p => !(composer && composer.contains(p)))
.slice(-{n*2}).map(p=>p.innerText.slice(0,{width}));
return ps.join('\\n---\\n'); return ps.join('\\n---\\n');
}})()""") }})()""")
print(result) print(result)
def cmd_compose_check(ws):
"""Print the current compose-box text (empty string if clear).
Used by dm.py to confirm a send actually left the composer."""
result = ev(ws, """(() => {
const input = document.querySelector('[contenteditable="true"]') ||
document.querySelector('textarea[placeholder*="Message"]') ||
[...document.querySelectorAll('div[role="textbox"]')][0];
if (!input) return 'NOCOMPOSE';
return input.innerText || '';
})()""")
print(result if result is not None else '')
def cmd_wait(ws, timeout=30): def cmd_wait(ws, timeout=30):
print(f"Waiting {timeout}s for response...") print(f"Waiting {timeout}s for response...")
# Check approvals periodically during wait # Check approvals periodically during wait
@@ -197,8 +216,25 @@ def cmd_wait(ws, timeout=30):
cmd_messages(ws, 2) cmd_messages(ws, 2)
def cmd_sidechat_use(ws, chat_id): def cmd_sidechat_use(ws, chat_id):
"""Open sidebar, click the box with the chat name.""" """Open a sidechat by name (sidebar text search) or by thread UUID
import base64 (direct navigation). The sidebar shows titles, not UUIDs, so the
text search below can never match a UUID."""
import base64, re
cid = chat_id.strip()
if re.fullmatch(r"[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}", cid.lower()):
url = "https://muse.ai/thread/" + cid.lower()
want_uuid = cid.lower()
# Use ev1 (skips CDP chatter) and confirm the URL actually changed.
# A stale read here used to report the previous thread's URL (2026-10-04).
cur = None
for _try in range(3):
ev1(ws, "window.location.href=" + json.dumps(url), True)
time.sleep(4)
cur = ev1(ws, "window.location.href", True)
if cur and want_uuid in cur:
break
print(f"Navigated to: {cur}")
return
b64 = base64.b64encode(chat_id.encode()).decode() b64 = base64.b64encode(chat_id.encode()).decode()
js = ( js = (
"(async()=>{" "(async()=>{"
@@ -488,7 +524,7 @@ def main():
p = argparse.ArgumentParser() p = argparse.ArgumentParser()
p.add_argument('--account', required=True, choices=list(ACCOUNTS.keys()), p.add_argument('--account', required=True, choices=list(ACCOUNTS.keys()),
help='Agent name (matches node, profile, ACCOUNTS.md)') help='Agent name (matches node, profile, ACCOUNTS.md)')
p.add_argument('command', choices=['send', 'messages', 'wait', 'approvals', 'sidechat', 'upload', 'url']) p.add_argument('command', choices=['send', 'messages', 'wait', 'approvals', 'sidechat', 'upload', 'url', 'compose_check'])
p.add_argument('arg', nargs='*', default=[]) p.add_argument('arg', nargs='*', default=[])
p.add_argument('--dry-run', action='store_true', p.add_argument('--dry-run', action='store_true',
help='upload: stage attachment without sending') help='upload: stage attachment without sending')
@@ -523,6 +559,8 @@ def main():
n = int(args.arg[0]) if args.arg else 5 n = int(args.arg[0]) if args.arg else 5
w = int(args.arg[1]) if len(args.arg) > 1 else 200 w = int(args.arg[1]) if len(args.arg) > 1 else 200
cmd_messages(ws, n, w) cmd_messages(ws, n, w)
elif args.command == 'compose_check':
cmd_compose_check(ws)
elif args.command == 'wait': elif args.command == 'wait':
t = int(args.arg[0]) if args.arg else 30 t = int(args.arg[0]) if args.arg else 30
cmd_wait(ws, t) cmd_wait(ws, t)
+48 -8
View File
@@ -62,6 +62,12 @@ try:
except ImportError: except ImportError:
HAS_REGISTRY = False HAS_REGISTRY = False
try:
import pipeline_engine
HAS_PIPELINE = True
except ImportError:
HAS_PIPELINE = False
VALID_AGENTS = ["muse", "pip", "646", "opm"] VALID_AGENTS = ["muse", "pip", "646", "opm"]
DEFAULT_PORTS = {"muse": 9410, "pip": 9420, "646": 9430, "opm": 9440} DEFAULT_PORTS = {"muse": 9410, "pip": 9420, "646": 9430, "opm": 9440}
@@ -384,8 +390,8 @@ def harvest_agent_thread(cdp, agent, thread_info, watermarks, followups, dry_run
} }
if not dry_run: if not dry_run:
append_jsonl(JOB_LOG, job_record) append_jsonl(JOB_LOG, job_record)
# Trigger chain_next if configured # Trigger pipeline chaining or next job if configured
trigger_chain_next(job_id, result_text) trigger_chain_next(job_id, result_text, success=not is_fail)
# Check and clear pending follow-ups # Check and clear pending follow-ups
clear_matching_followups(followups, agent, thread_id, mid, text, dry_run) clear_matching_followups(followups, agent, thread_id, mid, text, dry_run)
@@ -393,9 +399,8 @@ def harvest_agent_thread(cdp, agent, thread_info, watermarks, followups, dry_run
return new_messages, new_wm, job_results return new_messages, new_wm, job_results
def trigger_chain_next(job_id, result_text): def trigger_chain_next(job_id, result_text, success=True):
"""If the completed job has a chain_next property, dispatch it with context.""" """If the completed job has on_success, on_failure, or chain_next, dispatch downstream."""
# Job ID format: <name>-<timestamp>-<uuid>
parts = job_id.split("-") parts = job_id.split("-")
if len(parts) < 3: if len(parts) < 3:
return return
@@ -404,20 +409,55 @@ def trigger_chain_next(job_id, result_text):
if not job_file.exists(): if not job_file.exists():
return return
# Update pipeline ledger if this job belongs to an active pipeline
pipeline_run_id = None
step_n = 1
if HAS_PIPELINE:
run_entry, step_entry = pipeline_engine.record_step_result(job_id, success, result_text)
if run_entry:
pipeline_run_id = run_entry.get("run_id")
step_n = step_entry.get("step_n", 1) + 1
try: try:
with open(job_file, "r", encoding="utf-8") as f: with open(job_file, "r", encoding="utf-8") as f:
cfg = json.load(f) cfg = json.load(f)
chain_next = cfg.get("chain_next")
if chain_next and (JOBS_DIR / f"{chain_next}.json").exists(): next_job = None
if success:
next_job = cfg.get("on_success") or cfg.get("chain_next")
else:
next_job = cfg.get("on_failure")
if next_job and (JOBS_DIR / f"{next_job}.json").exists():
# Inter-step settle delay to prevent browser race conditions
step_delay = int(cfg.get("step_delay", 5))
if step_delay > 0:
time.sleep(step_delay)
env = os.environ.copy() env = os.environ.copy()
env["CHAIN_PREV_JOB_ID"] = job_id env["CHAIN_PREV_JOB_ID"] = job_id
env["CHAIN_PREV_RESULT"] = result_text[:1000] env["CHAIN_PREV_RESULT"] = result_text[:1000]
if pipeline_run_id:
env["CHAIN_PIPELINE_RUN_ID"] = pipeline_run_id
env["CHAIN_STEP_N"] = str(step_n)
cmd = [sys.executable, str(DISPATCH_PY), next_job]
if pipeline_run_id:
cmd.extend(["--pipeline-run", pipeline_run_id, "--step-n", str(step_n)])
subprocess.Popen( subprocess.Popen(
[sys.executable, str(DISPATCH_PY), chain_next], cmd,
env=env, env=env,
stdout=subprocess.DEVNULL, stdout=subprocess.DEVNULL,
stderr=subprocess.DEVNULL, stderr=subprocess.DEVNULL,
) )
else:
# End of chain for this pipeline run
if pipeline_run_id and HAS_PIPELINE:
if success:
pipeline_engine.complete_pipeline(pipeline_run_id)
else:
pipeline_engine.fail_pipeline(pipeline_run_id, "step_failed_without_fallback")
except Exception: except Exception:
pass pass
+277 -11
View File
@@ -55,6 +55,8 @@ FOLLOWUPS_FILE = NETVM_ROOT / "followups.json"
JOB_SIDECHATS_FILE = NETVM_ROOT / "job-sidechats.json" JOB_SIDECHATS_FILE = NETVM_ROOT / "job-sidechats.json"
RESPONSE_HARVESTER_PY = BIN_DIR / "response-harvester.py" RESPONSE_HARVESTER_PY = BIN_DIR / "response-harvester.py"
FOLLOWUP_SWEEPER_PY = BIN_DIR / "followup-sweeper.py" FOLLOWUP_SWEEPER_PY = BIN_DIR / "followup-sweeper.py"
PIPELINES_FILE = NETVM_ROOT / "pipelines.json"
JOB_DISPATCH_PY = BIN_DIR / "job-dispatch.py"
# Agent Constants # Agent Constants
VALID_NODES = ["muse", "pip", "646", "opm"] VALID_NODES = ["muse", "pip", "646", "opm"]
@@ -688,9 +690,14 @@ def cmd_dm_send(args):
message = args.message message = args.message
wait = getattr(args, "wait", False) wait = getattr(args, "wait", False)
timeout = getattr(args, "timeout", 60) or 60 timeout = getattr(args, "timeout", 60) or 60
allow_main = getattr(args, "allow_main_chat", False)
if target == "main": if target == "main":
print(c_yellow(" ⚠ [CHAT POLICY NOTE] Targeting Main Chat. Per CHAT_POLICY.md, avoid using Main Chat for routine tasks.")) if not allow_main:
print(c_red("ERROR: Targeting Main Chat is blocked by CHAT_POLICY.md to prevent chat degradation."), file=sys.stderr)
print(c_dim(" To send routine task/payloads, use a sidechat: --target '<sidechat>'.\n To override intentionally, pass --allow-main-chat."), file=sys.stderr)
sys.exit(1)
print(c_yellow(" ⚠ [CHAT POLICY OVERRIDE] Main Chat targeted with --allow-main-chat. Keep interaction minimal."))
should_sign = (sender == "super") and not getattr(args, "no_sign", False) should_sign = (sender == "super") and not getattr(args, "no_sign", False)
mid = None mid = None
@@ -738,13 +745,22 @@ def cmd_dm_send(args):
def cmd_dm_wo(args): def cmd_dm_wo(args):
sender = resolve_sender(args) sender = resolve_sender(args)
recipient = args.to recipient = args.to
target = getattr(args, "target", "main") target = getattr(args, "target", None)
if not target:
target = DEFAULT_AGENT_SIDECHATS.get(recipient, "main")
title = args.title title = args.title
body = args.body body = args.body
priority = getattr(args, "priority", "routine") priority = getattr(args, "priority", "routine")
wait = getattr(args, "wait", True)
timeout = getattr(args, "timeout", 60) or 60
allow_main = getattr(args, "allow_main_chat", False)
if target == "main": if target == "main":
print(c_yellow(" ⚠ [CHAT POLICY NOTE] Targeting Main Chat. Per CHAT_POLICY.md, prefer sidechats for work orders.")) if not allow_main:
print(c_red("ERROR: Work Orders must target sidechats per CHAT_POLICY.md."), file=sys.stderr)
print(c_dim(" Defaulting to sidechat recommended. To override, pass --allow-main-chat."), file=sys.stderr)
sys.exit(1)
print(c_yellow(" ⚠ [CHAT POLICY OVERRIDE] Work Order targeted to Main Chat with --allow-main-chat."))
wo_id = hashlib.sha256(f"{sender}-{recipient}-{title}-{time.time()}".encode()).hexdigest()[:8] wo_id = hashlib.sha256(f"{sender}-{recipient}-{title}-{time.time()}".encode()).hexdigest()[:8]
prefix = "[URGENT] " if priority == "urgent" else "" prefix = "[URGENT] " if priority == "urgent" else ""
@@ -762,9 +778,13 @@ def cmd_dm_wo(args):
"--raw", "--raw",
signed_wire signed_wire
] ]
print(f"Issuing Cryptographically Signed Work Order {c_green('[WO:' + wo_id + ']')} [{c_green('verified from:' + sender)}] -> [{c_bold(recipient)}/{target}]...") print(f"Issuing Cryptographically Signed Work Order {c_green('[WO:' + wo_id + ']')} [{c_green('verified from:' + sender)}] -> [{c_bold(recipient)}/{c_cyan(target)}]...")
res = subprocess.run(cmd) res = subprocess.run(cmd)
if res.returncode != 0:
sys.exit(res.returncode) sys.exit(res.returncode)
if wait:
wait_for_reply(recipient, target, msg_id=wo_id, timeout=timeout)
sys.exit(0)
except Exception as e: except Exception as e:
print(f"Signing notice: {e}, sending standard work order...", file=sys.stderr) print(f"Signing notice: {e}, sending standard work order...", file=sys.stderr)
@@ -777,9 +797,13 @@ def cmd_dm_wo(args):
wo_content wo_content
] ]
print(f"Issuing Work Order {c_cyan('[WO:' + wo_id + ']')} [{c_cyan('from:' + sender)}] -> [{c_bold(recipient)}/{target}]...") print(f"Issuing Work Order {c_cyan('[WO:' + wo_id + ']')} [{c_cyan('from:' + sender)}] -> [{c_bold(recipient)}/{c_cyan(target)}]...")
res = subprocess.run(cmd) res = subprocess.run(cmd)
if res.returncode != 0:
sys.exit(res.returncode) sys.exit(res.returncode)
if wait:
wait_for_reply(recipient, target, msg_id=wo_id, timeout=timeout)
sys.exit(0)
def cmd_dm_ack(args): def cmd_dm_ack(args):
sender = resolve_sender(args) sender = resolve_sender(args)
@@ -941,6 +965,22 @@ def cmd_dm_tail(args):
except KeyboardInterrupt: except KeyboardInterrupt:
print("\n" + c_dim("Tail stopped.")) print("\n" + c_dim("Tail stopped."))
TRANSFERS_REGISTRY_FILE = TRANSFERS_DIR / ".transfers-registry.json"
def _load_transfers_registry() -> dict:
if TRANSFERS_REGISTRY_FILE.exists():
try:
with open(TRANSFERS_REGISTRY_FILE, "r") as f:
return json.load(f)
except Exception:
return {}
return {}
def _save_transfers_registry(reg: dict):
TRANSFERS_DIR.mkdir(parents=True, exist_ok=True)
with open(TRANSFERS_REGISTRY_FILE, "w") as f:
json.dump(reg, f, indent=2)
def cmd_dm_send_file(args): def cmd_dm_send_file(args):
file_path = Path(args.file).resolve() file_path = Path(args.file).resolve()
if not file_path.exists() or not file_path.is_file(): if not file_path.exists() or not file_path.is_file():
@@ -983,6 +1023,24 @@ def cmd_dm_send_file(args):
# 3. Construct protocol payload pointer # 3. Construct protocol payload pointer
file_id = hashlib.sha256(f"{recipient}-{dest_filename}-{time.time()}".encode()).hexdigest()[:8] file_id = hashlib.sha256(f"{recipient}-{dest_filename}-{time.time()}".encode()).hexdigest()[:8]
# Save to transfers registry
reg = _load_transfers_registry()
reg[file_id] = {
"file_id": file_id,
"recipient": recipient,
"original_name": file_path.name,
"original_path": str(file_path),
"staged_path": str(dest_path),
"sha256": sha256_full,
"size_bytes": size_bytes,
"size_kb": f"{size_kb:.1f} KB",
"staged_at": datetime.now(timezone.utc).isoformat(),
"note": note,
"sender": "super",
"target": target
}
_save_transfers_registry(reg)
lines = [ lines = [
f"[FILE:{file_id}] {file_path.name} ({size_kb:.1f} KB, sha256:{sha256_short})", f"[FILE:{file_id}] {file_path.name} ({size_kb:.1f} KB, sha256:{sha256_short})",
f"Path: {dest_path}", f"Path: {dest_path}",
@@ -1016,21 +1074,53 @@ def cmd_dm_send_file(args):
wait_for_reply(recipient, target, msg_id=file_id, timeout=timeout) wait_for_reply(recipient, target, msg_id=file_id, timeout=timeout)
def cmd_dm_files(args): def cmd_dm_files(args):
subaction = getattr(args, "files_action", "list")
filter_agent = getattr(args, "agent", None) filter_agent = getattr(args, "agent", None)
if subaction == "clean":
older_than_days = getattr(args, "older_than", 7)
cutoff_sec = time.time() - (older_than_days * 86400)
reg = _load_transfers_registry()
removed_count = 0
retained_reg = {}
# Scan files on disk
for ad in (TRANSFERS_DIR.iterdir() if TRANSFERS_DIR.exists() else []):
if not ad.is_dir() or ad.name.startswith("."):
continue
if filter_agent and ad.name != filter_agent:
continue
for fp in list(ad.iterdir()):
if fp.is_file() and fp.stat().st_mtime < cutoff_sec:
try:
fp.unlink()
removed_count += 1
except Exception as e:
print(c_red(f"Error removing {fp}: {e}"))
for fid, rec in reg.items():
staged = Path(rec.get("staged_path", ""))
if staged.exists():
retained_reg[fid] = rec
_save_transfers_registry(retained_reg)
print(c_green(f"✔ Cleaned {removed_count} payload file(s) older than {older_than_days} day(s)."))
return
if not TRANSFERS_DIR.exists(): if not TRANSFERS_DIR.exists():
print(c_dim("No transferred files found (transfers/ directory is empty).")) print(c_dim("No transferred files found (transfers/ directory is empty)."))
return return
reg = _load_transfers_registry()
records = [] records = []
agent_dirs = [TRANSFERS_DIR / filter_agent] if filter_agent else sorted(TRANSFERS_DIR.iterdir()) agent_dirs = [TRANSFERS_DIR / filter_agent] if filter_agent else sorted(TRANSFERS_DIR.iterdir())
for ad in agent_dirs: for ad in agent_dirs:
if not ad.is_dir(): if not ad.is_dir() or ad.name.startswith("."):
continue continue
agent_name = ad.name agent_name = ad.name
for fp in sorted(ad.iterdir(), key=lambda p: p.stat().st_mtime, reverse=True): for fp in sorted(ad.iterdir(), key=lambda p: p.stat().st_mtime, reverse=True):
if not fp.is_file(): if not fp.is_file() or fp.name.startswith("."):
continue continue
st = fp.stat() st = fp.stat()
size_kb = f"{st.st_size / 1024:.1f} KB" size_kb = f"{st.st_size / 1024:.1f} KB"
@@ -1042,14 +1132,21 @@ def cmd_dm_files(args):
if m: if m:
clean_name = m.group(1) clean_name = m.group(1)
# Match in registry if available
reg_entry = next((r for r in reg.values() if r.get("staged_path") == str(fp)), {})
file_id = reg_entry.get("file_id", "-")
note = reg_entry.get("note", "")
records.append({ records.append({
"agent": agent_name, "agent": agent_name,
"file_id": file_id,
"filename": clean_name, "filename": clean_name,
"staged_name": fp.name, "staged_name": fp.name,
"path": str(fp), "path": str(fp),
"size_kb": size_kb, "size_kb": size_kb,
"mtime": mtime, "mtime": mtime,
"time_rel": t_rel "time_rel": t_rel,
"note": note
}) })
if args.json: if args.json:
@@ -1057,18 +1154,20 @@ def cmd_dm_files(args):
return return
print("\n" + c_bold(f"=== TRANSFERRED PAYLOADS & STAGED FILES ({len(records)}) ===") + "\n") print("\n" + c_bold(f"=== TRANSFERRED PAYLOADS & STAGED FILES ({len(records)}) ===") + "\n")
headers = ["AGENT", "FILENAME", "SIZE", "TRANSFERRED", "STAGED PATH"] headers = ["AGENT", "FILE ID", "FILENAME", "SIZE", "TRANSFERRED", "NOTE", "STAGED PATH"]
rows = [] rows = []
for r in records: for r in records:
rows.append([ rows.append([
c_cyan(r["agent"]), c_cyan(r["agent"]),
c_green(r["file_id"]) if r["file_id"] != "-" else c_dim("-"),
c_bold(r["filename"]), c_bold(r["filename"]),
r["size_kb"], r["size_kb"],
r["time_rel"], r["time_rel"],
r["note"][:20] if r["note"] else c_dim("-"),
c_dim(r["path"]) c_dim(r["path"])
]) ])
print_table(headers, rows) print_table(headers, rows)
print("\n" + c_dim(" Commands: super dm send-file --to <agent> <path> | super dm files --agent <agent>") + "\n") print("\n" + c_dim(" Commands: super dm send-file --to <agent> <path> | super dm files clean [--older-than N]") + "\n")
# --------------------------------------------------------------------------- # ---------------------------------------------------------------------------
# Domain: THREAD (Chat Oversight via box-chat.py) # Domain: THREAD (Chat Oversight via box-chat.py)
@@ -2033,6 +2132,145 @@ def cmd_followup_cancel(args):
os.replace(tmp_path, FOLLOWUPS_FILE) os.replace(tmp_path, FOLLOWUPS_FILE)
print(c_green(f"● Cancelled follow-up {match_key}.")) print(c_green(f"● Cancelled follow-up {match_key}."))
# ---------------------------------------------------------------------------
# Domain: PIPELINE (Multi-Agent Workflow Pipelines)
# ---------------------------------------------------------------------------
def cmd_pipeline_run(args):
name = args.name.strip()
job_file = JOBS_DIR / f"{name}.json"
if not job_file.exists():
print(c_red(f"Error: Pipeline root job '{name}.json' not found in jobs/"), file=sys.stderr)
sys.exit(1)
import pipeline_engine
run_id = pipeline_engine.generate_pipeline_run_id(name)
entry = pipeline_engine.create_pipeline(name, run_id=run_id)
target = entry.get("target")
print(c_bold(f"\n=== LAUNCHING MULTI-AGENT PIPELINE [{name}] ==="))
print(f" Run ID: {c_cyan(run_id)}")
print(f" Target: {c_yellow(target)}")
print(f" Step 1: {c_bold(name)}\n")
cmd = [sys.executable, str(JOB_DISPATCH_PY), name, "--pipeline-run", run_id, "--step-n", "1"]
if getattr(args, "dry_run", False):
cmd.append("--dry-run")
res = subprocess.run(cmd)
if res.returncode == 0:
print(c_green(f"\n✔ Pipeline '{run_id}' dispatched step 1 successfully."))
print(c_dim(" Track with: super pipeline status | super pipeline history\n"))
else:
print(c_red(f"\n✖ Step 1 dispatch failed with code {res.returncode}"), file=sys.stderr)
sys.exit(res.returncode)
def cmd_pipeline_status(args):
import pipeline_engine
active = pipeline_engine.list_active_pipelines()
if getattr(args, "json", False):
print(json.dumps(active, indent=2))
return
now_str = datetime.now().strftime("%H:%M:%S")
print(f"\n=== ACTIVE MULTI-AGENT PIPELINES === ({now_str} local)\n")
if not active:
print(c_dim(" (no active pipeline runs currently executing)"))
print(c_dim(" Launch one with: super pipeline run <job_name>\n"))
return
headers = ["RUN ID", "PIPELINE", "TARGET", "STEP", "CURRENT AGENT", "ACTIVE JOB ID", "STARTED", "STATUS"]
rows = []
for p in active:
run_id = p.get("run_id", "-")
p_name = p.get("pipeline_name", "-")
target = p.get("target", "-")
cur_step_n = p.get("current_step", 1)
steps = p.get("steps", [])
last_step = steps[-1] if steps else {}
agent = last_step.get("agent", "-")
job_id = last_step.get("job_id", "-")
created_at = p.get("created_at")
rows.append([
c_bold(run_id),
p_name,
c_yellow(target),
f"Step {cur_step_n}",
c_cyan(agent),
c_dim(job_id[:16] + "…") if len(job_id) > 16 else job_id,
parse_relative_time(created_at) if created_at else c_dim("-"),
badge_ok("RUNNING"),
])
print_table(headers, rows)
# Details on steps for active runs
for p in active:
print(c_bold(f"\n Active Pipeline Details [{p.get('run_id')}]:"))
for s in p.get("steps", []):
st = s.get("status", "unknown")
st_badge = badge_ok("COMPLETED") if st == "completed" else badge_warn("RUNNING") if st == "dispatched" else badge_err(st.upper())
res_snippet = s.get("result", "")
res_disp = f" → {c_dim(res_snippet[:70])}" if res_snippet else ""
print(f" • Step {s.get('step_n')}: {c_bold(s.get('job_name'))} ({s.get('agent')}) [{st_badge}]{res_disp}")
print(c_dim("\n Commands: super pipeline run <name> | super pipeline history\n"))
def cmd_pipeline_history(args):
import pipeline_engine
limit = getattr(args, "limit", 20) or 20
history = pipeline_engine.list_pipeline_history(limit=limit)
if getattr(args, "json", False):
print(json.dumps(history, indent=2))
return
now_str = datetime.now().strftime("%H:%M:%S")
print(f"\n=== PIPELINE RUN HISTORY (Recent {len(history)}) === ({now_str} local)\n")
if not history:
print(c_dim(" (no pipeline history found)"))
return
headers = ["RUN ID", "PIPELINE", "TARGET", "STEPS", "STARTED", "COMPLETED", "STATUS"]
rows = []
for p in history:
run_id = p.get("run_id", "-")
p_name = p.get("pipeline_name", "-")
target = p.get("target", "-")
steps = p.get("steps", [])
step_summary = f"{len(steps)} step{'s' if len(steps) != 1 else ''}"
created_at = p.get("created_at")
completed_at = p.get("completed_at")
st = p.get("status", "unknown")
if st == "completed":
badge = badge_ok("COMPLETED")
elif st == "running":
badge = badge_warn("RUNNING")
else:
badge = badge_err("FAILED")
rows.append([
c_bold(run_id),
p_name,
c_yellow(target),
step_summary,
parse_relative_time(created_at) if created_at else c_dim("-"),
parse_relative_time(completed_at) if completed_at else c_dim("active"),
badge,
])
print_table(headers, rows)
print(c_dim("\n Commands: super pipeline status | super pipeline run <name>\n"))
# --------------------------------------------------------------------------- # ---------------------------------------------------------------------------
# Domain: LOOP (Intrinsic Loop Strategy, Health, Break Taxonomy, and Control) # Domain: LOOP (Intrinsic Loop Strategy, Health, Break Taxonomy, and Control)
# --------------------------------------------------------------------------- # ---------------------------------------------------------------------------
@@ -2532,6 +2770,7 @@ def build_parser():
p_dm_send.add_argument("--to", required=True, choices=VALID_NODES, help="Recipient node") p_dm_send.add_argument("--to", required=True, choices=VALID_NODES, help="Recipient node")
p_dm_send.add_argument("--from", dest="from_agent", default=DEFAULT_SENDER, help="Sender identity (default: super)") p_dm_send.add_argument("--from", dest="from_agent", default=DEFAULT_SENDER, help="Sender identity (default: super)")
p_dm_send.add_argument("--target", default="main", help="Target chat/sidechat (default: main)") p_dm_send.add_argument("--target", default="main", help="Target chat/sidechat (default: main)")
p_dm_send.add_argument("--allow-main-chat", action="store_true", help="Explicit override allowing send directly to Main Chat")
p_dm_send.add_argument("-w", "--wait", action="store_true", help="Wait for recipient agent to reply") p_dm_send.add_argument("-w", "--wait", action="store_true", help="Wait for recipient agent to reply")
p_dm_send.add_argument("-t", "--timeout", type=int, default=60, help="Wait timeout in seconds (default: 60)") p_dm_send.add_argument("-t", "--timeout", type=int, default=60, help="Wait timeout in seconds (default: 60)")
p_dm_send.add_argument("--expect-reply", action="store_true", help="Arm follow-up tracking") p_dm_send.add_argument("--expect-reply", action="store_true", help="Arm follow-up tracking")
@@ -2556,7 +2795,9 @@ def build_parser():
# super dm files # super dm files
p_dm_files = dm_sub.add_parser("files", parents=[common], help="List transferred files and staged payloads") p_dm_files = dm_sub.add_parser("files", parents=[common], help="List transferred files and staged payloads")
p_dm_files.add_argument("files_action", nargs="?", default="list", choices=["list", "clean"], help="Action: list (default) or clean")
p_dm_files.add_argument("--agent", choices=VALID_NODES, default=None, help="Filter by recipient agent") p_dm_files.add_argument("--agent", choices=VALID_NODES, default=None, help="Filter by recipient agent")
p_dm_files.add_argument("--older-than", type=int, default=7, help="Retention cutoff in days for clean (default: 7)")
# super dm wo # super dm wo
p_dm_wo = dm_sub.add_parser("wo", parents=[common], help="Dispatch a structured Work Order ([WO:...])") p_dm_wo = dm_sub.add_parser("wo", parents=[common], help="Dispatch a structured Work Order ([WO:...])")
@@ -2564,7 +2805,10 @@ def build_parser():
p_dm_wo.add_argument("--title", required=True, help="Work Order title") p_dm_wo.add_argument("--title", required=True, help="Work Order title")
p_dm_wo.add_argument("--priority", choices=["routine", "urgent"], default="routine", help="Priority level") p_dm_wo.add_argument("--priority", choices=["routine", "urgent"], default="routine", help="Priority level")
p_dm_wo.add_argument("--from", dest="from_agent", default=DEFAULT_SENDER, help="Sender identity (default: super)") p_dm_wo.add_argument("--from", dest="from_agent", default=DEFAULT_SENDER, help="Sender identity (default: super)")
p_dm_wo.add_argument("--target", default="main", help="Target chat (default: main)") p_dm_wo.add_argument("--target", default=None, help="Target sidechat (default: agent primary task sidechat)")
p_dm_wo.add_argument("--allow-main-chat", action="store_true", help="Explicit override allowing work order in Main Chat")
p_dm_wo.add_argument("--no-wait", dest="wait", action="store_false", help="Do not wait for reply")
p_dm_wo.add_argument("-t", "--timeout", type=int, default=60, help="Wait timeout in seconds (default: 60)")
p_dm_wo.add_argument("--no-sign", action="store_true", help="Issue without cryptographic SSH signature") p_dm_wo.add_argument("--no-sign", action="store_true", help="Issue without cryptographic SSH signature")
p_dm_wo.add_argument("body", help="Work Order detailed description") p_dm_wo.add_argument("body", help="Work Order detailed description")
@@ -2763,6 +3007,18 @@ def build_parser():
p_l_harvest = loop_sub.add_parser("harvest", parents=[common], help="Run immediate harvest cycle") p_l_harvest = loop_sub.add_parser("harvest", parents=[common], help="Run immediate harvest cycle")
p_l_harvest.add_argument("--dry-run", action="store_true", help="Scrape without persisting watermarks") p_l_harvest.add_argument("--dry-run", action="store_true", help="Scrape without persisting watermarks")
# Domain: PIPELINE
p_pipe = subparsers.add_parser("pipeline", parents=[common], help="Multi-agent workflow pipelines and chained handoffs")
pipe_sub = p_pipe.add_subparsers(dest="action")
p_pipe_run = pipe_sub.add_parser("run", parents=[common], help="Start a new multi-agent pipeline run")
p_pipe_run.add_argument("name", help="Root pipeline job name (e.g. pipe-demo-step1)")
p_pipe_run.add_argument("--dry-run", action="store_true", help="Simulate pipeline without sending DMs")
p_pipe_status = pipe_sub.add_parser("status", parents=[common], help="Display active pipeline runs and step states")
p_pipe_hist = pipe_sub.add_parser("history", parents=[common], help="Show historical pipeline runs and outcomes")
p_pipe_hist.add_argument("-n", "--limit", type=int, default=20, help="Number of records to display")
# Top-level domain aliases: strat and vars # Top-level domain aliases: strat and vars
p_strat_alias = subparsers.add_parser("strat", parents=[common], help="Shortcut for 'loop strat'") p_strat_alias = subparsers.add_parser("strat", parents=[common], help="Shortcut for 'loop strat'")
strat_alias_sub = p_strat_alias.add_subparsers(dest="strat_action") strat_alias_sub = p_strat_alias.add_subparsers(dest="strat_action")
@@ -2922,6 +3178,16 @@ def main():
cmd_followup_cancel(args) cmd_followup_cancel(args)
else: else:
parser.print_help() parser.print_help()
elif args.domain == "pipeline":
act = getattr(args, "action", None)
if not act or act == "status":
cmd_pipeline_status(args)
elif act == "run":
cmd_pipeline_run(args)
elif act == "history":
cmd_pipeline_history(args)
else:
parser.print_help()
elif args.domain == "loop": elif args.domain == "loop":
act = getattr(args, "action", None) act = getattr(args, "action", None)
if not act or act == "status": if not act or act == "status":
+1 -1
View File
@@ -42,7 +42,7 @@ Examples:
## Lifecycle ## Lifecycle
1. **Create**: Operator or agent initiates via chat API 1. **Create**: Operator or agent initiates via chat API (IMPLEMENTED 2026-10-03)
- `muse-chat-api.py --account <agent> sidechat create --channel <id> --purpose <str>` - `muse-chat-api.py --account <agent> sidechat create --channel <id> --purpose <str>`
- Returns `side_chat_id` - Returns `side_chat_id`
+1 -1
View File
@@ -1,5 +1,5 @@
{ {
"heartbeat-opm": "5bd5b350-806c-4e04-939c-7ca7d1bd20df", "heartbeat-opm": "0077e918-9ac8-40ba-b0ad-65f9f78037c4",
"646-pip-coord": { "646-pip-coord": {
"thread_uuid": "4466d0c1-7961-4cf3-b99d-1ab7c38484c2", "thread_uuid": "4466d0c1-7961-4cf3-b99d-1ab7c38484c2",
"agent": "pip" "agent": "pip"
+1 -1
View File
@@ -1 +1 @@
{"agent":"opm","chain_next":null,"description":"Canary job to test the dispatcher - sends a simple DM","name":"canary-test","on_failure":"alert","prompt_template":"This is a canary test from the job dispatcher.\nJob ID: {job_id}\nDate: {date}\n\nPlease reply with [RESULT {job_id}] to confirm you received this.","schedule":"manual","sidechat":{"create":false},"timeout":300} {"agent":"opm","chain_next":null,"description":"Canary job to test the dispatcher - sends a simple DM","dm_target":"heartbeat","name":"canary-test","on_failure":"alert","prompt_template":"This is a canary test from the job dispatcher.\nJob ID: {job_id}\nDate: {date}\n\nPlease reply to confirm you received this.","schedule":"manual","sidechat":{"create":false},"timeout":300}