Files
box/bin/exec-constrained.py
T

1093 lines
37 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 HTTPServer, 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}$')
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=20, burst=5):
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):
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'],
'--limit', 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 _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 _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']
# 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)',
},
'job.run': {
'validate': _job_run_validate, 'build': _job_run_build,
'timeout': 300, 'side_effecting': True,
'desc': 'Run a job from the jobs directory',
},
'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)',
},
'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',
},
'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',
},
'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',
},
'exec.ping': {
'validate': _health_validate,
'build': lambda a: ['/bin/echo', 'PONG'],
'timeout': 10, 'side_effecting': False,
'desc': 'Canary no-op for watchdogs',
},
}
# 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.
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', 'chat.messages', 'health.check', 'thread.list', 'thread.view', 'exec.ping'}
def permitted(ident, op):
perms = PERMISSIONS.get(ident, DEFAULT_PERMS)
return op in perms
# ------------------------------------------------------- audit log
_audit_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
try:
spec = OPS[op]
clean = spec['validate'](args)
argv = spec['build'](clean)
except OpError as e:
self._send(400, {'error': f'bad args: {e}'})
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 = json.dumps(clean) if op.startswith(('files.', 'web.', 'service.')) else None
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:
self._send(200, {'rc': rc, 'stdout': out, 'stderr': err,
'duration_s': dt})
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')
args = ap.parse_args()
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 = HTTPServer((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()