fix(work): import hashlib and wire heal subparser into main CLI
This commit is contained in:
+167
-22
@@ -109,7 +109,18 @@ def iter_result_markers(text):
|
||||
"""Yield (job_id, result_text) for every [RESULT <job_id>] marker in text."""
|
||||
current_re = lookup_engine.get_result_regex() if HAS_LOOKUP_ENGINE else RESULT_RE
|
||||
for m in current_re.finditer(text or ""):
|
||||
yield m.group(1).strip(), m.group(2).strip()
|
||||
gd = m.groupdict()
|
||||
if "summary" in gd:
|
||||
# Engine shape: [RESULT <id>] [STATUS] <summary>. The
|
||||
# status word is optional (None for bare markers); keep
|
||||
# it when present so FAIL/ERROR still trips failure
|
||||
# detection downstream.
|
||||
status = (m.group("status") or "").strip()
|
||||
summary = (m.group("summary") or "").strip()
|
||||
result_text = f"{status} {summary}".strip() if status else summary
|
||||
yield m.group("job_id").strip(), result_text
|
||||
else:
|
||||
yield m.group(1).strip(), m.group(2).strip()
|
||||
|
||||
|
||||
# Verb markers for the digest response protocol:
|
||||
@@ -316,9 +327,19 @@ def result_has_evidence(result_text):
|
||||
return bool(_PROOF_EVIDENCE_RE.search(result_text or ""))
|
||||
|
||||
|
||||
# Automated in-thread proof requests disabled per fleet governance decision (2026-10-09)
|
||||
PROOF_REQUESTS_ENABLED = False
|
||||
|
||||
def maybe_request_proof(agent, thread_id, job_id, result_text, dry_run=False):
|
||||
"""Ask for checkable evidence when a success RESULT has none.
|
||||
|
||||
Disabled by default per fleet decision 2026-10-09: automated in-thread proof
|
||||
challenges trigger adversarial rejection loops and waste agent quota.
|
||||
"""
|
||||
if not PROOF_REQUESTS_ENABLED:
|
||||
return False
|
||||
"""Ask for checkable evidence when a success RESULT has none.
|
||||
|
||||
One-shot per (thread, job) via the nudge tracker. Returns True when a
|
||||
proof followup was scheduled.
|
||||
"""
|
||||
@@ -559,6 +580,42 @@ def format_tool_result_for_chat(op, raw_output):
|
||||
out = out[:900] + "\n…(truncated, refine the call for detail)"
|
||||
return f"box result:\n```\n{out}\n```"
|
||||
|
||||
if op == "flow.start" and isinstance(data, dict):
|
||||
if not data.get("ok"):
|
||||
return f"Flow start failed: {data.get('error')}"
|
||||
return f"Flow `{data.get('flow_id')}` started in pane `{data.get('session')}` (status: {data.get('status')})."
|
||||
|
||||
if op == "flow.read" and isinstance(data, dict):
|
||||
if not data.get("ok"):
|
||||
return f"Flow read failed: {data.get('error')}"
|
||||
st = data.get("status", "unknown")
|
||||
ec = data.get("exit_code")
|
||||
ec_str = f" (exit_code: {ec})" if ec is not None else ""
|
||||
pm = data.get("prompt_match")
|
||||
prompt_str = f"\nPrompt waiting: {pm.get('text', pm)}" if pm else ""
|
||||
delta = data.get("delta", "").strip()
|
||||
trunc = f" (last {data.get('lines_read')} lines)" if data.get("truncated") else ""
|
||||
body = f"\n```\n{delta}\n```" if delta else " (no new output)"
|
||||
return f"Flow `{data.get('flow_id')}` [{st}]{ec_str}{prompt_str}{trunc}:{body}"
|
||||
|
||||
if op == "flow.send" and isinstance(data, dict):
|
||||
if not data.get("ok"):
|
||||
return f"Flow send failed: {data.get('error')}"
|
||||
kind = "command" if data.get("is_command") else "keys"
|
||||
return f"Flow `{data.get('flow_id')}` sent {kind}: `{data.get('sent')}` (status: {data.get('status')})."
|
||||
|
||||
if op == "flow.list" and isinstance(data, dict):
|
||||
flows = data.get("flows", [])
|
||||
if not flows:
|
||||
return "No active flows."
|
||||
lines = [f"{len(flows)} flows:"]
|
||||
for f in flows[:8]:
|
||||
lines.append(f" • {f.get('flow_id')} [{f.get('status')}]: {f.get('session')} (cmd: {str(f.get('command', 'bash'))[:30]})")
|
||||
return "\n".join(lines)
|
||||
|
||||
if op == "flow.stop" and isinstance(data, dict):
|
||||
return f"Flow `{data.get('flow_id')}` stopped."
|
||||
|
||||
# General fallback: compact JSON capped to 400 chars
|
||||
s = json.dumps(data)
|
||||
return s[:400] + "..." if len(s) > 400 else s
|
||||
@@ -623,11 +680,59 @@ def is_fail_result(result_text):
|
||||
return t.startswith(FAIL_PREFIXES)
|
||||
|
||||
|
||||
RECENCY_WINDOW_SEC = 10800
|
||||
|
||||
_JOB_ID_RE = re.compile(r"^(.+)-(\d{8})-(\d{6})-([0-9a-f]{8})$")
|
||||
|
||||
|
||||
def dispatched_families_since(job_log_path, window_sec=RECENCY_WINDOW_SEC,
|
||||
now=None):
|
||||
"""Job families dispatched inside the window.
|
||||
|
||||
Scans job-log.jsonl for job_sent/job_dispatched events newer than
|
||||
``window_sec`` and returns their family names (the job id minus the
|
||||
trailing -YYYYMMDD-HHMMSS-<hash> run suffix). Missing, unreadable,
|
||||
or malformed input yields an empty set, never an exception.
|
||||
"""
|
||||
now = now or datetime.now(timezone.utc)
|
||||
cutoff = now.timestamp() - window_sec
|
||||
fams = set()
|
||||
try:
|
||||
handle = open(job_log_path, "r", encoding="utf-8")
|
||||
except OSError:
|
||||
return fams
|
||||
with handle:
|
||||
for line in handle:
|
||||
line = line.strip()
|
||||
if not line:
|
||||
continue
|
||||
try:
|
||||
event = json.loads(line)
|
||||
except Exception:
|
||||
continue
|
||||
if event.get("type") not in ("job_sent", "job_dispatched"):
|
||||
continue
|
||||
try:
|
||||
ts = datetime.fromisoformat(
|
||||
str(event.get("ts")).replace("Z", "+00:00")).timestamp()
|
||||
except Exception:
|
||||
continue
|
||||
if ts < cutoff:
|
||||
continue
|
||||
match = _JOB_ID_RE.match(str(event.get("job_id") or ""))
|
||||
if match:
|
||||
fams.add(match.group(1))
|
||||
return fams
|
||||
|
||||
|
||||
def get_monitored_threads(target_agent=None):
|
||||
"""
|
||||
Build dict of threads to monitor per agent:
|
||||
{ agent: [ {"id": "<uuid>", "name": "<alias>"} ] }
|
||||
Filters to permanent channels, threads with pending followups, or recent threads (< 3h).
|
||||
Filters to permanent channels, threads with pending followups,
|
||||
recently created threads (< 3h), or threads whose job family was
|
||||
dispatched recently (< 3h) so old persistent sidechats that still
|
||||
receive prompts stay monitored.
|
||||
"""
|
||||
agents = [target_agent] if target_agent else VALID_AGENTS
|
||||
threads_by_agent = {a: [] for a in agents}
|
||||
@@ -642,6 +747,7 @@ def get_monitored_threads(target_agent=None):
|
||||
|
||||
PERM_KEYWORDS = ("coord", "tasks", "task", "brain", "heartbeat", "sync", "audit", "main-loop")
|
||||
now = datetime.now(timezone.utc)
|
||||
recently_dispatched = dispatched_families_since(JOB_LOG, now=now)
|
||||
|
||||
state_files = [JOB_SIDECHATS_FILE, WAKE_SIDECHATS_FILE]
|
||||
for sf in state_files:
|
||||
@@ -684,7 +790,9 @@ def get_monitored_threads(target_agent=None):
|
||||
if isinstance(val, dict) and val.get("archived") and not is_pending:
|
||||
continue
|
||||
|
||||
if not (is_perm or is_pending or is_recent):
|
||||
is_dispatched = key in recently_dispatched
|
||||
|
||||
if not (is_perm or is_pending or is_recent or is_dispatched):
|
||||
continue
|
||||
|
||||
existing = [t["id"] for t in threads_by_agent[agent]]
|
||||
@@ -896,8 +1004,18 @@ def process_messages(raw_messages, agent, thread_id, thread_name, last_wm, follo
|
||||
append_jsonl(CHAT_HISTORY_LOG, record)
|
||||
|
||||
if author == "assistant":
|
||||
markers = list(iter_result_markers(text))
|
||||
verbs = list(iter_verb_markers(text))
|
||||
try:
|
||||
markers = list(iter_result_markers(text))
|
||||
verbs = list(iter_verb_markers(text))
|
||||
except Exception as e:
|
||||
# One poison message must not wedge the batch: without
|
||||
# this, the same crash repeats every cycle, the
|
||||
# watermark never advances past it, and the thread's
|
||||
# followups nag to escalation despite answered work.
|
||||
sys.stderr.write(
|
||||
"warning: marker extraction failed, treating as "
|
||||
f"plain reply: {e}\n")
|
||||
markers, verbs = [], []
|
||||
|
||||
# Synthesize [RESULT <job-id>] DECLINE if assistant explicitly refuses the task in plain text
|
||||
if not markers and not verbs and detect_explicit_refusal(text):
|
||||
@@ -946,19 +1064,27 @@ def process_messages(raw_messages, agent, thread_id, thread_name, last_wm, follo
|
||||
try:
|
||||
import muse_hybrid
|
||||
thread_url = f"https://box.muse-dev.online/thread/{thread_id}"
|
||||
tool_hint = (
|
||||
f"[Runtime Context: {thread_url}]\n"
|
||||
f"Tools: EMIT one [TOOL <op> <args>] line per action (you do not run it;"
|
||||
f" the runtime executes it and replies here). curl -sk -X POST"
|
||||
f" https://exec.muse-dev.online/exec works too.\n"
|
||||
f" • [TOOL tools.list {{}}] — discover every op dynamically\n"
|
||||
f" • [TOOL swarm.spawn {{\"count\": 1, \"task\": \"<task>\"}}] — spawn subagents\n"
|
||||
f" • [DM {{\"to\": \"<agent>\", \"target\": \"<sidechat>\", \"message\": \"<text>\"}}] — send a DM\n"
|
||||
f" • [TOOL box.exec {{\"action\": \"fleet-status\"}}] — call box (read-only actions)\n"
|
||||
f" • [TOOL followup.create {{\"in_m\": 5, \"prompt\": \"<reminder>\"}}]\n"
|
||||
f" • [TOOL health.check {{}}]\n\n"
|
||||
f"[Directive: Take next action or close with [RESULT <job_id>] <summary>]"
|
||||
)
|
||||
if op.startswith("flow."):
|
||||
flow_id = t_args.get("flow_id", "<flow_id>") if isinstance(t_args, dict) else "<flow_id>"
|
||||
tool_hint = (
|
||||
f"[Flow Directive: advance with [TOOL flow.send {{\"flow_id\": \"{flow_id}\", \"command\": \"...\"}}]"
|
||||
f" | read with [TOOL flow.read {{\"flow_id\": \"{flow_id}\"}}]"
|
||||
f" | close with [RESULT <job_id>] OK]"
|
||||
)
|
||||
else:
|
||||
tool_hint = (
|
||||
f"[Runtime Context: {thread_url}]\n"
|
||||
f"Tools: EMIT one [TOOL <op> <args>] line per action (you do not run it;"
|
||||
f" the runtime executes it and replies here). curl -sk -X POST"
|
||||
f" https://exec.muse-dev.online/exec works too.\n"
|
||||
f" • [TOOL tools.list {{}}] — discover every op dynamically\n"
|
||||
f" • [TOOL swarm.spawn {{\"count\": 1, \"task\": \"<task>\"}}] — spawn subagents\n"
|
||||
f" • [DM {{\"to\": \"<agent>\", \"target\": \"<sidechat>\", \"message\": \"<text>\"}}] — send a DM\n"
|
||||
f" • [TOOL box.exec {{\"action\": \"fleet-status\"}}] — call box (read-only actions)\n"
|
||||
f" • [TOOL followup.create {{\"in_m\": 5, \"prompt\": \"<reminder>\"}}]\n"
|
||||
f" • [TOOL health.check {{}}]\n\n"
|
||||
f"[Directive: Take next action or close with [RESULT <job_id>] <summary>]"
|
||||
)
|
||||
if t_ok:
|
||||
clean_msg = format_tool_result_for_chat(op, t_res)
|
||||
resp_text = f"Tool result (`{op}`):\n{clean_msg}\n\n{tool_hint}"
|
||||
@@ -969,7 +1095,12 @@ def process_messages(raw_messages, agent, thread_id, thread_name, last_wm, follo
|
||||
sys.stderr.write(f"warning: failed to post tool response back to thread: {te}\n")
|
||||
|
||||
if markers or verbs:
|
||||
seen_jobs = set()
|
||||
for job_id, result_text in markers:
|
||||
if job_id in seen_jobs:
|
||||
# Same verdict restated in one message: log once.
|
||||
continue
|
||||
seen_jobs.add(job_id)
|
||||
is_fail = is_fail_result(result_text)
|
||||
job_results += 1
|
||||
|
||||
@@ -983,6 +1114,11 @@ def process_messages(raw_messages, agent, thread_id, thread_name, last_wm, follo
|
||||
"thread_id": thread_id,
|
||||
"msg_id": mid,
|
||||
}
|
||||
if result_text.startswith("DECLINE:"):
|
||||
# Synthesized (or explicit) decline: still a
|
||||
# non-success (no chaining), but the auditor
|
||||
# buckets it as declined, not a failure.
|
||||
job_record["outcome"] = "declined"
|
||||
if not dry_run:
|
||||
append_jsonl(JOB_LOG, job_record)
|
||||
# Check if this is a swarm slot result: sw-YYYYMMDD-HHMMSS-xxxx/<slot>
|
||||
@@ -1000,10 +1136,12 @@ def process_messages(raw_messages, agent, thread_id, thread_name, last_wm, follo
|
||||
trigger_chain_next(job_id, result_text, success=not is_fail)
|
||||
if not is_fail:
|
||||
try:
|
||||
maybe_request_proof(agent, thread_id, job_id, result_text)
|
||||
maybe_request_proof(agent, thread_id, job_id, result_text,
|
||||
dry_run=dry_run)
|
||||
except Exception as pe:
|
||||
sys.stderr.write(f"warning: proof check failed: {pe}\n")
|
||||
archive_ephemeral_thread(agent, thread_id, job_id=job_id)
|
||||
archive_ephemeral_thread(agent, thread_id, job_id=job_id,
|
||||
dry_run=dry_run)
|
||||
clear_matching_followups(followups, agent, thread_id, mid, text,
|
||||
dry_run, job_id=job_id, verb="RESULT")
|
||||
for verb, job_id in verbs:
|
||||
@@ -1122,11 +1260,13 @@ def harvest_agent_thread(cdp, agent, thread_info, watermarks, followups, dry_run
|
||||
)
|
||||
|
||||
|
||||
def archive_ephemeral_thread(agent, thread_id, job_id=None):
|
||||
def archive_ephemeral_thread(agent, thread_id, job_id=None, dry_run=False):
|
||||
"""
|
||||
If thread_id belongs to an ephemeral job or one-off check,
|
||||
archive it via hybrid gateway and tag it as archived in job-sidechats.json.
|
||||
"""
|
||||
if dry_run:
|
||||
return
|
||||
if not thread_id or not re.fullmatch(r"[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}", thread_id.lower()):
|
||||
return
|
||||
|
||||
@@ -1699,10 +1839,15 @@ def clear_matching_followups(followups, agent, thread_id, mid, text, dry_run=Fal
|
||||
|
||||
# Match by job_id (from [RESULT <job_id>] or [VERB <job_id>]) --
|
||||
# works regardless of thread_uuid or target. This is an ADDITIONAL
|
||||
# path, not a replacement.
|
||||
# path, not a replacement. Also matches the followup's own key
|
||||
# (dm_id): agents quote the DM id from nudge text ([RESULT
|
||||
# <dm_id>]), which differs from job_id on DM-ordered followups
|
||||
# (observed live: [RESULT f4293153] vs job ml-muse-*).
|
||||
match_job = False
|
||||
if job_id and f_rec.get("job_id") and f_rec.get("job_id") == job_id:
|
||||
match_job = True
|
||||
elif job_id and job_id == f_id:
|
||||
match_job = True
|
||||
|
||||
# Non-RESULT verbs are job-scoped: they must not acknowledge/resolve
|
||||
# unrelated pending followups that merely share the thread. RESULT
|
||||
|
||||
Reference in New Issue
Block a user