|
| 1 | +"""Sandbox worker — runs ONE untrusted mock's code in a hardened subprocess. |
| 2 | +
|
| 3 | +Launched by :mod:`mockworld.sandbox` as ``python -m mockworld._sandbox_worker <dir>``. |
| 4 | +Everything untrusted (importing ``handlers.py``/``seed.py`` and running them) happens |
| 5 | +here, never in the parent. Before any untrusted import we neuter the network, |
| 6 | +subprocess/exec, ctypes, and file writes, and apply CPU/memory limits. |
| 7 | +
|
| 8 | +This is defense-in-depth, not a formal guarantee: in-process Python can't be made |
| 9 | +perfectly escape-proof. It removes the easy paths (exfiltration, spawning |
| 10 | +processes, trashing files) and contains crashes/hangs to a disposable child that |
| 11 | +holds none of the parent's state. For hard isolation, run mockworld in a container. |
| 12 | +
|
| 13 | +Protocol: length-prefixed JSON (4-byte big-endian length + UTF-8 body) on a private |
| 14 | +fd duped from stdout; the child's own stdout is redirected to stderr so handler |
| 15 | +`print()`s can't corrupt the stream. |
| 16 | +""" |
| 17 | + |
| 18 | +from __future__ import annotations |
| 19 | + |
| 20 | +import builtins |
| 21 | +import json |
| 22 | +import os |
| 23 | +import struct |
| 24 | +import sys |
| 25 | +from pathlib import Path |
| 26 | + |
| 27 | + |
| 28 | +def _harden() -> None: |
| 29 | + """Neuter dangerous capabilities on already-imported modules, then lock limits.""" |
| 30 | + def blocked(*_a, **_k): |
| 31 | + raise PermissionError("blocked in mockworld sandbox") |
| 32 | + |
| 33 | + import socket |
| 34 | + socket.socket = blocked |
| 35 | + socket.create_connection = blocked |
| 36 | + socket.create_server = blocked |
| 37 | + |
| 38 | + import subprocess |
| 39 | + for fn in ("Popen", "run", "call", "check_call", "check_output", "getoutput", "getstatusoutput"): |
| 40 | + if hasattr(subprocess, fn): |
| 41 | + setattr(subprocess, fn, blocked) |
| 42 | + |
| 43 | + for fn in ("system", "popen", "fork", "forkpty", "exec", "execv", "execve", "execvp", |
| 44 | + "execvpe", "execl", "execle", "execlp", "execlpe", "spawnv", "spawnve", |
| 45 | + "spawnl", "spawnlp", "remove", "unlink", "rmdir", "removedirs", "rename", |
| 46 | + "replace", "truncate", "kill", "killpg"): |
| 47 | + if hasattr(os, fn): |
| 48 | + setattr(os, fn, blocked) |
| 49 | + |
| 50 | + try: |
| 51 | + import ctypes |
| 52 | + for fn in ("CDLL", "PyDLL", "WinDLL", "OleDLL", "cdll", "pydll", "windll"): |
| 53 | + if hasattr(ctypes, fn): |
| 54 | + setattr(ctypes, fn, blocked) |
| 55 | + except Exception: |
| 56 | + pass |
| 57 | + |
| 58 | + real_open = builtins.open |
| 59 | + |
| 60 | + def safe_open(file, mode="r", *a, **k): |
| 61 | + if any(c in mode for c in "wax+"): |
| 62 | + raise PermissionError("file writes are blocked in the mockworld sandbox") |
| 63 | + return real_open(file, mode, *a, **k) |
| 64 | + |
| 65 | + builtins.open = safe_open |
| 66 | + |
| 67 | + try: |
| 68 | + import resource |
| 69 | + resource.setrlimit(resource.RLIMIT_CPU, (5, 6)) # ~5 CPU-seconds per worker |
| 70 | + mem = 512 * 1024 * 1024 |
| 71 | + for lim in ("RLIMIT_AS", "RLIMIT_DATA"): |
| 72 | + if hasattr(resource, lim): |
| 73 | + try: |
| 74 | + resource.setrlimit(getattr(resource, lim), (mem, mem)) |
| 75 | + except (ValueError, OSError): |
| 76 | + pass |
| 77 | + except Exception: |
| 78 | + pass |
| 79 | + |
| 80 | + |
| 81 | +def _read_msg(stream) -> dict | None: |
| 82 | + header = stream.read(4) |
| 83 | + if len(header) < 4: |
| 84 | + return None |
| 85 | + (length,) = struct.unpack(">I", header) |
| 86 | + return json.loads(stream.read(length)) |
| 87 | + |
| 88 | + |
| 89 | +def _write_msg(stream, obj: dict) -> None: |
| 90 | + body = json.dumps(obj).encode() |
| 91 | + stream.write(struct.pack(">I", len(body))) |
| 92 | + stream.write(body) |
| 93 | + stream.flush() |
| 94 | + |
| 95 | + |
| 96 | +def _result_payload(result) -> dict: |
| 97 | + if result.success: |
| 98 | + return {"success": True, "data": result.data, "meta": result.meta} |
| 99 | + err = result.err |
| 100 | + return {"success": False, "error": { |
| 101 | + "code": err.code, "message": err.message, "http_status": err.http_status, |
| 102 | + "body": err.body, "retry_after_s": err.retry_after_s}} |
| 103 | + |
| 104 | + |
| 105 | +def main() -> None: |
| 106 | + mock_dir = Path(sys.argv[1]) |
| 107 | + |
| 108 | + # Private protocol channel; keep the real stdout clean. |
| 109 | + proto_out = os.fdopen(os.dup(1), "wb") |
| 110 | + inp = sys.stdin.buffer |
| 111 | + sys.stdout = sys.stderr |
| 112 | + |
| 113 | + _harden() # <-- everything below runs under the neutered environment |
| 114 | + |
| 115 | + import yaml |
| 116 | + |
| 117 | + from mockworld.datagen import DataGen |
| 118 | + from mockworld.determinism import DeterministicContext |
| 119 | + from mockworld.errors import register_error |
| 120 | + from mockworld.handler_ctx import FaultHelper, HandlerCtx |
| 121 | + from mockworld.loader import SeedCtx, _import_module |
| 122 | + from mockworld.schema import MockDef |
| 123 | + from mockworld.state import _TOMBSTONE, StateView |
| 124 | + |
| 125 | + definition = MockDef.model_validate(yaml.safe_load((mock_dir / "mock.yaml").read_text())) |
| 126 | + for name, template in definition.errors.items(): |
| 127 | + register_error(name, template) |
| 128 | + handlers = _import_module(mock_dir / "handlers.py", "sbx_handlers") |
| 129 | + seed_module = _import_module(mock_dir / "seed.py", "sbx_seed") |
| 130 | + |
| 131 | + def handle(req: dict) -> dict: |
| 132 | + op = req["op"] |
| 133 | + dctx = DeterministicContext(req["seed"]) |
| 134 | + |
| 135 | + if op == "seed": |
| 136 | + seed_ctx = SeedCtx(rng=dctx.seed_rng(), ids=dctx.ids_for("__seed__", 0), |
| 137 | + fake=DataGen(dctx.seed_rng())) |
| 138 | + if definition.seed.generator.startswith("python:") and seed_module is not None: |
| 139 | + fn = getattr(seed_module, definition.seed.generator.split(".", 1)[-1]) |
| 140 | + snapshot = fn(seed_ctx, definition) |
| 141 | + else: |
| 142 | + snapshot = {} # builtin seed handled parent-side (trusted code) |
| 143 | + return {"ok": True, "snapshot": snapshot} |
| 144 | + |
| 145 | + if op == "call": |
| 146 | + tool = definition.tool(req["tool"]) |
| 147 | + fn = getattr(handlers, tool.handler_name, None) if handlers else None |
| 148 | + if fn is None: |
| 149 | + return {"ok": True, "result": {"success": False, "error": { |
| 150 | + "code": "internal_error", "message": f"handler {tool.handler_name!r} not found", |
| 151 | + "http_status": 500, "body": {}, "retry_after_s": None}}, |
| 152 | + "mutations": {}, "deletes": {}} |
| 153 | + |
| 154 | + view = StateView(req["state"], {}, list(req["state"].keys())) |
| 155 | + ctx = HandlerCtx(state=view, clock=dctx.clock_for(req["step"]), |
| 156 | + ids=dctx.ids_for(req["tool"], req["idx"]), |
| 157 | + rng=dctx.rng_for(req["tool"], req["idx"]), |
| 158 | + tool=req["tool"], faults=FaultHelper()) |
| 159 | + result = fn(ctx, req["params"]) |
| 160 | + |
| 161 | + mutations: dict = {} |
| 162 | + deletes: dict = {} |
| 163 | + for coll, entries in view._scratch.items(): |
| 164 | + for key, value in entries.items(): |
| 165 | + if value is _TOMBSTONE: |
| 166 | + deletes.setdefault(coll, []).append(key) |
| 167 | + else: |
| 168 | + mutations.setdefault(coll, {})[key] = value |
| 169 | + return {"ok": True, "result": _result_payload(result), |
| 170 | + "mutations": mutations, "deletes": deletes} |
| 171 | + |
| 172 | + return {"ok": False, "error": f"unknown op {op!r}"} |
| 173 | + |
| 174 | + while True: |
| 175 | + req = _read_msg(inp) |
| 176 | + if req is None: |
| 177 | + break |
| 178 | + try: |
| 179 | + resp = handle(req) |
| 180 | + except Exception as exc: # never crash the worker on handler error |
| 181 | + resp = {"ok": False, "error": f"{type(exc).__name__}: {exc}"} |
| 182 | + _write_msg(proto_out, resp) |
| 183 | + |
| 184 | + |
| 185 | +if __name__ == "__main__": |
| 186 | + main() |
0 commit comments