feat(sidechats): auto-archive completed ephemeral job threads via hybrid gateway

This commit is contained in:
operator
2026-10-05 03:57:06 +00:00
parent 2132d15cb4
commit d09e1c4872
3 changed files with 79 additions and 0 deletions
+41
View File
@@ -381,6 +381,10 @@ def get_monitored_threads(target_agent=None):
if not (key.startswith("pipe-") or key.startswith("test-") or key.startswith("onboarding-")):
is_recent = True
# If already archived, exclude from active monitoring unless pending followup
if isinstance(val, dict) and val.get("archived") and not is_pending:
continue
if not (is_perm or is_pending or is_recent):
continue
@@ -625,6 +629,8 @@ def process_messages(raw_messages, agent, thread_id, thread_name, last_wm, follo
if not dry_run:
append_jsonl(JOB_LOG, job_record)
trigger_chain_next(job_id, result_text, success=not is_fail)
if not is_fail:
archive_ephemeral_thread(agent, thread_id, job_id=job_id)
clear_matching_followups(followups, agent, thread_id, mid, text,
dry_run, job_id=job_id, verb="RESULT")
for verb, job_id in verbs:
@@ -741,6 +747,41 @@ def harvest_agent_thread(cdp, agent, thread_info, watermarks, followups, dry_run
)
def archive_ephemeral_thread(agent, thread_id, job_id=None):
"""
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 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
# Check if thread is marked as persistent in job-sidechats.json
try:
if JOB_SIDECHATS_FILE.exists():
sc_state = json.loads(JOB_SIDECHATS_FILE.read_text(encoding="utf-8"))
for key, val in sc_state.items():
if isinstance(val, dict) and val.get("thread_uuid") == thread_id:
if val.get("type") == "persistent":
return # Do not archive persistent coordinator channels
val["archived"] = True
val["archived_at"] = utcnow()
val["archived_by_job"] = job_id
JOB_SIDECHATS_FILE.write_text(json.dumps(sc_state, indent=2), encoding="utf-8")
except Exception as e:
sys.stderr.write(f"warning: failed to update job-sidechats state: {e}\n")
# Call muse-threads.py archive via subprocess (runs in node's netns)
try:
helper = BIN_DIR / "muse-threads.py"
subprocess.run(
[sys.executable, str(helper), "archive", "--agent", agent, "--thread", thread_id],
capture_output=True, text=True, timeout=15
)
print(f"[{agent}] Archived ephemeral thread {thread_id[:8]} upon completion of job {job_id or 'unknown'}")
except Exception as ae:
sys.stderr.write(f"warning: archive_ephemeral_thread failed: {ae}\n")
# Chain deduplication
CHAINED_JOBS_FILE = Path(__file__).parent / "chained-jobs.json"