NEXORA / nexora /tools.py
devildasdf's picture
Release validated NEXORA research prototype, tiny weights and evidence
12496fc verified
Raw History Blame Contribute Delete
10.1 kB
"""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