diff --git a/bin/response-harvester.py b/bin/response-harvester.py index 2a2e236..8f3365a 100755 --- a/bin/response-harvester.py +++ b/bin/response-harvester.py @@ -1042,36 +1042,56 @@ def reconcile_and_dispatch_swarms(dry_run=False): f"Reply with [RESULT {sid}/{slot_idx}] OK or FAIL ." ) - # Verified dispatch: confirm the send landed BEFORE - # marking the slot running. Fire-and-forget used to - # freeze slots from birth -- the slot read "running" - # while no worker was ever notified (2026-10-05: dev's - # wedged WireGuard data path silently killed every - # gateway dispatch). + # Native Subagent Execution Bridge: + # Spawns an interactive child session directly on the target worker + # via muse_hybrid (Option 3), ensuring an active consumer actually executes + # the prompt and returns results. dispatch_ok = False + subagent_sid = None try: - dm_cmd = [ - sys.executable, str(BIN_DIR / "dm.py"), "send", - "--agent", creator, - "--to", worker, - "--target", slot_target, - prompt, - ] - proc = subprocess.run( - dm_cmd, stdout=subprocess.PIPE, stderr=subprocess.PIPE, - text=True, timeout=DISPATCH_VERIFY_TIMEOUT) - _dout = proc.stdout or "" - dispatch_ok = (proc.returncode == 0 and "SENT" in _dout) + import muse_hybrid + sub_res, sub_err = muse_hybrid.start_session(worker, title=f"swarm-{sid[:8]}-s{slot_idx}") + if not sub_err and sub_res and sub_res.get("session_id"): + subagent_sid = sub_res["session_id"] + send_res, send_err = muse_hybrid.send_message(worker, prompt, thread_id=subagent_sid, wait=0) + if not send_err: + dispatch_ok = True + try: + import subagent_tracker + subagent_tracker.register_session(creator, subagent_sid, title=f"{sid}-s{slot_idx}", prompt=prompt) + except Exception: + pass + print(f"[swarm] Dispatched slot {slot_idx} of {sid} to native subagent session {subagent_sid} on {worker}") if not dispatch_ok: + sys.stderr.write(f"warning: native subagent dispatch failed for {sid} slot {slot_idx} -> {worker} (err: {sub_err or send_err})\n") + except Exception as se: + sys.stderr.write(f"warning: muse_hybrid subagent dispatch exception: {se}\n") + + # Fallback to dm.py send if native session creation fails + if not dispatch_ok: + try: + dm_cmd = [ + sys.executable, str(BIN_DIR / "dm.py"), "send", + "--agent", creator, + "--to", worker, + "--target", slot_target, + prompt, + ] + proc = subprocess.run( + dm_cmd, stdout=subprocess.PIPE, stderr=subprocess.PIPE, + text=True, timeout=DISPATCH_VERIFY_TIMEOUT) + _dout = proc.stdout or "" + dispatch_ok = (proc.returncode == 0 and "SENT" in _dout) + if not dispatch_ok: + sys.stderr.write( + "warning: swarm fallback dispatch unverified %s slot %d -> %s " + "(rc=%d out=%.100s err=%.100s)\n" + % (sid, slot_idx, worker, proc.returncode, + _dout, proc.stderr or "")) + except Exception as de: sys.stderr.write( - "warning: swarm dispatch unverified %s slot %d -> %s " - "(rc=%d out=%.100s err=%.100s)\n" - % (sid, slot_idx, worker, proc.returncode, - _dout, proc.stderr or "")) - except Exception as de: - sys.stderr.write( - "warning: swarm dispatch error %s slot %d -> %s: %s\n" - % (sid, slot_idx, worker, de)) + "warning: swarm dispatch error %s slot %d -> %s: %s\n" + % (sid, slot_idx, worker, de)) if not dispatch_ok: # Leave pending so the next harvester pass retries @@ -1101,6 +1121,8 @@ def reconcile_and_dispatch_swarms(dry_run=False): s["agent_id"] = worker s["status"] = "running" s["updated_ts"] = now + if subagent_sid: + s["subagent_session_id"] = subagent_sid busy_workers.add(worker) modified = True print(f"[swarm] Dispatched slot {slot_idx} of {sid} to {worker} in sidechat {slot_target}")