feat(box): enable real-time agent tool triggers and capabilities expansion
This commit is contained in:
@@ -240,6 +240,56 @@ case "$cmd" in
|
||||
;;
|
||||
esac
|
||||
;;
|
||||
cron)
|
||||
sub="${1:-runs}"
|
||||
shift || true
|
||||
case "$sub" in
|
||||
runs|list)
|
||||
call_exec "cron.runs" "{}"
|
||||
;;
|
||||
status)
|
||||
name="${1:-heartbeat}"
|
||||
args=$(python3 -c "import json, sys; print(json.dumps({'name': sys.argv[1]}))" "$name")
|
||||
call_exec "cron.status" "$args"
|
||||
;;
|
||||
view)
|
||||
name="${1:?usage: box cron view <name>}"
|
||||
args=$(python3 -c "import json, sys; print(json.dumps({'name': sys.argv[1]}))" "$name")
|
||||
call_exec "cron.view" "$args"
|
||||
;;
|
||||
run)
|
||||
name="${1:?usage: box cron run <name>}"
|
||||
args=$(python3 -c "import json, sys; print(json.dumps({'job': sys.argv[1]}))" "$name")
|
||||
call_exec "cron.run" "$args"
|
||||
;;
|
||||
*)
|
||||
echo "Usage: box cron runs|status|view|run ..."
|
||||
;;
|
||||
esac
|
||||
;;
|
||||
vars)
|
||||
sub="${1:-list}"
|
||||
shift || true
|
||||
case "$sub" in
|
||||
list)
|
||||
call_exec "vars.list" "{}"
|
||||
;;
|
||||
get)
|
||||
name="${1:?usage: box vars get <name>}"
|
||||
args=$(python3 -c "import json, sys; print(json.dumps({'name': sys.argv[1]}))" "$name")
|
||||
call_exec "vars.get" "$args"
|
||||
;;
|
||||
set)
|
||||
name="${1:?usage: box vars set <name> <value>}"
|
||||
val="${2:?usage: box vars set <name> <value>}"
|
||||
args=$(python3 -c "import json, sys; print(json.dumps({'name': sys.argv[1], 'value': sys.argv[2]}))" "$name" "$val")
|
||||
call_exec "vars.set" "$args"
|
||||
;;
|
||||
*)
|
||||
echo "Usage: box vars list|get|set ..."
|
||||
;;
|
||||
esac
|
||||
;;
|
||||
health)
|
||||
call_exec "health.check" "{}"
|
||||
;;
|
||||
@@ -261,6 +311,13 @@ Usage:
|
||||
box dm read [<target=main>] [<limit=10>]
|
||||
box thread list [<agent>]
|
||||
box thread view <thread_id> [<limit=15>]
|
||||
box cron runs
|
||||
box cron status [<name=heartbeat>]
|
||||
box cron view <name>
|
||||
box cron run <name>
|
||||
box vars list
|
||||
box vars get <name>
|
||||
box vars set <name> <value>
|
||||
box health
|
||||
box ping
|
||||
box ops
|
||||
|
||||
+124
-1
@@ -67,6 +67,7 @@ SIG_MAX_SKEW = 300
|
||||
AGENTS = ('muse', 'pip', '646', 'opm', 'def', 'dev')
|
||||
TARGET_RE = re.compile(r'^[A-Za-z0-9][A-Za-z0-9/_.-]{0,63}$')
|
||||
JOB_RE = re.compile(r'^[a-z0-9][a-z0-9-]{0,63}$')
|
||||
NAME_RE = re.compile(r'^[A-Za-z0-9_.-]{1,64}$')
|
||||
IDENT_RE = re.compile(r'^[a-z0-9-]+$')
|
||||
NONCE_RE = re.compile(r'^[0-9a-fA-F]{16,128}$')
|
||||
HEX_RE = re.compile(r'^[0-9a-f]{8,128}$')
|
||||
@@ -523,6 +524,91 @@ def _pipeline_run_build(a):
|
||||
return [sys.executable, os.path.join(BIN_DIR, 'super-cli.py'), 'deploy', 'pipeline', a['name']]
|
||||
|
||||
|
||||
def _cron_status_validate(raw):
|
||||
if not isinstance(raw, dict):
|
||||
raise OpError('args must be an object')
|
||||
allowed = {'name'}
|
||||
for k in raw:
|
||||
if k not in allowed:
|
||||
raise OpError(f'unknown arg: {k}')
|
||||
return {'name': _job_name(raw.get('name')) if raw.get('name') else 'heartbeat'}
|
||||
|
||||
|
||||
def _cron_status_build(a):
|
||||
return [sys.executable, os.path.join(BIN_DIR, 'box-ctl.py'), 'timer-status', a['name']]
|
||||
|
||||
|
||||
def _cron_view_validate(raw):
|
||||
if not isinstance(raw, dict):
|
||||
raise OpError('args must be an object')
|
||||
allowed = {'name'}
|
||||
for k in raw:
|
||||
if k not in allowed:
|
||||
raise OpError(f'unknown arg: {k}')
|
||||
return {'name': _job_name(raw.get('name'))}
|
||||
|
||||
|
||||
def _cron_view_build(a):
|
||||
return [sys.executable, os.path.join(BIN_DIR, 'box-ctl.py'), 'job-get', a['name']]
|
||||
|
||||
|
||||
def _cron_runs_validate(raw):
|
||||
if raw not in ({}, None):
|
||||
raise OpError('cron.runs takes no required args')
|
||||
return {}
|
||||
|
||||
|
||||
def _cron_runs_build(a):
|
||||
return [sys.executable, os.path.join(BIN_DIR, 'box-ctl.py'), 'job-list']
|
||||
|
||||
|
||||
def _vars_list_validate(raw):
|
||||
if raw not in ({}, None):
|
||||
raise OpError('vars.list takes no required args')
|
||||
return {}
|
||||
|
||||
|
||||
def _vars_list_build(a):
|
||||
return [sys.executable, os.path.join(BIN_DIR, 'box-ctl.py'), 'vars-list']
|
||||
|
||||
|
||||
def _vars_get_validate(raw):
|
||||
if not isinstance(raw, dict):
|
||||
raise OpError('args must be an object')
|
||||
allowed = {'name'}
|
||||
for k in raw:
|
||||
if k not in allowed:
|
||||
raise OpError(f'unknown arg: {k}')
|
||||
name = raw.get('name')
|
||||
if not isinstance(name, str) or not NAME_RE.fullmatch(name):
|
||||
raise OpError('name must match safe identifier')
|
||||
return {'name': name}
|
||||
|
||||
|
||||
def _vars_get_build(a):
|
||||
return [sys.executable, os.path.join(BIN_DIR, 'box-ctl.py'), 'vars-get', a['name']]
|
||||
|
||||
|
||||
def _vars_set_validate(raw):
|
||||
if not isinstance(raw, dict):
|
||||
raise OpError('args must be an object')
|
||||
allowed = {'name', 'value'}
|
||||
for k in raw:
|
||||
if k not in allowed:
|
||||
raise OpError(f'unknown arg: {k}')
|
||||
name = raw.get('name')
|
||||
if not isinstance(name, str) or not NAME_RE.fullmatch(name):
|
||||
raise OpError('name must match safe identifier')
|
||||
val = raw.get('value')
|
||||
if val is None or not isinstance(val, (str, int, float, bool)):
|
||||
raise OpError('value must be a scalar')
|
||||
return {'name': name, 'value': str(val)}
|
||||
|
||||
|
||||
def _vars_set_build(a):
|
||||
return [sys.executable, os.path.join(BIN_DIR, 'box-ctl.py'), 'vars-set', a['name'], a['value']]
|
||||
|
||||
|
||||
# op -> {validate, build, timeout, side_effecting, description}
|
||||
OPS = {
|
||||
'dm.send': {
|
||||
@@ -580,6 +666,41 @@ OPS = {
|
||||
'timeout': 120, 'side_effecting': True,
|
||||
'desc': 'Dispatch a multi-step pipeline across agents',
|
||||
},
|
||||
'cron.status': {
|
||||
'validate': _cron_status_validate, 'build': _cron_status_build,
|
||||
'timeout': 30, 'side_effecting': False,
|
||||
'desc': 'Check systemd timer / scheduled check status',
|
||||
},
|
||||
'cron.view': {
|
||||
'validate': _cron_view_validate, 'build': _cron_view_build,
|
||||
'timeout': 30, 'side_effecting': False,
|
||||
'desc': 'View scheduled job definition and prompt',
|
||||
},
|
||||
'cron.runs': {
|
||||
'validate': _cron_runs_validate, 'build': _cron_runs_build,
|
||||
'timeout': 30, 'side_effecting': False,
|
||||
'desc': 'List all scheduled crons/jobs across the fleet',
|
||||
},
|
||||
'cron.run': {
|
||||
'validate': _job_run_validate, 'build': _job_run_build,
|
||||
'timeout': 300, 'side_effecting': True,
|
||||
'desc': 'Trigger on-demand execution of a scheduled job',
|
||||
},
|
||||
'vars.list': {
|
||||
'validate': _vars_list_validate, 'build': _vars_list_build,
|
||||
'timeout': 30, 'side_effecting': False,
|
||||
'desc': 'List intrinsic loop strategy and timing variables',
|
||||
},
|
||||
'vars.get': {
|
||||
'validate': _vars_get_validate, 'build': _vars_get_build,
|
||||
'timeout': 30, 'side_effecting': False,
|
||||
'desc': 'Get specific intrinsic loop variable',
|
||||
},
|
||||
'vars.set': {
|
||||
'validate': _vars_set_validate, 'build': _vars_set_build,
|
||||
'timeout': 30, 'side_effecting': True,
|
||||
'desc': 'Update intrinsic loop variable',
|
||||
},
|
||||
'exec.ping': {
|
||||
'validate': _health_validate,
|
||||
'build': lambda a: ['/bin/echo', 'PONG'],
|
||||
@@ -757,7 +878,9 @@ class Handler(BaseHTTPRequestHandler):
|
||||
self._send(400, {'error': f'bad args: {e}'})
|
||||
return
|
||||
except Exception as e:
|
||||
self._send(500, {'error': 'internal'})
|
||||
sys.stderr.write(f'Validation/build exception for {op}: {e}\n')
|
||||
import traceback; traceback.print_exc()
|
||||
self._send(500, {'error': f'internal: {e}'})
|
||||
return
|
||||
|
||||
# --- execute (NO shell, argv only)
|
||||
|
||||
+2
-1
@@ -486,12 +486,13 @@ 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] = {
|
||||
"thread_uuid": new_uuid,
|
||||
"agent": agent,
|
||||
"title": channel_title,
|
||||
"type": "persistent" if is_persistent else "ephemeral",
|
||||
"created_at": datetime.now(timezone.utc).isoformat()
|
||||
}
|
||||
save_sidechat_state(sc_state)
|
||||
|
||||
@@ -28,6 +28,7 @@ import hashlib
|
||||
import json
|
||||
import os
|
||||
import re
|
||||
import ssl
|
||||
import subprocess
|
||||
import sys
|
||||
import time
|
||||
@@ -105,6 +106,86 @@ def iter_verb_markers(text):
|
||||
yield m.group(1), m.group(2).strip()
|
||||
|
||||
|
||||
def parse_tool_calls(text):
|
||||
"""
|
||||
Extract structured tool/exec calls from assistant messages.
|
||||
Supports:
|
||||
1. [TOOL <op> <json_args>] or [EXEC <op> <json_args>]
|
||||
2. ```box / ```tool / ```exec JSON blocks
|
||||
"""
|
||||
calls = []
|
||||
for m in re.finditer(r"\[(?:TOOL|EXEC)\s+([a-zA-Z0-9_.-]+)(?:\s+(.*?))?\]", text or ""):
|
||||
op = m.group(1).strip()
|
||||
raw_args = (m.group(2) or "").strip()
|
||||
args = {}
|
||||
if raw_args:
|
||||
try:
|
||||
args = json.loads(raw_args)
|
||||
except Exception:
|
||||
args = {"raw": raw_args}
|
||||
calls.append((op, args))
|
||||
|
||||
for m in re.finditer(r"```(?:box|tool|exec)\s*\n(.*?)```", text or "", re.DOTALL):
|
||||
block = m.group(1).strip()
|
||||
try:
|
||||
d = json.loads(block)
|
||||
if isinstance(d, dict) and "op" in d:
|
||||
calls.append((d["op"], d.get("args", {})))
|
||||
except Exception:
|
||||
pass
|
||||
return calls
|
||||
|
||||
|
||||
def execute_agent_tool(agent, op, args, timeout=30):
|
||||
"""
|
||||
Execute tool call via local exec-constrained HTTPS daemon over Tailscale.
|
||||
Returns (ok, result_or_error_string).
|
||||
"""
|
||||
token = None
|
||||
token_file = os.path.expanduser(f"~/.exec-tokens/{agent}")
|
||||
if not os.path.exists(token_file):
|
||||
token_file = os.path.expanduser("~/.exec-server-token")
|
||||
if os.path.exists(token_file):
|
||||
try:
|
||||
with open(token_file) as f:
|
||||
token = f.read().strip()
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
url = "https://100.123.153.75:8444/exec"
|
||||
payload = json.dumps({"op": op, "args": args}).encode("utf-8")
|
||||
headers = {
|
||||
"Content-Type": "application/json",
|
||||
"User-Agent": f"Box-Harvester/{agent}",
|
||||
}
|
||||
if token:
|
||||
headers["Authorization"] = f"Bearer {token}"
|
||||
|
||||
ctx = ssl.create_default_context()
|
||||
ctx.check_hostname = False
|
||||
ctx.verify_mode = ssl.CERT_NONE
|
||||
|
||||
req = urllib.request.Request(url, data=payload, headers=headers, method="POST")
|
||||
try:
|
||||
with urllib.request.urlopen(req, context=ctx, timeout=timeout) as resp:
|
||||
data = resp.read().decode("utf-8")
|
||||
res = json.loads(data)
|
||||
if res.get("rc") == 0:
|
||||
return True, res.get("stdout", "").strip() or "OK"
|
||||
elif "stdout" in res:
|
||||
return False, (res.get("stdout") or res.get("stderr") or "failed").strip()
|
||||
elif "error" in res:
|
||||
return False, res["error"]
|
||||
return True, json.dumps(res)
|
||||
except urllib.error.HTTPError as he:
|
||||
try:
|
||||
err_body = he.read().decode("utf-8")
|
||||
return False, f"HTTP {he.code}: {err_body}"
|
||||
except Exception:
|
||||
return False, f"HTTP {he.code}"
|
||||
except Exception as e:
|
||||
return False, str(e)
|
||||
|
||||
|
||||
def utcnow():
|
||||
return datetime.now(timezone.utc).isoformat()
|
||||
@@ -416,6 +497,35 @@ def process_messages(raw_messages, agent, thread_id, thread_name, last_wm, follo
|
||||
if author == "assistant":
|
||||
markers = list(iter_result_markers(text))
|
||||
verbs = list(iter_verb_markers(text))
|
||||
# Check for structured tool / command calls from the agent
|
||||
tool_calls = parse_tool_calls(text)
|
||||
if tool_calls and not dry_run:
|
||||
for op, args in tool_calls:
|
||||
print(f"[{agent}] Executing tool '{op}' from message {mid[:8]} in thread {thread_name or thread_id[:8]}")
|
||||
t_ok, t_res = execute_agent_tool(agent, op, args)
|
||||
tool_record = {
|
||||
"ts": utcnow(),
|
||||
"type": "tool_exec",
|
||||
"agent": agent,
|
||||
"thread_id": thread_id,
|
||||
"op": op,
|
||||
"args": args,
|
||||
"success": t_ok,
|
||||
"result_snippet": str(t_res)[:300],
|
||||
"msg_id": mid,
|
||||
}
|
||||
append_jsonl(JOB_LOG, tool_record)
|
||||
|
||||
# Post 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}"
|
||||
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")
|
||||
|
||||
if markers or verbs:
|
||||
for job_id, result_text in markers:
|
||||
is_fail = is_fail_result(result_text)
|
||||
|
||||
Reference in New Issue
Block a user