04339bad14
- bin/box-ctl.py: wire tmux-tally, tmux-auto-status, tmux-auto-toggle, tmux-auto-once, and onboard-connects actions with idempotent allowlists - bin/exec-constrained.py: register tmux.tally, tmux.auto_status, and onboard.connects ops for HTTPS execution - bin/super-cli.py: wire box tmux dispatch to approver, add cmd_run/cmd_watch, and json unread formatting - .agents/skills/box/SKILL.md: document tmux worker tally and auto-approval capabilities
2639 lines
90 KiB
Python
Executable File
2639 lines
90 KiB
Python
Executable File
#!/usr/bin/env python3
|
|
"""
|
|
exec-constrained.py — Constrained HTTPS script-execution endpoint for bl.
|
|
|
|
SSH-fallback: when SSH to bl is down, operators can still run a constrained
|
|
set of scripts (dm.py, job-dispatch, chromebox control) over HTTPS.
|
|
|
|
This is the hardened replacement for exec-server.py's /exec endpoint, which
|
|
executes ARBITRARY shell commands (subprocess.run(cmd, shell=True)). This
|
|
server NEVER takes a command string. Clients request a named OPERATION with
|
|
validated ARGUMENTS; the server maps op -> fixed argv. No shell. No
|
|
interpolation. No RCE.
|
|
|
|
Usage:
|
|
python3 exec-constrained.py --port 8444
|
|
|
|
Client:
|
|
curl -sk -X POST https://127.0.0.1:8444/exec \\
|
|
-H 'Authorization: Bearer <token>' \\
|
|
-H 'Content-Type: application/json' \\
|
|
-d '{"op": "dm.send", "args": {"agent": "opm", "to": "646",
|
|
"target": "main", "message": "hello"}}'
|
|
|
|
Or signature auth (no secret in transit):
|
|
payload=$(python3 -c "import json,time,secrets; print(json.dumps({
|
|
'op': 'dm.send',
|
|
'args': {'agent':'opm','to':'646','target':'main','message':'hi'},
|
|
'ts': time.time(), 'nonce': secrets.token_hex(16)}))")
|
|
sig=$(printf '%s' "$payload" | ssh-keygen -Y sign -f ~/.ssh/id_frontdoor -n exec-constrained)
|
|
curl -sk -X POST https://127.0.0.1:8444/exec \\
|
|
-H 'Content-Type: application/json' \\
|
|
-d "$(python3 -c "import json,sys; print(json.dumps({
|
|
'identity': 'operator-main',
|
|
'payload': sys.argv[1], 'signature': sys.argv[2]}))" "$payload" "$sig")"
|
|
|
|
Auth model: master bearer token + per-agent bearer tokens (same files as
|
|
exec-server.py, so migration is drop-in) OR ssh-keygen -Y signatures over
|
|
the op envelope (namespace 'exec-constrained'). See SECURITY.md.
|
|
"""
|
|
import argparse
|
|
import hashlib
|
|
import hmac
|
|
import json
|
|
import os
|
|
import re
|
|
import secrets
|
|
import subprocess
|
|
import sys
|
|
import tempfile
|
|
import threading
|
|
import time
|
|
from http.server import ThreadingHTTPServer, BaseHTTPRequestHandler
|
|
import ssl
|
|
|
|
# ---------------------------------------------------------------- config
|
|
|
|
BIN_DIR = '/home/super/Projects/NetVM/bin'
|
|
JOBS_DIR = '/home/super/Projects/NetVM/jobs'
|
|
TOKEN_FILE = '/home/super/.exec-server-token'
|
|
TOKEN_DIR = '/home/super/.exec-tokens'
|
|
SIGNERS_FILE = '/home/super/.exec-signers'
|
|
NONCE_FILE = '/home/super/.exec-constrained-nonces'
|
|
AUDIT_LOG = '/home/super/.exec-constrained-audit.jsonl'
|
|
SIG_NAMESPACE = 'exec-constrained'
|
|
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}$')
|
|
DM_ID_RE = re.compile(r'^[0-9a-fA-F]{6,64}$')
|
|
TEST_MODULE_RE = re.compile(r'^tests\.[a-z0-9_]+$')
|
|
JOB_DISPATCH_ID_RE = re.compile(r'^[a-z0-9][a-z0-9-]{0,63}-\d{8}-\d{6}-[a-f0-9]{8}$')
|
|
STRAT_TYPES = frozenset({'wake', 'job', 'siphon', 'manual', 'health', 'heartbeat'})
|
|
STRAT_PRIORITIES = frozenset({'routine', 'normal', 'important'})
|
|
SUBTYPE_RE = re.compile(r'^[A-Za-z0-9_.-]{1,64}$')
|
|
MD_ACCOUNT_RE = re.compile(r'^[A-Za-z0-9][A-Za-z0-9_-]{0,31}$')
|
|
MD_FILENAME_RE = re.compile(r'^[A-Za-z0-9_.-]{1,128}$')
|
|
MD_SUBPATH_RE = re.compile(r'^[A-Za-z0-9_.-]+(/[A-Za-z0-9_.-]+)*$')
|
|
MD_TEMPLATE_FILES = frozenset({'SOUL.md', 'PROACTIVE_PREFERENCES.md',
|
|
'HEARTBEAT.md', 'AGENTS.md', 'MEMORY.md',
|
|
'USER.md', 'TOOLS.md', 'IDENTITY.md'})
|
|
MD_MAX_AMEND = 256 * 1024
|
|
MD_MAX_APPEND = 64 * 1024
|
|
|
|
MAX_BODY = 64 * 1024 # 64KB request cap
|
|
MAX_MESSAGE = 2000 # dm.py message cap
|
|
MAX_OUTPUT = 256 * 1024 # per-stream output cap
|
|
WORK_DIR = '/home/super'
|
|
|
|
|
|
# ------------------------------------------------------- rate limiting
|
|
|
|
class RateLimiter:
|
|
"""Token bucket per identity. Thread-safe."""
|
|
def __init__(self, rate_per_min=60, burst=30):
|
|
self.rate = rate_per_min / 60.0
|
|
self.burst = burst
|
|
self._buckets = {}
|
|
self._lock = threading.Lock()
|
|
|
|
def allow(self, ident):
|
|
now = time.monotonic()
|
|
with self._lock:
|
|
tokens, last = self._buckets.get(ident, (self.burst, now))
|
|
tokens = min(self.burst, tokens + (now - last) * self.rate)
|
|
if tokens >= 1.0:
|
|
self._buckets[ident] = (tokens - 1.0, now)
|
|
return True
|
|
self._buckets[ident] = (tokens, now)
|
|
return False
|
|
|
|
|
|
LIMITER = RateLimiter()
|
|
|
|
|
|
# ------------------------------------------------------------- auth
|
|
|
|
def get_token():
|
|
try:
|
|
with open(TOKEN_FILE) as f:
|
|
return f.read().strip()
|
|
except FileNotFoundError:
|
|
return None
|
|
|
|
|
|
def check_token(token):
|
|
"""Bearer token -> identity label, or None. Never logs the token."""
|
|
master = get_token()
|
|
if token and master and hmac.compare_digest(token, master):
|
|
return 'master'
|
|
try:
|
|
names = os.listdir(TOKEN_DIR)
|
|
except FileNotFoundError:
|
|
return None
|
|
for name in names:
|
|
if not IDENT_RE.fullmatch(name):
|
|
continue
|
|
p = os.path.join(TOKEN_DIR, name)
|
|
if not os.path.isfile(p):
|
|
continue
|
|
try:
|
|
with open(p) as f:
|
|
t = f.read().strip()
|
|
except OSError:
|
|
continue
|
|
if t and hmac.compare_digest(token, t):
|
|
return name
|
|
return None
|
|
|
|
|
|
def check_nonce(nonce):
|
|
with _nonce_lock:
|
|
now = time.time()
|
|
fresh = []
|
|
try:
|
|
with open(NONCE_FILE) as f:
|
|
for line in f:
|
|
parts = line.split()
|
|
if len(parts) != 2:
|
|
continue
|
|
n, t = parts
|
|
try:
|
|
if now - float(t) < 2 * SIG_MAX_SKEW:
|
|
fresh.append((n, t))
|
|
except ValueError:
|
|
pass
|
|
except FileNotFoundError:
|
|
pass
|
|
if any(n == nonce for n, _ in fresh):
|
|
return False
|
|
fresh.append((nonce, str(now)))
|
|
try:
|
|
with open(NONCE_FILE, 'w') as f:
|
|
for n, t in fresh:
|
|
f.write(f'{n} {t}\n')
|
|
os.chmod(NONCE_FILE, 0o600)
|
|
except OSError:
|
|
return False
|
|
return True
|
|
|
|
|
|
def check_signature(identity, payload, signature):
|
|
"""ssh-keygen -Y signature over {"op","args","ts","nonce"} -> identity/None."""
|
|
if not IDENT_RE.fullmatch(identity or ''):
|
|
return None
|
|
try:
|
|
data = json.loads(payload)
|
|
except Exception:
|
|
return None
|
|
if not isinstance(data, dict):
|
|
return None
|
|
op = data.get('op')
|
|
args = data.get('args')
|
|
ts = data.get('ts')
|
|
nonce = data.get('nonce')
|
|
if not isinstance(op, str) or op not in OPS:
|
|
return None
|
|
if not isinstance(args, dict):
|
|
return None
|
|
if not isinstance(nonce, str) or not NONCE_RE.fullmatch(nonce):
|
|
return None
|
|
try:
|
|
ts = float(ts)
|
|
except (TypeError, ValueError):
|
|
return None
|
|
if abs(time.time() - ts) > SIG_MAX_SKEW:
|
|
return None
|
|
if not check_nonce(nonce):
|
|
return None
|
|
sig_path = None
|
|
try:
|
|
with tempfile.NamedTemporaryFile('w', delete=False, suffix='.sig') as f:
|
|
f.write(signature if signature.endswith('\n') else signature + '\n')
|
|
sig_path = f.name
|
|
p = subprocess.run(
|
|
['ssh-keygen', '-Y', 'verify', '-f', SIGNERS_FILE, '-I', identity,
|
|
'-n', SIG_NAMESPACE, '-s', sig_path],
|
|
input=payload.encode(), capture_output=True, timeout=15)
|
|
return identity if p.returncode == 0 else None
|
|
except Exception:
|
|
return None
|
|
finally:
|
|
if sig_path:
|
|
try:
|
|
os.unlink(sig_path)
|
|
except OSError:
|
|
pass
|
|
|
|
# ------------------------------------------------------- allowlist
|
|
#
|
|
# Each op maps to a FIXED script with VALIDATED arguments. The client can
|
|
# never influence: the executable path, the subcommand, or any flag name.
|
|
# Only whitelisted argument VALUES flow through, each checked below.
|
|
#
|
|
# To add an op: add an entry here with a build() function. build() receives
|
|
# the validated args dict and returns an argv list (no shell). Arg spec is
|
|
# enforced by validate() before build() runs.
|
|
|
|
class OpError(Exception):
|
|
pass
|
|
|
|
|
|
def _clean_message(s):
|
|
if not isinstance(s, str) or not s.strip():
|
|
raise OpError('message must be a non-empty string')
|
|
if len(s) > MAX_MESSAGE:
|
|
raise OpError(f'message too long (max {MAX_MESSAGE})')
|
|
if any(ord(c) < 32 and c not in '\n\t' for c in s):
|
|
raise OpError('message contains control characters')
|
|
return s
|
|
|
|
|
|
def _agent(v):
|
|
if v not in AGENTS:
|
|
raise OpError(f'agent must be one of {AGENTS}')
|
|
return v
|
|
|
|
|
|
def _target(v):
|
|
if not isinstance(v, str) or not TARGET_RE.fullmatch(v):
|
|
raise OpError('target must match ^[A-Za-z0-9][A-Za-z0-9/_.-]{0,63}$')
|
|
return v
|
|
|
|
|
|
def _opt_int(v, lo, hi, name):
|
|
if v is None:
|
|
return None
|
|
if isinstance(v, bool) or not isinstance(v, int):
|
|
raise OpError(f'{name} must be an integer')
|
|
if not (lo <= v <= hi):
|
|
raise OpError(f'{name} must be {lo}..{hi}')
|
|
return v
|
|
|
|
|
|
def _job_name(v):
|
|
# Must exist in JOBS_DIR and match the safe pattern (no path traversal).
|
|
if not isinstance(v, str) or not JOB_RE.fullmatch(v):
|
|
raise OpError('job must match ^[a-z0-9][a-z0-9-]{0,63}$')
|
|
path = os.path.join(JOBS_DIR, v + '.json')
|
|
if not os.path.isfile(path):
|
|
raise OpError('unknown job')
|
|
return v
|
|
|
|
|
|
def _dm_send_build(a):
|
|
argv = [sys.executable, os.path.join(BIN_DIR, 'dm.py'), 'send',
|
|
'--agent', a['agent'], '--target', a['target']]
|
|
if a.get('to'):
|
|
argv += ['--to', a['to']]
|
|
if a.get('expect_reply'):
|
|
argv += ['--expect-reply']
|
|
if a.get('reply_timeout') is not None:
|
|
argv += ['--reply-timeout', str(a['reply_timeout'])]
|
|
if a.get('reply_nudges') is not None:
|
|
argv += ['--reply-nudges', str(a['reply_nudges'])]
|
|
if a.get('reply_escalate'):
|
|
argv += ['--reply-escalate', a['reply_escalate']]
|
|
if a.get('route'):
|
|
argv += ['--route', a['route']]
|
|
for t in a.get('tags', []):
|
|
argv += ['--tag', t]
|
|
argv.append(a['message'])
|
|
return argv
|
|
|
|
|
|
def _dm_send_validate(raw):
|
|
if not isinstance(raw, dict):
|
|
raise OpError('args must be an object')
|
|
allowed = {'agent', 'to', 'target', 'message', 'expect_reply',
|
|
'reply_timeout', 'reply_nudges', 'reply_escalate',
|
|
'route', 'tags'}
|
|
for k in raw:
|
|
if k not in allowed:
|
|
raise OpError(f'unknown arg: {k}')
|
|
a = {}
|
|
a['agent'] = _agent(raw.get('agent'))
|
|
if 'to' in raw and raw['to'] is not None:
|
|
a['to'] = _agent(raw['to'])
|
|
a['target'] = _target(raw.get('target'))
|
|
a['message'] = _clean_message(raw.get('message'))
|
|
a['expect_reply'] = bool(raw.get('expect_reply', False))
|
|
a['reply_timeout'] = _opt_int(raw.get('reply_timeout'), 60, 604800,
|
|
'reply_timeout')
|
|
a['reply_nudges'] = _opt_int(raw.get('reply_nudges'), 0, 10,
|
|
'reply_nudges')
|
|
esc = raw.get('reply_escalate')
|
|
if esc is not None:
|
|
a['reply_escalate'] = _agent(esc)
|
|
route = raw.get('route')
|
|
if route is not None:
|
|
if not isinstance(route, str) or not TARGET_RE.fullmatch(route):
|
|
raise OpError('route must match target pattern')
|
|
a['route'] = route
|
|
tags = raw.get('tags', [])
|
|
if not isinstance(tags, list) or len(tags) > 8:
|
|
raise OpError('tags must be a list of at most 8 strings')
|
|
clean_tags = []
|
|
for t in tags:
|
|
if not isinstance(t, str) or len(t) > 120:
|
|
raise OpError('tag must be a string <= 120 chars')
|
|
# Canonical tag vocabulary only (see DEPLOY-DECISIONS.md).
|
|
if not re.fullmatch(r'[a-z0-9_:-]+=[a-zA-Z0-9_.:/-]*', t) and \
|
|
not re.fullmatch(r'[a-z0-9_:-]+', t):
|
|
raise OpError(f'malformed tag: {t}')
|
|
clean_tags.append(t)
|
|
a['tags'] = clean_tags
|
|
return a
|
|
|
|
|
|
def _dm_thread_validate(raw):
|
|
# Same shape as dm.send but uses the thread subcommand.
|
|
return _dm_send_validate(raw)
|
|
|
|
|
|
def _dm_thread_build(a):
|
|
argv = _dm_send_build(a)
|
|
argv[2] = 'thread'
|
|
return argv
|
|
|
|
|
|
def _dm_read_validate(raw):
|
|
if not isinstance(raw, dict):
|
|
raise OpError('args must be an object')
|
|
allowed = {'agent', 'target', 'limit'}
|
|
for k in raw:
|
|
if k not in allowed:
|
|
raise OpError(f'unknown arg: {k}')
|
|
return {
|
|
'agent': _agent(raw.get('agent')),
|
|
'target': _target(raw.get('target')),
|
|
'limit': _opt_int(raw.get('limit', 20), 1, 100, 'limit') or 20,
|
|
}
|
|
|
|
|
|
def _dm_read_build(a):
|
|
return [sys.executable, os.path.join(BIN_DIR, 'dm.py'), 'read',
|
|
'--agent', a['agent'], '--target', a['target'],
|
|
'--n', str(a['limit'])]
|
|
|
|
|
|
def _job_run_validate(raw):
|
|
if not isinstance(raw, dict):
|
|
raise OpError('args must be an object')
|
|
allowed = {'job'}
|
|
for k in raw:
|
|
if k not in allowed:
|
|
raise OpError(f'unknown arg: {k}')
|
|
return {'job': _job_name(raw.get('job'))}
|
|
|
|
|
|
def _job_run_build(a):
|
|
return [sys.executable, os.path.join(BIN_DIR, 'job-dispatch.py'), a['job']]
|
|
|
|
|
|
def _chat_messages_validate(raw):
|
|
if not isinstance(raw, dict):
|
|
raise OpError('args must be an object')
|
|
allowed = {'account', 'limit'}
|
|
for k in raw:
|
|
if k not in allowed:
|
|
raise OpError(f'unknown arg: {k}')
|
|
return {
|
|
'account': _agent(raw.get('account')),
|
|
'limit': _opt_int(raw.get('limit', 20), 1, 100, 'limit') or 20,
|
|
}
|
|
|
|
|
|
def _chat_messages_build(a):
|
|
return [sys.executable, os.path.join(BIN_DIR, 'muse-chat-api.py'),
|
|
'--account', a['account'], 'messages', '--limit', str(a['limit'])]
|
|
|
|
|
|
def _chat_send_validate(raw):
|
|
if not isinstance(raw, dict):
|
|
raise OpError('args must be an object')
|
|
allowed = {'account', 'message', 'thread'}
|
|
for k in raw:
|
|
if k not in allowed:
|
|
raise OpError(f'unknown arg: {k}')
|
|
a = {
|
|
'account': _agent(raw.get('account')),
|
|
'message': _clean_message(raw.get('message')),
|
|
}
|
|
if raw.get('thread') is not None:
|
|
th = raw['thread']
|
|
if not isinstance(th, str) or not HEX_RE.fullmatch(th):
|
|
raise OpError('thread must be a hex uuid')
|
|
a['thread'] = th
|
|
return a
|
|
|
|
|
|
def _chat_send_build(a):
|
|
argv = [sys.executable, os.path.join(BIN_DIR, 'muse-chat-api.py'),
|
|
'--account', a['account'], 'send', '--message', a['message']]
|
|
if a.get('thread'):
|
|
argv += ['--thread', a['thread']]
|
|
return argv
|
|
|
|
|
|
def _safe_str(v, max_len=120, name='string'):
|
|
if v is None:
|
|
return ''
|
|
if not isinstance(v, str):
|
|
raise OpError(f'{name} must be a string')
|
|
if len(v) > max_len:
|
|
raise OpError(f'{name} exceeds max length {max_len}')
|
|
if any(ord(c) < 32 and c not in '\n\t' for c in v):
|
|
raise OpError(f'{name} contains control characters')
|
|
return v.strip()
|
|
|
|
|
|
def _health_validate(raw):
|
|
if raw not in ({}, None):
|
|
raise OpError('health.check takes no args')
|
|
return {}
|
|
|
|
|
|
def _health_build(a):
|
|
return [sys.executable, os.path.join(BIN_DIR, 'super-cli.py'), 'fleet', 'status', '--json']
|
|
|
|
|
|
def _cdp_latency_build(a):
|
|
return [sys.executable, os.path.join(BIN_DIR, 'box-ctl.py'), 'cdp-latency']
|
|
|
|
|
|
def _chrome_errors_build(a):
|
|
return [sys.executable, os.path.join(BIN_DIR, 'box-ctl.py'), 'chrome-errors']
|
|
|
|
|
|
def _watchdog_alerts_build(a):
|
|
return [sys.executable, os.path.join(BIN_DIR, 'box-ctl.py'), 'watchdog-alerts']
|
|
|
|
|
|
def _relay_health_build(a):
|
|
return [sys.executable, os.path.join(BIN_DIR, 'box-ctl.py'), 'relay-health']
|
|
|
|
|
|
def _quality_check_validate(raw):
|
|
if raw not in ({}, None):
|
|
if isinstance(raw, dict):
|
|
allowed = {'scope'}
|
|
for k in raw:
|
|
if k not in allowed:
|
|
raise OpError(f'unknown arg: {k}')
|
|
else:
|
|
raise OpError('quality.check args must be an object or empty')
|
|
return {}
|
|
|
|
|
|
def _quality_check_build(a):
|
|
return [sys.executable, os.path.join(BIN_DIR, 'box-ctl.py'), 'quality-check']
|
|
|
|
|
|
def _subagent_spawn_validate(raw):
|
|
if not isinstance(raw, dict):
|
|
raise OpError('args must be an object')
|
|
allowed = {'agent', 'title', 'prompt', 'wait'}
|
|
for k in raw:
|
|
if k not in allowed:
|
|
raise OpError(f'unknown arg: {k}')
|
|
return {
|
|
'agent': _agent(raw.get('agent')),
|
|
'title': _safe_str(raw.get('title', 'subagent'), 120, 'title') or 'subagent',
|
|
'prompt': _clean_message(raw.get('prompt')),
|
|
'wait': _opt_int(raw.get('wait', 0), 0, 60, 'wait') or 0,
|
|
}
|
|
|
|
|
|
def _subagent_spawn_build(a):
|
|
return [sys.executable, os.path.join(BIN_DIR, 'super-cli.py'),
|
|
'deploy', 'subagent',
|
|
'--agent', a['agent'],
|
|
'--title', a['title'],
|
|
'--wait', str(a['wait']),
|
|
a['prompt']]
|
|
|
|
|
|
def _thread_list_validate(raw):
|
|
if not isinstance(raw, dict):
|
|
raise OpError('args must be an object')
|
|
allowed = {'agent'}
|
|
for k in raw:
|
|
if k not in allowed:
|
|
raise OpError(f'unknown arg: {k}')
|
|
return {'agent': _agent(raw.get('agent'))}
|
|
|
|
|
|
def _thread_list_build(a):
|
|
return [sys.executable, os.path.join(BIN_DIR, 'super-cli.py'),
|
|
'thread', 'list', a['agent'], '--json']
|
|
|
|
|
|
def _thread_view_validate(raw):
|
|
if not isinstance(raw, dict):
|
|
raise OpError('args must be an object')
|
|
allowed = {'agent', 'thread', 'limit'}
|
|
for k in raw:
|
|
if k not in allowed:
|
|
raise OpError(f'unknown arg: {k}')
|
|
th = raw.get('thread')
|
|
if not isinstance(th, str) or not TARGET_RE.fullmatch(th):
|
|
raise OpError('thread must match target pattern')
|
|
return {
|
|
'agent': _agent(raw.get('agent')),
|
|
'thread': th,
|
|
'limit': _opt_int(raw.get('limit', 15), 1, 100, 'limit') or 15,
|
|
}
|
|
|
|
|
|
def _thread_view_build(a):
|
|
return [sys.executable, os.path.join(BIN_DIR, 'super-cli.py'),
|
|
'thread', 'view', a['agent'], a['thread'],
|
|
'--limit', str(a['limit']), '--json']
|
|
|
|
|
|
def _pipeline_run_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 _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 _cron_timer_start_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_timer_start_build(a):
|
|
return [sys.executable, os.path.join(BIN_DIR, 'box-ctl.py'), 'timer-start', a['name']]
|
|
|
|
|
|
def _cron_timer_create_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_timer_create_build(a):
|
|
return [sys.executable, os.path.join(BIN_DIR, 'box-ctl.py'), 'timer-create', a['name']]
|
|
|
|
|
|
# ------------------------------------------------------- tmux ops
|
|
def _tmux_session_name(s):
|
|
if not isinstance(s, str) or not re.fullmatch(r'^[a-zA-Z0-9_.-]{1,64}$', s):
|
|
raise OpError('session must match ^[a-zA-Z0-9_.-]{1,64}$')
|
|
return s
|
|
|
|
|
|
def _tmux_send_validate(raw):
|
|
if not isinstance(raw, dict):
|
|
raise OpError('args must be an object')
|
|
allowed = {'session', 'keys', 'no_enter'}
|
|
for k in raw:
|
|
if k not in allowed:
|
|
raise OpError(f'unknown arg: {k}')
|
|
if not raw.get('session'):
|
|
raise OpError('session is required')
|
|
if 'keys' not in raw:
|
|
raise OpError('keys is required')
|
|
return {
|
|
'session': _tmux_session_name(raw['session']),
|
|
'keys': str(raw['keys']),
|
|
'no_enter': bool(raw.get('no_enter', False)),
|
|
}
|
|
|
|
|
|
def _tmux_send_build(a):
|
|
cmd = [sys.executable, os.path.join(BIN_DIR, 'muse-tmux.py'), 'send', a['session'], a['keys']]
|
|
if a.get('no_enter'):
|
|
cmd.append('--no-enter')
|
|
return cmd
|
|
|
|
|
|
def _tmux_capture_validate(raw):
|
|
if not isinstance(raw, dict):
|
|
raise OpError('args must be an object')
|
|
allowed = {'session', 'lines'}
|
|
for k in raw:
|
|
if k not in allowed:
|
|
raise OpError(f'unknown arg: {k}')
|
|
if not raw.get('session'):
|
|
raise OpError('session is required')
|
|
lines = raw.get('lines', 30)
|
|
try:
|
|
lines = int(lines)
|
|
if lines < 1 or lines > 500:
|
|
lines = 30
|
|
except Exception:
|
|
lines = 30
|
|
return {
|
|
'session': _tmux_session_name(raw['session']),
|
|
'lines': lines,
|
|
}
|
|
|
|
|
|
def _tmux_capture_build(a):
|
|
return [sys.executable, os.path.join(BIN_DIR, 'muse-tmux.py'), 'capture', a['session'], '--lines', str(a['lines'])]
|
|
|
|
|
|
def _tmux_list_validate(raw):
|
|
return {}
|
|
|
|
|
|
def _tmux_list_build(a):
|
|
return [sys.executable, os.path.join(BIN_DIR, 'muse-tmux.py'), 'list']
|
|
|
|
|
|
def _tmux_new_validate(raw):
|
|
if not isinstance(raw, dict):
|
|
raise OpError('args must be an object')
|
|
allowed = {'session', 'window', 'command'}
|
|
for k in raw:
|
|
if k not in allowed:
|
|
raise OpError(f'unknown arg: {k}')
|
|
if not raw.get('session'):
|
|
raise OpError('session is required')
|
|
return {
|
|
'session': _tmux_session_name(raw['session']),
|
|
'window': str(raw.get('window', '')) if raw.get('window') else None,
|
|
'command': str(raw.get('command', '')) if raw.get('command') else None,
|
|
}
|
|
|
|
|
|
def _tmux_new_build(a):
|
|
cmd = [sys.executable, os.path.join(BIN_DIR, 'muse-tmux.py'), 'new', a['session']]
|
|
if a.get('window'):
|
|
cmd.extend(['--window', a['window']])
|
|
if a.get('command'):
|
|
cmd.extend(['--command', a['command']])
|
|
return cmd
|
|
|
|
|
|
def _tmux_kill_validate(raw):
|
|
if not isinstance(raw, dict):
|
|
raise OpError('args must be an object')
|
|
allowed = {'session'}
|
|
for k in raw:
|
|
if k not in allowed:
|
|
raise OpError(f'unknown arg: {k}')
|
|
if not raw.get('session'):
|
|
raise OpError('session is required')
|
|
return {'session': _tmux_session_name(raw['session'])}
|
|
|
|
|
|
def _tmux_kill_build(a):
|
|
return [sys.executable, os.path.join(BIN_DIR, 'muse-tmux.py'), 'kill', a['session']]
|
|
|
|
|
|
def _tmux_prune_validate(raw):
|
|
if not isinstance(raw, dict):
|
|
raise OpError('args must be an object')
|
|
allowed = {'ttl'}
|
|
for k in raw:
|
|
if k not in allowed:
|
|
raise OpError(f'unknown arg: {k}')
|
|
ttl = raw.get('ttl', 7200)
|
|
try:
|
|
ttl = int(ttl)
|
|
if ttl < 60:
|
|
ttl = 60
|
|
except Exception:
|
|
ttl = 7200
|
|
return {'ttl': ttl}
|
|
|
|
|
|
def _tmux_prune_build(a):
|
|
return [sys.executable, os.path.join(BIN_DIR, 'muse-tmux.py'), 'prune', '--ttl', str(a['ttl'])]
|
|
|
|
|
|
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']]
|
|
|
|
|
|
def _files_read_validate(raw):
|
|
if not isinstance(raw, dict):
|
|
raise OpError('args must be an object')
|
|
allowed = {'path', 'lines'}
|
|
for k in raw:
|
|
if k not in allowed:
|
|
raise OpError(f'unknown arg: {k}')
|
|
p = raw.get('path')
|
|
if not isinstance(p, str) or not p.strip() or '..' in p:
|
|
raise OpError('path must be a safe relative or repo path without ".."')
|
|
return {
|
|
'path': p.strip(),
|
|
'lines': _opt_int(raw.get('lines', 100), 1, 1000, 'lines') or 100
|
|
}
|
|
|
|
|
|
def _files_read_build(a):
|
|
# Pass JSON args via stdin to box-sys-op.py
|
|
return [sys.executable, os.path.join(BIN_DIR, 'box-sys-op.py'), 'files.read']
|
|
|
|
|
|
def _files_write_validate(raw):
|
|
if not isinstance(raw, dict):
|
|
raise OpError('args must be an object')
|
|
allowed = {'path', 'content'}
|
|
for k in raw:
|
|
if k not in allowed:
|
|
raise OpError(f'unknown arg: {k}')
|
|
p = raw.get('path')
|
|
if not isinstance(p, str) or not p.strip() or '..' in p:
|
|
raise OpError('path must be a safe relative or repo path without ".."')
|
|
content = raw.get('content')
|
|
if not isinstance(content, str):
|
|
raise OpError('content must be a string')
|
|
if len(content.encode('utf-8')) > 64 * 1024:
|
|
raise OpError('content exceeds max size 64KB')
|
|
return {'path': p.strip(), 'content': content}
|
|
|
|
|
|
def _files_write_build(a):
|
|
return [sys.executable, os.path.join(BIN_DIR, 'box-sys-op.py'), 'files.write']
|
|
|
|
|
|
def _web_fetch_validate(raw):
|
|
if not isinstance(raw, dict):
|
|
raise OpError('args must be an object')
|
|
allowed = {'url'}
|
|
for k in raw:
|
|
if k not in allowed:
|
|
raise OpError(f'unknown arg: {k}')
|
|
u = raw.get('url')
|
|
if not isinstance(u, str) or not u.strip():
|
|
raise OpError('url must be a string')
|
|
if not (u.startswith('http://') or u.startswith('https://')):
|
|
raise OpError('url must start with http:// or https://')
|
|
return {'url': u.strip()}
|
|
|
|
|
|
def _web_fetch_build(a):
|
|
return [sys.executable, os.path.join(BIN_DIR, 'box-sys-op.py'), 'web.fetch']
|
|
|
|
|
|
def _service_status_validate(raw):
|
|
if not isinstance(raw, dict):
|
|
raise OpError('args must be an object')
|
|
allowed = {'unit', 'name'}
|
|
for k in raw:
|
|
if k not in allowed:
|
|
raise OpError(f'unknown arg: {k}')
|
|
unit = raw.get('unit') or raw.get('name')
|
|
if not isinstance(unit, str) or not NAME_RE.fullmatch(unit):
|
|
raise OpError('unit must match safe identifier')
|
|
return {'unit': unit}
|
|
|
|
|
|
def _service_status_build(a):
|
|
return [sys.executable, os.path.join(BIN_DIR, 'box-sys-op.py'), 'service.status']
|
|
|
|
|
|
def _service_restart_validate(raw):
|
|
if not isinstance(raw, dict):
|
|
raise OpError('args must be an object')
|
|
allowed = {'unit', 'name'}
|
|
for k in raw:
|
|
if k not in allowed:
|
|
raise OpError(f'unknown arg: {k}')
|
|
unit = raw.get('unit') or raw.get('name')
|
|
if not isinstance(unit, str) or not NAME_RE.fullmatch(unit):
|
|
raise OpError('unit must match safe identifier')
|
|
return {'unit': unit}
|
|
|
|
|
|
def _service_restart_build(a):
|
|
return [sys.executable, os.path.join(BIN_DIR, 'box-sys-op.py'), 'service.restart']
|
|
|
|
|
|
def _followup_create_validate(raw):
|
|
if not isinstance(raw, dict):
|
|
raise OpError('args must be an object')
|
|
allowed = {'agent', 'in_m', 'prompt', 'thread', 'sidechat', 'sender'}
|
|
for k in raw:
|
|
if k not in allowed:
|
|
raise OpError(f'unknown arg: {k}')
|
|
agent = _agent(raw.get('agent', 'opm'))
|
|
sender = raw.get('sender')
|
|
if sender:
|
|
sender = _agent(sender)
|
|
try:
|
|
in_m = float(raw.get('in_m', 1))
|
|
except (TypeError, ValueError):
|
|
raise OpError('in_m must be a number')
|
|
if in_m < 0.05 or in_m > 1440:
|
|
raise OpError('in_m must be between 0.05 and 1440 minutes')
|
|
prompt = _clean_message(raw.get('prompt'))
|
|
sidechat = raw.get('thread') or raw.get('sidechat')
|
|
if sidechat and not TARGET_RE.fullmatch(str(sidechat)):
|
|
raise OpError('thread/sidechat must match safe identifier')
|
|
return {
|
|
'agent': agent,
|
|
'sender': sender or agent,
|
|
'in_m': in_m,
|
|
'prompt': prompt,
|
|
'thread': str(sidechat) if sidechat else None
|
|
}
|
|
|
|
|
|
def _followup_create_build(a):
|
|
return [sys.executable, os.path.join(BIN_DIR, 'box-sys-op.py'), 'followup.create']
|
|
|
|
|
|
def _swarm_spawn_validate(raw):
|
|
if not isinstance(raw, dict):
|
|
raise OpError('args must be an object')
|
|
allowed = {'count', 'task', 'label'}
|
|
for k in raw:
|
|
if k not in allowed:
|
|
raise OpError(f'unknown arg: {k}')
|
|
count = _opt_int(raw.get('count', 1), 1, 50, 'count') or 1
|
|
task = _clean_message(raw.get('task'))
|
|
label = raw.get('label')
|
|
if label and not TARGET_RE.fullmatch(str(label)):
|
|
raise OpError('label must match safe identifier')
|
|
return {'count': count, 'task': task, 'label': str(label) if label else None}
|
|
|
|
|
|
def _swarm_spawn_build(a):
|
|
cmd = [sys.executable, os.path.join(BIN_DIR, 'box-ctl.py'), 'swarm-spawn', str(a['count']), a['task']]
|
|
if a.get('label'):
|
|
cmd.extend(['--label', a['label']])
|
|
return cmd
|
|
|
|
|
|
def _swarm_status_validate(raw):
|
|
if not isinstance(raw, dict):
|
|
raise OpError('args must be an object')
|
|
allowed = {'swarm_id', 'id'}
|
|
for k in raw:
|
|
if k not in allowed:
|
|
raise OpError(f'unknown arg: {k}')
|
|
sid = raw.get('swarm_id') or raw.get('id')
|
|
if not isinstance(sid, str) or not TARGET_RE.fullmatch(sid):
|
|
raise OpError('swarm_id must match safe identifier')
|
|
return {'swarm_id': sid}
|
|
|
|
|
|
def _swarm_status_build(a):
|
|
return [sys.executable, os.path.join(BIN_DIR, 'box-ctl.py'), 'swarm-status', a['swarm_id']]
|
|
|
|
|
|
def _swarm_list_validate(raw):
|
|
if raw not in ({}, None):
|
|
raise OpError('swarm.list takes no required args')
|
|
return {}
|
|
|
|
|
|
def _swarm_list_build(a):
|
|
return [sys.executable, os.path.join(BIN_DIR, 'box-ctl.py'), 'swarm-list']
|
|
|
|
|
|
def _swarm_results_validate(raw):
|
|
if not isinstance(raw, dict):
|
|
raise OpError('args must be an object')
|
|
allowed = {'swarm_id', 'id'}
|
|
for k in raw:
|
|
if k not in allowed:
|
|
raise OpError(f'unknown arg: {k}')
|
|
sid = raw.get('swarm_id') or raw.get('id')
|
|
if not isinstance(sid, str) or not TARGET_RE.fullmatch(sid):
|
|
raise OpError('swarm_id must match safe identifier')
|
|
return {'swarm_id': sid}
|
|
|
|
|
|
def _swarm_results_build(a):
|
|
return [sys.executable, os.path.join(BIN_DIR, 'box-ctl.py'), 'swarm-results', a['swarm_id']]
|
|
|
|
|
|
def _fleet_unread_validate(raw):
|
|
if not isinstance(raw, dict):
|
|
raise OpError('args must be an object')
|
|
allowed = {'agent'}
|
|
for k in raw:
|
|
if k not in allowed:
|
|
raise OpError(f'unknown arg: {k}')
|
|
agent = raw.get('agent')
|
|
return {'agent': _agent(agent) if agent else None}
|
|
|
|
|
|
def _fleet_unread_build(a):
|
|
argv = [sys.executable, os.path.join(BIN_DIR, 'box-ctl.py'), 'unread']
|
|
if a.get('agent'):
|
|
argv += ['--agent', a['agent']]
|
|
return argv
|
|
|
|
|
|
def _dm_log_validate(raw):
|
|
if not isinstance(raw, dict):
|
|
raise OpError('args must be an object')
|
|
allowed = {'limit', 'agent'}
|
|
for k in raw:
|
|
if k not in allowed:
|
|
raise OpError(f'unknown arg: {k}')
|
|
agent = raw.get('agent')
|
|
return {
|
|
'limit': _opt_int(raw.get('limit', 20), 1, 100, 'limit') or 20,
|
|
'agent': _agent(agent) if agent else None,
|
|
}
|
|
|
|
|
|
def _dm_log_build(a):
|
|
argv = [sys.executable, os.path.join(BIN_DIR, 'box-ctl.py'),
|
|
'dm-log', str(a['limit'])]
|
|
if a.get('agent'):
|
|
argv += ['--agent', a['agent']]
|
|
return argv
|
|
|
|
|
|
def _repo_path(v):
|
|
# Repo-relative path, no escapes. Mirrors box-ctl.py _git_path.
|
|
if not isinstance(v, str) or not v.strip():
|
|
raise OpError('path must be a non-empty string')
|
|
if '..' in v or v.startswith('/') or \
|
|
any(ord(c) < 32 or ord(c) == 127 for c in v):
|
|
raise OpError("path must be repo-relative without '..'")
|
|
return v
|
|
|
|
|
|
def _sidechat_name(v):
|
|
# Sidechat names may contain spaces ("646 tasks"); reject only
|
|
# control characters and enforce length. Mirrors box-ctl.py.
|
|
if not isinstance(v, str) or not v or len(v) > 64:
|
|
raise OpError('sidechat must be 1-64 chars')
|
|
if any(ord(c) < 32 or ord(c) == 127 for c in v):
|
|
raise OpError('sidechat contains control characters')
|
|
return v
|
|
|
|
|
|
def _git_diff_validate(raw):
|
|
if not isinstance(raw, dict):
|
|
raise OpError('args must be an object')
|
|
allowed = {'path', 'stat'}
|
|
for k in raw:
|
|
if k not in allowed:
|
|
raise OpError(f'unknown arg: {k}')
|
|
path = raw.get('path')
|
|
return {
|
|
'path': _repo_path(path) if path is not None else None,
|
|
'stat': bool(raw.get('stat', False)),
|
|
}
|
|
|
|
|
|
def _git_diff_build(a):
|
|
argv = [sys.executable, os.path.join(BIN_DIR, 'box-ctl.py'), 'git-diff']
|
|
if a.get('stat'):
|
|
argv.append('--stat')
|
|
if a.get('path'):
|
|
argv += ['--path', a['path']]
|
|
return argv
|
|
|
|
|
|
def _git_log_validate(raw):
|
|
if not isinstance(raw, dict):
|
|
raise OpError('args must be an object')
|
|
allowed = {'limit', 'path'}
|
|
for k in raw:
|
|
if k not in allowed:
|
|
raise OpError(f'unknown arg: {k}')
|
|
path = raw.get('path')
|
|
return {
|
|
'limit': _opt_int(raw.get('limit', 10), 1, 50, 'limit') or 10,
|
|
'path': _repo_path(path) if path is not None else None,
|
|
}
|
|
|
|
|
|
def _git_log_build(a):
|
|
argv = [sys.executable, os.path.join(BIN_DIR, 'box-ctl.py'),
|
|
'git-log', '--limit', str(a['limit'])]
|
|
if a.get('path'):
|
|
argv += ['--path', a['path']]
|
|
return argv
|
|
|
|
|
|
def _kfilter(v):
|
|
# unittest -k pattern: plain string, passed as argv (no shell).
|
|
if not isinstance(v, str) or not v.strip() or len(v) > 200:
|
|
raise OpError('filter must be 1-200 chars')
|
|
if any(ord(c) < 32 or ord(c) == 127 for c in v):
|
|
raise OpError('filter contains control characters')
|
|
return v
|
|
|
|
|
|
def _tests_run_validate(raw):
|
|
if not isinstance(raw, dict):
|
|
raise OpError('args must be an object')
|
|
allowed = {'test', 'filter'}
|
|
for k in raw:
|
|
if k not in allowed:
|
|
raise OpError(f'unknown arg: {k}')
|
|
test = raw.get('test')
|
|
filt = raw.get('filter')
|
|
if test is not None and \
|
|
(not isinstance(test, str) or not TEST_MODULE_RE.fullmatch(test)):
|
|
raise OpError('test must match ^tests\\.[a-z0-9_]+$')
|
|
return {
|
|
'test': test,
|
|
'filter': _kfilter(filt) if filt is not None else None,
|
|
}
|
|
|
|
|
|
def _tests_run_build(a):
|
|
argv = [sys.executable, os.path.join(BIN_DIR, 'box-ctl.py'), 'tests-run']
|
|
if a.get('test'):
|
|
argv.append(a['test'])
|
|
if a.get('filter'):
|
|
argv += ['--filter', a['filter']]
|
|
return argv
|
|
|
|
|
|
def _notify_send_validate(raw):
|
|
if not isinstance(raw, dict):
|
|
raise OpError('args must be an object')
|
|
allowed = {'agent', 'message', 'sidechat', 'sender'}
|
|
for k in raw:
|
|
if k not in allowed:
|
|
raise OpError(f'unknown arg: {k}')
|
|
msg = _clean_message(raw.get('message'))
|
|
if len(msg) > 1000:
|
|
raise OpError('message too long (max 1000)')
|
|
sidechat = raw.get('sidechat')
|
|
sender = raw.get('sender')
|
|
return {
|
|
'agent': _agent(raw.get('agent')),
|
|
'message': msg,
|
|
'sidechat': _sidechat_name(sidechat) if sidechat is not None else None,
|
|
'sender': _agent(sender) if sender is not None else None,
|
|
}
|
|
|
|
|
|
def _notify_send_build(a):
|
|
argv = [sys.executable, os.path.join(BIN_DIR, 'box-ctl.py'),
|
|
'notify', a['agent'], a['message']]
|
|
if a.get('sidechat'):
|
|
argv += ['--sidechat', a['sidechat']]
|
|
if a.get('sender'):
|
|
argv += ['--sender', a['sender']]
|
|
return argv
|
|
|
|
|
|
def _dm_ack_validate(raw):
|
|
if not isinstance(raw, dict):
|
|
raise OpError('args must be an object')
|
|
allowed = {'id', 'to', 'sender', 'sidechat', 'allow_main_chat'}
|
|
for k in raw:
|
|
if k not in allowed:
|
|
raise OpError(f'unknown arg: {k}')
|
|
ref_id = raw.get('id')
|
|
if not isinstance(ref_id, str) or not DM_ID_RE.fullmatch(ref_id):
|
|
raise OpError('id must be 6-64 hex chars')
|
|
sender = raw.get('sender')
|
|
if not sender:
|
|
raise OpError('sender is required')
|
|
sidechat = raw.get('sidechat')
|
|
return {
|
|
'id': ref_id,
|
|
'to': _agent(raw.get('to')),
|
|
'sender': _agent(sender),
|
|
'sidechat': _sidechat_name(sidechat) if sidechat is not None else None,
|
|
'allow_main_chat': bool(raw.get('allow_main_chat', False)),
|
|
}
|
|
|
|
|
|
def _dm_ack_build(a):
|
|
argv = [sys.executable, os.path.join(BIN_DIR, 'box-ctl.py'),
|
|
'ack', a['id'], '--to', a['to'], '--sender', a['sender']]
|
|
if a.get('sidechat'):
|
|
argv += ['--sidechat', a['sidechat']]
|
|
if a.get('allow_main_chat'):
|
|
argv.append('--allow-main-chat')
|
|
return argv
|
|
|
|
|
|
def _job_put_validate(raw):
|
|
if not isinstance(raw, dict):
|
|
raise OpError('args must be an object')
|
|
allowed = {'name', 'definition'}
|
|
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 JOB_RE.fullmatch(name):
|
|
raise OpError('name must match ^[a-z0-9][a-z0-9-]{0,63}$')
|
|
definition = raw.get('definition')
|
|
if not isinstance(definition, dict):
|
|
raise OpError('definition must be an object')
|
|
if definition.get('name') != name:
|
|
raise OpError('definition name must match name')
|
|
# Full schema is enforced by box-ctl.py validate_job (single copy).
|
|
return {'name': name, 'definition': definition}
|
|
|
|
|
|
def _job_put_build(a):
|
|
return [sys.executable, os.path.join(BIN_DIR, 'box-ctl.py'),
|
|
'job-put', a['name']]
|
|
|
|
|
|
def _job_trigger_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 _job_trigger_build(a):
|
|
return [sys.executable, os.path.join(BIN_DIR, 'box-ctl.py'),
|
|
'job-trigger', a['name']]
|
|
|
|
|
|
def _job_chain_validate(raw):
|
|
if not isinstance(raw, dict):
|
|
raise OpError('args must be an object')
|
|
allowed = {'from', 'to', 'on_failure'}
|
|
for k in raw:
|
|
if k not in allowed:
|
|
raise OpError(f'unknown arg: {k}')
|
|
# Self-chain and cycle guards live in box-ctl.py (single copy).
|
|
return {
|
|
'from': _job_name(raw.get('from')),
|
|
'to': _job_name(raw.get('to')),
|
|
'on_failure': bool(raw.get('on_failure', False)),
|
|
}
|
|
|
|
|
|
def _job_chain_build(a):
|
|
argv = [sys.executable, os.path.join(BIN_DIR, 'box-ctl.py'),
|
|
'job-chain', a['from'], a['to']]
|
|
if a.get('on_failure'):
|
|
argv.append('--on-failure')
|
|
return argv
|
|
|
|
|
|
def _job_next_validate(raw):
|
|
if not isinstance(raw, dict):
|
|
raise OpError('args must be an object')
|
|
allowed = {'job_id', 'success'}
|
|
for k in raw:
|
|
if k not in allowed:
|
|
raise OpError(f'unknown arg: {k}')
|
|
job_id = raw.get('job_id')
|
|
if not isinstance(job_id, str) or not JOB_DISPATCH_ID_RE.fullmatch(job_id):
|
|
raise OpError('job_id must look like <name>-YYYYMMDD-HHMMSS-<8hex>')
|
|
success = raw.get('success')
|
|
if success is not None and not isinstance(success, bool):
|
|
raise OpError('success must be a boolean')
|
|
return {'job_id': job_id, 'success': success}
|
|
|
|
|
|
def _job_next_build(a):
|
|
argv = [sys.executable, os.path.join(BIN_DIR, 'box-ctl.py'),
|
|
'job-next', a['job_id']]
|
|
if a.get('success') is True:
|
|
argv.append('--success')
|
|
elif a.get('success') is False:
|
|
argv.append('--fail')
|
|
return argv
|
|
|
|
|
|
def _timer_stop_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 _timer_stop_build(a):
|
|
return [sys.executable, os.path.join(BIN_DIR, 'box-ctl.py'),
|
|
'timer-stop', a['name']]
|
|
|
|
|
|
def _timer_disable_build(a):
|
|
return [sys.executable, os.path.join(BIN_DIR, 'box-ctl.py'),
|
|
'timer-disable', a['name']]
|
|
|
|
|
|
def _loop_remediate_validate(raw):
|
|
if not isinstance(raw, dict):
|
|
raise OpError('args must be an object')
|
|
allowed = {'dry_run'}
|
|
for k in raw:
|
|
if k not in allowed:
|
|
raise OpError(f'unknown arg: {k}')
|
|
return {'dry_run': bool(raw.get('dry_run', False))}
|
|
|
|
|
|
def _loop_remediate_build(a):
|
|
argv = [sys.executable, os.path.join(BIN_DIR, 'box-ctl.py'),
|
|
'loop-remediate']
|
|
if a.get('dry_run'):
|
|
argv.append('--dry-run')
|
|
return argv
|
|
|
|
|
|
def _loop_resolve_validate(raw):
|
|
if not isinstance(raw, dict):
|
|
raise OpError('args must be an object')
|
|
allowed = {'dm_id', 'note'}
|
|
for k in raw:
|
|
if k not in allowed:
|
|
raise OpError(f'unknown arg: {k}')
|
|
dm_id = raw.get('dm_id')
|
|
if not isinstance(dm_id, str) or not DM_ID_RE.fullmatch(dm_id):
|
|
raise OpError('dm_id must be 6-64 hex chars')
|
|
note = raw.get('note')
|
|
return {
|
|
'dm_id': dm_id,
|
|
'note': _clean_message(note) if note is not None else None,
|
|
}
|
|
|
|
|
|
def _loop_resolve_build(a):
|
|
argv = [sys.executable, os.path.join(BIN_DIR, 'box-ctl.py'),
|
|
'loop-resolve', a['dm_id']]
|
|
if a.get('note'):
|
|
argv.append(a['note'])
|
|
return argv
|
|
|
|
|
|
def _strat_type(v):
|
|
if not isinstance(v, str) or v.lower() not in STRAT_TYPES:
|
|
raise OpError('type must be one of %s' % sorted(STRAT_TYPES))
|
|
return v.lower()
|
|
|
|
|
|
def _strat_common(raw, allowed):
|
|
for k in raw:
|
|
if k not in allowed:
|
|
raise OpError(f'unknown arg: {k}')
|
|
subtype = raw.get('subtype')
|
|
if subtype is not None and \
|
|
(not isinstance(subtype, str) or not SUBTYPE_RE.fullmatch(subtype)):
|
|
raise OpError('subtype must match ^[A-Za-z0-9_.-]{1,64}$')
|
|
agent = raw.get('agent')
|
|
return subtype, _agent(agent) if agent is not None else None
|
|
|
|
|
|
def _strat_int(v, field):
|
|
if isinstance(v, bool) or not isinstance(v, int):
|
|
raise OpError(f'{field} must be an integer')
|
|
return v
|
|
|
|
|
|
def _strat_set_validate(raw):
|
|
if not isinstance(raw, dict):
|
|
raise OpError('args must be an object')
|
|
allowed = {'type', 'subtype', 'agent', 'track', 'priority',
|
|
'timeout_s', 'nudges', 'escalate'}
|
|
subtype, agent = _strat_common(raw, allowed)
|
|
track = raw.get('track')
|
|
if track is not None and not isinstance(track, bool):
|
|
raise OpError('track must be a boolean')
|
|
priority = raw.get('priority')
|
|
if priority is not None:
|
|
if not isinstance(priority, str) or \
|
|
priority.lower() not in STRAT_PRIORITIES:
|
|
raise OpError('priority must be one of %s'
|
|
% sorted(STRAT_PRIORITIES))
|
|
priority = priority.lower()
|
|
timeout_s = raw.get('timeout_s')
|
|
if timeout_s is not None:
|
|
timeout_s = _strat_int(timeout_s, 'timeout_s')
|
|
nudges = raw.get('nudges')
|
|
if nudges is not None:
|
|
nudges = _strat_int(nudges, 'nudges')
|
|
escalate = raw.get('escalate')
|
|
if escalate is not None:
|
|
if not isinstance(escalate, str) or not escalate.strip() or \
|
|
len(escalate) > 64 or \
|
|
any(ord(c) < 32 or ord(c) == 127 for c in escalate):
|
|
raise OpError('escalate must be 1-64 chars, no controls')
|
|
return {
|
|
'type': _strat_type(raw.get('type')),
|
|
'subtype': subtype, 'agent': agent, 'track': track,
|
|
'priority': priority, 'timeout_s': timeout_s, 'nudges': nudges,
|
|
'escalate': escalate,
|
|
}
|
|
|
|
|
|
def _strat_set_build(a):
|
|
payload = {}
|
|
for k in ('subtype', 'agent', 'track', 'priority', 'timeout_s',
|
|
'nudges', 'escalate'):
|
|
if a.get(k) is not None:
|
|
payload[k] = a[k]
|
|
return [sys.executable, os.path.join(BIN_DIR, 'box-ctl.py'),
|
|
'strat-set', a['type'], json.dumps(payload)]
|
|
|
|
|
|
def _strat_reset_validate(raw):
|
|
if not isinstance(raw, dict):
|
|
raise OpError('args must be an object')
|
|
allowed = {'type', 'subtype', 'agent'}
|
|
subtype, agent = _strat_common(raw, allowed)
|
|
return {'type': _strat_type(raw.get('type')),
|
|
'subtype': subtype, 'agent': agent}
|
|
|
|
|
|
def _strat_reset_build(a):
|
|
argv = [sys.executable, os.path.join(BIN_DIR, 'box-ctl.py'),
|
|
'strat-reset', a['type']]
|
|
if a.get('subtype'):
|
|
argv.append(a['subtype'])
|
|
if a.get('agent'):
|
|
argv += ['--agent', a['agent']]
|
|
return argv
|
|
|
|
|
|
def _vars_name_validate(raw, allowed):
|
|
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
|
|
|
|
|
|
def _vars_reset_validate(raw):
|
|
if not isinstance(raw, dict):
|
|
raise OpError('args must be an object')
|
|
return {'name': _vars_name_validate(raw, {'name'})}
|
|
|
|
|
|
def _vars_reset_build(a):
|
|
return [sys.executable, os.path.join(BIN_DIR, 'box-ctl.py'),
|
|
'vars-reset', a['name']]
|
|
|
|
|
|
def _vars_rollback_validate(raw):
|
|
if not isinstance(raw, dict):
|
|
raise OpError('args must be an object')
|
|
name = _vars_name_validate(raw, {'name', 'revision'})
|
|
revision = raw.get('revision')
|
|
if revision is not None:
|
|
if isinstance(revision, bool):
|
|
raise OpError('revision must be an int step or timestamp')
|
|
if isinstance(revision, int):
|
|
if revision < 1:
|
|
raise OpError('revision step must be >= 1')
|
|
elif not isinstance(revision, str) or not revision.strip() or \
|
|
len(revision) > 64 or \
|
|
any(ord(c) < 32 or ord(c) == 127 for c in revision):
|
|
raise OpError('revision must be an int step or timestamp')
|
|
return {'name': name, 'revision': revision}
|
|
|
|
|
|
def _vars_rollback_build(a):
|
|
argv = [sys.executable, os.path.join(BIN_DIR, 'box-ctl.py'),
|
|
'vars-rollback', a['name']]
|
|
if a.get('revision') is not None:
|
|
argv.append(str(a['revision']))
|
|
return argv
|
|
|
|
|
|
def _md_check_keys(raw, allowed):
|
|
for k in raw:
|
|
if k not in allowed:
|
|
raise OpError(f'unknown arg: {k}')
|
|
|
|
|
|
def _md_account(v):
|
|
if not isinstance(v, str) or not MD_ACCOUNT_RE.fullmatch(v):
|
|
raise OpError('account must match ^[A-Za-z0-9][A-Za-z0-9_-]{0,31}$')
|
|
return v
|
|
|
|
|
|
def _md_filename(v, template_only=False):
|
|
if template_only:
|
|
if v not in MD_TEMPLATE_FILES:
|
|
raise OpError('filename must be one of %s'
|
|
% sorted(MD_TEMPLATE_FILES))
|
|
return v
|
|
if not isinstance(v, str) or v in ('.', '..') \
|
|
or not MD_FILENAME_RE.fullmatch(v):
|
|
raise OpError('filename must be a plain basename (no directories)')
|
|
return v
|
|
|
|
|
|
def _md_subpath(v):
|
|
if v in (None, ''):
|
|
return ''
|
|
if not isinstance(v, str) or not MD_SUBPATH_RE.fullmatch(v) \
|
|
or '..' in v.split('/'):
|
|
raise OpError('path must be a subdir without ..')
|
|
return v
|
|
|
|
|
|
def _md_audit_validate(raw):
|
|
if not isinstance(raw, dict):
|
|
raise OpError('args must be an object')
|
|
_md_check_keys(raw, {'accounts'})
|
|
accounts = raw.get('accounts')
|
|
if accounts is None:
|
|
return {'accounts': None}
|
|
if not isinstance(accounts, list) or not accounts:
|
|
raise OpError('accounts must be a non-empty list')
|
|
return {'accounts': [_md_account(a) for a in accounts]}
|
|
|
|
|
|
def _md_audit_build(a):
|
|
argv = [sys.executable, os.path.join(BIN_DIR, 'box-ctl.py'), 'md-audit']
|
|
if a.get('accounts'):
|
|
argv += a['accounts']
|
|
return argv
|
|
|
|
|
|
def _md_list_validate(raw):
|
|
if not isinstance(raw, dict):
|
|
raise OpError('args must be an object')
|
|
_md_check_keys(raw, {'account', 'path'})
|
|
return {'account': _md_account(raw.get('account')),
|
|
'path': _md_subpath(raw.get('path'))}
|
|
|
|
|
|
def _md_list_build(a):
|
|
argv = [sys.executable, os.path.join(BIN_DIR, 'box-ctl.py'),
|
|
'md-list', a['account']]
|
|
if a.get('path'):
|
|
argv.append(a['path'])
|
|
return argv
|
|
|
|
|
|
def _md_read_validate(raw):
|
|
if not isinstance(raw, dict):
|
|
raise OpError('args must be an object')
|
|
_md_check_keys(raw, {'account', 'filename'})
|
|
return {'account': _md_account(raw.get('account')),
|
|
'filename': _md_filename(raw.get('filename'))}
|
|
|
|
|
|
def _md_read_build(a):
|
|
return [sys.executable, os.path.join(BIN_DIR, 'box-ctl.py'),
|
|
'md-read', a['account'], a['filename']]
|
|
|
|
|
|
def _md_diff_validate(raw):
|
|
if not isinstance(raw, dict):
|
|
raise OpError('args must be an object')
|
|
_md_check_keys(raw, {'account', 'filename'})
|
|
return {'account': _md_account(raw.get('account')),
|
|
'filename': _md_filename(raw.get('filename'), template_only=True)}
|
|
|
|
|
|
def _md_diff_build(a):
|
|
return [sys.executable, os.path.join(BIN_DIR, 'box-ctl.py'),
|
|
'md-diff', a['account'], a['filename']]
|
|
|
|
|
|
def _md_pull_validate(raw):
|
|
if not isinstance(raw, dict):
|
|
raise OpError('args must be an object')
|
|
_md_check_keys(raw, {'account', 'filename'})
|
|
return {'account': _md_account(raw.get('account')),
|
|
'filename': _md_filename(raw.get('filename'), template_only=True)}
|
|
|
|
|
|
def _md_pull_build(a):
|
|
return [sys.executable, os.path.join(BIN_DIR, 'box-ctl.py'),
|
|
'md-pull', a['account'], a['filename']]
|
|
|
|
|
|
def _md_inject_drive_validate(raw):
|
|
if not isinstance(raw, dict):
|
|
raise OpError('args must be an object')
|
|
_md_check_keys(raw, {'account'})
|
|
return {'account': _md_account(raw.get('account'))}
|
|
|
|
|
|
def _md_inject_drive_build(a):
|
|
return [sys.executable, os.path.join(BIN_DIR, 'box-ctl.py'),
|
|
'md-inject-drive', a['account']]
|
|
|
|
|
|
def _md_sync_all_validate(raw):
|
|
if not isinstance(raw, dict):
|
|
raise OpError('args must be an object')
|
|
_md_check_keys(raw, set())
|
|
return {}
|
|
|
|
|
|
def _md_sync_all_build(a):
|
|
return [sys.executable, os.path.join(BIN_DIR, 'box-ctl.py'), 'md-sync-all']
|
|
|
|
|
|
def _md_author(v):
|
|
if not isinstance(v, str) or not v.strip() or len(v) > 64 or \
|
|
any(ord(c) < 32 or ord(c) == 127 for c in v):
|
|
raise OpError('author must be 1-64 chars, no controls')
|
|
return v
|
|
|
|
|
|
def _md_reason(v):
|
|
if not isinstance(v, str) or len(v) > 256 or \
|
|
any(ord(c) < 32 or ord(c) == 127 for c in v):
|
|
raise OpError('reason must be 0-256 chars, no controls')
|
|
return v
|
|
|
|
|
|
def _md_amend_validate(raw):
|
|
if not isinstance(raw, dict):
|
|
raise OpError('args must be an object')
|
|
_md_check_keys(raw, {'filename', 'content', 'author', 'reason'})
|
|
content = raw.get('content')
|
|
if not isinstance(content, str) or not content.strip() \
|
|
or len(content) > MD_MAX_AMEND:
|
|
raise OpError('content must be 1-%d chars' % MD_MAX_AMEND)
|
|
return {'filename': _md_filename(raw.get('filename'), template_only=True),
|
|
'content': content,
|
|
'author': _md_author(raw.get('author', 'operator')),
|
|
'reason': _md_reason(raw.get('reason', ''))}
|
|
|
|
|
|
def _md_amend_build(a):
|
|
argv = [sys.executable, os.path.join(BIN_DIR, 'box-ctl.py'),
|
|
'md-amend', a['filename'], '--stdin',
|
|
'--author', a['author']]
|
|
if a.get('reason'):
|
|
argv += ['--reason', a['reason']]
|
|
return argv
|
|
|
|
|
|
def _md_append_validate(raw):
|
|
if not isinstance(raw, dict):
|
|
raise OpError('args must be an object')
|
|
_md_check_keys(raw, {'filename', 'text', 'author', 'section'})
|
|
text = raw.get('text')
|
|
if not isinstance(text, str) or not text.strip() \
|
|
or len(text) > MD_MAX_APPEND:
|
|
raise OpError('text must be 1-%d chars' % MD_MAX_APPEND)
|
|
section = raw.get('section')
|
|
if section is not None:
|
|
if not isinstance(section, str) or not section.strip() \
|
|
or len(section) > 128 or \
|
|
any(ord(c) < 32 or ord(c) == 127 for c in section):
|
|
raise OpError('section must be 1-128 chars, no controls')
|
|
return {'filename': _md_filename(raw.get('filename'), template_only=True),
|
|
'text': text,
|
|
'author': _md_author(raw.get('author', 'operator')),
|
|
'section': section}
|
|
|
|
|
|
def _md_append_build(a):
|
|
argv = [sys.executable, os.path.join(BIN_DIR, 'box-ctl.py'),
|
|
'md-append', a['filename'], '--stdin',
|
|
'--author', a['author']]
|
|
if a.get('section'):
|
|
argv += ['--section', a['section']]
|
|
return argv
|
|
|
|
|
|
def _approval_keys(raw, allowed):
|
|
for k in raw:
|
|
if k not in allowed:
|
|
raise OpError(f'unknown arg: {k}')
|
|
|
|
|
|
def _approval_opt_node(raw):
|
|
node = raw.get('node')
|
|
if node is None:
|
|
return None
|
|
return _agent(node)
|
|
|
|
|
|
def _approval_main_chat(raw):
|
|
v = raw.get('allow_main_chat', False)
|
|
if not isinstance(v, bool):
|
|
raise OpError('allow_main_chat must be a boolean')
|
|
return v
|
|
|
|
|
|
def _approval_check_validate(raw):
|
|
if not isinstance(raw, dict):
|
|
raise OpError('args must be an object')
|
|
_approval_keys(raw, {'node'})
|
|
return {'node': _approval_opt_node(raw)}
|
|
|
|
|
|
def _approval_check_build(a):
|
|
argv = [sys.executable, os.path.join(BIN_DIR, 'box-ctl.py'),
|
|
'approval-check']
|
|
if a.get('node'):
|
|
argv.append(a['node'])
|
|
return argv
|
|
|
|
|
|
def _approval_deny_validate(raw):
|
|
if not isinstance(raw, dict):
|
|
raise OpError('args must be an object')
|
|
_approval_keys(raw, {'node', 'message', 'allow_main_chat'})
|
|
return {'node': _agent(raw.get('node')),
|
|
'message': _clean_message(raw.get('message')),
|
|
'allow_main_chat': _approval_main_chat(raw)}
|
|
|
|
|
|
def _approval_deny_build(a):
|
|
argv = [sys.executable, os.path.join(BIN_DIR, 'box-ctl.py'),
|
|
'approval-deny', a['node'], '--message', a['message']]
|
|
if a.get('allow_main_chat'):
|
|
argv.append('--allow-main-chat')
|
|
return argv
|
|
|
|
|
|
def _approval_auto_validate(raw):
|
|
if not isinstance(raw, dict):
|
|
raise OpError('args must be an object')
|
|
_approval_keys(raw, {'node'})
|
|
return {'node': _approval_opt_node(raw)}
|
|
|
|
|
|
def _approval_auto_build(a):
|
|
argv = [sys.executable, os.path.join(BIN_DIR, 'box-ctl.py'),
|
|
'approval-auto']
|
|
if a.get('node'):
|
|
argv.append(a['node'])
|
|
return argv
|
|
|
|
|
|
def _approval_allow_validate(raw):
|
|
if not isinstance(raw, dict):
|
|
raise OpError('args must be an object')
|
|
_approval_keys(raw, {'node', 'message', 'allow_main_chat'})
|
|
return {'node': _agent(raw.get('node')),
|
|
'message': _clean_message(raw.get('message')),
|
|
'allow_main_chat': _approval_main_chat(raw)}
|
|
|
|
|
|
def _approval_allow_build(a):
|
|
# One-shot only: --always/--force are never passed remotely.
|
|
argv = [sys.executable, os.path.join(BIN_DIR, 'box-ctl.py'),
|
|
'approval-allow', a['node'], '--message', a['message']]
|
|
if a.get('allow_main_chat'):
|
|
argv.append('--allow-main-chat')
|
|
return argv
|
|
|
|
|
|
def _tmux_tally_validate(raw):
|
|
if not isinstance(raw, dict):
|
|
raise OpError('args must be an object')
|
|
return {}
|
|
|
|
|
|
def _tmux_tally_build(a):
|
|
return [sys.executable, os.path.join(BIN_DIR, 'tmux_auto_approver.py'), 'tally', '--json']
|
|
|
|
|
|
def _tmux_auto_status_validate(raw):
|
|
if not isinstance(raw, dict):
|
|
raise OpError('args must be an object')
|
|
return {}
|
|
|
|
|
|
def _tmux_auto_status_build(a):
|
|
return [sys.executable, os.path.join(BIN_DIR, 'tmux_auto_approver.py'), 'status', '--json']
|
|
|
|
|
|
def _onboard_connects_validate(raw):
|
|
if not isinstance(raw, dict):
|
|
raise OpError('args must be an object')
|
|
return {}
|
|
|
|
|
|
def _onboard_connects_build(a):
|
|
return [sys.executable, os.path.join(BIN_DIR, 'onboard_pipeline.py'), 'connects', '--json']
|
|
|
|
|
|
def _stdin_body(op, clean):
|
|
"""Request body piped to the backend's stdin (or None).
|
|
|
|
Most ops pass everything via argv. box-sys-op.py file/web/service/
|
|
followup ops consume the validated envelope; box-ctl.py job-put
|
|
consumes the raw job definition (box-ctl matches definition.name
|
|
against the argv name itself); box-ctl.py md-amend/md-append consume
|
|
the raw markdown content (argv carries filename, flags, --stdin).
|
|
"""
|
|
if op == 'job.put':
|
|
return json.dumps(clean['definition'])
|
|
if op == 'md.amend':
|
|
return clean['content']
|
|
if op == 'md.append':
|
|
return clean['text']
|
|
if op.startswith(('files.', 'web.', 'service.', 'followup.')):
|
|
return json.dumps(clean)
|
|
return None
|
|
|
|
|
|
# box.exec allowlist: read-only box-ctl actions only (mirrors the read-only
|
|
# subset of box-ctl.py IDEMPOTENT_ACTIONS). Actions with side-effecting
|
|
# subverbs (main-loop enable/disable), path args (git-diff/git-log), or
|
|
# complex argv shapes (job-next, policy-*, quality-validate, job-result)
|
|
# are deliberately excluded.
|
|
BOX_EXEC_NOARG_ACTIONS = frozenset({
|
|
'fleet-status', 'watchdog-alerts', 'relay-health', 'cdp-latency',
|
|
'chrome-errors', 'identity-audit', 'timer-list', 'job-list',
|
|
'vars-list', 'strat-list', 'loop-status', 'loop-health', 'loop-breaks',
|
|
'thread-list', 'dm-log', 'unread', 'swarm-list', 'quality-check',
|
|
'git-status',
|
|
})
|
|
BOX_EXEC_ONEARG_ACTIONS = frozenset({
|
|
'job-get', 'job-status', 'timer-status', 'vars-get', 'vars-history',
|
|
'strat-get', 'swarm-status', 'swarm-results',
|
|
})
|
|
_BOX_ARG_RE = re.compile(r'^[A-Za-z0-9][A-Za-z0-9/_.-]{0,127}$')
|
|
|
|
|
|
def _box_exec_validate(raw):
|
|
if not isinstance(raw, dict):
|
|
raise OpError('args must be an object')
|
|
for k in raw:
|
|
if k not in {'action', 'arg', 'agent'}:
|
|
raise OpError(f'unknown arg: {k}')
|
|
action = raw.get('action')
|
|
if action in BOX_EXEC_NOARG_ACTIONS:
|
|
if raw.get('arg') is not None:
|
|
raise OpError(f'{action} takes no arg')
|
|
return {'action': action}
|
|
if action in BOX_EXEC_ONEARG_ACTIONS:
|
|
arg = raw.get('arg')
|
|
if arg is None:
|
|
return {'action': action}
|
|
if not isinstance(arg, str) or not _BOX_ARG_RE.fullmatch(arg):
|
|
raise OpError('arg must match safe token')
|
|
return {'action': action, 'arg': arg}
|
|
raise OpError(f'unknown or non-read-only box action: {action}')
|
|
|
|
|
|
def _box_exec_build(a):
|
|
argv = [sys.executable, os.path.join(BIN_DIR, 'box-ctl.py'), a['action']]
|
|
if a.get('arg'):
|
|
argv.append(a['arg'])
|
|
return argv
|
|
|
|
|
|
def _tools_list_validate(raw):
|
|
if not isinstance(raw, dict):
|
|
raise OpError('args must be an object')
|
|
for k in raw:
|
|
if k not in {'agent'}:
|
|
raise OpError(f'unknown arg: {k}')
|
|
return {}
|
|
|
|
|
|
def _tools_list_build(a):
|
|
return [sys.executable, os.path.join(BIN_DIR, 'exec-constrained.py'),
|
|
'--list-ops']
|
|
|
|
|
|
# op -> {validate, build, timeout, side_effecting, description}
|
|
OPS = {
|
|
'dm.send': {
|
|
'validate': _dm_send_validate, 'build': _dm_send_build,
|
|
'timeout': 120, 'side_effecting': True,
|
|
'desc': 'Send a DM via dm.py (verified delivery)',
|
|
},
|
|
'dm.thread': {
|
|
'validate': _dm_thread_validate, 'build': _dm_thread_build,
|
|
'timeout': 120, 'side_effecting': True,
|
|
'desc': 'Send a threaded DM via dm.py',
|
|
},
|
|
'dm.read': {
|
|
'validate': _dm_read_validate, 'build': _dm_read_build,
|
|
'timeout': 60, 'side_effecting': False,
|
|
'desc': 'Read recent DMs (read-only)',
|
|
},
|
|
'dm.log': {
|
|
'validate': _dm_log_validate, 'build': _dm_log_build,
|
|
'timeout': 30, 'side_effecting': False,
|
|
'desc': 'Tail the DM send log, optionally agent-scoped (read-only)',
|
|
},
|
|
'dm.ack': {
|
|
'validate': _dm_ack_validate, 'build': _dm_ack_build,
|
|
'timeout': 150, 'side_effecting': True,
|
|
'desc': 'Acknowledge a DM or work order as [ACK:id] (sidechat-first)',
|
|
},
|
|
'notify.send': {
|
|
'validate': _notify_send_validate, 'build': _notify_send_build,
|
|
'timeout': 150, 'side_effecting': True,
|
|
'desc': 'Notify an agent in its sidechat (sidechat-first)',
|
|
},
|
|
'job.run': {
|
|
'validate': _job_run_validate, 'build': _job_run_build,
|
|
'timeout': 300, 'side_effecting': True,
|
|
'desc': 'Run a job from the jobs directory',
|
|
},
|
|
'job.put': {
|
|
'validate': _job_put_validate, 'build': _job_put_build,
|
|
'timeout': 60, 'side_effecting': True,
|
|
'desc': 'Create or update a job definition (schema-validated, committed)',
|
|
},
|
|
'job.trigger': {
|
|
'validate': _job_trigger_validate, 'build': _job_trigger_build,
|
|
'timeout': 330, 'side_effecting': True,
|
|
'desc': 'Trigger a job dispatch now (audited JSON contract)',
|
|
},
|
|
'job.chain': {
|
|
'validate': _job_chain_validate, 'build': _job_chain_build,
|
|
'timeout': 60, 'side_effecting': True,
|
|
'desc': 'Wire chain_next (or on_failure) between two jobs',
|
|
},
|
|
'job.next': {
|
|
'validate': _job_next_validate, 'build': _job_next_build,
|
|
'timeout': 30, 'side_effecting': False,
|
|
'desc': 'Dry-run: what would dispatch next for a job id (read-only)',
|
|
},
|
|
'loop.remediate': {
|
|
'validate': _loop_remediate_validate, 'build': _loop_remediate_build,
|
|
'timeout': 300, 'side_effecting': True,
|
|
'desc': 'Auto-heal soft loop breaks (supports dry_run preview)',
|
|
},
|
|
'loop.resolve': {
|
|
'validate': _loop_resolve_validate, 'build': _loop_resolve_build,
|
|
'timeout': 30, 'side_effecting': True,
|
|
'desc': 'Mark a followup loop resolved with an optional note',
|
|
},
|
|
'strat.set': {
|
|
'validate': _strat_set_validate, 'build': _strat_set_build,
|
|
'timeout': 30, 'side_effecting': True,
|
|
'desc': 'Set a strategy override (type-scoped, validated)',
|
|
},
|
|
'strat.reset': {
|
|
'validate': _strat_reset_validate, 'build': _strat_reset_build,
|
|
'timeout': 30, 'side_effecting': True,
|
|
'desc': 'Reset a strategy override to built-in default',
|
|
},
|
|
'chat.messages': {
|
|
'validate': _chat_messages_validate, 'build': _chat_messages_build,
|
|
'timeout': 60, 'side_effecting': False,
|
|
'desc': 'Read recent chat messages (read-only)',
|
|
},
|
|
'chat.send': {
|
|
'validate': _chat_send_validate, 'build': _chat_send_build,
|
|
'timeout': 120, 'side_effecting': True,
|
|
'desc': 'Send a chat message via muse-chat-api.py',
|
|
},
|
|
'health.check': {
|
|
'validate': _health_validate, 'build': _health_build,
|
|
'timeout': 30, 'side_effecting': False,
|
|
'desc': 'Run fleet status health check (read-only)',
|
|
},
|
|
'fleet.unread': {
|
|
'validate': _fleet_unread_validate, 'build': _fleet_unread_build,
|
|
'timeout': 60, 'side_effecting': False,
|
|
'desc': 'Fleet unread/activity counts, optionally agent-scoped (read-only)',
|
|
},
|
|
'git.status': {
|
|
'validate': _health_validate,
|
|
'build': lambda a: [sys.executable, os.path.join(BIN_DIR, 'box-ctl.py'), 'git-status'],
|
|
'timeout': 30, 'side_effecting': False,
|
|
'desc': 'Git working-tree status short + branch (read-only)',
|
|
},
|
|
'git.diff': {
|
|
'validate': _git_diff_validate, 'build': _git_diff_build,
|
|
'timeout': 30, 'side_effecting': False,
|
|
'desc': 'Git diff, capped at 64KB, optional path/stat (read-only)',
|
|
},
|
|
'git.log': {
|
|
'validate': _git_log_validate, 'build': _git_log_build,
|
|
'timeout': 30, 'side_effecting': False,
|
|
'desc': 'Recent commits as sha/subject, optional path (read-only)',
|
|
},
|
|
'tests.run': {
|
|
'validate': _tests_run_validate, 'build': _tests_run_build,
|
|
'timeout': 600, 'side_effecting': True,
|
|
'desc': 'Run repo unit tests: one tests.<module> or full suite, optional -k filter',
|
|
},
|
|
'md.audit': {
|
|
'validate': _md_audit_validate, 'build': _md_audit_build,
|
|
'timeout': 300, 'side_effecting': False,
|
|
'desc': 'Audit agent .md files + drive score across fleet (read-only)',
|
|
},
|
|
'md.list': {
|
|
'validate': _md_list_validate, 'build': _md_list_build,
|
|
'timeout': 120, 'side_effecting': False,
|
|
'desc': 'List agent container files, capped at 200 entries (read-only)',
|
|
},
|
|
'md.read': {
|
|
'validate': _md_read_validate, 'build': _md_read_build,
|
|
'timeout': 120, 'side_effecting': False,
|
|
'desc': 'Read an agent container file, capped at 64KB (read-only)',
|
|
},
|
|
'md.diff': {
|
|
'validate': _md_diff_validate, 'build': _md_diff_build,
|
|
'timeout': 120, 'side_effecting': False,
|
|
'desc': 'Diff agent file against shared template, capped at 64KB (read-only)',
|
|
},
|
|
'md.pull': {
|
|
'validate': _md_pull_validate, 'build': _md_pull_build,
|
|
'timeout': 120, 'side_effecting': True,
|
|
'desc': 'Pull canonical shared template into an agent container',
|
|
},
|
|
'md.inject_drive': {
|
|
'validate': _md_inject_drive_validate, 'build': _md_inject_drive_build,
|
|
'timeout': 300, 'side_effecting': True,
|
|
'desc': 'Inject high-drive operator templates into one agent',
|
|
},
|
|
'md.sync_all': {
|
|
'validate': _md_sync_all_validate, 'build': _md_sync_all_build,
|
|
'timeout': 600, 'side_effecting': True,
|
|
'desc': 'Inject high-drive operator templates across all agents',
|
|
},
|
|
'md.amend': {
|
|
'validate': _md_amend_validate, 'build': _md_amend_build,
|
|
'timeout': 120, 'side_effecting': True,
|
|
'desc': 'Amend a shared operator template (validated, git-committed)',
|
|
},
|
|
'md.append': {
|
|
'validate': _md_append_validate, 'build': _md_append_build,
|
|
'timeout': 120, 'side_effecting': True,
|
|
'desc': 'Append a note to a shared operator template (git-committed)',
|
|
},
|
|
'approval.check': {
|
|
'validate': _approval_check_validate, 'build': _approval_check_build,
|
|
'timeout': 180, 'side_effecting': False,
|
|
'desc': 'Fleet browser approval status, optionally node-scoped (read-only)',
|
|
},
|
|
'approval.deny': {
|
|
'validate': _approval_deny_validate, 'build': _approval_deny_build,
|
|
'timeout': 120, 'side_effecting': True,
|
|
'desc': 'Deny a node pending approval (message required)',
|
|
},
|
|
'approval.auto': {
|
|
'validate': _approval_auto_validate, 'build': _approval_auto_build,
|
|
'timeout': 300, 'side_effecting': True,
|
|
'desc': 'Auto-allow TRUSTED non-key prompts, fleet or one node',
|
|
},
|
|
'approval.allow': {
|
|
'validate': _approval_allow_validate, 'build': _approval_allow_build,
|
|
'timeout': 120, 'side_effecting': True,
|
|
'desc': 'One-shot allow of a node pending approval (message required; no always/force)',
|
|
},
|
|
'tmux.tally': {
|
|
'validate': _tmux_tally_validate, 'build': _tmux_tally_build,
|
|
'timeout': 30, 'side_effecting': False,
|
|
'desc': 'Tally all tmux sessions, workers, and panes across fleet sockets (JSON, read-only)',
|
|
},
|
|
'tmux.auto_status': {
|
|
'validate': _tmux_auto_status_validate, 'build': _tmux_auto_status_build,
|
|
'timeout': 30, 'side_effecting': False,
|
|
'desc': 'Tmux auto-approval runtime status, toggle state, and rule inventory (JSON, read-only)',
|
|
},
|
|
'onboard.connects': {
|
|
'validate': _onboard_connects_validate, 'build': _onboard_connects_build,
|
|
'timeout': 30, 'side_effecting': False,
|
|
'desc': 'Consolidated active fleet and client onboard connects (JSON, read-only)',
|
|
},
|
|
'cdp.latency': {
|
|
'validate': _health_validate, 'build': _cdp_latency_build,
|
|
'timeout': 30, 'side_effecting': False,
|
|
'desc': 'Run CDP latency probe across fleet profiles (read-only)',
|
|
},
|
|
'chrome.errors': {
|
|
'validate': _health_validate, 'build': _chrome_errors_build,
|
|
'timeout': 30, 'side_effecting': False,
|
|
'desc': 'Scan chrome log files for FATAL/crash/OOM errors since watermark (read-only)',
|
|
},
|
|
'watchdog.alerts': {
|
|
'validate': _health_validate, 'build': _watchdog_alerts_build,
|
|
'timeout': 30, 'side_effecting': False,
|
|
'desc': 'Check for new browser relaunch watchdog alerts since watermark (read-only)',
|
|
},
|
|
'relay.health': {
|
|
'validate': _health_validate, 'build': _relay_health_build,
|
|
'timeout': 30, 'side_effecting': False,
|
|
'desc': 'Check systemd relay service status across all node ports (read-only)',
|
|
},
|
|
'quality.check': {
|
|
'validate': _quality_check_validate, 'build': _quality_check_build,
|
|
'timeout': 30, 'side_effecting': False,
|
|
'desc': 'Run box quality gates and self-diagnostics suite (read-only)',
|
|
},
|
|
'subagent.spawn': {
|
|
'validate': _subagent_spawn_validate, 'build': _subagent_spawn_build,
|
|
'timeout': 120, 'side_effecting': True,
|
|
'desc': 'Spawn an autonomous subagent session via muse-cli gateway',
|
|
},
|
|
'thread.list': {
|
|
'validate': _thread_list_validate, 'build': _thread_list_build,
|
|
'timeout': 60, 'side_effecting': False,
|
|
'desc': 'List threads/sessions for an agent',
|
|
},
|
|
'thread.view': {
|
|
'validate': _thread_view_validate, 'build': _thread_view_build,
|
|
'timeout': 60, 'side_effecting': False,
|
|
'desc': 'View messages in a thread/session',
|
|
},
|
|
'pipeline.run': {
|
|
'validate': _pipeline_run_validate, 'build': _pipeline_run_build,
|
|
'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',
|
|
},
|
|
'cron.timer_create': {
|
|
'validate': _cron_timer_create_validate, 'build': _cron_timer_create_build,
|
|
'timeout': 30, 'side_effecting': True,
|
|
'desc': 'Create systemd user timer unit for a job',
|
|
},
|
|
'cron.timer_start': {
|
|
'validate': _cron_timer_start_validate, 'build': _cron_timer_start_build,
|
|
'timeout': 30, 'side_effecting': True,
|
|
'desc': 'Enable and start systemd user timer unit for a job',
|
|
},
|
|
'cron.timer_stop': {
|
|
'validate': _timer_stop_validate, 'build': _timer_stop_build,
|
|
'timeout': 30, 'side_effecting': True,
|
|
'desc': 'Stop the systemd user timer unit for a job',
|
|
},
|
|
'cron.timer_disable': {
|
|
'validate': _timer_stop_validate, 'build': _timer_disable_build,
|
|
'timeout': 30, 'side_effecting': True,
|
|
'desc': 'Disable the systemd user timer unit for a 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',
|
|
},
|
|
'vars.reset': {
|
|
'validate': _vars_reset_validate, 'build': _vars_reset_build,
|
|
'timeout': 30, 'side_effecting': True,
|
|
'desc': 'Restore a loop variable to its default value',
|
|
},
|
|
'vars.rollback': {
|
|
'validate': _vars_rollback_validate, 'build': _vars_rollback_build,
|
|
'timeout': 30, 'side_effecting': True,
|
|
'desc': 'Roll a loop variable back to a previous revision',
|
|
},
|
|
'files.read': {
|
|
'validate': _files_read_validate, 'build': _files_read_build,
|
|
'timeout': 30, 'side_effecting': False,
|
|
'desc': 'Read sandboxed file within repository tree',
|
|
},
|
|
'files.write': {
|
|
'validate': _files_write_validate, 'build': _files_write_build,
|
|
'timeout': 30, 'side_effecting': True,
|
|
'desc': 'Write sandboxed file within repository tree',
|
|
},
|
|
'web.fetch': {
|
|
'validate': _web_fetch_validate, 'build': _web_fetch_build,
|
|
'timeout': 30, 'side_effecting': False,
|
|
'desc': 'Perform safe HTTP/HTTPS GET request with SSRF guard',
|
|
},
|
|
'service.status': {
|
|
'validate': _service_status_validate, 'build': _service_status_build,
|
|
'timeout': 30, 'side_effecting': False,
|
|
'desc': 'Inspect allowlisted fleet systemd service status',
|
|
},
|
|
'service.restart': {
|
|
'validate': _service_restart_validate, 'build': _service_restart_build,
|
|
'timeout': 30, 'side_effecting': True,
|
|
'desc': 'Restart allowlisted fleet systemd service',
|
|
},
|
|
'swarm.spawn': {
|
|
'validate': _swarm_spawn_validate, 'build': _swarm_spawn_build,
|
|
'timeout': 30, 'side_effecting': True,
|
|
'desc': 'Spawn autonomous subagent swarm slots managed by Box',
|
|
},
|
|
'swarm.status': {
|
|
'validate': _swarm_status_validate, 'build': _swarm_status_build,
|
|
'timeout': 30, 'side_effecting': False,
|
|
'desc': 'Inspect subagent swarm status and progress',
|
|
},
|
|
'swarm.list': {
|
|
'validate': _swarm_list_validate, 'build': _swarm_list_build,
|
|
'timeout': 30, 'side_effecting': False,
|
|
'desc': 'List all subagent swarms and their counts',
|
|
},
|
|
'swarm.results': {
|
|
'validate': _swarm_results_validate, 'build': _swarm_results_build,
|
|
'timeout': 30, 'side_effecting': False,
|
|
'desc': 'Retrieve results and summaries for a completed swarm',
|
|
},
|
|
'followup.create': {
|
|
'validate': _followup_create_validate, 'build': _followup_create_build,
|
|
'timeout': 30, 'side_effecting': True,
|
|
'desc': 'Schedule a delayed autonomous self-followup reminder via transient systemd timer',
|
|
},
|
|
'tmux.send': {
|
|
'validate': _tmux_send_validate, 'build': _tmux_send_build,
|
|
'timeout': 15, 'side_effecting': True,
|
|
'desc': 'Send keystrokes to a shared Muse tmux session (/tmp/tmux-muse.sock)',
|
|
},
|
|
'tmux.capture': {
|
|
'validate': _tmux_capture_validate, 'build': _tmux_capture_build,
|
|
'timeout': 15, 'side_effecting': False,
|
|
'desc': 'Capture pane output from a shared Muse tmux session (/tmp/tmux-muse.sock)',
|
|
},
|
|
'tmux.list': {
|
|
'validate': _tmux_list_validate, 'build': _tmux_list_build,
|
|
'timeout': 10, 'side_effecting': False,
|
|
'desc': 'List sessions on shared Muse tmux socket (/tmp/tmux-muse.sock)',
|
|
},
|
|
'tmux.new': {
|
|
'validate': _tmux_new_validate, 'build': _tmux_new_build,
|
|
'timeout': 15, 'side_effecting': True,
|
|
'desc': 'Create a new session on shared Muse tmux socket (/tmp/tmux-muse.sock)',
|
|
},
|
|
'tmux.kill': {
|
|
'validate': _tmux_kill_validate, 'build': _tmux_kill_build,
|
|
'timeout': 15, 'side_effecting': True,
|
|
'desc': 'Kill a session on shared Muse tmux socket (/tmp/tmux-muse.sock)',
|
|
},
|
|
'tmux.prune': {
|
|
'validate': _tmux_prune_validate, 'build': _tmux_prune_build,
|
|
'timeout': 15, 'side_effecting': True,
|
|
'desc': 'Reap stale unattached sessions inactive for >TTL (default 2h)',
|
|
},
|
|
'exec.ping': {
|
|
'validate': _health_validate,
|
|
'build': lambda a: ['/bin/echo', 'PONG'],
|
|
'timeout': 10, 'side_effecting': False,
|
|
'desc': 'Canary no-op for watchdogs',
|
|
},
|
|
'box.exec': {
|
|
'validate': _box_exec_validate, 'build': _box_exec_build,
|
|
'timeout': 60, 'side_effecting': False,
|
|
'desc': 'Call a read-only box-ctl action (fleet-status, dm-log, job-get, ...)',
|
|
},
|
|
'tools.list': {
|
|
'validate': _tools_list_validate, 'build': _tools_list_build,
|
|
'timeout': 15, 'side_effecting': False,
|
|
'desc': 'List all exec ops with descriptions (dynamic discovery)',
|
|
},
|
|
}
|
|
|
|
# identity -> set of ops. 'master' may invoke everything. Unknown identities
|
|
# get the read-only subset. Per-agent tokens inherit their agent name as the
|
|
# identity; tighten per agent here as needed.
|
|
# ---- read-only tolerance + hyphenated aliases used by job templates ----
|
|
_ARG_SYNONYMS = {'job': 'name', 'id': 'name', 'job_name': 'name', 'unit': 'name'}
|
|
|
|
|
|
def _tolerant_args(spec, args):
|
|
"""Read-only ops: drop unknown args (mapping common synonyms) instead of failing.
|
|
Returns (args, dropped_names)."""
|
|
if args is None:
|
|
return {}, []
|
|
if not isinstance(args, dict):
|
|
return args, []
|
|
args, dropped = dict(args), []
|
|
for _ in range(len(args) + 1):
|
|
try:
|
|
spec['validate'](args)
|
|
break
|
|
except OpError as e:
|
|
m = re.match(r'unknown arg: (\S+)', str(e))
|
|
if m and m.group(1) in args:
|
|
k = m.group(1)
|
|
v = args.pop(k)
|
|
syn = _ARG_SYNONYMS.get(k)
|
|
if syn and syn not in args:
|
|
args[syn] = v
|
|
else:
|
|
dropped.append(k)
|
|
continue
|
|
if 'takes no' in str(e):
|
|
dropped.extend(args)
|
|
args = {}
|
|
continue
|
|
break
|
|
return args, dropped
|
|
|
|
|
|
def _register_aliases():
|
|
alias = {'job-list': 'cron.runs', 'quality.validate': 'quality.check',
|
|
'vars-list': 'vars.list', 'fleet-status': 'health.check'}
|
|
for new, old in alias.items():
|
|
if new not in OPS and old in OPS:
|
|
OPS[new] = dict(OPS[old], desc=OPS[old]['desc'] + f' (alias of {old})')
|
|
def _ctl(*words):
|
|
return lambda a: [sys.executable, os.path.join(BIN_DIR, 'box-ctl.py'), *words]
|
|
def _noargs(raw):
|
|
return {}
|
|
for new, words, desc in (('loop-status', ('loop-status',), 'Fleet loop health (read-only)'),
|
|
('timer-list', ('timer', 'list'), 'List box timers (read-only)')):
|
|
if new not in OPS:
|
|
OPS[new] = {'validate': _noargs, 'build': _ctl(*words), 'timeout': 30,
|
|
'side_effecting': False, 'desc': desc}
|
|
|
|
|
|
_register_aliases()
|
|
|
|
PERMISSIONS = {
|
|
'master': set(OPS),
|
|
'operator-main': set(OPS),
|
|
'operator-646': set(OPS),
|
|
'operator-muse': set(OPS),
|
|
'operator-pip': set(OPS),
|
|
'operator-opm': set(OPS),
|
|
'operator-dev': set(OPS),
|
|
'operator-def': set(OPS),
|
|
'operator-super': set(OPS),
|
|
'super': set(OPS),
|
|
'646': set(OPS),
|
|
'pip': set(OPS),
|
|
'muse': set(OPS),
|
|
'opm': set(OPS),
|
|
'dev': set(OPS),
|
|
'def': set(OPS),
|
|
'exec-canary': {'exec.ping'},
|
|
}
|
|
DEFAULT_PERMS = {'dm.read', 'dm.log', 'chat.messages', 'health.check', 'fleet.unread',
|
|
'thread.list', 'thread.view', 'exec.ping',
|
|
'git.status', 'git.diff', 'git.log', 'job.next',
|
|
'md.audit', 'md.list', 'md.read', 'md.diff',
|
|
'approval.check', 'tmux.tally', 'tmux.auto_status', 'onboard.connects'}
|
|
|
|
|
|
def permitted(ident, op):
|
|
perms = PERMISSIONS.get(ident, DEFAULT_PERMS)
|
|
return op in perms
|
|
|
|
# ------------------------------------------------------- audit log
|
|
|
|
_audit_lock = threading.Lock()
|
|
_nonce_lock = threading.Lock()
|
|
|
|
|
|
def audit(entry):
|
|
"""Append-only JSONL audit. Never includes tokens or signatures."""
|
|
entry = dict(entry)
|
|
entry['ts'] = time.strftime('%Y-%m-%dT%H:%M:%SZ', time.gmtime())
|
|
line = json.dumps(entry, separators=(',', ':')) + '\n'
|
|
with _audit_lock:
|
|
try:
|
|
with open(AUDIT_LOG, 'a') as f:
|
|
f.write(line)
|
|
os.chmod(AUDIT_LOG, 0o600)
|
|
except OSError:
|
|
pass
|
|
|
|
|
|
# ------------------------------------------------------- HTTP handler
|
|
|
|
class Handler(BaseHTTPRequestHandler):
|
|
server_version = 'exec-constrained/1.0'
|
|
|
|
def log_message(self, format, *args):
|
|
sys.stderr.write(f'{self.client_address[0]} - {format % args}\n')
|
|
|
|
def _send(self, code, obj):
|
|
body = json.dumps(obj).encode()
|
|
self.send_response(code)
|
|
self.send_header('Content-Type', 'application/json')
|
|
self.send_header('Content-Length', str(len(body)))
|
|
self.end_headers()
|
|
self.wfile.write(body)
|
|
|
|
def do_GET(self):
|
|
if self.path == '/health':
|
|
self._send(200, {'status': 'ok', 'mode': 'constrained'})
|
|
elif self.path == '/ops':
|
|
# List ops (names + descriptions only; no arg details needed).
|
|
self._send(200, {'ops': [
|
|
{'op': k, 'desc': v['desc'],
|
|
'side_effecting': v['side_effecting']}
|
|
for k, v in sorted(OPS.items())]})
|
|
elif self.path in ('/box', '/box-relay.sh'):
|
|
relay_path = os.path.join(BIN_DIR, 'box-relay.sh')
|
|
if os.path.isfile(relay_path):
|
|
try:
|
|
with open(relay_path, 'rb') as f:
|
|
data = f.read()
|
|
self.send_response(200)
|
|
self.send_header('Content-Type', 'text/x-shellscript')
|
|
self.send_header('Content-Length', str(len(data)))
|
|
self.end_headers()
|
|
self.wfile.write(data)
|
|
return
|
|
except Exception:
|
|
pass
|
|
self._send(404, {'error': 'not found'})
|
|
else:
|
|
self._send(404, {'error': 'not found'})
|
|
|
|
def do_POST(self):
|
|
if self.path == '/exec':
|
|
self._handle_exec()
|
|
else:
|
|
self._send(404, {'error': 'not found'})
|
|
|
|
def _handle_exec(self):
|
|
peer = self.client_address[0]
|
|
try:
|
|
length = int(self.headers.get('Content-Length', 0))
|
|
except ValueError:
|
|
length = 0
|
|
if length <= 0 or length > MAX_BODY:
|
|
self._send(413 if length > MAX_BODY else 400,
|
|
{'error': 'bad content length'})
|
|
return
|
|
try:
|
|
data = json.loads(self.rfile.read(length))
|
|
except Exception:
|
|
self._send(400, {'error': 'invalid json'})
|
|
return
|
|
if not isinstance(data, dict):
|
|
self._send(400, {'error': 'body must be an object'})
|
|
return
|
|
|
|
# --- auth: Bearer <redacted> header (preferred) or legacy body token,
|
|
# --- or ssh-keygen signature envelope.
|
|
ident = None
|
|
op = data.get('op')
|
|
args = data.get('args', {})
|
|
authz = self.headers.get('Authorization', '')
|
|
if authz.startswith('Bearer '):
|
|
token = authz[7:].strip()
|
|
ident = check_token(token)
|
|
elif data.get('token'):
|
|
ident = check_token(data['token'])
|
|
elif data.get('identity') and data.get('payload') \
|
|
and data.get('signature'):
|
|
# Signature envelope carries op/args; ignore body's op/args.
|
|
ident = check_signature(data['identity'], data['payload'],
|
|
data['signature'])
|
|
if ident:
|
|
try:
|
|
env = json.loads(data['payload'])
|
|
op, args = env.get('op'), env.get('args', {})
|
|
except Exception:
|
|
op, args = None, {}
|
|
if not ident:
|
|
audit({'event': 'auth_failed', 'peer': peer,
|
|
'op': str(op)[:64]})
|
|
self._send(401, {'error': 'unauthorized'})
|
|
return
|
|
|
|
# --- rate limit
|
|
if not LIMITER.allow(ident):
|
|
audit({'event': 'rate_limited', 'ident': ident, 'peer': peer,
|
|
'op': str(op)[:64]})
|
|
self._send(429, {'error': 'rate limited'})
|
|
return
|
|
|
|
# --- op + permission + args
|
|
if not isinstance(op, str) or op not in OPS:
|
|
self._send(400, {'error': 'unknown op',
|
|
'hint': 'GET /ops for the allowlist'})
|
|
return
|
|
if not permitted(ident, op):
|
|
audit({'event': 'forbidden', 'ident': ident, 'peer': peer,
|
|
'op': op})
|
|
self._send(403, {'error': 'forbidden for this identity'})
|
|
return
|
|
dropped = []
|
|
try:
|
|
spec = OPS[op]
|
|
if not spec.get('side_effecting'):
|
|
args, dropped = _tolerant_args(spec, args)
|
|
clean = spec['validate'](args)
|
|
argv = spec['build'](clean)
|
|
except OpError as e:
|
|
self._send(400, {'error': f'bad args: {e}', 'op': op,
|
|
'desc': OPS[op].get('desc', ''),
|
|
'hint': 'GET /ops for the allowlist; see desc for usage'})
|
|
return
|
|
except Exception as e:
|
|
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)
|
|
audit({'event': 'exec_start', 'ident': ident, 'peer': peer,
|
|
'op': op,
|
|
'args_sha': hashlib.sha256(
|
|
json.dumps(clean, sort_keys=True).encode()
|
|
).hexdigest()[:16]})
|
|
t0 = time.monotonic()
|
|
try:
|
|
stdin_input = _stdin_body(op, clean)
|
|
p = subprocess.run(argv, input=stdin_input, capture_output=True, text=True,
|
|
timeout=spec['timeout'], cwd=WORK_DIR)
|
|
rc = p.returncode
|
|
out = p.stdout[-MAX_OUTPUT:]
|
|
err = p.stderr[-MAX_OUTPUT:]
|
|
ok = True
|
|
except subprocess.TimeoutExpired:
|
|
rc, out, err, ok = -1, '', 'timeout', False
|
|
except Exception as e:
|
|
rc, out, err, ok = -1, '', str(e)[:500], False
|
|
dt = round(time.monotonic() - t0, 2)
|
|
audit({'event': 'exec_done', 'ident': ident, 'peer': peer,
|
|
'op': op, 'rc': rc, 'duration_s': dt, 'ok': ok})
|
|
if ok:
|
|
resp = {'rc': rc, 'stdout': out, 'stderr': err, 'duration_s': dt}
|
|
if dropped:
|
|
resp['dropped_args'] = dropped
|
|
self._send(200, resp)
|
|
else:
|
|
self._send(500, {'error': err, 'rc': rc})
|
|
|
|
|
|
|
|
# ------------------------------------------------------- main
|
|
|
|
def main():
|
|
global TOKEN_FILE, TOKEN_DIR, SIGNERS_FILE, NONCE_FILE, AUDIT_LOG
|
|
global BIN_DIR, JOBS_DIR, WORK_DIR
|
|
ap = argparse.ArgumentParser()
|
|
ap.add_argument('--port', type=int, default=8444)
|
|
ap.add_argument('--host', default='127.0.0.1')
|
|
ap.add_argument('--token-file', default=TOKEN_FILE)
|
|
ap.add_argument('--token-dir', default=TOKEN_DIR)
|
|
ap.add_argument('--signers-file', default=SIGNERS_FILE)
|
|
ap.add_argument('--nonce-file', default=NONCE_FILE)
|
|
ap.add_argument('--audit-log', default=AUDIT_LOG)
|
|
ap.add_argument('--bin-dir', default=BIN_DIR)
|
|
ap.add_argument('--jobs-dir', default=JOBS_DIR)
|
|
ap.add_argument('--cert-file',
|
|
default='/home/super/.exec-constrained-cert.pem')
|
|
ap.add_argument('--key-file',
|
|
default='/home/super/.exec-constrained-key.pem')
|
|
ap.add_argument('--work-dir', default='/home/super')
|
|
ap.add_argument('--list-ops', action='store_true',
|
|
help='Print the op registry as JSON and exit (backs tools.list)')
|
|
args = ap.parse_args()
|
|
|
|
if args.list_ops:
|
|
print(json.dumps({
|
|
'ok': True,
|
|
'ops': [{'op': k, 'desc': v.get('desc', ''),
|
|
'side_effecting': bool(v.get('side_effecting', False)),
|
|
'timeout': v.get('timeout', 30)}
|
|
for k, v in sorted(OPS.items())],
|
|
}))
|
|
return
|
|
|
|
TOKEN_FILE, TOKEN_DIR = args.token_file, args.token_dir
|
|
SIGNERS_FILE, NONCE_FILE = args.signers_file, args.nonce_file
|
|
AUDIT_LOG = args.audit_log
|
|
BIN_DIR, JOBS_DIR = args.bin_dir, args.jobs_dir
|
|
WORK_DIR = args.work_dir
|
|
|
|
server = ThreadingHTTPServer((args.host, args.port), Handler)
|
|
|
|
cert_file = args.cert_file
|
|
key_file = args.key_file
|
|
if not os.path.exists(cert_file):
|
|
print('Generating self-signed cert...', file=sys.stderr)
|
|
subprocess.run([
|
|
'openssl', 'req', '-x509', '-newkey', 'rsa:2048',
|
|
'-keyout', key_file, '-out', cert_file,
|
|
'-days', '3650', '-nodes',
|
|
'-subj', '/CN=bl-exec-constrained',
|
|
], check=True, capture_output=True)
|
|
os.chmod(key_file, 0o600)
|
|
|
|
ctx = ssl.SSLContext(ssl.PROTOCOL_TLS_SERVER)
|
|
ctx.load_cert_chain(cert_file, key_file)
|
|
server.socket = ctx.wrap_socket(server.socket, server_side=True)
|
|
|
|
print(f'Constrained exec listening on https://{args.host}:{args.port}/exec',
|
|
file=sys.stderr)
|
|
print(f'Allowlist: {", ".join(sorted(OPS))}', file=sys.stderr)
|
|
try:
|
|
server.serve_forever()
|
|
except KeyboardInterrupt:
|
|
print('\nShutting down', file=sys.stderr)
|
|
|
|
|
|
if __name__ == '__main__':
|
|
main()
|