From 581bbfe3b48fba4479f56c42cd4e0acffe113893 Mon Sep 17 00:00:00 2001 From: operator Date: Mon, 5 Oct 2026 03:47:18 +0000 Subject: [PATCH] fix(harvester): format tool outputs into concise human-readable messages --- bin/response-harvester.py | 54 ++++++++++++++++++++++++++++++++++++--- 1 file changed, 51 insertions(+), 3 deletions(-) diff --git a/bin/response-harvester.py b/bin/response-harvester.py index fa5c0c9..3f21a3b 100755 --- a/bin/response-harvester.py +++ b/bin/response-harvester.py @@ -187,6 +187,51 @@ def execute_agent_tool(agent, op, args, timeout=30): return False, str(e) +def format_tool_result_for_chat(op, raw_output): + """ + Format tool output into a clean, human-readable chat message + instead of spewing a raw wall of unformatted JSON. + """ + if not raw_output: + return "OK" + + try: + data = json.loads(raw_output) + except Exception: + s = str(raw_output).strip() + return s[:500] if len(s) > 500 else s + + if op == "health.check" and isinstance(data, dict) and "fleet" in data: + nodes = data.get("fleet", []) + all_ok = all(n.get("cdp_ok") and n.get("proc_alive") for n in nodes) + up_count = sum(1 for n in nodes if n.get("cdp_ok") and n.get("proc_alive")) + lines = [f"Fleet Health: {'ALL GREEN' if all_ok else 'DEGRADED'} ({up_count}/{len(nodes)} nodes online)"] + for n in nodes: + st = "OK" if n.get("cdp_ok") and n.get("proc_alive") else "FAIL" + lat = n.get("latency_ms", 0) + lines.append(f" • {n.get('node')}: {st} ({lat}ms)") + return "\n".join(lines) + + if op == "cron.runs" and isinstance(data, dict) and "jobs" in data: + jobs = data.get("jobs", []) + return f"{len(jobs)} scheduled jobs configured (e.g. {', '.join(j.get('name') for j in jobs[:6])})" + + if op == "cron.status" and isinstance(data, dict): + return f"Cron '{data.get('name')}': {'ACTIVE' if data.get('active') else 'INACTIVE'} (last result: {data.get('last_result', 'unknown')})" + + if op == "vars.list" and isinstance(data, dict) and "variables" in data: + vars_dict = data.get("variables", {}) + sample = ", ".join(f"{k}={v}" for k, v in list(vars_dict.items())[:5]) + return f"{len(vars_dict)} variables: {sample}..." + + if op == "vars.get" and isinstance(data, dict): + return f"{data.get('name')} = {data.get('value')}" + + # General fallback: compact JSON capped to 400 chars + s = json.dumps(data) + return s[:400] + "..." if len(s) > 400 else s + + def utcnow(): return datetime.now(timezone.utc).isoformat() @@ -516,12 +561,15 @@ def process_messages(raw_messages, agent, thread_id, thread_name, last_wm, follo } append_jsonl(JOB_LOG, tool_record) - # Post tool output back to the originating thread via fast headless gateway + # Post clean tool output back to the originating thread via fast headless gateway if thread_id and 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()): try: import muse_hybrid - status_label = "OUTPUT" if t_ok else "ERROR" - resp_text = f"[{status_label} {op}]\n{t_res}" + if t_ok: + clean_msg = format_tool_result_for_chat(op, t_res) + resp_text = f"Tool result (`{op}`):\n{clean_msg}" + else: + resp_text = f"Tool error (`{op}`): {t_res}" muse_hybrid.send_message(agent, resp_text, thread_id=thread_id, wait=0) except Exception as te: sys.stderr.write(f"warning: failed to post tool response back to thread: {te}\n")