fix: placement verification and sidechat policy enforcement
- dm.py: placement-aware verify (verify_placement), checked_uuid logging, placement_mismatch events - main-chat-watchdog.py: P0 alerts on placement_mismatch - box-ctl.py: notify routes to sidechat, box policy command - jobs: ops-audit and pipe-demo use sidechats - response-harvester.py: chain deduplication Session: sidechat/ops-restore
This commit is contained in:
+225
-9
@@ -36,9 +36,23 @@ CTL_LOG = NETVM_ROOT / "box-ctl.jsonl"
|
||||
JOB_LOG = NETVM_ROOT / "job-log.jsonl"
|
||||
DISPATCHER = BIN / "job-dispatch.py"
|
||||
DM_PY = BIN / "dm.py"
|
||||
POLICY_FILE = NETVM_ROOT / "CHAT_POLICY.md"
|
||||
DM_LOG = NETVM_ROOT / "dm-log.jsonl"
|
||||
WATCHDOG_STATE = NETVM_ROOT / "main-chat-watchdog.state"
|
||||
|
||||
NAME_RE = re.compile(r"^[a-z0-9-]{1,64}$")
|
||||
VALID_AGENTS = {"muse", "pip", "646", "opm"}
|
||||
|
||||
# Default sidechat per agent for `box notify`. Mirrors super-cli.py
|
||||
# DEFAULT_AGENT_SIDECHATS (kept in sync manually; super-cli is the
|
||||
# canonical copy). Sidechat-first policy: notify never targets main
|
||||
# chat unless --allow-main-chat is passed explicitly.
|
||||
NOTIFY_SIDECHATS = {
|
||||
"646": "646 tasks",
|
||||
"opm": "heartbeat",
|
||||
"pip": "646-pip-coord",
|
||||
"muse": "muse tasks",
|
||||
}
|
||||
VALID_ON_FAILURE = {"retry", "alert", "ignore"}
|
||||
KNOWN_PLACEHOLDERS = {"job_id", "job_name", "datetime", "date", "last_run"}
|
||||
PROTOCOL_LITERALS = ("[REQ", "[CONFIRM", "[JOB", "[RESULT")
|
||||
@@ -602,19 +616,38 @@ def act_job_trigger(name):
|
||||
out(True, name=name, triggered=True)
|
||||
|
||||
|
||||
def act_notify(agent, message):
|
||||
def _valid_sidechat_name(name):
|
||||
# Sidechat names may contain spaces ("646 tasks"); reject only
|
||||
# control characters and enforce a sane length. dm.py resolves
|
||||
# the name (alias -> job-sidechats.json -> autoprovision).
|
||||
return bool(name) and len(name) <= 64 and not any(
|
||||
ord(c) < 32 or ord(c) == 127 for c in name)
|
||||
|
||||
|
||||
def act_notify(agent, message, sidechat=None, allow_main_chat=False):
|
||||
if agent not in VALID_AGENTS:
|
||||
fail("BAD_NAME", f"agent must be one of {sorted(VALID_AGENTS)}")
|
||||
if not message or len(message) > 1000:
|
||||
fail("INVALID_JOB", "message must be 1–1000 chars")
|
||||
# dm.py send --agent opm --to <agent> --target main "<message>"
|
||||
fail("INVALID_JOB", "message must be 1\u20131000 chars")
|
||||
# Sidechat-first policy: default to the agent's sidechat, never main.
|
||||
if allow_main_chat:
|
||||
target = "main"
|
||||
extra = ["--allow-main-chat"]
|
||||
else:
|
||||
target = sidechat if sidechat else NOTIFY_SIDECHATS.get(agent)
|
||||
if not target:
|
||||
fail("BAD_NAME", f"no default sidechat for agent {agent}; pass --sidechat <name>")
|
||||
if sidechat and not _valid_sidechat_name(sidechat):
|
||||
fail("BAD_NAME", "sidechat name must be 1-64 chars, no control characters")
|
||||
extra = []
|
||||
r = run([sys.executable, str(DM_PY), "send",
|
||||
"--agent", "opm", "--to", agent, "--target", "main", message],
|
||||
"--agent", "opm", "--to", agent, "--target", target,
|
||||
*extra, message],
|
||||
timeout=120)
|
||||
if r.returncode != 0:
|
||||
fail("DISPATCH_FAILED", f"dm.py send failed: {(r.stderr or r.stdout).strip()[-500:]}")
|
||||
audit("notify", agent)
|
||||
out(True, agent=agent, sent=True)
|
||||
out(True, agent=agent, target=target, sent=True)
|
||||
|
||||
|
||||
def act_fleet_status():
|
||||
@@ -646,6 +679,156 @@ def act_dm_log(limit=50):
|
||||
fail("DM_LOG_ERROR", "failed to read dm log", {"stderr": r.stderr})
|
||||
|
||||
|
||||
|
||||
def _policy_meta():
|
||||
"""Parse Version:/Date: from CHAT_POLICY.md."""
|
||||
version, date = "unknown", "unknown"
|
||||
try:
|
||||
for line in POLICY_FILE.read_text(encoding="utf-8", errors="replace").splitlines():
|
||||
s = line.strip()
|
||||
if s.startswith("**Version:**"):
|
||||
version = s.split("**Version:**", 1)[1].strip().rstrip("\\")
|
||||
elif s.startswith("**Date:**"):
|
||||
date = s.split("**Date:**", 1)[1].strip().rstrip("\\")
|
||||
except OSError:
|
||||
pass
|
||||
return version, date
|
||||
|
||||
|
||||
def _watchdog_state():
|
||||
"""Read the main-chat watchdog watermark (None if missing)."""
|
||||
try:
|
||||
return json.loads(WATCHDOG_STATE.read_text(encoding="utf-8"))
|
||||
except Exception:
|
||||
return None
|
||||
|
||||
|
||||
def _policy_scan():
|
||||
"""Classify dm-log.jsonl events for the sidechat-first policy.
|
||||
|
||||
Mirrors main-chat-watchdog.py classification:
|
||||
VIOLATION : type=sent, target=main, no tags.allow_main_chat
|
||||
AUTHORIZED: type=sent, target=main, tags.allow_main_chat set
|
||||
BLOCKED : type=main_chat_blocked (gate worked)
|
||||
Events before the allow_main_chat audit marker was adopted cannot be
|
||||
classified, so the compliance window starts at marker adoption.
|
||||
"""
|
||||
per_agent = {}
|
||||
legacy_untagged = 0
|
||||
malformed = 0
|
||||
adoption_ts = None
|
||||
total_scanned = 0
|
||||
|
||||
def bucket(agent):
|
||||
return per_agent.setdefault(agent, {
|
||||
"blocked": 0,
|
||||
"authorized_main": 0,
|
||||
"violations": 0,
|
||||
"total_sends": 0,
|
||||
})
|
||||
|
||||
try:
|
||||
with open(DM_LOG, "r", encoding="utf-8", errors="replace") as f:
|
||||
lines = f.readlines()
|
||||
except OSError:
|
||||
return None, {"error": f"cannot read {DM_LOG}"}
|
||||
|
||||
for line in lines:
|
||||
line = line.strip()
|
||||
if not line:
|
||||
continue
|
||||
try:
|
||||
ev = json.loads(line)
|
||||
except json.JSONDecodeError:
|
||||
continue
|
||||
tags = ev.get("tags")
|
||||
if isinstance(tags, dict) and "allow_main_chat" in tags:
|
||||
ts = ev.get("ts") or ""
|
||||
if adoption_ts is None or ts < adoption_ts:
|
||||
adoption_ts = ts
|
||||
|
||||
for line in lines:
|
||||
line = line.strip()
|
||||
if not line:
|
||||
continue
|
||||
total_scanned += 1
|
||||
try:
|
||||
ev = json.loads(line)
|
||||
except json.JSONDecodeError:
|
||||
malformed += 1
|
||||
continue
|
||||
ts = ev.get("ts") or ""
|
||||
if adoption_ts is not None and ts < adoption_ts:
|
||||
if ev.get("type") == "sent" and ev.get("target") == "main":
|
||||
legacy_untagged += 1
|
||||
continue
|
||||
etype = ev.get("type")
|
||||
agent = ev.get("agent") or "unknown"
|
||||
if etype == "main_chat_blocked":
|
||||
bucket(agent)["blocked"] += 1
|
||||
elif etype == "sent":
|
||||
bucket(agent)["total_sends"] += 1
|
||||
if ev.get("target") == "main":
|
||||
tags = ev.get("tags") or {}
|
||||
if tags.get("allow_main_chat"):
|
||||
bucket(agent)["authorized_main"] += 1
|
||||
else:
|
||||
bucket(agent)["violations"] += 1
|
||||
|
||||
return per_agent, {
|
||||
"window_start": adoption_ts,
|
||||
"events_scanned": total_scanned,
|
||||
"malformed": malformed,
|
||||
"legacy_untagged_main_sends": legacy_untagged,
|
||||
}
|
||||
|
||||
|
||||
def act_policy():
|
||||
audit("policy")
|
||||
version, date = _policy_meta()
|
||||
per_agent, meta = _policy_scan()
|
||||
if per_agent is None:
|
||||
fail("POLICY_ERROR", "failed to read dm log", meta)
|
||||
totals = {"blocked": 0, "authorized_main": 0, "violations": 0, "total_sends": 0}
|
||||
for counts in per_agent.values():
|
||||
for k in totals:
|
||||
totals[k] += counts[k]
|
||||
out(True,
|
||||
version=version,
|
||||
date=date,
|
||||
rule=("No DM naturally lands in main chat. Main requires explicit "
|
||||
"opt-in (--allow-main-chat / allow_main_chat)."),
|
||||
compliance_window=meta,
|
||||
watchdog_state=_watchdog_state(),
|
||||
agents=per_agent,
|
||||
totals=totals,
|
||||
status="clean" if totals["violations"] == 0 else "violations found")
|
||||
|
||||
|
||||
def act_policy_check(agent):
|
||||
if agent not in VALID_AGENTS:
|
||||
fail("BAD_NAME", f"agent must be one of {sorted(VALID_AGENTS)}")
|
||||
audit("policy-check", agent)
|
||||
per_agent, meta = _policy_scan()
|
||||
if per_agent is None:
|
||||
fail("POLICY_ERROR", "failed to read dm log", meta)
|
||||
counts = per_agent.get(agent, {
|
||||
"blocked": 0, "authorized_main": 0, "violations": 0, "total_sends": 0})
|
||||
out(True, agent=agent,
|
||||
compliance_window=meta, **counts,
|
||||
status="clean" if counts["violations"] == 0 else "violations found")
|
||||
|
||||
|
||||
def act_policy_show():
|
||||
audit("policy-show")
|
||||
try:
|
||||
text = POLICY_FILE.read_text(encoding="utf-8", errors="replace")
|
||||
except OSError as e:
|
||||
fail("POLICY_ERROR", f"cannot read {POLICY_FILE}: {e}")
|
||||
version, date = _policy_meta()
|
||||
out(True, version=version, date=date, path=str(POLICY_FILE), policy=text)
|
||||
|
||||
|
||||
def act_vars_list():
|
||||
try:
|
||||
from variables import Variables
|
||||
@@ -904,7 +1087,14 @@ loop actions:
|
||||
loop-remediate [--dry-run]
|
||||
|
||||
notify:
|
||||
notify <agent> <message>"""
|
||||
notify <agent> <message> [--sidechat <name>] [--allow-main-chat]
|
||||
Send a DM to an agent's default sidechat (never main chat unless
|
||||
--allow-main-chat is passed explicitly).
|
||||
|
||||
policy:
|
||||
policy chat policy version, rule, and compliance summary
|
||||
policy check <agent> per-agent compliance detail
|
||||
policy show print the full CHAT_POLICY.md"""
|
||||
|
||||
|
||||
def main(argv):
|
||||
@@ -1055,9 +1245,35 @@ def main(argv):
|
||||
dry = "--dry-run" in rest
|
||||
act_loop_remediate(dry_run=dry)
|
||||
elif action == "notify":
|
||||
if len(rest) != 2:
|
||||
fail("BAD_NAME", "usage: notify <agent> <message>")
|
||||
act_notify(rest[0], rest[1])
|
||||
# notify <agent> <message> [--sidechat <name>] [--allow-main-chat]
|
||||
args = list(rest)
|
||||
allow_main = False
|
||||
sidechat = None
|
||||
if "--allow-main-chat" in args:
|
||||
allow_main = True
|
||||
args.remove("--allow-main-chat")
|
||||
if "--sidechat" in args:
|
||||
i = args.index("--sidechat")
|
||||
if i + 1 >= len(args):
|
||||
fail("BAD_NAME", "usage: notify <agent> <message> [--sidechat <name>] [--allow-main-chat]")
|
||||
sidechat = args[i + 1]
|
||||
del args[i:i + 2]
|
||||
if len(args) != 2:
|
||||
fail("BAD_NAME", "usage: notify <agent> <message> [--sidechat <name>] [--allow-main-chat]")
|
||||
act_notify(args[0], args[1], sidechat=sidechat, allow_main_chat=allow_main)
|
||||
elif action == "policy":
|
||||
if not rest:
|
||||
act_policy()
|
||||
elif rest[0] == "check":
|
||||
if len(rest) != 2:
|
||||
fail("BAD_NAME", "usage: policy check <agent>")
|
||||
act_policy_check(rest[1])
|
||||
elif rest[0] == "show":
|
||||
if len(rest) != 1:
|
||||
fail("BAD_NAME", "usage: policy show")
|
||||
act_policy_show()
|
||||
else:
|
||||
fail("BAD_NAME", "usage: policy [check <agent>|show]")
|
||||
elif action == "fleet-status":
|
||||
act_fleet_status()
|
||||
elif action == "dm-log":
|
||||
|
||||
@@ -17,14 +17,18 @@ after 3 attempts. `send --raw` transmits verbatim (for pre-signed
|
||||
messages): no tagging, no truncation — the signed payload must survive
|
||||
byte-identical.
|
||||
|
||||
Policy (2026-10-04): sidechat-first. Main-chat sends are refused unless
|
||||
--allow-main-chat is passed explicitly. Never saturate main threads by
|
||||
default; use a sidechat target instead.
|
||||
|
||||
Usage:
|
||||
dm.py send --agent opm --to 646 --target main "message"
|
||||
dm.py send --agent opm --target main --raw "$(dm-sign.sh --from operator-main 'hi')"
|
||||
dm.py verify-sig --agent opm --target main # scan recent reads for signed DMs
|
||||
dm.py send --agent opm --to 646 --target "646 tasks" "message"
|
||||
dm.py send --agent opm --target main --allow-main-chat --raw "$(dm-sign.sh --from operator-main 'hi')"
|
||||
dm.py verify-sig --agent opm --target "646 tasks" # scan recent reads for signed DMs
|
||||
dm.py verify-sig "$(dm-sign.sh --from operator-main 'hi')" # verify text directly
|
||||
dm.py read --agent 646 --target main [n]
|
||||
dm.py read --agent 646 --target "646 tasks" [n]
|
||||
dm.py log [--n 20]
|
||||
dm.py thread --from opm --to 646 --target main "message"
|
||||
dm.py thread --from opm --to 646 --target "646 tasks" "message"
|
||||
"""
|
||||
import argparse
|
||||
import os
|
||||
@@ -203,12 +207,16 @@ LOG_FILE = "/home/super/Projects/NetVM/dm-log.jsonl"
|
||||
SIDCHAT_ALIASES = {
|
||||
"646-opm-work": "d410b9ad-f667-465f-a103-43fabc0f69fe",
|
||||
"646 tasks": "1e75a740-d08f-443d-a0f9-793db196e24f",
|
||||
"heartbeat": "1e75a740-d08f-443d-a0f9-793db196e24f",
|
||||
"heartbeat": "0077e918-9ac8-40ba-b0ad-65f9f78037c4", # dedicated heartbeat sidechat (was colliding with 646 tasks)
|
||||
}
|
||||
|
||||
def resolve_sidechat_target(target):
|
||||
"""Resolve target alias or name to UUID dynamically from job-sidechats.json."""
|
||||
if not target or target == "main":
|
||||
if not target or not str(target).strip():
|
||||
raise ValueError(
|
||||
"DM target must not be empty: specify a sidechat name/UUID, "
|
||||
"or main (main requires --allow-main-chat)")
|
||||
if 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()
|
||||
@@ -229,6 +237,32 @@ def resolve_sidechat_target(target):
|
||||
return target
|
||||
|
||||
|
||||
def resolve_sidechat_source(target):
|
||||
"""Return where a target resolved from, for placement auditing.
|
||||
|
||||
Values: main / direct_uuid / static_alias / dynamic_mapping / passthrough.
|
||||
Mirrors resolve_sidechat_target() logic without changing its signature.
|
||||
"""
|
||||
if not target or not str(target).strip():
|
||||
return "invalid"
|
||||
if target == "main":
|
||||
return "main"
|
||||
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 "direct_uuid"
|
||||
if target in SIDCHAT_ALIASES:
|
||||
return "static_alias"
|
||||
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)
|
||||
if sc_data.get(target):
|
||||
return "dynamic_mapping"
|
||||
except Exception:
|
||||
pass
|
||||
return "passthrough"
|
||||
|
||||
|
||||
def log_event(event):
|
||||
"""Append to JSONL log."""
|
||||
event["ts"] = datetime.now(timezone.utc).isoformat()
|
||||
@@ -250,8 +284,81 @@ 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 verify_placement(recipient, msg_id, target, thread_uuid):
|
||||
"""Authoritative post-send placement verification (2026-10-04).
|
||||
|
||||
The send/verify loop only proves the message exists in whichever chat
|
||||
the browser was parked on at read time. A silent navigation drift (the
|
||||
SPA restoring main chat after a thread-URL navigation, a stale DOM read)
|
||||
used to produce verified:true for messages that actually landed in main.
|
||||
|
||||
This re-navigates to the target by direct URL, asserts the post-nav URL
|
||||
matches the target, then reads the chat back and asserts the message ID
|
||||
is present IN THAT CHAT. Returns (ok, detail). Read-only retries only --
|
||||
never resends: a misplaced send must fail loudly, not be duplicated.
|
||||
"""
|
||||
detail = {"target": target, "thread_uuid": thread_uuid}
|
||||
try:
|
||||
if target == "main":
|
||||
_rc, _out, _err = run_full(
|
||||
f"{NETVM_EXEC} {recipient} -- python3 {API} --account {recipient} sidechat main")
|
||||
if _rc != 0:
|
||||
detail.update({"reason": "nav_failed", "rc": _rc, "err": _err[:200]})
|
||||
return False, detail
|
||||
time.sleep(2)
|
||||
_u_rc, _u_out, _u_err = run_full(
|
||||
f"{NETVM_EXEC} {recipient} -- python3 {API} --account {recipient} url")
|
||||
detail["actual_url"] = _u_out[:200]
|
||||
if "/thread/" in _u_out:
|
||||
detail["reason"] = "url_mismatch"
|
||||
detail["expected"] = "main chat (no /thread/ in URL)"
|
||||
return False, detail
|
||||
else:
|
||||
expected_uuid = (thread_uuid or "").lower()
|
||||
detail["expected_uuid"] = expected_uuid or None
|
||||
if not expected_uuid:
|
||||
detail["reason"] = "no_thread_uuid"
|
||||
return False, detail
|
||||
# Direct-URL navigation: sidechat use <uuid> sets
|
||||
# window.location.href to /thread/<uuid> and confirms it.
|
||||
_rc, _out, _err = run_full(
|
||||
f"{NETVM_EXEC} {recipient} -- python3 {API} --account {recipient} sidechat use {expected_uuid}")
|
||||
if _rc != 0 or "NOTFOUND" in _out or expected_uuid not in _out.lower():
|
||||
detail.update({"reason": "nav_failed", "rc": _rc,
|
||||
"out": _out[:200], "err": _err[:200]})
|
||||
return False, detail
|
||||
time.sleep(2)
|
||||
# Independent URL assertion: the nav command's own confirm read
|
||||
# the same window.location.href, so sample it again here.
|
||||
_u_rc, _u_out, _u_err = run_full(
|
||||
f"{NETVM_EXEC} {recipient} -- python3 {API} --account {recipient} url")
|
||||
detail["actual_url"] = _u_out[:200]
|
||||
if expected_uuid not in _u_out.lower():
|
||||
detail["reason"] = "url_mismatch"
|
||||
return False, detail
|
||||
# Read-back with read-only retries: the SPA may still be rendering
|
||||
# the thread after navigation. A miss here must NOT trigger a resend.
|
||||
for _r in range(3):
|
||||
try:
|
||||
check_msgs = run(
|
||||
f"{NETVM_EXEC} {recipient} -- python3 {API} --account {recipient} messages 5 200")
|
||||
except Exception as e:
|
||||
detail["read_error"] = str(e)[:100]
|
||||
time.sleep(2)
|
||||
continue
|
||||
if msg_id in check_msgs:
|
||||
return True, detail
|
||||
time.sleep(2)
|
||||
detail["reason"] = "msg_not_in_target_chat"
|
||||
return False, detail
|
||||
except Exception as e:
|
||||
detail["reason"] = "exception"
|
||||
detail["error"] = str(e)[:200]
|
||||
return False, detail
|
||||
|
||||
def dm_send(agent, target, message, verify=True, raw=False,
|
||||
to_agent=None, tags=None, nudge_meta=None):
|
||||
to_agent=None, tags=None, nudge_meta=None,
|
||||
allow_main_chat=False):
|
||||
"""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>]."""
|
||||
@@ -275,9 +382,39 @@ def dm_send(agent, target, message, verify=True, raw=False,
|
||||
# Single unified attribution format (matches verify-sig's regex).
|
||||
tagged = f"[from:{agent}] [id:{msg_id}] {message}"
|
||||
|
||||
# Empty target is an error, not a silent main redirect (2026-10-04).
|
||||
# Only explicit --allow-main-chat reaches main.
|
||||
if not target or not str(target).strip():
|
||||
print(f"ERROR: DM target must not be empty (agent={agent} to={recipient}). "
|
||||
f"Specify a sidechat name/UUID, or main with --allow-main-chat.",
|
||||
file=sys.stderr)
|
||||
sys.exit(2)
|
||||
# Sidechat-first policy (2026-10-04): refuse main-chat sends unless the
|
||||
# caller explicitly opted in. Main threads must never be saturated by
|
||||
# default; use a sidechat target instead.
|
||||
if target == "main" and not allow_main_chat:
|
||||
log_event({"type": "main_chat_blocked", "id": msg_id, "agent": agent,
|
||||
"to": recipient, "target": target})
|
||||
print(f"ERROR: Refusing main-chat send by policy (agent={agent} to={recipient}). "
|
||||
f"Use a sidechat target, or pass --allow-main-chat for explicit main-chat sends.",
|
||||
file=sys.stderr)
|
||||
sys.exit(2)
|
||||
|
||||
# Canonical follow-up tags: trailing [bracket] tokens are metadata,
|
||||
# stripped from the delivered text, recorded on the log events.
|
||||
tags = tags or {}
|
||||
# Sidechat-first policy audit marker (2026-10-04): record explicit
|
||||
# main-chat opt-in so the main-chat watchdog can distinguish
|
||||
# authorized main sends from policy violations.
|
||||
if allow_main_chat and target == "main":
|
||||
tags["allow_main_chat"] = True
|
||||
# Extract job_id from [JOB <id>] marker for followup correlation.
|
||||
# This lets the harvester resolve followups by job_id when a
|
||||
# [RESULT <job_id>] reply arrives, even if thread_uuid is null.
|
||||
if "job_id" not in tags:
|
||||
_jm = re.search(r'\[JOB\s+([A-Za-z0-9_-]+)\]', message)
|
||||
if _jm:
|
||||
tags["job_id"] = _jm.group(1)
|
||||
if not raw:
|
||||
message, text_tags = extract_trailing_tags(message)
|
||||
if text_tags:
|
||||
@@ -301,7 +438,8 @@ def dm_send(agent, target, message, verify=True, raw=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})
|
||||
log_event({"type": "alias_resolved", "id": msg_id, "target": target, "thread_uuid": nav_target,
|
||||
"source": resolve_sidechat_source(target)})
|
||||
if target == "main":
|
||||
_rc, _out, _err = run_full(f"{NETVM_EXEC} {recipient} -- python3 {API} --account {recipient} sidechat main")
|
||||
else:
|
||||
@@ -333,7 +471,8 @@ def dm_send(agent, target, message, verify=True, raw=False,
|
||||
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"})
|
||||
"target": target, "status": "sidechat_created_pending_uuid",
|
||||
"browser_url": thread_url})
|
||||
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:
|
||||
@@ -358,7 +497,8 @@ def dm_send(agent, target, message, verify=True, raw=False,
|
||||
file=sys.stderr)
|
||||
sys.exit(1)
|
||||
log_event({"type": "nav_ok", "id": msg_id, "agent": agent, "to": recipient,
|
||||
"target": target, "thread_uuid": thread_uuid})
|
||||
"target": target, "thread_uuid": thread_uuid,
|
||||
"browser_url": thread_url})
|
||||
|
||||
time.sleep(2)
|
||||
|
||||
@@ -372,6 +512,7 @@ def dm_send(agent, target, message, verify=True, raw=False,
|
||||
# so we verify by reading the recipient's chat for our message ID.
|
||||
max_retries = 3
|
||||
delivered = False
|
||||
verified_attempt = None
|
||||
for attempt in range(max_retries):
|
||||
send_out = run(f'{NETVM_EXEC} {recipient} -- python3 {API} --account {recipient} send "{safe}"',
|
||||
priority="high")
|
||||
@@ -408,8 +549,28 @@ def dm_send(agent, target, message, verify=True, raw=False,
|
||||
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,
|
||||
"thread_uuid": thread_uuid, "attempt": attempt + 1})
|
||||
verified_attempt = attempt + 1
|
||||
# Capture the actual browser URL at verify time. `messages`
|
||||
# reads whichever chat the browser is currently parked on; if
|
||||
# that isn't the thread we navigated to, the send went to the
|
||||
# wrong chat (placement-blindness, 2026-10-04). This is only
|
||||
# an early signal -- the authoritative placement check runs
|
||||
# after this loop, and verified:true is only logged there.
|
||||
_v_rc, _v_out, _v_err = run_full(f'{NETVM_EXEC} {recipient} -- python3 {API} --account {recipient} url')
|
||||
_v_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})", _v_out, re.I)
|
||||
checked_uuid = _v_m.group(1).lower() if _v_m else None
|
||||
if target == "main":
|
||||
if checked_uuid:
|
||||
# Browser was parked in a sidechat thread while
|
||||
# sending to main: the message likely went to the
|
||||
# sidechat, not main.
|
||||
log_event({"type": "placement_mismatch", "id": msg_id, "agent": agent, "to": recipient,
|
||||
"target": target, "expected": "main", "actual_uuid": checked_uuid,
|
||||
"attempt": attempt + 1})
|
||||
elif checked_uuid and thread_uuid and checked_uuid != thread_uuid:
|
||||
log_event({"type": "placement_mismatch", "id": msg_id, "agent": agent, "to": recipient,
|
||||
"target": target, "expected_uuid": thread_uuid, "actual_uuid": checked_uuid,
|
||||
"attempt": attempt + 1})
|
||||
break
|
||||
else:
|
||||
log_event({"type": "retry", "id": msg_id, "agent": agent, "to": recipient, "attempt": attempt + 1})
|
||||
@@ -458,6 +619,29 @@ def dm_send(agent, target, message, verify=True, raw=False,
|
||||
if delivered and thread_uuid and "thread" not in tags:
|
||||
tags["thread"] = thread_uuid
|
||||
|
||||
if delivered:
|
||||
# Authoritative placement verification (2026-10-04): the loop above
|
||||
# only proved the message exists in whichever chat the browser was
|
||||
# parked on. Re-navigate to the target by direct URL, assert the
|
||||
# post-nav URL, and read the message back IN THAT CHAT before
|
||||
# logging verified:true. A failure here means the message landed in
|
||||
# the wrong chat: fail loudly (do NOT resend -- that would
|
||||
# duplicate the misplaced message).
|
||||
placed, pdetail = verify_placement(recipient, msg_id, target, thread_uuid)
|
||||
if not placed:
|
||||
log_event({"type": "placement_failed", "id": msg_id, "agent": agent,
|
||||
"to": recipient, "loop_attempt": verified_attempt, **pdetail})
|
||||
log_event({"type": "failed", "id": msg_id, "agent": agent, "to": recipient,
|
||||
"target": target, "verified": False, "reason": "placement_failed"})
|
||||
print(f"DM {msg_id} from {agent} to {recipient}/{target}: FAILED "
|
||||
f"(placement check failed: {pdetail.get('reason')}; "
|
||||
f"expected={pdetail.get('expected_uuid') or pdetail.get('expected')}; "
|
||||
f"actual_url={pdetail.get('actual_url')})", file=sys.stderr)
|
||||
sys.exit(1)
|
||||
log_event({"type": "verified", "id": msg_id, "agent": agent, "to": recipient,
|
||||
"target": target, "thread_uuid": thread_uuid,
|
||||
"placement": "confirmed", "attempt": verified_attempt})
|
||||
|
||||
# Note: Leave recipient browser parked in the sidechat thread to preserve Main Chat DOM
|
||||
log_event({"type": "send_done", "id": msg_id, "agent": agent, "to": recipient, "target": target})
|
||||
|
||||
@@ -646,20 +830,24 @@ def dm_log(n=20):
|
||||
e = json.loads(line)
|
||||
print(f"{e['ts'][:19]} {e.get('type','?'):12} {e.get('id','-'):8} {e.get('agent','-')}/{e.get('target','-')}")
|
||||
|
||||
def dm_thread(from_agent, to_agent, target, message):
|
||||
def dm_thread(from_agent, to_agent, target, message, allow_main_chat=False):
|
||||
"""Thread from one agent to another. Attribution is applied exactly once,
|
||||
in the unified [from:X] [id:Y] format, by dm_send."""
|
||||
return dm_send(from_agent, target, message, to_agent=to_agent)
|
||||
return dm_send(from_agent, target, message, to_agent=to_agent,
|
||||
allow_main_chat=allow_main_chat)
|
||||
|
||||
def main():
|
||||
p = argparse.ArgumentParser(description="DM: Headless Direct Messages (tagged [from:X] [id:Y], logged; delivery confirmed by recipient read-back)")
|
||||
p = argparse.ArgumentParser(description="DM: Headless Direct Messages (tagged [from:X] [id:Y], logged; delivery confirmed by recipient read-back). Sidechat-first policy: --target main requires --allow-main-chat.")
|
||||
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_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('--target', required=True,
|
||||
help='Conversation: sidechat name or thread UUID. Sidechat-first policy: --target main requires --allow-main-chat.')
|
||||
ps.add_argument('--allow-main-chat', action='store_true',
|
||||
help='Explicit opt-in for main-chat sends (refused by default per sidechat-first policy)')
|
||||
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)')
|
||||
@@ -698,7 +886,8 @@ def main():
|
||||
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)
|
||||
nudge_meta=nudge_meta,
|
||||
allow_main_chat=a.allow_main_chat)
|
||||
ps.set_defaults(func=_send_with_tags)
|
||||
|
||||
psel = sub.add_parser('select', help='Select active conversation')
|
||||
@@ -726,9 +915,13 @@ def main():
|
||||
pt = sub.add_parser('thread', help='Thread from one agent to another')
|
||||
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('--target', required=True,
|
||||
help='Conversation: sidechat name or thread UUID. Sidechat-first policy: main requires --allow-main-chat.')
|
||||
pt.add_argument('--allow-main-chat', action='store_true',
|
||||
help='Explicit opt-in for main-chat sends (refused by default per sidechat-first policy)')
|
||||
pt.add_argument('message')
|
||||
pt.set_defaults(func=lambda a: dm_thread(a.from_agent, a.to_agent, a.target, a.message))
|
||||
pt.set_defaults(func=lambda a: dm_thread(a.from_agent, a.to_agent, a.target, a.message,
|
||||
allow_main_chat=a.allow_main_chat))
|
||||
|
||||
args = p.parse_args()
|
||||
args.func(args)
|
||||
@@ -771,7 +964,7 @@ def _register_followup(msg_id, agent, recipient, target, tags):
|
||||
if tags.get("reply:expected"):
|
||||
tags_obj["reply:expected"] = True
|
||||
for key in ("reply:timeout", "reply:nudges", "reply:escalate",
|
||||
"route", "thread"):
|
||||
"route", "thread", "job_id"):
|
||||
if key in tags:
|
||||
tags_obj[key] = tags[key]
|
||||
|
||||
@@ -824,6 +1017,7 @@ def _register_followup(msg_id, agent, recipient, target, tags):
|
||||
"recipient": recipient,
|
||||
"target": target,
|
||||
"thread_uuid": tags.get("thread"),
|
||||
"job_id": tags.get("job_id"),
|
||||
"route": tags.get("route"),
|
||||
"sent_at": now_dt.isoformat(),
|
||||
"deadline": deadline_dt.isoformat(),
|
||||
|
||||
+70
-1
@@ -82,7 +82,9 @@ def append_job_log(entry):
|
||||
|
||||
|
||||
def send_dm(sender, recipient, target, text):
|
||||
"""Dispatch a DM via dm.py."""
|
||||
"""Dispatch a DM via dm.py. The sweeper's delivery target is explicit
|
||||
followup state (or explicit in-code escalation routing), so pass the
|
||||
main-chat opt-in when the target is main (sidechat-first policy)."""
|
||||
cmd = [
|
||||
sys.executable,
|
||||
str(DM_PY),
|
||||
@@ -90,6 +92,7 @@ def send_dm(sender, recipient, target, text):
|
||||
"--agent", sender,
|
||||
"--to", recipient,
|
||||
"--target", target,
|
||||
] + (["--allow-main-chat"] if target == "main" else []) + [
|
||||
text,
|
||||
]
|
||||
try:
|
||||
@@ -125,6 +128,34 @@ def sweep_cycle(dry_run=False):
|
||||
orig_target = rec.get("target", "main")
|
||||
thread_uuid = rec.get("thread_uuid")
|
||||
|
||||
# Ghost-followup fail-fast (2026-10-04, operator-main): a null
|
||||
# thread_uuid means the sidechat was never provisioned, so the
|
||||
# harvester can never match a reply. Nudging is pointless -- flag
|
||||
# for manual triage once instead of burning the nudge budget and
|
||||
# escalating a ghost. Main-chat followups are unaffected (the
|
||||
# harvester matches those by target).
|
||||
if thread_uuid is None and orig_target != "main" and not rec.get("needs_review"):
|
||||
rec["needs_review"] = True
|
||||
rec["status"] = "needs_review"
|
||||
rec["review_reason"] = (
|
||||
"ghost: unresolvable, manual triage "
|
||||
f"(thread_uuid null, target={orig_target}; "
|
||||
"reply can never auto-resolve)"
|
||||
)
|
||||
modified = True
|
||||
append_job_log({
|
||||
"ts": utcnow_str(),
|
||||
"type": "followup_ghost_suppressed",
|
||||
"dm_id": dm_id,
|
||||
"recipient": recipient,
|
||||
"target": orig_target,
|
||||
"reason": "ghost: unresolvable, manual triage",
|
||||
})
|
||||
print(f"Sweeper: suppressing ghost followup {dm_id} "
|
||||
f"(thread_uuid null, target={orig_target}) -> needs_review",
|
||||
file=sys.stderr)
|
||||
continue
|
||||
|
||||
if nudges_sent < nudges_allowed:
|
||||
# Deliver next nudge
|
||||
nudge_num = nudges_sent + 1
|
||||
@@ -161,6 +192,44 @@ def sweep_cycle(dry_run=False):
|
||||
})
|
||||
else:
|
||||
print(f"Sweeper WARNING: nudge send failed: {out}", file=sys.stderr)
|
||||
# Failure-mode fix (2026-10-04): advance state on send
|
||||
# failure so a failing nudge is never re-fired every
|
||||
# timer tick. failed_sends is tracked separately from
|
||||
# nudges_sent so a delivery failure does not consume a
|
||||
# real nudge. Backoff: 5 min base, doubling per failure.
|
||||
failed = rec.get("failed_sends", 0) + 1
|
||||
rec["failed_sends"] = failed
|
||||
rec["last_failure_at"] = utcnow_str()
|
||||
backoff_s = 300 * (2 ** min(failed - 1, 4)) # 5m,10m,20m,40m,80m cap
|
||||
rec["deadline"] = (now + timedelta(seconds=backoff_s)).isoformat()
|
||||
modified = True
|
||||
append_job_log({
|
||||
"ts": utcnow_str(),
|
||||
"type": "followup_nudge_failed",
|
||||
"dm_id": dm_id,
|
||||
"nudge_num": nudge_num,
|
||||
"recipient": recipient,
|
||||
"target": delivery_target,
|
||||
"failed_sends": failed,
|
||||
"backoff_s": backoff_s,
|
||||
"error": str(out)[:200],
|
||||
})
|
||||
# After 3 consecutive failures, stop retrying blindly and
|
||||
# flag for manual review (delivery may be uncertain or the
|
||||
# target may be permanently broken).
|
||||
if failed >= 3:
|
||||
rec["needs_review"] = True
|
||||
rec["review_reason"] = (
|
||||
f"nudge send failed {failed} times consecutively "
|
||||
f"(last: {str(out)[:120]})"
|
||||
)
|
||||
append_job_log({
|
||||
"ts": utcnow_str(),
|
||||
"type": "followup_needs_review",
|
||||
"dm_id": dm_id,
|
||||
"recipient": recipient,
|
||||
"failed_sends": failed,
|
||||
})
|
||||
else:
|
||||
nudges_count += 1
|
||||
|
||||
|
||||
+27
-3
@@ -9,6 +9,10 @@ Reads /home/super/Projects/NetVM/jobs/<job_name>.json,
|
||||
renders the prompt template, sends DM via dm.py, logs to job-log.jsonl.
|
||||
|
||||
Part of the JOB system (see docs/JOB-SPEC.md).
|
||||
|
||||
Sidechat-first policy (2026-10-04): a job that resolves to target "main"
|
||||
without explicit opt-in fails loudly instead of saturating main threads.
|
||||
Opt in via job JSON "allow_main_chat": true, or --allow-main-chat.
|
||||
"""
|
||||
|
||||
import sys
|
||||
@@ -202,9 +206,11 @@ def build_followup_tags(followup):
|
||||
return args
|
||||
|
||||
|
||||
def send_dm(agent, target, message, dry_run=False, followup_tags=None):
|
||||
def send_dm(agent, target, message, dry_run=False, followup_tags=None,
|
||||
allow_main_chat=False):
|
||||
"""Send DM via dm.py. followup_tags: flat ['--tag', 'k=v', ...] list
|
||||
from build_followup_tags(), or None."""
|
||||
from build_followup_tags(), or None. allow_main_chat passes the explicit
|
||||
main-chat opt-in through to dm.py (sidechat-first policy)."""
|
||||
if dry_run:
|
||||
print(f"[DRY RUN] Would send to {agent} ({target}):")
|
||||
if followup_tags:
|
||||
@@ -218,6 +224,7 @@ def send_dm(agent, target, message, dry_run=False, followup_tags=None):
|
||||
|
||||
cmd = ([str(DM_PY), "send", "--agent", "opm", "--to", agent,
|
||||
"--target", target]
|
||||
+ (["--allow-main-chat"] if (target == "main" and allow_main_chat) else [])
|
||||
+ (followup_tags or []) + [message])
|
||||
result = subprocess.run(cmd, capture_output=True, text=True, timeout=120)
|
||||
|
||||
@@ -286,6 +293,8 @@ def main():
|
||||
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")
|
||||
p.add_argument("--allow-main-chat", action="store_true",
|
||||
help="Explicit opt-in: allow this job to dispatch to main chat (refused by default per sidechat-first policy)")
|
||||
args = p.parse_args()
|
||||
|
||||
job_name = args.job_name
|
||||
@@ -370,6 +379,20 @@ def main():
|
||||
sc_name = render_prompt(sc_cfg.get("name_template", "job-{job_name}-{date}"), variables)
|
||||
target = sc_name
|
||||
|
||||
# Sidechat-first policy (2026-10-04): refuse to dispatch to main chat
|
||||
# unless the job explicitly opts in. Never fall back to main silently.
|
||||
allow_main = bool(job.get("allow_main_chat")) or args.allow_main_chat
|
||||
if target == "main" and not allow_main:
|
||||
print(f"ERROR: Job '{job_name}' resolves to main chat; refusing by sidechat-first policy. "
|
||||
f"Set a sidechat target (dm_target/target/sidechat.create) in the job JSON, "
|
||||
f"set \"allow_main_chat\": true, or pass --allow-main-chat.", file=sys.stderr)
|
||||
log_event("job_failed", {
|
||||
"job_id": job_id,
|
||||
"error": "main_chat_blocked_by_policy",
|
||||
"pipeline_run_id": pipeline_run_id,
|
||||
})
|
||||
sys.exit(1)
|
||||
|
||||
# Log job_sent
|
||||
log_event("job_sent", {
|
||||
"job_id": job_id,
|
||||
@@ -383,7 +406,8 @@ def main():
|
||||
|
||||
# Dispatch via dm.py (handles main or sidechat with auto-provisioning and verification)
|
||||
msg_id = send_dm(agent, target, dm_message,
|
||||
dry_run=dry_run, followup_tags=followup_tags)
|
||||
dry_run=dry_run, followup_tags=followup_tags,
|
||||
allow_main_chat=allow_main)
|
||||
|
||||
if msg_id and not dry_run:
|
||||
print(f"Dispatched job {job_id} to {agent}/{target} (DM: {msg_id})")
|
||||
|
||||
Executable
+194
@@ -0,0 +1,194 @@
|
||||
#!/usr/bin/env python3
|
||||
"""main-chat-watchdog.py — enforce the sidechat-first DM policy.
|
||||
|
||||
Tails dm-log.jsonl (read-only) and flags any DM actually SENT to main chat
|
||||
without an explicit allow_main_chat marker, plus any placement_mismatch /
|
||||
placement_failed events (message landed in the wrong chat). Distinguishes:
|
||||
|
||||
VIOLATION : type=sent, target=main, no tags.allow_main_chat -> policy broken (P0)
|
||||
AUTHORIZED: type=sent, target=main, tags.allow_main_chat=true -> explicit opt-in
|
||||
BLOCKED : type=main_chat_blocked -> policy working
|
||||
MISMATCH : type=placement_mismatch -> placed in wrong chat (P0)
|
||||
PFAILED : type=placement_failed -> placement hard-failed (P0)
|
||||
|
||||
State: byte-offset watermark in main-chat-watchdog.state (handles log
|
||||
rotation via inode check). First run starts at EOF (no historical backfill —
|
||||
old entries predate the allow_main_chat audit marker).
|
||||
|
||||
Violations are appended to logs/main-chat-violations.jsonl and printed to
|
||||
stdout (journal). Exit 0 = clean, 1 = violations/P0s found, 2 = error.
|
||||
|
||||
Read-only: never modifies dm-log.jsonl.
|
||||
"""
|
||||
|
||||
import json
|
||||
import os
|
||||
import sys
|
||||
from datetime import datetime, timezone
|
||||
|
||||
BASE = "/home/super/Projects/NetVM"
|
||||
DM_LOG = os.path.join(BASE, "dm-log.jsonl")
|
||||
STATE_FILE = os.path.join(BASE, "main-chat-watchdog.state")
|
||||
VIOLATIONS_LOG = os.path.join(BASE, "logs", "main-chat-violations.jsonl")
|
||||
|
||||
|
||||
def utcnow():
|
||||
return datetime.now(timezone.utc).isoformat()
|
||||
|
||||
|
||||
def load_state():
|
||||
try:
|
||||
with open(STATE_FILE) as f:
|
||||
return json.load(f)
|
||||
except (FileNotFoundError, json.JSONDecodeError, ValueError):
|
||||
return {}
|
||||
|
||||
|
||||
def save_state(state):
|
||||
tmp = STATE_FILE + ".tmp"
|
||||
with open(tmp, "w") as f:
|
||||
json.dump(state, f)
|
||||
os.replace(tmp, STATE_FILE)
|
||||
|
||||
|
||||
def main():
|
||||
try:
|
||||
st = os.stat(DM_LOG)
|
||||
except FileNotFoundError:
|
||||
print(f"watchdog ERROR: {DM_LOG} not found", file=sys.stderr)
|
||||
return 2
|
||||
|
||||
state = load_state()
|
||||
# First run (or rotation): start at EOF, don't backfill history that
|
||||
# predates the allow_main_chat audit marker.
|
||||
if state.get("inode") != st.st_ino:
|
||||
offset = st.st_size
|
||||
if state:
|
||||
print(f"watchdog: log rotated or first run (inode {st.st_ino}), "
|
||||
f"starting at EOF offset {offset}")
|
||||
else:
|
||||
offset = min(state.get("offset", 0), st.st_size)
|
||||
|
||||
violations = [] # P0: gate bypass (sent to main, no opt-in)
|
||||
mismatches = [] # P0: placement_mismatch / placement_failed
|
||||
blocked = 0
|
||||
authorized = 0
|
||||
scanned = 0
|
||||
malformed = 0
|
||||
|
||||
try:
|
||||
with open(DM_LOG, "r", encoding="utf-8", errors="replace") as f:
|
||||
f.seek(offset)
|
||||
for line in f:
|
||||
line = line.strip()
|
||||
if not line:
|
||||
continue
|
||||
scanned += 1
|
||||
try:
|
||||
ev = json.loads(line)
|
||||
except json.JSONDecodeError:
|
||||
malformed += 1
|
||||
continue
|
||||
etype = ev.get("type")
|
||||
if etype == "main_chat_blocked":
|
||||
# Policy working: the dm.py gate refused a main send.
|
||||
blocked += 1
|
||||
elif etype == "placement_mismatch":
|
||||
# P0: verify-time URL check found the browser parked in a
|
||||
# different chat than the intended target. The send
|
||||
# completed but landed in the wrong place. Schema has two
|
||||
# variants: sidechat ("expected_uuid") and main-drift
|
||||
# ("expected": "main").
|
||||
mismatches.append({
|
||||
"ts": utcnow(),
|
||||
"kind": "placement_mismatch",
|
||||
"severity": "P0",
|
||||
"dm_id": ev.get("id"),
|
||||
"agent": ev.get("agent"),
|
||||
"to": ev.get("to"),
|
||||
"target": ev.get("target"),
|
||||
"expected_uuid": ev.get("expected_uuid") or ev.get("expected"),
|
||||
"actual_uuid": ev.get("actual_uuid"),
|
||||
"attempt": ev.get("attempt"),
|
||||
"event_ts": ev.get("ts"),
|
||||
})
|
||||
elif etype == "placement_failed":
|
||||
# P0: the authoritative post-send placement check failed
|
||||
# hard (sender exited 1, send NOT marked verified).
|
||||
d = ev.get("detail") or {}
|
||||
mismatches.append({
|
||||
"ts": utcnow(),
|
||||
"kind": "placement_failed",
|
||||
"severity": "P0",
|
||||
"dm_id": ev.get("id"),
|
||||
"agent": ev.get("agent"),
|
||||
"to": ev.get("to"),
|
||||
"target": ev.get("target"),
|
||||
"reason": ev.get("reason") or d.get("reason"),
|
||||
"expected_uuid": ev.get("expected_uuid") or d.get("expected_uuid"),
|
||||
"actual_url": ev.get("actual_url") or d.get("actual_url"),
|
||||
"loop_attempt": ev.get("loop_attempt"),
|
||||
"event_ts": ev.get("ts"),
|
||||
})
|
||||
elif etype == "sent" and ev.get("target") == "main":
|
||||
tags = ev.get("tags") or {}
|
||||
if tags.get("allow_main_chat"):
|
||||
authorized += 1
|
||||
else:
|
||||
violations.append({
|
||||
"ts": utcnow(),
|
||||
"kind": "main_chat_send",
|
||||
"severity": "P0",
|
||||
"dm_id": ev.get("id"),
|
||||
"agent": ev.get("agent"),
|
||||
"to": ev.get("to"),
|
||||
"target": "main",
|
||||
"msg_preview": (ev.get("msg") or "")[:120],
|
||||
"event_ts": ev.get("ts"),
|
||||
})
|
||||
new_offset = f.tell()
|
||||
except OSError as e:
|
||||
print(f"watchdog ERROR reading log: {e}", file=sys.stderr)
|
||||
return 2
|
||||
|
||||
save_state({"offset": new_offset, "inode": st.st_ino})
|
||||
|
||||
findings = violations + mismatches
|
||||
summary = (f"watchdog: scanned={scanned} blocked={blocked} "
|
||||
f"authorized_main={authorized} "
|
||||
f"violations={len(violations)} "
|
||||
f"placement_mismatches={len(mismatches)} "
|
||||
f"malformed={malformed}")
|
||||
print(summary)
|
||||
|
||||
if findings:
|
||||
try:
|
||||
os.makedirs(os.path.dirname(VIOLATIONS_LOG), exist_ok=True)
|
||||
with open(VIOLATIONS_LOG, "a", encoding="utf-8") as vf:
|
||||
for v in findings:
|
||||
vf.write(json.dumps(v) + "\n")
|
||||
except OSError as e:
|
||||
print(f"watchdog ERROR writing violations log: {e}", file=sys.stderr)
|
||||
return 2
|
||||
for v in violations:
|
||||
print(f"VIOLATION main-chat send without opt-in: "
|
||||
f"id={v['dm_id']} agent={v['agent']} to={v['to']} "
|
||||
f"at={v['event_ts']} preview={v['msg_preview']!r}")
|
||||
for m in mismatches:
|
||||
if m["kind"] == "placement_mismatch":
|
||||
print(f"P0 PLACEMENT_MISMATCH message landed in wrong chat: "
|
||||
f"id={m['dm_id']} agent={m['agent']} to={m['to']} "
|
||||
f"target={m['target']} expected={m['expected_uuid']} "
|
||||
f"actual={m['actual_uuid']} at={m['event_ts']}")
|
||||
else:
|
||||
print(f"P0 PLACEMENT_FAILED placement check failed: "
|
||||
f"id={m['dm_id']} agent={m['agent']} to={m['to']} "
|
||||
f"target={m['target']} reason={m['reason']} "
|
||||
f"expected={m['expected_uuid']} "
|
||||
f"actual_url={m['actual_url']} at={m['event_ts']}")
|
||||
return 1
|
||||
return 0
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
sys.exit(main())
|
||||
@@ -119,6 +119,18 @@ def append_jsonl(path, record):
|
||||
f.write(json.dumps(record) + "\n")
|
||||
|
||||
|
||||
# Case-insensitive failure/decline detection for [RESULT] text.
|
||||
# "declined"/"decline"/"reject" count as failures so declines do NOT
|
||||
# trigger onward pipeline chaining (previously only uppercase
|
||||
# FAILED/UNABLE/FAIL matched, so a lowercase "declined" was logged as success).
|
||||
FAIL_PREFIXES = ("failed", "unable", "fail", "declined", "decline", "reject", "error")
|
||||
|
||||
|
||||
def is_fail_result(result_text):
|
||||
t = (result_text or "").lstrip().lower()
|
||||
return t.startswith(FAIL_PREFIXES)
|
||||
|
||||
|
||||
def get_monitored_threads(target_agent=None):
|
||||
"""
|
||||
Build dict of threads to monitor per agent:
|
||||
@@ -342,10 +354,12 @@ def harvest_main_feed(cdp, agent, watermarks, followups, dry_run=False):
|
||||
|
||||
if author == "assistant":
|
||||
m_res = re.search(r"\[RESULT\s+([A-Za-z0-9_-]+)\]\s*(.*)", text, re.S)
|
||||
result_job_id = None
|
||||
if m_res:
|
||||
job_id = m_res.group(1).strip()
|
||||
result_job_id = job_id
|
||||
result_text = m_res.group(2).strip()
|
||||
is_fail = result_text.startswith("FAILED") or result_text.startswith("UNABLE") or result_text.startswith("FAIL")
|
||||
is_fail = is_fail_result(result_text)
|
||||
job_results += 1
|
||||
job_record = {
|
||||
"ts": utcnow(),
|
||||
@@ -361,7 +375,8 @@ def harvest_main_feed(cdp, agent, watermarks, followups, dry_run=False):
|
||||
append_jsonl(JOB_LOG, job_record)
|
||||
trigger_chain_next(job_id, result_text, success=not is_fail)
|
||||
|
||||
clear_matching_followups(followups, agent, "main", mid, text, dry_run)
|
||||
clear_matching_followups(followups, agent, "main", mid, text, dry_run,
|
||||
job_id=result_job_id)
|
||||
|
||||
return new_messages, new_wm, job_results
|
||||
|
||||
@@ -462,10 +477,12 @@ def harvest_agent_thread(cdp, agent, thread_info, watermarks, followups, dry_run
|
||||
# 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)
|
||||
result_job_id = None
|
||||
if m_res:
|
||||
job_id = m_res.group(1).strip()
|
||||
result_job_id = job_id
|
||||
result_text = m_res.group(2).strip()
|
||||
is_fail = result_text.startswith("FAILED") or result_text.startswith("UNABLE") or result_text.startswith("FAIL")
|
||||
is_fail = is_fail_result(result_text)
|
||||
job_results += 1
|
||||
|
||||
job_record = {
|
||||
@@ -484,13 +501,42 @@ def harvest_agent_thread(cdp, agent, thread_info, watermarks, followups, dry_run
|
||||
trigger_chain_next(job_id, result_text, success=not is_fail)
|
||||
|
||||
# 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,
|
||||
job_id=result_job_id)
|
||||
|
||||
return new_messages, new_wm, job_results
|
||||
|
||||
|
||||
# Chain deduplication
|
||||
CHAINED_JOBS_FILE = Path(__file__).parent / "chained-jobs.json"
|
||||
|
||||
def _has_chained(job_id):
|
||||
try:
|
||||
if CHAINED_JOBS_FILE.exists():
|
||||
import json as _j
|
||||
with open(CHAINED_JOBS_FILE) as f:
|
||||
return job_id in _j.load(f)
|
||||
except: pass
|
||||
return False
|
||||
|
||||
def _mark_chained(job_id):
|
||||
try:
|
||||
import json as _j
|
||||
c = []
|
||||
if CHAINED_JOBS_FILE.exists():
|
||||
with open(CHAINED_JOBS_FILE) as f: c = _j.load(f)
|
||||
if job_id not in c:
|
||||
c.append(job_id)
|
||||
c = c[-1000:]
|
||||
with open(CHAINED_JOBS_FILE, "w") as f: _j.dump(c, f)
|
||||
except: pass
|
||||
|
||||
def trigger_chain_next(job_id, result_text, success=True):
|
||||
"""If the completed job has on_success, on_failure, or chain_next, dispatch downstream."""
|
||||
if _has_chained(job_id):
|
||||
return
|
||||
_mark_chained(job_id)
|
||||
|
||||
m = re.match(r"^(.*)-(\d{8}-\d{6}-[a-f0-9]{8})$", job_id)
|
||||
if m:
|
||||
job_name = m.group(1)
|
||||
@@ -557,8 +603,14 @@ def trigger_chain_next(job_id, result_text, success=True):
|
||||
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."""
|
||||
def clear_matching_followups(followups, agent, thread_id, mid, text, dry_run=False,
|
||||
job_id=None):
|
||||
"""Resolve follow-up records if an assistant message is detected in the thread.
|
||||
|
||||
Matches on thread identity (thread_uuid or target='main') OR on job_id
|
||||
(from a [RESULT <job_id>] reply). The job_id path works regardless of
|
||||
thread_uuid or target, fixing ghost followups with null thread_uuid.
|
||||
"""
|
||||
if not followups:
|
||||
return
|
||||
|
||||
@@ -576,7 +628,13 @@ def clear_matching_followups(followups, agent, thread_id, mid, text, dry_run=Fal
|
||||
elif f_rec.get("thread_uuid") and f_rec.get("thread_uuid") == thread_id:
|
||||
match_thread = True
|
||||
|
||||
if match_thread:
|
||||
# Match by job_id (from [RESULT <job_id>]) -- works regardless of
|
||||
# thread_uuid or target. This is an ADDITIONAL path, not a replacement.
|
||||
match_job = False
|
||||
if job_id and f_rec.get("job_id") and f_rec.get("job_id") == job_id:
|
||||
match_job = True
|
||||
|
||||
if match_thread or match_job:
|
||||
f_rec["status"] = "resolved"
|
||||
f_rec["resolved_at"] = utcnow()
|
||||
f_rec["resolved_by_mid"] = mid
|
||||
|
||||
+66
-1
@@ -65,7 +65,7 @@ DEFAULT_AGENT_SIDECHATS = {
|
||||
"646": "646 tasks",
|
||||
"opm": "heartbeat",
|
||||
"pip": "646-pip-coord",
|
||||
"muse": "main",
|
||||
"muse": "muse tasks",
|
||||
}
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
@@ -2202,6 +2202,71 @@ def cmd_followup_cancel(args):
|
||||
os.replace(tmp_path, FOLLOWUPS_FILE)
|
||||
print(c_green(f"● Cancelled follow-up {match_key}."))
|
||||
|
||||
# Also stand down the VM-side nudge controller (request-sweeper.py),
|
||||
# which tracks dm_followup records independently of bl's local file.
|
||||
# Best-effort: the local cancel already succeeded, so a VM failure
|
||||
# (unreachable, no record, bad response) must not fail the command.
|
||||
_vm_followup_cancel(match_key)
|
||||
|
||||
|
||||
def _vm_followup_cancel(dm_id):
|
||||
"""POST /api/box/followups/cancel so the VM-side sweeper stops nudging.
|
||||
|
||||
Best-effort: never raises. Prints a one-line outcome.
|
||||
"""
|
||||
import json as _json
|
||||
import os as _os
|
||||
import subprocess as _sp
|
||||
import time as _time
|
||||
import urllib.parse as _up
|
||||
import urllib.request as _ur
|
||||
import urllib.error as _ue
|
||||
|
||||
box_api = _os.environ.get("BOX_API_BASE", "https://box.muse-dev.online")
|
||||
key = _os.environ.get("BOX_SIGN_KEY",
|
||||
_os.path.expanduser("~/.ssh/id_ed25519"))
|
||||
endpoint = "followups/cancel"
|
||||
ts_now = str(int(_time.time()))
|
||||
payload = ("%s\n%s" % (ts_now, endpoint)).encode()
|
||||
try:
|
||||
pr = _sp.run(["ssh-keygen", "-Y", "sign", "-f", key, "-n", "box"],
|
||||
input=payload, capture_output=True, timeout=15)
|
||||
sig = pr.stdout.decode().strip()
|
||||
if not sig:
|
||||
raise RuntimeError("empty signature")
|
||||
except Exception as e:
|
||||
print(c_warn(" VM follow-up cancel skipped (signing failed: %s)" % e))
|
||||
return
|
||||
query = _up.urlencode({"identity": "bl", "ts": ts_now, "sig": sig})
|
||||
url = "%s/api/box/followups/cancel?%s" % (box_api, query)
|
||||
try:
|
||||
req = _ur.Request(url,
|
||||
data=_json.dumps({"dm_id": dm_id}).encode(),
|
||||
headers={"Content-Type": "application/json",
|
||||
"User-Agent": "super-cli/1.0 (bl)"},
|
||||
method="POST")
|
||||
with _ur.urlopen(req, timeout=30) as r:
|
||||
resp = _json.load(r)
|
||||
except _ue.HTTPError as e:
|
||||
if e.code == 404:
|
||||
print(c_dim(" VM: no dm_followup record (bl-only follow-up)."))
|
||||
else:
|
||||
try:
|
||||
detail = e.read().decode("utf-8", errors="ignore")[:120]
|
||||
except Exception:
|
||||
detail = ""
|
||||
print(c_warn(" VM cancel failed: HTTP %s %s" % (e.code, detail)))
|
||||
return
|
||||
except Exception as e:
|
||||
print(c_warn(" VM cancel failed: %s" % e))
|
||||
return
|
||||
if resp.get("canceled"):
|
||||
print(c_green(" VM: follow-up canceled (%s)."
|
||||
% resp.get("request_id")))
|
||||
else:
|
||||
print(c_dim(" VM: %s." % resp.get("reason", "not canceled")))
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Domain: PIPELINE (Multi-Agent Workflow Pipelines)
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
Reference in New Issue
Block a user