fix(harvester): format tool outputs into concise human-readable messages
This commit is contained in:
@@ -187,6 +187,51 @@ def execute_agent_tool(agent, op, args, timeout=30):
|
|||||||
return False, str(e)
|
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():
|
def utcnow():
|
||||||
return datetime.now(timezone.utc).isoformat()
|
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)
|
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()):
|
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:
|
try:
|
||||||
import muse_hybrid
|
import muse_hybrid
|
||||||
status_label = "OUTPUT" if t_ok else "ERROR"
|
if t_ok:
|
||||||
resp_text = f"[{status_label} {op}]\n{t_res}"
|
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)
|
muse_hybrid.send_message(agent, resp_text, thread_id=thread_id, wait=0)
|
||||||
except Exception as te:
|
except Exception as te:
|
||||||
sys.stderr.write(f"warning: failed to post tool response back to thread: {te}\n")
|
sys.stderr.write(f"warning: failed to post tool response back to thread: {te}\n")
|
||||||
|
|||||||
Reference in New Issue
Block a user