From ec547b28acdf72cc6fbcd1ef2dc414f28aebf3e8 Mon Sep 17 00:00:00 2001 From: operator Date: Mon, 5 Oct 2026 02:27:23 +0000 Subject: [PATCH] feat(sidechats): headless gateway new channel auto-provisioning and dedicated heartbeat channel --- bin/dm.py | 54 +++++++++++++++++++++++++++++++ bin/job-dispatch.py | 51 +++++++++++++++++++++++++++--- job-sidechats.json | 61 +++++++++++++++++++++++++++++++++--- jobs/box-deep-health.json | 5 +-- jobs/box-http-health.json | 5 +-- jobs/box-service-health.json | 23 +++++++++++++- jobs/canary-test.json | 15 ++++++++- 7 files changed, 200 insertions(+), 14 deletions(-) diff --git a/bin/dm.py b/bin/dm.py index ae54d03..180265d 100755 --- a/bin/dm.py +++ b/bin/dm.py @@ -598,6 +598,60 @@ def dm_send(agent, target, message, verify=True, raw=False, except Exception as _e: log_event({"type": "gateway_fallback", "id": msg_id, "error": str(_e)[:100]}) + # Fast gateway auto-provisioning for named sidechat targets: + if not is_uuid and recipient in VALID_AGENTS and target != "main": + try: + import muse_hybrid + matched_uuid = None + gw_threads, gw_t_err = muse_hybrid.get_threads(recipient) + if not gw_t_err and gw_threads: + for t in gw_threads: + t_title = (t.get("title") or "").strip().lower() + if t_title == target.strip().lower(): + matched_uuid = t.get("session_id") + break + + if not matched_uuid: + # Spawn a brand new sidechat channel via fast headless gateway + res, err = muse_hybrid.start_session(recipient, title=target) + if res and not err and res.get("session_id"): + matched_uuid = res.get("session_id") + log_event({"type": "sidechat_autoprovisioned", "id": msg_id, "target": target, + "thread_uuid": matched_uuid, "agent": recipient, "transport": "gateway"}) + + if matched_uuid: + # Save mapping to job-sidechats.json + sc_data = load_sidechat_map() + save_key = f"{target}@{recipient}" if f"{target}@{recipient}" in sc_data else target + sc_data[save_key] = { + "thread_uuid": matched_uuid, + "agent": recipient, + "title": target, + "created_at": datetime.now(timezone.utc).isoformat() + } + tmp_sc = f"{SIDCHAT_MAP_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, SIDCHAT_MAP_FILE) + + # Send immediately via gateway + gw_res, gw_err = muse_hybrid.send_message(recipient, tagged, thread_id=matched_uuid, wait=0) + if gw_res and not gw_err: + thread_uuid = matched_uuid + tags["thread"] = thread_uuid + log_event({"type": "verified", "id": msg_id, "agent": agent, "to": recipient, + "target": target, "thread_uuid": thread_uuid, "transport": "gateway", + "placement": "confirmed"}) + log_event({"type": "send_done", "id": msg_id, "agent": agent, "to": recipient, "target": target}) + log_event({"type": "sent", "id": msg_id, "agent": agent, "to": recipient, "target": target, + "verified": True, "transport": "gateway", "tags": tags}) + print(f"DM {msg_id} from {agent} to {recipient}/{target}: SENT and VERIFIED thread={thread_uuid} (gateway, auto-provisioned)") + if tags.get("reply:expected"): + _register_followup(msg_id, agent, recipient, target, tags) + return msg_id + except Exception as _ge: + log_event({"type": "gateway_autoprovision_fallback", "id": msg_id, "error": str(_ge)[:100]}) + if target == "main": _rc, _out, _err = run_full(f"{NETVM_EXEC} {recipient} -- python3 {API} --account {recipient} sidechat main") else: diff --git a/bin/job-dispatch.py b/bin/job-dispatch.py index 9c7f9cf..7d9adaa 100755 --- a/bin/job-dispatch.py +++ b/bin/job-dispatch.py @@ -455,14 +455,57 @@ def main(): target = p_entry.get("target") or f"pipe-{pipeline_run_id.split('-')[-1]}" else: target = f"pipe-{pipeline_run_id.split('-')[-1]}" + elif job.get("sidechat", {}).get("create"): + sc_cfg = job.get("sidechat", {}) + sc_name = render_prompt(sc_cfg.get("name_template", "job-{job_name}-{date}"), variables) + reuse_key = sc_cfg.get("reuse_key") + + # Check if reuse_key exists in job-sidechats.json and thread is still alive + sc_state = load_sidechat_state() + reused_uuid = None + if reuse_key and reuse_key in sc_state: + val = sc_state[reuse_key] + cand_uuid = val.get("thread_uuid") if isinstance(val, dict) else val + if cand_uuid: + try: + import muse_hybrid + threads, err = muse_hybrid.get_threads(agent) + if not err and threads: + thread_ids = [t.get("session_id") for t in threads] + if cand_uuid in thread_ids: + reused_uuid = cand_uuid + except Exception: + pass + + if reused_uuid: + target = reused_uuid + else: + # Spawn a brand new sidechat/channel via fast headless gateway! + channel_title = sc_name + try: + 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 + sc_state[key_to_save] = { + "thread_uuid": new_uuid, + "agent": agent, + "title": channel_title, + "created_at": datetime.now(timezone.utc).isoformat() + } + save_sidechat_state(sc_state) + target = new_uuid + print(f"Spawned new sidechat channel '{channel_title}' ({new_uuid}) for {agent}") + else: + target = sc_name + except Exception as e: + sys.stderr.write(f"warning: fast gateway session-start exception ({e}), falling back to name {sc_name}\n") + target = sc_name elif job.get("dm_target"): target = job.get("dm_target").strip() elif job.get("target"): target = job.get("target").strip() - elif job.get("sidechat", {}).get("create"): - sc_cfg = job.get("sidechat", {}) - 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. diff --git a/job-sidechats.json b/job-sidechats.json index f9e8b7e..f09408e 100644 --- a/job-sidechats.json +++ b/job-sidechats.json @@ -92,14 +92,16 @@ "purpose": "Main loop coordination: health reports, digests, alerts" }, "heartbeat": { - "thread_uuid": "5f18476d-8994-49e7-a9e0-4838732363fe", + "thread_uuid": "51695451-ffba-49e2-8c6d-ffa4f511ce5a", "agent": "opm", - "description": "Alias for main-loop brain heartbeat" + "description": "Dedicated heartbeat thread (provisioned 2026-10-05 via headless gateway)", + "updated_at": "2026-10-05T02:22:19+00:00" }, "heartbeat-opm": { - "thread_uuid": "5f18476d-8994-49e7-a9e0-4838732363fe", + "thread_uuid": "51695451-ffba-49e2-8c6d-ffa4f511ce5a", "agent": "opm", - "description": "Reuse key for OPM heartbeat" + "description": "Reuse key for OPM heartbeat -> dedicated heartbeat thread", + "updated_at": "2026-10-05T02:22:19+00:00" }, "646-pip@646": { "thread_uuid": "e3442b78-0300-4034-8cbe-7a0ab6c5764a", @@ -156,5 +158,56 @@ "agent": "dev", "created_at": "2026-10-04T23:37:01+00:00", "note": "onboarding-alpha-probe" + }, + "onboarding-def": { + "thread_uuid": "466005b1-7065-4798-bf31-bbe785d65f3a", + "agent": "def", + "created_at": "2026-10-04T23:45:37+00:00", + "note": "onboarding-def-probe" + }, + "main-loop-brain": { + "thread_uuid": "74d05d6b-af95-423f-8999-f733ab647785", + "agent": "646", + "created_at": "2026-10-05T00:53:28.473884+00:00" + }, + "dev-audit-channel": { + "thread_uuid": "729d31e0-5154-44ad-97c8-748d36daa448", + "agent": "dev", + "title": "dev-audit-channel", + "created_at": "2026-10-05T02:23:45.941926+00:00" + }, + "canary-test": { + "thread_uuid": "bcda3950-9966-4edd-a014-f931d8b17701", + "agent": "opm", + "title": "canary-test", + "created_at": "2026-10-05T02:24:25.475505+00:00" + }, + "heartbeat-thread": { + "thread_uuid": "51695451-ffba-49e2-8c6d-ffa4f511ce5a", + "agent": "opm", + "description": "Alias for heartbeat (dedicated heartbeat channel provisioned 2026-10-05)" + }, + "main-loop": { + "thread_uuid": "5f18476d-8994-49e7-a9e0-4838732363fe", + "agent": "opm", + "description": "Alias for main-loop brain; added 2026-10-05 to stop autoprovision junk" + }, + "box-http-health": { + "thread_uuid": "d9c596c4-5ace-44dc-9619-9862cf56968d", + "agent": "646", + "title": "box-http-health", + "created_at": "2026-10-05T02:25:41.573556+00:00" + }, + "box-service-health": { + "thread_uuid": "94ab9f9f-012c-41e5-ad35-dc44f2b24281", + "agent": "646", + "title": "box-service-health", + "created_at": "2026-10-05T02:26:11.075842+00:00" + }, + "box-deep-health": { + "thread_uuid": "6628c035-4413-4d9f-863c-56ee861c8c83", + "agent": "646", + "title": "box-deep-health", + "created_at": "2026-10-05T02:26:37.043639+00:00" } } \ No newline at end of file diff --git a/jobs/box-deep-health.json b/jobs/box-deep-health.json index dbfaf98..5769efd 100644 --- a/jobs/box-deep-health.json +++ b/jobs/box-deep-health.json @@ -2,7 +2,6 @@ "agent": "646", "chain_next": null, "description": "Box deep health \u2014 daily comprehensive check", - "dm_target": "646 tasks", "followup": { "escalate": "opm", "expect_reply": true, @@ -15,7 +14,9 @@ "prompt_template": "Box deep health check (daily).\nJob ID: {job_id}\nTime: {datetime}\n\nRun the FULL Box health suite on the VM:\n ssh dev-operator-646@34.139.37.135 \"/srv/box/bin/box-health-check.sh all\"\n\nThen verify these endpoints and services:\n1. https://box.muse-dev.online/ loads (Fleet Console UI, assets box.js and box.css 200 OK)\n2. https://box.muse-dev.online/api/box/fleet returns live fleet status\n3. https://box.muse-dev.online/api/box/dm/log returns live DM log entries\n4. https://box.muse-dev.online/followups loads (operator view, follow-up dashboard)\n5. https://box.muse-dev.online/api/box/followups/summary returns valid JSON\n6. Check /srv/box/dm-log.jsonl tail for errors in the last 24h\n7. Confirm the request-sweeper processed follow-ups in the last hour\n (grep dm_followup /srv/box/box_requests.jsonl | tail -5)\n\nReport any degradation, even if checks pass (slow responses, log warnings).\n\nReply with OK or FAIL , plus any observations worth tracking.", "schedule": "0 9 * * *", "sidechat": { - "create": false + "create": true, + "name_template": "box-deep-health", + "reuse_key": "box-deep-health" }, "timeout": 1800 } diff --git a/jobs/box-http-health.json b/jobs/box-http-health.json index f61f8cb..2883730 100644 --- a/jobs/box-http-health.json +++ b/jobs/box-http-health.json @@ -2,7 +2,6 @@ "agent": "646", "chain_next": null, "description": "Box HTTP endpoint health \u2014 every 15 minutes", - "dm_target": "646 tasks", "followup": { "escalate": "opm", "expect_reply": true, @@ -15,7 +14,9 @@ "prompt_template": "Box HTTP health check.\nJob ID: {job_id}\nTime: {datetime}\n\nRun the Box HTTP checks on the VM:\n ssh dev-operator-646@34.139.37.135 \"/srv/box/bin/box-health-check.sh http\"\n(or run /srv/box/bin/box-health-check.sh http via any VM access you have)\n\nAlso verify the new live Box Console endpoints:\n1. https://box.muse-dev.online/ (Front-Door Console loads, assets box.js and box.css return 200)\n2. https://box.muse-dev.online/api/box/fleet (Returns 200 with live fleet array)\n3. https://box.muse-dev.online/api/box/dm/log (Returns 200 with recent DM log entries)\n\nExpected: all OK. If any FAIL, investigate immediately (check board.service,\nCaddy, recent deploys) and include the failure lines in your reply.\n\nReply with OK or FAIL .", "schedule": "*/15 * * * *", "sidechat": { - "create": false + "create": true, + "name_template": "box-http-health", + "reuse_key": "box-http-health" }, "timeout": 600 } diff --git a/jobs/box-service-health.json b/jobs/box-service-health.json index c4435c4..2950f2d 100644 --- a/jobs/box-service-health.json +++ b/jobs/box-service-health.json @@ -1 +1,22 @@ -{"agent":"646","chain_next":null,"description":"Box service + data health — every 15 minutes (offset 7m)","dm_target":"646 tasks","followup":{"escalate":"opm","expect_reply":true,"nudges":2,"route":"box-health","timeout":"30m"},"name":"box-service-health","on_failure":"alert","prompt_template":"Box service health check.\nJob ID: {job_id}\nTime: {datetime}\n\nRun the Box service and data checks on the VM:\n ssh dev-operator-646@34.139.37.135 \"/srv/box/bin/box-health-check.sh services\"\n ssh dev-operator-646@34.139.37.135 \"/srv/box/bin/box-health-check.sh data\"\n(or run /srv/box/bin/box-health-check.sh via any VM access you have)\n\nExpected: board.service active, caddy active, sweeper timer firing,\n/srv/box writable. If any FAIL, investigate and include failure lines.\n\nReply with [RESULT {job_id}] OK or [RESULT {job_id}] FAIL .","schedule":"7,22,37,52 * * * *","sidechat":{"create":false},"timeout":600} \ No newline at end of file +{ + "agent": "646", + "chain_next": null, + "description": "Box service + data health — every 15 minutes (offset 7m)", + "followup": { + "escalate": "opm", + "expect_reply": true, + "nudges": 2, + "route": "box-health", + "timeout": "30m" + }, + "name": "box-service-health", + "on_failure": "alert", + "prompt_template": "Box service health check.\nJob ID: {job_id}\nTime: {datetime}\n\nRun the Box service and data checks on the VM:\n ssh dev-operator-646@34.139.37.135 \"/srv/box/bin/box-health-check.sh services\"\n ssh dev-operator-646@34.139.37.135 \"/srv/box/bin/box-health-check.sh data\"\n(or run /srv/box/bin/box-health-check.sh via any VM access you have)\n\nExpected: board.service active, caddy active, sweeper timer firing,\n/srv/box writable. If any FAIL, investigate and include failure lines.\n\nReply with [RESULT {job_id}] OK or [RESULT {job_id}] FAIL .", + "schedule": "7,22,37,52 * * * *", + "sidechat": { + "create": true, + "name_template": "box-service-health", + "reuse_key": "box-service-health" + }, + "timeout": 600 +} \ No newline at end of file diff --git a/jobs/canary-test.json b/jobs/canary-test.json index f3f5b25..331b397 100644 --- a/jobs/canary-test.json +++ b/jobs/canary-test.json @@ -1 +1,14 @@ -{"agent":"opm","chain_next":null,"description":"Canary job to test the dispatcher - sends a simple DM","dm_target":"heartbeat","name":"canary-test","on_failure":"alert","prompt_template":"This is a canary test from the job dispatcher.\nJob ID: {job_id}\nDate: {date}\n\nPlease reply to confirm you received this.","schedule":"manual","sidechat":{"create":false},"timeout":300} \ No newline at end of file +{ + "agent": "opm", + "chain_next": null, + "description": "Canary job to test the dispatcher - sends a simple DM", + "name": "canary-test", + "on_failure": "alert", + "prompt_template": "This is a canary test from the job dispatcher.\nJob ID: {job_id}\nDate: {date}\n\nPlease reply to confirm you received this.", + "schedule": "manual", + "sidechat": { + "create": true, + "name_template": "canary-test" + }, + "timeout": 300 +} \ No newline at end of file