From d09e1c4872c90fe14be9309ba55a992fd996b97c Mon Sep 17 00:00:00 2001 From: operator Date: Mon, 5 Oct 2026 03:57:06 +0000 Subject: [PATCH] feat(sidechats): auto-archive completed ephemeral job threads via hybrid gateway --- bin/job-dispatch.py | 1 + bin/response-harvester.py | 41 +++++++++++++++++++++++++++++++++++++++ job-sidechats.json | 37 +++++++++++++++++++++++++++++++++++ 3 files changed, 79 insertions(+) diff --git a/bin/job-dispatch.py b/bin/job-dispatch.py index e3b8549..b42c5af 100755 --- a/bin/job-dispatch.py +++ b/bin/job-dispatch.py @@ -486,6 +486,7 @@ def main(): import muse_hybrid res, err = muse_hybrid.start_session(agent, title=channel_title) if res and not err and res.get("session_id"): + new_uuid = res.get("session_id") key_to_save = reuse_key or sc_name is_persistent = bool(reuse_key) sc_state[key_to_save] = { diff --git a/bin/response-harvester.py b/bin/response-harvester.py index 172b251..7eed868 100755 --- a/bin/response-harvester.py +++ b/bin/response-harvester.py @@ -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" diff --git a/job-sidechats.json b/job-sidechats.json index a12b957..b7241b6 100644 --- a/job-sidechats.json +++ b/job-sidechats.json @@ -312,5 +312,42 @@ "title": "box-service-health-2026-10-05T03:37:00.096660+00:00", "created_at": "2026-10-05T03:37:02.331936+00:00", "type": "ephemeral" + }, + "box-http-health-2026-10-05T03:45:00.595432+00:00": { + "thread_uuid": "39036562-762b-4a70-a90c-a8466e862aa6", + "agent": "646", + "title": "box-http-health-2026-10-05T03:45:00.595432+00:00", + "created_at": "2026-10-05T03:45:07.455504+00:00" + }, + "box-service-health-2026-10-05T03:52:00.107905+00:00": { + "thread_uuid": "1aa6544d-c0de-4eac-9d49-aeccb8785b9b", + "agent": "646", + "title": "box-service-health-2026-10-05T03:52:00.107905+00:00", + "created_at": "2026-10-05T03:52:07.285636+00:00" + }, + "box-service-health-2026-10-05T03:55:00.424082+00:00": { + "thread_uuid": "4974823b-5ca6-4c82-9ec2-046238aade49", + "agent": "646", + "title": "box-service-health-2026-10-05T03:55:00.424082+00:00", + "created_at": "2026-10-05T03:55:26.693139+00:00", + "archived": true, + "archived_at": "2026-10-05T03:56:24.794310+00:00", + "archived_by_job": "box-service-health-20261005-035500-1e649722" + }, + "canary-2026-10-05T03:55:20.212700+00:00": { + "thread_uuid": "9875b9ed-233c-45db-9a29-bce44d4a3169", + "agent": "646", + "title": "canary-2026-10-05T03:55:20.212700+00:00", + "created_at": "2026-10-05T03:55:30.069381+00:00" + }, + "canary-2026-10-05T03:55:55.378058+00:00": { + "thread_uuid": "3472ea41-3c3e-43e1-ad10-74e5474d1352", + "agent": "646", + "title": "canary-2026-10-05T03:55:55.378058+00:00", + "type": "ephemeral", + "created_at": "2026-10-05T03:55:57.750978+00:00", + "archived": true, + "archived_at": "2026-10-05T03:56:25.108240+00:00", + "archived_by_job": "canary-test-20261005-035555-885b4732" } } \ No newline at end of file