File size: 10,146 Bytes
12496fc | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 201 202 203 204 205 206 207 208 209 210 211 212 213 214 215 216 217 | """Explicit capabilities and receipts. Host execution is NOT a security sandbox."""
from dataclasses import asdict, dataclass, field
from hashlib import sha256
from pathlib import Path
import json
import os
import subprocess
import tempfile
import threading
import time
import uuid
import jsonschema
def schema(properties, required):
return {"type": "object", "properties": properties, "required": required, "additionalProperties": False}
STR = {"type": "string", "minLength": 1}
SCHEMAS = {
"filesystem.read": schema({"path": STR}, ["path"]),
"filesystem.write": schema({"path": STR, "text": {"type": "string"}, "expected_sha256": STR}, ["path", "text"]),
"filesystem.list": schema({"path": STR}, ["path"]),
"shell.exec": schema({"command": STR}, ["command"]),
"git.diff": schema({}, []),
}
CLASSES = {"filesystem.read": {"READ"}, "filesystem.list": {"READ"},
"filesystem.write": {"WRITE"}, "shell.exec": {"EXECUTE"}, "git.diff": {"READ", "EXECUTE"}}
@dataclass
class Policy:
root: str
permissions: list[str] = field(default_factory=lambda: ["READ"])
commands: dict[str, list[str]] = field(default_factory=dict)
timeout_seconds: float = 30
output_limit: int = 16000
max_file_bytes: int = 1000000
allow_host_execution: bool = False
def __post_init__(self):
if self.timeout_seconds <= 0 or self.output_limit < 1 or self.max_file_bytes < 1:
raise ValueError("Policy limits must be positive")
if set(self.permissions) - {"READ", "WRITE", "EXECUTE", "NETWORK", "PRIVILEGED", "IRREVERSIBLE"}:
raise ValueError("Unknown permission class")
if any(not v or not all(isinstance(a, str) and a for a in v) for v in self.commands.values()):
raise ValueError("Commands must be nonempty argument arrays")
@dataclass
class Receipt:
id: str
tool: str
ok: bool
output: str
elapsed_ms: float
truncated: bool = False
exit_code: int | None = None
error: str | None = None
output_sha256: str | None = None
class Executor:
def __init__(self, policy: Policy, audit_path=None):
self.policy = policy
self.root = Path(policy.root).resolve(strict=True)
if not self.root.is_dir():
raise ValueError("Workspace must be a directory")
self.audit_path = Path(audit_path) if audit_path else None
self.cache = {}
self.lock = threading.RLock()
def available_tools(self):
return {name: spec for name, spec in SCHEMAS.items()
if CLASSES[name] <= set(self.policy.permissions)
and (name not in {"shell.exec", "git.diff"} or self.policy.allow_host_execution)
and (name != "shell.exec" or self.policy.commands)}
def path(self, name):
candidate = Path(name)
if candidate.is_absolute() or candidate.drive or ":" in name:
raise PermissionError("Only relative workspace paths are accepted")
target = (self.root / candidate).resolve()
if not target.is_relative_to(self.root):
raise PermissionError("Path escapes workspace")
rel = target.relative_to(self.root)
if any(p.lower() in {".git", ".ssh", ".aws", ".env", "private", ".cache"} or p.lower().startswith(".env.") for p in rel.parts):
raise PermissionError("Private or control path is excluded")
return target
def _run(self, argv):
if not self.policy.allow_host_execution:
raise PermissionError("Host execution is disabled; use an isolated workspace before enabling")
env = {k: v for k, v in os.environ.items() if k.upper() in {"PATH", "SYSTEMROOT", "WINDIR", "TEMP", "TMP", "PATHEXT"}}
# Drain continuously; bounded capture prevents output flooding from exhausting RAM.
flags = subprocess.CREATE_NEW_PROCESS_GROUP | subprocess.CREATE_NO_WINDOW if os.name == "nt" else 0
proc = subprocess.Popen(argv, cwd=self.root, env=env, stdin=subprocess.DEVNULL,
stdout=subprocess.PIPE, stderr=subprocess.STDOUT,
shell=False, creationflags=flags, start_new_session=os.name != "nt")
data, total = bytearray(), [0]
def drain():
while chunk := proc.stdout.read(4096):
total[0] += len(chunk)
remaining = self.policy.output_limit - len(data)
if remaining > 0:
data.extend(chunk[:remaining])
reader = threading.Thread(target=drain, daemon=True)
reader.start()
timed_out = False
try:
proc.wait(timeout=self.policy.timeout_seconds)
except subprocess.TimeoutExpired:
timed_out = True
if os.name == "nt":
subprocess.run(["taskkill", "/PID", str(proc.pid), "/T", "/F"], capture_output=True, timeout=10)
else:
import signal
os.killpg(proc.pid, signal.SIGKILL)
proc.wait(timeout=10)
reader.join(timeout=2)
if reader.is_alive():
# Detached descendants are not contained by this host runner.
raise RuntimeError("Output pipe remained open after command; isolated execution required")
proc.stdout.close()
return data.decode("utf-8", errors="replace"), proc.returncode, total[0] > len(data), timed_out
def execute(self, name, arguments, *, call_id=None):
with self.lock:
return self._execute(name, arguments, call_id=call_id)
def _execute(self, name, arguments, *, call_id=None):
call_id = call_id or uuid.uuid4().hex
fingerprint = sha256(json.dumps([name, arguments], sort_keys=True).encode()).hexdigest()
if call_id in self.cache:
old, receipt = self.cache[call_id]
if old != fingerprint:
raise ValueError("Idempotency key reused for different operation")
return receipt
start = time.perf_counter()
output, code, truncated, error, ok = "", None, False, None, False
try:
if name not in SCHEMAS:
raise ValueError("Unknown tool")
jsonschema.validate(arguments, SCHEMAS[name])
if not CLASSES[name] <= set(self.policy.permissions):
raise PermissionError("Required permission is disabled")
if name == "filesystem.read":
path = self.path(arguments["path"])
with path.open("rb") as f:
raw = f.read(self.policy.max_file_bytes + 1)
if len(raw) > self.policy.max_file_bytes:
raise ValueError("File exceeds configured read limit")
output = raw.decode("utf-8")
elif name == "filesystem.list":
path = self.path(arguments["path"])
names = []
for p in sorted(path.iterdir()):
try:
self.path(str(p.relative_to(self.root)))
names.append(p.name + ("/" if p.is_dir() else ""))
except PermissionError:
pass
if len(names) >= 1000:
truncated = True
break
output = json.dumps(names)
elif name == "filesystem.write":
path = self.path(arguments["path"])
raw = arguments["text"].encode("utf-8")
if len(raw) > self.policy.max_file_bytes:
raise ValueError("Write exceeds configured limit")
old = arguments.get("expected_sha256")
if path.exists():
if old is None or sha256(path.read_bytes()).hexdigest() != old:
raise ValueError("Existing files require matching expected_sha256")
elif old is not None:
raise ValueError("Expected an existing file")
path.parent.mkdir(parents=True, exist_ok=True)
fd, tmp = tempfile.mkstemp(dir=path.parent, prefix=".nexora-")
try:
with os.fdopen(fd, "wb") as f:
f.write(raw)
f.flush()
os.fsync(f.fileno())
os.replace(tmp, path)
finally:
if os.path.exists(tmp):
os.unlink(tmp)
output = json.dumps({"path": arguments["path"], "sha256": sha256(path.read_bytes()).hexdigest(), "bytes": len(raw)})
else:
if name == "git.diff":
argv = ["git", "--no-pager", "diff", "--no-ext-diff", "--no-textconv"]
else:
key = arguments["command"]
if key not in self.policy.commands:
raise PermissionError("Command is not configured by the owner")
argv = self.policy.commands[key]
output, code, truncated, timed_out = self._run(argv)
if timed_out:
error = "Command timed out"
elif code != 0:
error = f"Command exited {code}"
ok = error is None
except Exception as exc:
error = f"{type(exc).__name__}: {str(exc)[:500]}"
digest = sha256(output.encode()).hexdigest()
truncated = truncated or len(output) > self.policy.output_limit
receipt = Receipt(call_id, name, ok, output[:self.policy.output_limit],
(time.perf_counter()-start)*1000, truncated, code, error, digest)
if self.audit_path:
self.audit_path.parent.mkdir(parents=True, exist_ok=True)
# Do not persist tool arguments/content in the default audit log.
with self.audit_path.open("a", encoding="utf-8") as f:
f.write(json.dumps({**asdict(receipt), "output": "[omitted]", "error": None if not error else error.split(":")[0], "request_sha256": fingerprint}) + "\n")
self.cache[call_id] = (fingerprint, receipt)
return receipt
|