feat(automation): sidechat auto-provisioning, response harvester, followup sweeper, and super CLI

This commit is contained in:
operator
2026-10-04 16:23:10 +00:00
parent a214355f16
commit 4a935bd0c7
7 changed files with 4417 additions and 45 deletions
+8
View File
@@ -0,0 +1,8 @@
# NetVM Transfers & Local State
transfers/
*.log
__pycache__/
*.pyc
*.bak
*.bak-*
*.orig
+494 -18
View File
@@ -35,13 +35,198 @@ import sys
import uuid
import json
import os
from datetime import datetime, timezone
import urllib.request
import urllib.parse
from datetime import datetime, timezone, timedelta
API = "/home/super/Projects/NetVM/bin/muse-chat-api.py"
NETVM_EXEC = "/home/super/Projects/NetVM/bin/netvm-exec.sh"
VALID_AGENTS = ["muse", "pip", "646", "opm"]
VALID_SENDERS = ["muse", "pip", "646", "opm", "super"]
VALID_RECIPIENTS = ["muse", "pip", "646", "opm"]
# ---- Canonical follow-up tags (2026-10-04, DEPLOY-DECISIONS.md) ----
# Tags declare follow-up policy at send time. They are metadata only:
# stripped from the delivered text, recorded on send_start/sent log
# events. The request-store sweeper (VM side) guarantees nudges,
# retries, closure, and escalation. Untagged DMs behave exactly as before.
TAG_TIMEOUT_MIN = 60
TAG_TIMEOUT_MAX = 604800
TAG_TIMEOUT_DEFAULT = 3600 # 1h, per deploy decisions
TAG_NUDGES_MIN = 0
TAG_NUDGES_MAX = 10
TAG_NUDGES_DEFAULT = 2
TAG_ROUTE_RE = re.compile(r"^[A-Za-z0-9][A-Za-z0-9_-]{0,63}$")
TAG_THREAD_RE = re.compile(r"^[A-Za-z0-9_-]{1,64}$")
# Trailing-token form: [reply:expected] [reply:timeout=7200] ...
TAG_TOKEN_RE = re.compile(r"\[([A-Za-z][A-Za-z0-9:]*)"
r"(?:=([^\]]*))?\]\s*$")
def parse_canonical_tag(item):
"""Parse one canonical tag token. Returns (key, value).
Bare 'reply:expected' -> ('reply:expected', True).
'route:X' / 'thread:X' use colon form; 'reply:timeout=N' etc use '='.
Raises ValueError with a human message on any problem."""
item = item.strip()
if "=" in item:
k, v = item.split("=", 1)
k, v = k.strip(), v.strip()
elif item.startswith("route:"):
k, v = "route", item[len("route:"):].strip()
elif item.startswith("thread:"):
k, v = "thread", item[len("thread:"):].strip()
else:
k, v = item, True
if k == "reply:expected":
# Bare 'reply:expected' or 'reply:expected=' (empty value from
# --tag flag generation); anything else is an error.
if v is not True and v != "":
raise ValueError("reply:expected takes no value")
return k, True
if k == "reply:timeout":
try:
n = int(v)
except (TypeError, ValueError):
raise ValueError("reply:timeout must be integer seconds")
if not TAG_TIMEOUT_MIN <= n <= TAG_TIMEOUT_MAX:
raise ValueError("reply:timeout must be %d..%d"
% (TAG_TIMEOUT_MIN, TAG_TIMEOUT_MAX))
return k, n
if k == "reply:nudges":
try:
n = int(v)
except (TypeError, ValueError):
raise ValueError("reply:nudges must be an integer")
if not TAG_NUDGES_MIN <= n <= TAG_NUDGES_MAX:
raise ValueError("reply:nudges must be %d..%d"
% (TAG_NUDGES_MIN, TAG_NUDGES_MAX))
return k, n
if k == "reply:escalate":
if not v or not TAG_ROUTE_RE.fullmatch(str(v)):
raise ValueError("reply:escalate must be an identity")
return k, str(v)
if k == "route":
if not TAG_ROUTE_RE.fullmatch(str(v)):
raise ValueError("route must match ^[A-Za-z0-9][A-Za-z0-9_-]{0,63}$")
return k, str(v)
if k == "thread":
if not TAG_THREAD_RE.fullmatch(str(v)):
raise ValueError("thread must match ^[A-Za-z0-9_-]{1,64}$")
return k, str(v)
raise ValueError("unknown tag %r; known: reply:expected, reply:timeout=, "
"reply:nudges=, reply:escalate=, route:, thread:" % (k,))
def parse_tags(tag_list):
"""Validate --tag items. Returns canonical dict. Raises ValueError."""
tags = {}
for item in tag_list or []:
k, v = parse_canonical_tag(item)
if k in tags:
raise ValueError("duplicate tag %r" % (k,))
tags[k] = v
# Cross-field consistency: policy tags require reply:expected.
if "reply:expected" not in tags:
for k in ("reply:timeout", "reply:nudges", "reply:escalate"):
if k in tags:
raise ValueError("tag %s requires reply:expected" % (k,))
# Defaults when reply:expected without explicit policy.
if "reply:expected" in tags:
tags.setdefault("reply:timeout", TAG_TIMEOUT_DEFAULT)
tags.setdefault("reply:nudges", TAG_NUDGES_DEFAULT)
return tags
def extract_trailing_tags(message):
"""Pull trailing [tag] tokens off message text. Returns
(stripped_message, tags_dict). Unknown brackets are left alone."""
tags = {}
text = message.rstrip()
while True:
m = TAG_TOKEN_RE.search(text)
if not m:
break
key = m.group(1)
# Only consume known canonical keys; anything else stays.
probe = key + ("=" + m.group(2) if m.group(2) is not None else "")
try:
k, v = parse_canonical_tag(probe)
except ValueError:
break
if k in tags:
break # duplicate; leave it in the text
tags[k] = v
text = text[:m.start()].rstrip()
return text, tags
def merge_tags(flag_tags, cli_tags, text_tags):
"""Precedence: --tag > dedicated flags > trailing text tokens."""
merged = dict(text_tags)
merged.update(flag_tags)
merged.update(cli_tags)
# Defaults when reply:expected without explicit policy.
if "reply:expected" in merged:
merged.setdefault("reply:timeout", TAG_TIMEOUT_DEFAULT)
merged.setdefault("reply:nudges", TAG_NUDGES_DEFAULT)
return merged
def tags_from_flags(a):
"""Build canonical tags from argparse --expect-reply et al."""
tags = {}
if getattr(a, "expect_reply", False):
tags["reply:expected"] = True
if getattr(a, "reply_timeout", None) is not None:
k, v = parse_canonical_tag("reply:timeout=%s" % a.reply_timeout)
tags[k] = v
if getattr(a, "reply_nudges", None) is not None:
k, v = parse_canonical_tag("reply:nudges=%s" % a.reply_nudges)
tags[k] = v
if getattr(a, "reply_escalate", None):
k, v = parse_canonical_tag("reply:escalate=%s" % a.reply_escalate)
tags[k] = v
if getattr(a, "route", None):
k, v = parse_canonical_tag("route:%s" % a.route)
tags[k] = v
if getattr(a, "thread", None):
k, v = parse_canonical_tag("thread:%s" % a.thread)
tags[k] = v
return tags
LOG_FILE = "/home/super/Projects/NetVM/dm-log.jsonl"
# Well-known sidechat name -> thread UUID aliases.
# These bypass fuzzy name matching for reliable placement.
# Resolved 2026-10-04. Add new entries as sidechats are created.
SIDCHAT_ALIASES = {
"646-opm-work": "d410b9ad-f667-465f-a103-43fabc0f69fe",
}
def resolve_sidechat_target(target):
"""Resolve target alias or name to UUID dynamically from job-sidechats.json."""
if not target or target == "main":
return target
if re.fullmatch(r"[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}", target.lower()):
return target.lower()
if target in SIDCHAT_ALIASES:
return SIDCHAT_ALIASES[target]
sc_file = "/home/super/Projects/NetVM/job-sidechats.json"
if os.path.exists(sc_file):
try:
with open(sc_file, "r", encoding="utf-8") as f:
sc_data = json.load(f)
val = sc_data.get(target)
if isinstance(val, dict):
return val.get("thread_uuid") or val.get("uuid") or target
elif isinstance(val, str):
return val
except Exception:
pass
return target
def log_event(event):
"""Append to JSONL log."""
event["ts"] = datetime.now(timezone.utc).isoformat()
@@ -63,7 +248,8 @@ def run_full(cmd, timeout=60):
result = subprocess.run(cmd, shell=True, capture_output=True, text=True, timeout=timeout)
return result.returncode, result.stdout.strip(), result.stderr.strip()
def dm_send(agent, target, message, verify=True, raw=False, to_agent=None):
def dm_send(agent, target, message, verify=True, raw=False,
to_agent=None, tags=None, nudge_meta=None):
"""Send a DM. raw=True sends verbatim (for pre-signed messages from
dm-sign.sh, which already carry [from:X] [id:Y]): no tagging, no
truncation. Non-raw messages are tagged [from:<agent>] [id:<uuid8>]."""
@@ -71,10 +257,10 @@ def dm_send(agent, target, message, verify=True, raw=False, to_agent=None):
# agent is the sender (for attribution). If to_agent is None, send to own chat.
recipient = to_agent if to_agent else agent
if agent not in VALID_AGENTS:
if agent not in VALID_SENDERS:
print(f"ERROR: Unknown agent {agent}", file=sys.stderr)
sys.exit(1)
if recipient not in VALID_AGENTS:
if recipient not in VALID_RECIPIENTS:
print(f"ERROR: Unknown recipient {recipient}", file=sys.stderr)
sys.exit(1)
@@ -87,15 +273,85 @@ def dm_send(agent, target, message, verify=True, raw=False, to_agent=None):
# Single unified attribution format (matches verify-sig's regex).
tagged = f"[from:{agent}] [id:{msg_id}] {message}"
log_event({"type": "send_start", "id": msg_id, "agent": agent, "to": recipient, "target": target, "msg": message[:100]})
# Canonical follow-up tags: trailing [bracket] tokens are metadata,
# stripped from the delivered text, recorded on the log events.
tags = tags or {}
if not raw:
message, text_tags = extract_trailing_tags(message)
if text_tags:
# Re-tag with the stripped body so the wire text is clean.
tagged = f"[from:{agent}] [id:{msg_id}] {message}"
for k, v in text_tags.items():
tags.setdefault(k, v)
_log = {"type": "send_start", "id": msg_id, "agent": agent, "to": recipient, "target": target, "msg": message[:100], "tags": tags}
if nudge_meta:
_log["nudge_meta"] = nudge_meta
log_event(_log)
# Navigate the RECIPIENT's browser to the target chat (the send and the
# read-back both happen there; navigating the sender's browser was a bug
# for cross-operator side-chat targets).
# Use run_full: a silent navigation failure used to send the message to
# whatever chat the browser happened to be parked on (2026-10-04).
thread_uuid = None
thread_url = None
is_new_sidechat = False
# Resolve well-known aliases or dynamic thread mappings to UUIDs.
nav_target = resolve_sidechat_target(target)
if nav_target != target:
log_event({"type": "alias_resolved", "id": msg_id, "target": target, "thread_uuid": nav_target})
if target == "main":
run(f"{NETVM_EXEC} {recipient} -- python3 {API} --account {recipient} sidechat main")
_rc, _out, _err = run_full(f"{NETVM_EXEC} {recipient} -- python3 {API} --account {recipient} sidechat main")
else:
run(f"{NETVM_EXEC} {recipient} -- python3 {API} --account {recipient} sidechat use {target}")
_rc, _out, _err = run_full(f"{NETVM_EXEC} {recipient} -- python3 {API} --account {recipient} sidechat use {nav_target}")
if _rc != 0:
log_event({"type": "nav_failed", "id": msg_id, "agent": agent, "to": recipient,
"target": target, "rc": _rc, "err": _err[:200]})
print(f"DM {msg_id} from {agent} to {recipient}/{target}: FAILED (navigation rc={_rc}: {_err[:120]})",
file=sys.stderr)
sys.exit(1)
# cmd_sidechat_use exits 0 even on NOTFOUND (it prints "Navigated to: NOTFOUND").
# If sidechat is not found, auto-provision a new thread via `sidechat create`!
if target != "main":
if "NOTFOUND" in _out:
log_event({"type": "sidechat_autoprovision_start", "id": msg_id, "agent": agent,
"to": recipient, "target": target})
_c_rc, _c_out, _c_err = run_full(f"{NETVM_EXEC} {recipient} -- python3 {API} --account {recipient} sidechat create")
if _c_rc != 0 or "Created:" not in _c_out:
log_event({"type": "nav_failed", "id": msg_id, "agent": agent, "to": recipient,
"target": target, "reason": "autoprovision_create_failed",
"rc": _c_rc, "err": _c_err[:200], "out": _c_out[:200]})
print(f"DM {msg_id} from {agent} to {recipient}/{target}: FAILED (sidechat auto-creation failed: {_c_err[:120]})",
file=sys.stderr)
sys.exit(1)
is_new_sidechat = True
log_event({"type": "nav_ok", "id": msg_id, "agent": agent, "to": recipient,
"target": target, "status": "sidechat_created_pending_uuid"})
else:
m = re.search(r"/thread/([0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12})", _out, re.I)
if m:
thread_uuid = m.group(1).lower()
thread_url = "https://muse.ai/thread/" + thread_uuid
else:
# Navigation "succeeded" but we can't confirm where we landed.
# Don't trust it: fail instead of sending to an unknown chat.
log_event({"type": "nav_failed", "id": msg_id, "agent": agent, "to": recipient,
"target": target, "reason": "no_thread_url", "out": _out[:200]})
print(f"DM {msg_id} from {agent} to {recipient}/{target}: FAILED (could not confirm thread URL)",
file=sys.stderr)
sys.exit(1)
# If we resolved via alias, assert the browser actually landed on
# the expected thread. A mismatch means the navigation didn't take
# (stale URL, blocked nav) -- fail instead of misdelivering.
if nav_target != target and thread_uuid != nav_target.lower():
log_event({"type": "nav_failed", "id": msg_id, "agent": agent, "to": recipient,
"target": target, "reason": "uuid_mismatch",
"expected": nav_target.lower(), "got": thread_uuid})
print(f"DM {msg_id} from {agent} to {recipient}/{target}: FAILED (navigation landed on wrong thread)",
file=sys.stderr)
sys.exit(1)
log_event({"type": "nav_ok", "id": msg_id, "agent": agent, "to": recipient,
"target": target, "thread_uuid": thread_uuid})
time.sleep(2)
@@ -110,16 +366,43 @@ def dm_send(agent, target, message, verify=True, raw=False, to_agent=None):
max_retries = 3
delivered = False
for attempt in range(max_retries):
run(f'{NETVM_EXEC} {recipient} -- python3 {API} --account {recipient} send "{safe}"',
send_out = run(f'{NETVM_EXEC} {recipient} -- python3 {API} --account {recipient} send "{safe}"',
priority="high")
time.sleep(3) # Wait for message to propagate
# Verify by reading recipient's chat (independent check, not local echo)
# The send command reports 'sent' / 'enter-sent' / 'NOINPUT'.
# 'NOINPUT' (or empty output) means the compose box was never found:
# retrying verification would be meaningless, so fail fast.
if send_out.strip() not in ("sent", "enter-sent"):
log_event({"type": "send_failed", "id": msg_id, "agent": agent, "to": recipient,
"target": target, "attempt": attempt + 1, "send_out": send_out[:80]})
if attempt < max_retries - 1:
time.sleep(2)
continue
# Assert the composer cleared: if our text (with the [id:...] tag) is
# still sitting in the compose box, the send click never fired and any
# later "verification" would be a false positive (2026-10-04 bug).
compose = run(f'{NETVM_EXEC} {recipient} -- python3 {API} --account {recipient} compose_check')
if compose and compose != "NOCOMPOSE" and msg_id in compose:
log_event({"type": "compose_stuck", "id": msg_id, "agent": agent, "to": recipient,
"target": target, "attempt": attempt + 1})
# Clear the stuck draft so a retry starts clean.
run(f'{NETVM_EXEC} {recipient} -- python3 {API} --account {recipient} send ""',
priority="high")
if attempt < max_retries - 1:
time.sleep(2)
continue
# Verify by reading recipient's chat (independent check, not local echo).
# cmd_messages now excludes the composer subtree, so a hit here means
# the message is actually in the chat history.
try:
check_msgs = run(f'{NETVM_EXEC} {recipient} -- python3 {API} --account {recipient} messages 5 200')
if msg_id in check_msgs:
delivered = True
log_event({"type": "verified", "id": msg_id, "agent": agent, "to": recipient, "target": target, "attempt": attempt + 1})
log_event({"type": "verified", "id": msg_id, "agent": agent, "to": recipient, "target": target,
"thread_uuid": thread_uuid, "attempt": attempt + 1})
break
else:
log_event({"type": "retry", "id": msg_id, "agent": agent, "to": recipient, "attempt": attempt + 1})
@@ -129,14 +412,61 @@ def dm_send(agent, target, message, verify=True, raw=False, to_agent=None):
if attempt < max_retries - 1:
time.sleep(2) # Brief pause before retry
if delivered and is_new_sidechat:
# We sent to /thread/new; muse.ai now assigns a permanent UUID.
# Capture the current URL BEFORE navigating away to main.
for _ in range(10):
_u_rc, _u_out, _u_err = run_full(f"{NETVM_EXEC} {recipient} -- python3 {API} --account {recipient} url")
m_uuid = re.search(r"/thread/([0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12})", _u_out, re.I)
if m_uuid:
thread_uuid = m_uuid.group(1).lower()
thread_url = "https://muse.ai/thread/" + thread_uuid
break
time.sleep(1)
if thread_uuid:
sc_file = "/home/super/Projects/NetVM/job-sidechats.json"
try:
sc_data = {}
if os.path.exists(sc_file):
with open(sc_file, "r", encoding="utf-8") as f:
sc_data = json.load(f)
sc_data[target] = {
"thread_uuid": thread_uuid,
"agent": recipient,
"created_at": datetime.now(timezone.utc).isoformat()
}
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)
log_event({"type": "sidechat_autoprovisioned", "id": msg_id, "target": target,
"thread_uuid": thread_uuid, "agent": recipient})
except Exception as e:
log_event({"type": "sidechat_persist_error", "id": msg_id, "error": str(e)[:100]})
tags["thread"] = thread_uuid
else:
log_event({"type": "sidechat_uuid_capture_failed", "id": msg_id, "out": _u_out[:200]})
if delivered and thread_uuid and "thread" not in tags:
tags["thread"] = thread_uuid
# Park the recipient's browser back on main
run(f"{NETVM_EXEC} {recipient} -- python3 {API} --account {recipient} sidechat main")
log_event({"type": "send_done", "id": msg_id, "agent": agent, "to": recipient, "target": target})
if delivered:
log_event({"type": "sent", "id": msg_id, "agent": agent, "to": recipient, "target": target, "verified": True})
_sent = {"type": "sent", "id": msg_id, "agent": agent, "to": recipient, "target": target, "verified": True, "tags": tags}
if nudge_meta:
_sent["nudge_meta"] = nudge_meta
log_event(_sent)
print(f"DM {msg_id} from {agent} to {recipient}/{target}: SENT and VERIFIED")
# Follow-up record creation: if tags request tracking, POST to the
# VM's /api/box/followups endpoint. Fail gracefully -- the DM already
# sent, so a record-creation failure is logged, not fatal.
if tags.get("reply:expected"):
_register_followup(msg_id, agent, recipient, target, tags)
else:
log_event({"type": "failed", "id": msg_id, "agent": agent, "to": recipient, "target": target, "verified": False})
print(f"DM {msg_id} from {agent} to {recipient}/{target}: FAILED (not found in recipient chat after {max_retries} attempts)", file=sys.stderr)
@@ -150,10 +480,11 @@ def dm_read(agent, target, n=5, quiet=False, width=200):
print(f"ERROR: Unknown agent {agent}", file=sys.stderr)
sys.exit(1)
nav_target = resolve_sidechat_target(target)
if target == "main":
run(f"{NETVM_EXEC} {agent} -- python3 {API} --account {agent} sidechat main")
else:
run(f"{NETVM_EXEC} {agent} -- python3 {API} --account {agent} sidechat use {target}")
run(f"{NETVM_EXEC} {agent} -- python3 {API} --account {agent} sidechat use {nav_target}")
time.sleep(2)
rc, msgs, err = run_full(f"{NETVM_EXEC} {agent} -- python3 {API} --account {agent} messages {n} {width}")
@@ -187,11 +518,12 @@ def dm_select(agent, target=None):
return r.stdout.strip()
if target:
nav_target = resolve_sidechat_target(target)
# Switch to target conversation
if target == "main":
run(f"{NETVM_EXEC} {agent} -- python3 {API} --account {agent} sidechat main")
else:
run(f"{NETVM_EXEC} {agent} -- python3 {API} --account {agent} sidechat use {target}")
run(f"{NETVM_EXEC} {agent} -- python3 {API} --account {agent} sidechat use {nav_target}")
# Store selection
with open(state_file, "w") as f:
json.dump({"agent": agent, "target": target}, f)
@@ -310,15 +642,50 @@ def main():
sub = p.add_subparsers(dest='cmd', required=True)
ps = sub.add_parser('send', help='Send a DM (tagged; delivery confirmed by recipient read-back, up to 3 attempts)')
ps.add_argument('--agent', required=True, choices=VALID_AGENTS)
ps.add_argument('--to', required=False, choices=VALID_AGENTS, default=None,
ps.add_argument('--agent', required=True, choices=VALID_SENDERS)
ps.add_argument('--to', required=False, choices=VALID_RECIPIENTS, default=None,
help='Recipient operator (for cross-operator DMs). Uses recipient\'s browser/chat.')
ps.add_argument('--target', required=True)
ps.add_argument('--no-verify', action='store_true',
help='Accepted for compatibility but ignored: recipient-side read-back verification always runs.')
ps.add_argument('--raw', action='store_true', help='Send verbatim: no tagging, no truncation (for pre-signed messages from dm-sign.sh)')
ps.add_argument('--expect-reply', action='store_true',
help='Tag reply:expected: create a follow-up record (nudge/escalate per policy)')
ps.add_argument('--reply-timeout', default=None,
help='reply:timeout= seconds until nudge/escalation (60..604800, default 3600)')
ps.add_argument('--reply-nudges', default=None,
help='reply:nudges= max nudges before escalation (0..10, default 2)')
ps.add_argument('--reply-escalate', default=None,
help='reply:escalate= identity to alert on timeout')
ps.add_argument('--route', default=None, help='route: attach follow-up to a conversation route')
ps.add_argument('--thread', default=None, help='thread: associate DM with a thread UUID')
ps.add_argument('--tag', action='append', default=[],
help='Raw canonical tag (repeatable): reply:expected, reply:timeout=N, reply:nudges=N, reply:escalate=X, route:R, thread:U')
ps.add_argument('--nudge-meta', default=None,
help='JSON string of nudge metadata (followup_id, nudge_n, dm_id) for audit logging')
ps.add_argument('message')
ps.set_defaults(func=lambda a: dm_send(a.agent, a.target, a.message, verify=not a.no_verify, raw=a.raw, to_agent=a.to))
def _send_with_tags(a):
try:
cli_tags = parse_tags(a.tag)
flag_tags = tags_from_flags(a)
tags = merge_tags(flag_tags, cli_tags, {})
except ValueError as e:
print(f"TAG_ERROR: {e}", file=sys.stderr)
sys.exit(2)
nudge_meta = None
if a.nudge_meta:
try:
import json as _json
nudge_meta = _json.loads(a.nudge_meta)
if not isinstance(nudge_meta, dict):
raise ValueError("nudge-meta must be a JSON object")
except Exception as e:
print(f"NUDGE_META_ERROR: {e}", file=sys.stderr)
sys.exit(2)
return dm_send(a.agent, a.target, a.message, verify=not a.no_verify,
raw=a.raw, to_agent=a.to, tags=tags,
nudge_meta=nudge_meta)
ps.set_defaults(func=_send_with_tags)
psel = sub.add_parser('select', help='Select active conversation')
psel.add_argument('--agent', required=True, choices=VALID_AGENTS)
@@ -343,8 +710,8 @@ def main():
pl.set_defaults(func=lambda a: dm_log(a.n))
pt = sub.add_parser('thread', help='Thread from one agent to another')
pt.add_argument('--from', dest='from_agent', required=True, choices=VALID_AGENTS)
pt.add_argument('--to', dest='to_agent', required=True, choices=VALID_AGENTS)
pt.add_argument('--from', dest='from_agent', required=True, choices=VALID_SENDERS)
pt.add_argument('--to', dest='to_agent', required=True, choices=VALID_RECIPIENTS)
pt.add_argument('--target', required=True)
pt.add_argument('message')
pt.set_defaults(func=lambda a: dm_thread(a.from_agent, a.to_agent, a.target, a.message))
@@ -352,5 +719,114 @@ def main():
args = p.parse_args()
args.func(args)
BOX_API_BASE = os.environ.get("BOX_API_BASE", "https://box.muse-dev.online")
BOX_SIGN_KEY = os.environ.get("BOX_SIGN_KEY", os.path.expanduser("~/.ssh/id_ed25519"))
def _box_sign(identity, endpoint):
"""Sign a box API request. Returns (ts, sig_armored) or (None, None)."""
ts = str(int(time.time()))
payload = f"{ts}\n{endpoint}".encode()
try:
p = subprocess.run(
["ssh-keygen", "-Y", "sign", "-f", BOX_SIGN_KEY, "-n", "box"],
input=payload, capture_output=True, timeout=15)
if p.returncode != 0:
return None, None
return ts, p.stdout.decode()
except Exception:
return None, None
def _register_followup(msg_id, agent, recipient, target, tags):
"""POST a follow-up record to the VM. Never raises -- logs and returns."""
identity_map = {
"bl": "bl",
"opm": "operator-main",
"646": "operator-646",
}
identity = "bl"
ts, sig = _box_sign(identity, "followups")
if not ts or not sig:
log_event({"type": "followup_register_failed", "id": msg_id,
"error": "signing failed"})
return
tags_obj = {}
if tags.get("reply:expected"):
tags_obj["reply:expected"] = True
for key in ("reply:timeout", "reply:nudges", "reply:escalate",
"route", "thread"):
if key in tags:
tags_obj[key] = tags[key]
body = {
"dm_id": msg_id,
"to": recipient,
"from": agent,
"target": target,
"tags": tags_obj,
"dm_sent_ts": datetime.now(timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ"),
}
query = urllib.parse.urlencode({
"identity": identity,
"ts": ts,
"sig": sig,
})
url = f"{BOX_API_BASE}/api/box/followups?{query}"
try:
req = urllib.request.Request(
url, data=json.dumps(body).encode(),
headers={"Content-Type": "application/json",
# Cloudflare 1010-blocks python-urllib's default UA;
# use an identifiable custom UA instead.
"User-Agent": "NetVM-dm.py/1.0 (bl)"},
method="POST")
with urllib.request.urlopen(req, timeout=30) as r:
resp = json.load(r)
log_event({"type": "followup_registered", "id": msg_id,
"request_id": resp.get("request_id"),
"tracked": resp.get("tracked")})
except Exception as e:
log_event({"type": "followup_register_failed", "id": msg_id,
"error": str(e)[:200]})
# Also persist to local followups.json for bl autonomy
try:
f_path = "/home/super/Projects/NetVM/followups.json"
followups = {}
if os.path.exists(f_path):
with open(f_path, "r", encoding="utf-8") as f:
followups = json.load(f)
now_dt = datetime.now(timezone.utc)
timeout_s = int(tags.get("reply:timeout", 3600))
nudges_n = int(tags.get("reply:nudges", 2))
deadline_dt = now_dt + timedelta(seconds=timeout_s)
followups[msg_id] = {
"dm_id": msg_id,
"sender": agent,
"recipient": recipient,
"target": target,
"thread_uuid": tags.get("thread"),
"route": tags.get("route"),
"sent_at": now_dt.isoformat(),
"deadline": deadline_dt.isoformat(),
"timeout_s": timeout_s,
"nudges_allowed": nudges_n,
"nudges_sent": 0,
"escalate_to": tags.get("reply:escalate", "opm"),
"status": "pending"
}
tmp = f"{f_path}.tmp.{os.getpid()}"
with open(tmp, "w", encoding="utf-8") as f:
json.dump(followups, f, indent=2)
os.replace(tmp, f_path)
except Exception as e:
log_event({"type": "local_followup_save_failed", "id": msg_id, "error": str(e)[:200]})
if __name__ == '__main__':
main()
+219
View File
@@ -0,0 +1,219 @@
#!/usr/bin/env python3
"""
followup-sweeper.py — Autonomous follow-up deadline tracking and nudge sweeper.
Monitors pending follow-ups in followups.json, delivers progressive nudges to
recipients when deadlines expire (in-thread first, Main Chat on final nudge),
and executes terminal escalations to opm when all nudges are exhausted.
Usage:
python3 followup-sweeper.py --once
python3 followup-sweeper.py --loop --interval 30
"""
import argparse
import json
import os
import subprocess
import sys
import time
from datetime import datetime, timezone, timedelta
from pathlib import Path
# Paths
NETVM_ROOT = Path("/home/super/Projects/NetVM")
BIN_DIR = NETVM_ROOT / "bin"
FOLLOWUPS_FILE = NETVM_ROOT / "followups.json"
JOB_LOG = NETVM_ROOT / "job-log.jsonl"
DM_PY = BIN_DIR / "dm.py"
def utcnow_dt():
return datetime.now(timezone.utc)
def utcnow_str():
return utcnow_dt().isoformat()
def parse_iso(ts_str):
if not ts_str:
return None
try:
ts_clean = ts_str.replace("Z", "+00:00")
dt = datetime.fromisoformat(ts_clean)
if dt.tzinfo is None:
dt = dt.replace(tzinfo=timezone.utc)
return dt
except Exception:
return None
def load_followups():
if not FOLLOWUPS_FILE.exists():
return {}
try:
with open(FOLLOWUPS_FILE, "r", encoding="utf-8") as f:
return json.load(f)
except Exception:
return {}
def save_followups(data):
tmp_path = f"{FOLLOWUPS_FILE}.tmp.{os.getpid()}"
with open(tmp_path, "w", encoding="utf-8") as f:
json.dump(data, f, indent=2)
os.replace(tmp_path, FOLLOWUPS_FILE)
def append_job_log(entry):
os.makedirs(os.path.dirname(os.path.abspath(JOB_LOG)), exist_ok=True)
with open(JOB_LOG, "a", encoding="utf-8") as f:
f.write(json.dumps(entry) + "\n")
def send_dm(sender, recipient, target, text):
"""Dispatch a DM via dm.py."""
cmd = [
sys.executable,
str(DM_PY),
"send",
"--agent", sender,
"--to", recipient,
"--target", target,
text,
]
try:
res = subprocess.run(cmd, capture_output=True, text=True, timeout=90)
return res.returncode == 0, res.stdout.strip() or res.stderr.strip()
except Exception as e:
return False, str(e)
def sweep_cycle(dry_run=False):
followups = load_followups()
if not followups:
return {"status": "ok", "pending": 0, "nudges_sent": 0, "escalations": 0}
now = utcnow_dt()
nudges_count = 0
escalations_count = 0
modified = False
for dm_id, rec in list(followups.items()):
if rec.get("status") != "pending":
continue
deadline_dt = parse_iso(rec.get("deadline"))
if not deadline_dt or now < deadline_dt:
continue
# Deadline has expired!
nudges_sent = rec.get("nudges_sent", 0)
nudges_allowed = rec.get("nudges_allowed", 2)
sender = rec.get("sender", "opm")
recipient = rec.get("recipient")
orig_target = rec.get("target", "main")
thread_uuid = rec.get("thread_uuid")
if nudges_sent < nudges_allowed:
# Deliver next nudge
nudge_num = nudges_sent + 1
is_final = (nudge_num == nudges_allowed)
# Routing: In-thread first, Main on final nudge
delivery_target = "main" if is_final else orig_target
nudge_text = (
f"[nudge {nudge_num}/{nudges_allowed}] [ref:{dm_id}] "
f"Reminder: awaiting reply to request sent at {rec.get('sent_at', 'earlier')}."
)
if is_final and orig_target != "main":
nudge_text += f" (Origin thread: {orig_target})"
print(f"Sweeper: Sending nudge {nudge_num}/{nudges_allowed} to {recipient}/{delivery_target}...")
if not dry_run:
ok, out = send_dm(sender, recipient, delivery_target, nudge_text)
if ok:
nudges_count += 1
rec["nudges_sent"] = nudge_num
rec["last_nudge_at"] = utcnow_str()
# Calculate interval for next nudge: proportional to timeout or default 10m
timeout_s = rec.get("timeout_s", 1800)
step_s = max(300, timeout_s // (nudges_allowed + 1))
rec["deadline"] = (now + timedelta(seconds=step_s)).isoformat()
modified = True
append_job_log({
"ts": utcnow_str(),
"type": "followup_nudged",
"dm_id": dm_id,
"nudge_num": nudge_num,
"recipient": recipient,
"target": delivery_target,
})
else:
print(f"Sweeper WARNING: nudge send failed: {out}", file=sys.stderr)
else:
nudges_count += 1
else:
# All nudges exhausted: Terminal escalation
escalate_to = rec.get("escalate_to", "opm")
esc_text = (
f"[ESCALATION] Agent {recipient} failed to reply to DM {dm_id} "
f"after {nudges_allowed} nudges. Target was: {orig_target} "
f"(thread: {thread_uuid or 'n/a'}). Request sent: {rec.get('sent_at')}."
)
print(f"Sweeper: Escalating expired follow-up {dm_id} to {escalate_to}...")
if not dry_run:
ok, out = send_dm("bl", escalate_to, "main", esc_text)
rec["status"] = "escalated"
rec["escalated_at"] = utcnow_str()
escalations_count += 1
modified = True
append_job_log({
"ts": utcnow_str(),
"type": "followup_escalated",
"dm_id": dm_id,
"recipient": recipient,
"escalated_to": escalate_to,
})
else:
escalations_count += 1
if modified and not dry_run:
save_followups(followups)
pending_count = sum(1 for r in followups.values() if r.get("status") == "pending")
return {
"status": "ok",
"pending": pending_count,
"nudges_sent": nudges_count,
"escalations": escalations_count,
}
def main():
parser = argparse.ArgumentParser(description="Autonomous follow-up deadline tracker and sweeper")
parser.add_argument("--once", action="store_true", help="Run once and exit (default)")
parser.add_argument("--loop", action="store_true", help="Run continuously in a daemon loop")
parser.add_argument("--interval", type=int, default=60, help="Interval in seconds for loop (default 60)")
parser.add_argument("--dry-run", action="store_true", help="Inspect without sending nudges or updating records")
args = parser.parse_args()
if not args.loop:
stats = sweep_cycle(dry_run=args.dry_run)
print(f"[{datetime.now(timezone.utc).strftime('%H:%M:%SZ')}] Sweep cycle: {stats['pending']} pending, {stats['nudges_sent']} nudges, {stats['escalations']} escalations.")
return
print(f"Starting follow-up sweeper loop (interval={args.interval}s)...")
while True:
try:
stats = sweep_cycle(dry_run=args.dry_run)
print(f"[{datetime.now(timezone.utc).strftime('%H:%M:%SZ')}] Sweep cycle: {stats['pending']} pending, {stats['nudges_sent']} nudges, {stats['escalations']} escalations.")
except Exception as e:
print(f"ERROR in sweeper loop: {e}", file=sys.stderr)
time.sleep(args.interval)
if __name__ == "__main__":
main()
+135 -5
View File
@@ -109,10 +109,98 @@ def render_prompt(template, variables):
result = result.replace(f"{{{key}}}", str(value))
return result
def send_dm(agent, target, message, dry_run=False):
"""Send DM via dm.py"""
# ---- follow-up tracking (DM follow-up system integration) ----------------
# Jobs opt in via a "followup" block in the job JSON:
#
# "followup": {
# "expect_reply": true, # required: enables tracking
# "timeout": "1h", # duration ("30s","15m","2h","1d") or seconds
# # int; default "1h" (3600s)
# "nudges": 2, # 0..10, default 2
# "escalate": "opm", # identity string, default "opm"
# "route": "646-pip-coord" # optional route_id
# }
#
# The dispatcher translates this into dm.py --tag flags using the canonical
# vocabulary (box-threads/DEPLOY-DECISIONS.md). dm.py strips the tags from
# delivered text and creates a dm_followup request-store record after
# SENT+VERIFIED. Jobs without a followup block behave exactly as today.
#
# LIMITATIONS (v1):
# - The heartbeat job NEVER gets follow-ups (loopback health check).
# Hardcoded guard below; a followup block on heartbeat is ignored loudly.
# - Sidechat sends (muse-chat-api.py direct path) do not go through dm.py,
# so --tag flags cannot attach. v2 needs a record-creation path that does
# not send (e.g. POST /api/box/followups, or a bl->VM queue; bl cannot
# currently SSH to the VM). The dispatcher logs a warning when a
# sidechat-targeted job has followup enabled.
HEARTBEAT_JOB_NAME = "heartbeat"
def parse_followup_duration(value):
"""Parse a followup timeout into seconds. Accepts int (seconds) or
strings like '30s', '15m', '2h', '1d'. Returns int seconds.
Raises ValueError on bad input."""
if isinstance(value, int) and not isinstance(value, bool):
s = value
elif isinstance(value, str):
m = re.fullmatch(r"(\d+)\s*([smhd])?", value.strip().lower())
if not m:
raise ValueError("bad duration %r" % (value,))
n = int(m.group(1))
unit = m.group(2) or "s"
s = n * {"s": 1, "m": 60, "h": 3600, "d": 86400}[unit]
else:
raise ValueError("timeout must be int seconds or duration string")
if not 60 <= s <= 604800:
raise ValueError("timeout must be 60..604800s (1m..7d), got %d" % s)
return s
def build_followup_tags(followup):
"""Translate a job's followup block into dm.py --tag arguments.
Returns a flat list like ['--tag', 'reply:timeout=3600', ...].
Returns [] if followup is falsy or expect_reply is not true.
Raises ValueError on invalid config (caller logs a warning and sends
the DM untagged -- the job itself must never fail over this)."""
if not followup or not followup.get("expect_reply"):
return []
args = []
# Bare trigger. dm.py's parse_tags splits each --tag on '='; an empty
# value means "present". If the deployed dm.py requires a non-empty
# value for this key, use 'reply:expected=true' instead.
args += ["--tag", "reply:expected="]
if "timeout" in followup:
s = parse_followup_duration(followup["timeout"])
args += ["--tag", "reply:timeout=%d" % s]
if "nudges" in followup:
n = followup["nudges"]
if not isinstance(n, int) or isinstance(n, bool) or not 0 <= n <= 10:
raise ValueError("nudges must be int 0..10")
args += ["--tag", "reply:nudges=%d" % n]
if "escalate" in followup:
e = followup["escalate"]
if not isinstance(e, str) or not re.fullmatch(r"[a-z0-9_-]{1,64}", e):
raise ValueError("escalate must be an identity string")
args += ["--tag", "reply:escalate=%s" % e]
if "route" in followup:
r = followup["route"]
if not isinstance(r, str) or not re.fullmatch(r"[a-z0-9_-]{1,64}", r):
raise ValueError("route must be a route_id string")
args += ["--tag", "route:%s" % r]
# 'thread' is intentionally not settable from job JSON; it names a
# specific existing thread and is filled by the dispatcher when known.
return args
def send_dm(agent, target, message, dry_run=False, followup_tags=None):
"""Send DM via dm.py. followup_tags: flat ['--tag', 'k=v', ...] list
from build_followup_tags(), or None."""
if dry_run:
print(f"[DRY RUN] Would send to {agent} ({target}):")
if followup_tags:
print(f"[DRY RUN] With follow-up tags: {' '.join(followup_tags)}")
print(message[:200] + "..." if len(message) > 200 else message)
return "dry-run-id"
@@ -120,8 +208,9 @@ def send_dm(agent, target, message, dry_run=False):
if HAS_RATE_LIMITER:
rate_limit_wait(agent)
cmd = [str(DM_PY), "send", "--agent", "opm", "--to", agent,
"--target", target, message]
cmd = ([str(DM_PY), "send", "--agent", "opm", "--to", agent,
"--target", target]
+ (followup_tags or []) + [message])
result = subprocess.run(cmd, capture_output=True, text=True, timeout=60)
if result.returncode != 0:
@@ -202,8 +291,33 @@ def main():
"job_name": job_name,
"date": datetime.now(timezone.utc).strftime("%Y-%m-%d"),
"datetime": datetime.now(timezone.utc).isoformat(),
"prev_job_id": os.environ.get("CHAIN_PREV_JOB_ID", ""),
"prev_result": os.environ.get("CHAIN_PREV_RESULT", ""),
}
# Follow-up tracking (opt-in via job JSON "followup" block; see helpers).
# The heartbeat job is a loopback health check and must never be tracked.
followup_cfg = job.get("followup")
followup_tags = []
if followup_cfg:
if job_name == HEARTBEAT_JOB_NAME:
print(f"Warning: job '{job_name}' must not use follow-up "
f"tracking (loopback); ignoring followup block",
file=sys.stderr)
log_event("job_followup_skipped",
{"job_id": job_id, "reason": "heartbeat_loopback"})
else:
try:
followup_tags = build_followup_tags(followup_cfg)
if followup_tags:
log_event("job_followup_armed",
{"job_id": job_id, "tags": followup_tags})
except ValueError as e:
print(f"Warning: invalid followup block: {e}; "
f"sending untagged", file=sys.stderr)
log_event("job_followup_invalid",
{"job_id": job_id, "error": str(e)})
# Render prompt
prompt_template = job.get("prompt_template", "")
if not prompt_template:
@@ -262,6 +376,11 @@ def main():
capture_uuid = False
else:
target = "main"
# dm_target override: job JSON can specify a dm.py --target
# (sidechat name/UUID) for tracked sends to a thread.
_dt = job.get("dm_target")
if _dt and isinstance(_dt, str) and _dt.strip():
target = _dt.strip()
# Log job_sent
log_event("job_sent", {
@@ -297,11 +416,22 @@ def main():
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, dry_run=dry_run)
msg_id = send_dm(agent, target, dm_message,
dry_run=dry_run, followup_tags=followup_tags)
if msg_id and not dry_run:
print(f"Dispatched job {job_id} to {agent} (DM: {msg_id})")
+566
View File
@@ -0,0 +1,566 @@
#!/usr/bin/env python3
"""
response-harvester.py — Fleet agent readback and response harvesting daemon.
Monitors Chromebox agents (muse, pip, 646, opm), harvests incoming messages from
Main Chat and registered sidechats, maintains persistent watermarks, appends to
chat-history.jsonl, resolves pending follow-ups, and records [RESULT] completions
in job-log.jsonl.
Features:
- Direct CDP over host veth interfaces (fast, no sudo needed).
- cdp_queue integration with PRIORITY_LOW (never blocks operator/DMs).
- URL state preservation (restores browser to initial thread/main via Ctrl+J).
- Bounded scroll-back for virtualized DOM (#hatch-chat-scroll).
- Per-node fault isolation (CDP errors on one node do not abort the cycle).
- Dual-mode execution (--once for systemd timers/CLI, --loop for daemon).
Usage:
python3 response-harvester.py --once
python3 response-harvester.py --loop --interval 30
python3 response-harvester.py --agent 646 --once
"""
import argparse
import hashlib
import json
import os
import re
import subprocess
import sys
import time
import urllib.error
import urllib.request
import websocket
from datetime import datetime, timezone
from pathlib import Path
# Paths
NETVM_ROOT = Path("/home/super/Projects/NetVM")
BIN_DIR = NETVM_ROOT / "bin"
LOGS_DIR = NETVM_ROOT / "logs"
CHAT_HISTORY_LOG = LOGS_DIR / "chat-history.jsonl"
WATERMARKS_FILE = NETVM_ROOT / "siphon-watermarks.json"
JOB_LOG = NETVM_ROOT / "job-log.jsonl"
FOLLOWUPS_FILE = NETVM_ROOT / "followups.json"
JOB_SIDECHATS_FILE = NETVM_ROOT / "job-sidechats.json"
WAKE_SIDECHATS_FILE = Path("/home/super/sidechat-wake/wake-sidechats.json")
JOBS_DIR = NETVM_ROOT / "jobs"
DISPATCH_PY = BIN_DIR / "job-dispatch.py"
# Ensure bin is in sys.path
sys.path.insert(0, str(BIN_DIR))
try:
from cdp_queue import cdp_slot, PRIORITY_LOW
HAS_CDP_QUEUE = True
except ImportError:
HAS_CDP_QUEUE = False
try:
import netvm_registry
HAS_REGISTRY = True
except ImportError:
HAS_REGISTRY = False
VALID_AGENTS = ["muse", "pip", "646", "opm"]
DEFAULT_PORTS = {"muse": 9410, "pip": 9420, "646": 9430, "opm": 9440}
def utcnow():
return datetime.now(timezone.utc).isoformat()
def get_node_network(node):
"""Derive veth peer IP and CDP port from node identity."""
tag = hashlib.sha256(node.encode()).hexdigest()[:8]
idx = int(tag[:3], 16) % 200 + 10
peer_ip = f"10.201.{idx}.2"
port = None
if HAS_REGISTRY:
try:
port = netvm_registry.port_for(node)
except Exception:
pass
if not port:
port = DEFAULT_PORTS.get(node, 9410)
return peer_ip, port
def load_json_file(path, default=None):
if default is None:
default = {}
if not os.path.exists(path):
return default
try:
with open(path, "r", encoding="utf-8") as f:
return json.load(f)
except Exception:
return default
def save_json_file(path, data):
tmp_path = f"{path}.tmp.{os.getpid()}"
with open(tmp_path, "w", encoding="utf-8") as f:
json.dump(data, f, indent=2)
os.replace(tmp_path, path)
def append_jsonl(path, record):
os.makedirs(os.path.dirname(os.path.abspath(path)), exist_ok=True)
with open(path, "a", encoding="utf-8") as f:
f.write(json.dumps(record) + "\n")
def get_monitored_threads(target_agent=None):
"""
Build dict of threads to monitor per agent:
{ agent: [ {"id": "main", "name": "main"}, {"id": "<uuid>", "name": "<alias>"} ] }
"""
agents = [target_agent] if target_agent else VALID_AGENTS
threads_by_agent = {a: [{"id": "main", "name": "Main Chat"}] for a in agents}
state_files = [JOB_SIDECHATS_FILE, WAKE_SIDECHATS_FILE]
for sf in state_files:
if not sf.exists():
continue
try:
data = json.loads(sf.read_text(encoding="utf-8"))
for key, val in data.items():
if key.startswith("_"):
continue
if isinstance(val, dict):
uuid = val.get("thread_uuid") or val.get("uuid")
agent = val.get("agent", "opm")
elif isinstance(val, str):
uuid = val
agent = "opm"
else:
continue
if uuid and agent in threads_by_agent:
# Avoid duplicate threads
existing = [t["id"] for t in threads_by_agent[agent]]
if uuid not in existing:
threads_by_agent[agent].append({"id": uuid, "name": key})
except Exception:
continue
return threads_by_agent
class CDPClient:
"""Lightweight direct CDP client over host veth."""
def __init__(self, node, peer_ip, port, timeout=10):
self.node = node
self.peer_ip = peer_ip
self.port = port
self.timeout = timeout
self.ws = None
self.msg_id = 0
def connect(self):
url = f"http://{self.peer_ip}:{self.port}/json/list"
req = urllib.request.Request(url)
with urllib.request.urlopen(req, timeout=self.timeout) as resp:
targets = json.load(resp)
pages = [t for t in targets if t.get("type") == "page"]
if not pages:
raise RuntimeError(f"No page target found on CDP for {self.node}")
ws_url = pages[0]["webSocketDebuggerUrl"]
self.ws = websocket.create_connection(ws_url, timeout=self.timeout)
def send_cmd(self, method, params=None):
self.msg_id += 1
cid = self.msg_id
payload = {"id": cid, "method": method, "params": params or {}}
self.ws.send(json.dumps(payload))
while True:
raw = self.ws.recv()
data = json.loads(raw)
if data.get("id") == cid:
return data
def evaluate(self, expr, await_promise=False):
res = self.send_cmd(
"Runtime.evaluate",
{"expression": expr, "returnByValue": True, "awaitPromise": await_promise},
)
result = res.get("result", {}).get("result", {})
if res.get("result", {}).get("exceptionDetails"):
desc = res["result"]["exceptionDetails"].get("text", "JS exception")
raise RuntimeError(f"CDP eval error: {desc}")
return result.get("value")
def dispatch_key(self, key, code, modifiers=0):
self.send_cmd(
"Input.dispatchKeyEvent",
{
"type": "rawKeyDown",
"key": key,
"code": code,
"modifiers": modifiers,
"windowsVirtualKeyCode": 74 if code == "KeyJ" else 0,
},
)
self.send_cmd(
"Input.dispatchKeyEvent",
{
"type": "keyUp",
"key": key,
"code": code,
"modifiers": modifiers,
"windowsVirtualKeyCode": 74 if code == "KeyJ" else 0,
},
)
def close(self):
if self.ws:
try:
self.ws.close()
except Exception:
pass
self.ws = None
DOM_EXTRACT_JS = """(() => {
const els = [...document.querySelectorAll('[data-message-id]')];
return els.map(m => {
const id = m.getAttribute('data-message-id');
const ps = [...m.querySelectorAll('p')].map(p => (p.innerText || '').trim()).filter(Boolean);
let text = ps.join('\\n');
if (!text) {
text = (m.innerText || '').replace(/^(Assistant message:|User message:)\\s*/i, '').trim();
}
const t = m.querySelector('time');
return {
id: id,
author: id.startsWith('assistant-msg') ? 'assistant' : 'user',
text: text,
ts: t ? (t.getAttribute('datetime') || t.innerText || null) : null
};
});
})()"""
def scrape_thread_messages(cdp, thread_id, watermark, max_scrollbacks=3):
"""Scrape messages with bounded scroll-back if watermark is out of view."""
messages = cdp.evaluate(DOM_EXTRACT_JS) or []
# If watermark exists and is already in view, or no watermark, no scroll-back needed
seen_ids = {m["id"] for m in messages if m.get("id")}
if watermark and watermark not in seen_ids and max_scrollbacks > 0:
# Bounded scroll-back loop
for _ in range(max_scrollbacks):
cdp.evaluate("""(() => {
const sc = document.getElementById('hatch-chat-scroll');
if (sc) sc.scrollTop = 0;
})()""")
time.sleep(0.8)
older = cdp.evaluate(DOM_EXTRACT_JS) or []
for m in older:
if m.get("id") and m["id"] not in seen_ids:
messages.insert(0, m)
seen_ids.add(m["id"])
if watermark in seen_ids:
break
return messages
def harvest_agent_thread(cdp, agent, thread_info, watermarks, followups, dry_run=False):
"""
Harvests new messages for a single thread, preserves URL state,
and returns (new_messages, new_watermark, job_results_count).
"""
thread_id = thread_info["id"]
thread_name = thread_info["name"]
wm_key = f"{agent}:{thread_id}"
last_wm = watermarks.get(wm_key, "")
# 1. Capture current URL before navigating
initial_url = cdp.evaluate("window.location.href") or "https://muse.ai/"
is_init_main = "/thread/" not in initial_url
# 2. Navigate to target thread if not already there
try:
if thread_id == "main":
if not is_init_main:
# Dispatch Ctrl+J (modifier 2 = Control)
cdp.dispatch_key("j", "KeyJ", modifiers=2)
time.sleep(2.0)
else:
target_url = f"https://muse.ai/thread/{thread_id}"
if initial_url.strip() != target_url:
cdp.evaluate(f"window.location.href = {json.dumps(target_url)}")
# Settle wait
time.sleep(2.5)
# 3. Scrape messages
raw_messages = scrape_thread_messages(cdp, thread_id, last_wm)
finally:
# 4. State preservation: restore browser back to initial state
try:
curr_url = cdp.evaluate("window.location.href") or ""
if is_init_main:
if "/thread/" in curr_url:
cdp.dispatch_key("j", "KeyJ", modifiers=2)
else:
if curr_url.strip() != initial_url.strip():
cdp.evaluate(f"window.location.href = {json.dumps(initial_url)}")
except Exception:
pass
if not raw_messages:
return [], last_wm, 0
# 5. Filter for new messages based on watermark
new_messages = []
if not last_wm:
# Initial run on this thread: watermark at current latest message to avoid flooding backlog
new_wm = raw_messages[-1]["id"]
return [], new_wm, 0
else:
# Find index of last_wm
wm_idx = -1
for i, m in enumerate(raw_messages):
if m["id"] == last_wm:
wm_idx = i
break
if wm_idx >= 0:
new_messages = raw_messages[wm_idx + 1 :]
else:
# Watermark not found in loaded window (older than scroll limit)
# Process all visible messages that are newer than timestamp or just unread tail
new_messages = raw_messages
if not new_messages:
return [], last_wm, 0
new_wm = new_messages[-1]["id"]
job_results = 0
# 6. Ingest new messages
for msg in new_messages:
mid = msg.get("id", "")
author = msg.get("author", "unknown")
text = msg.get("text", "")
msg_ts = msg.get("ts") or utcnow()
# Append to chat-history.jsonl
record = {
"ts": utcnow(),
"agent": agent,
"thread_id": thread_id,
"thread_name": thread_name,
"msg_id": mid,
"author": author,
"text": text,
"source_ts": msg_ts,
}
if not dry_run:
append_jsonl(CHAT_HISTORY_LOG, record)
# Check for [RESULT <job_id>] in assistant messages
if author == "assistant":
m_res = re.search(r"\[RESULT\s+([A-Za-z0-9_-]+)\]\s*(.*)", text, re.S)
if m_res:
job_id = m_res.group(1).strip()
result_text = m_res.group(2).strip()
is_fail = result_text.startswith("FAILED") or result_text.startswith("UNABLE") or result_text.startswith("FAIL")
job_results += 1
job_record = {
"ts": utcnow(),
"type": "job_result",
"job_id": job_id,
"agent": agent,
"success": not is_fail,
"result_snippet": result_text[:300],
"thread_id": thread_id,
"msg_id": mid,
}
if not dry_run:
append_jsonl(JOB_LOG, job_record)
# Trigger chain_next if configured
trigger_chain_next(job_id, result_text)
# Check and clear pending follow-ups
clear_matching_followups(followups, agent, thread_id, mid, text, dry_run)
return new_messages, new_wm, job_results
def trigger_chain_next(job_id, result_text):
"""If the completed job has a chain_next property, dispatch it with context."""
# Job ID format: <name>-<timestamp>-<uuid>
parts = job_id.split("-")
if len(parts) < 3:
return
job_name = "-".join(parts[:-2])
job_file = JOBS_DIR / f"{job_name}.json"
if not job_file.exists():
return
try:
with open(job_file, "r", encoding="utf-8") as f:
cfg = json.load(f)
chain_next = cfg.get("chain_next")
if chain_next and (JOBS_DIR / f"{chain_next}.json").exists():
env = os.environ.copy()
env["CHAIN_PREV_JOB_ID"] = job_id
env["CHAIN_PREV_RESULT"] = result_text[:1000]
subprocess.Popen(
[sys.executable, str(DISPATCH_PY), chain_next],
env=env,
stdout=subprocess.DEVNULL,
stderr=subprocess.DEVNULL,
)
except Exception:
pass
def clear_matching_followups(followups, agent, thread_id, mid, text, dry_run=False):
"""Resolve follow-up records if an assistant message is detected in the thread."""
if not followups:
return
modified = False
for f_id, f_rec in followups.items():
if f_rec.get("status") != "pending":
continue
if f_rec.get("recipient") != agent:
continue
# Match either exact thread_uuid, or target alias 'main'
match_thread = False
if f_rec.get("target") == "main" and thread_id == "main":
match_thread = True
elif f_rec.get("thread_uuid") and f_rec.get("thread_uuid") == thread_id:
match_thread = True
if match_thread:
f_rec["status"] = "resolved"
f_rec["resolved_at"] = utcnow()
f_rec["resolved_by_mid"] = mid
f_rec["resolved_snippet"] = text[:150]
modified = True
if modified and not dry_run:
save_json_file(FOLLOWUPS_FILE, followups)
def harvest_cycle(target_agent=None, dry_run=False, output_json=False):
"""Execute one full harvest cycle across agents and threads."""
watermarks = load_json_file(WATERMARKS_FILE)
followups = load_json_file(FOLLOWUPS_FILE)
monitored = get_monitored_threads(target_agent)
cycle_stats = {
"timestamp": utcnow(),
"agents": {},
"total_new_messages": 0,
"total_job_results": 0,
}
for agent, thread_list in monitored.items():
agent_stats = {"status": "ok", "threads": {}, "new_messages": 0, "job_results": 0}
peer_ip, port = get_node_network(agent)
# Wrap in CDP slot if available
slot_ctx = (
cdp_slot(agent, priority=PRIORITY_LOW, timeout=10)
if HAS_CDP_QUEUE
else None
)
try:
if slot_ctx:
with slot_ctx:
cdp = CDPClient(agent, peer_ip, port, timeout=8)
cdp.connect()
try:
for t_info in thread_list:
new_msgs, new_wm, j_res = harvest_agent_thread(
cdp, agent, t_info, watermarks, followups, dry_run
)
wm_key = f"{agent}:{t_info['id']}"
if new_wm and not dry_run:
watermarks[wm_key] = new_wm
agent_stats["threads"][t_info["name"]] = len(new_msgs)
agent_stats["new_messages"] += len(new_msgs)
agent_stats["job_results"] += j_res
finally:
cdp.close()
else:
cdp = CDPClient(agent, peer_ip, port, timeout=8)
cdp.connect()
try:
for t_info in thread_list:
new_msgs, new_wm, j_res = harvest_agent_thread(
cdp, agent, t_info, watermarks, followups, dry_run
)
wm_key = f"{agent}:{t_info['id']}"
if new_wm and not dry_run:
watermarks[wm_key] = new_wm
agent_stats["threads"][t_info["name"]] = len(new_msgs)
agent_stats["new_messages"] += len(new_msgs)
agent_stats["job_results"] += j_res
finally:
cdp.close()
except Exception as e:
agent_stats["status"] = "error"
agent_stats["error"] = str(e)[:200]
cycle_stats["agents"][agent] = agent_stats
cycle_stats["total_new_messages"] += agent_stats["new_messages"]
cycle_stats["total_job_results"] += agent_stats["job_results"]
if not dry_run:
save_json_file(WATERMARKS_FILE, watermarks)
# Output formatting
if output_json:
print(json.dumps(cycle_stats))
else:
ts_short = cycle_stats["timestamp"].split("T")[1][:8]
summary_parts = []
for ag, st in cycle_stats["agents"].items():
if st["status"] == "ok":
summary_parts.append(f"{ag}: {st['new_messages']} msgs ({st['job_results']} results)")
else:
summary_parts.append(f"{ag}: [UNREACHABLE: {st.get('error', 'err')[:40]}]")
print(f"[{ts_short}Z] Harvest cycle: {', '.join(summary_parts)}")
return cycle_stats
def main():
parser = argparse.ArgumentParser(description="Fleet agent readback and response harvester")
parser.add_argument("--once", action="store_true", help="Run once and exit (default)")
parser.add_argument("--loop", action="store_true", help="Run continuously in a daemon loop")
parser.add_argument("--interval", type=int, default=30, help="Interval in seconds for --loop (default 30)")
parser.add_argument("--agent", choices=VALID_AGENTS, default=None, help="Harvest only specific agent")
parser.add_argument("--dry-run", action="store_true", help="Scrape without persisting watermarks or logs")
parser.add_argument("--json", action="store_true", help="Output summary as JSON")
args = parser.parse_args()
# Default to --once if --loop is not provided
if not args.loop:
harvest_cycle(target_agent=args.agent, dry_run=args.dry_run, output_json=args.json)
return
print(f"Starting response-harvester daemon (interval={args.interval}s, agent={args.agent or 'all'})...")
while True:
try:
harvest_cycle(target_agent=args.agent, dry_run=args.dry_run, output_json=args.json)
except Exception as e:
print(f"ERROR in harvest loop: {e}", file=sys.stderr)
time.sleep(args.interval)
if __name__ == "__main__":
main()
+2957
View File
File diff suppressed because it is too large Load Diff
+16
View File
@@ -0,0 +1,16 @@
{
"heartbeat-opm": "5bd5b350-806c-4e04-939c-7ca7d1bd20df",
"646-pip-coord": {
"thread_uuid": "4466d0c1-7961-4cf3-b99d-1ab7c38484c2",
"agent": "pip"
},
"646-opm-coord": {
"thread_uuid": "4139dd4e-96fe-4222-a99f-82a59a7b0eeb",
"agent": "opm"
},
"test-auto-prov": {
"thread_uuid": "b96dd020-b429-4f5f-9e52-dbe8f805ac6a",
"agent": "646",
"created_at": "2026-10-04T16:20:33.919952+00:00"
}
}