"""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