From a20951d2094ffcead261a12e135f210b1248edc0 Mon Sep 17 00:00:00 2001 From: Claude Date: Thu, 17 Sep 2026 13:01:57 +0000 Subject: [PATCH] Add backend="fabric": cells in worker processes that can be killed The sub-interpreter backend leaves one thing unsolved that a multi-tenant deployment cannot live without: a running cell cannot be reclaimed. close() refuses while the guest executes, there is no kill, and an async exception aimed at the thread does not reach the interpreter running in it. In-process, a runaway guest is permanent -- the cell is abandoned and its thread stays pinned for the life of the process. The fabric puts cells somewhere killable. Three levels, each a different kind of boundary: * a cell separates namespaces -- its own sys.modules, builtins and globals; * a worker process is the kill domain, and the only place a memory cap means anything, because sys.getallocatedblocks() is process-global rather than per-interpreter so there is no per-cell figure to limit; * the fabric decides which worker a tenant's cells land in, which is how a deployment chooses its blast radius. A deadline that expires now kills the worker and the sandbox raises saying so. Measured end to end: a runaway cell is reclaimed in its deadline, the worker process is gone, another tenant's cells are untouched, and the wedged tenant gets a fresh worker on its next spawn. Placement is the blast-radius decision, so it is explicit. tenant_isolation defaults to on, meaning a worker only ever hosts one tenant's cells and a kill costs that tenant alone; turning it off packs tenants together for density and makes them share a fate. Callers state which they want rather than getting one silently. Cost of the kill domain, measured on free-threaded 3.14: 1.20 ms for an exec+recv round trip against 0.82 ms in-process, plus a 160 ms worker spawn per tenant that WorkerPool.prewarm moves off the request path. Two defects found and fixed while building it, both in the worker lifecycle: * A worker killed from outside stayed a zombie until something happened to poll it, so a fabric recycling under load would accumulate one per recycle. The reader thread now reaps on EOF, which is the earliest reliable signal the process is finished. * Popen.poll() returns None when it cannot take the waitpid lock, so while the reader thread was inside wait() a dead worker still reported is_alive() -- long enough for placement to hand it a new cell. The worker is now marked unusable when its channel hits EOF, before the reap. Workers are spawned, never forked: forking a process that already hosts sub-interpreters and their threads segfaults. Exceptions crossing back are rebuilt from a fixed name map rather than a dynamic lookup, so a worker running guest code cannot name an arbitrary supervisor class and have it constructed. The fabric is a fault boundary, not a security one -- a guest that escapes its cell owns its worker, and a worker is an ordinary process with the supervisor's privileges. backend="process" remains the boundary mode. Kernel confinement of workers is now a roadmap item rather than a gap with no plan: tenant_isolation makes a worker's policy unambiguous, so seccomp/Landlock could be applied at spawn. Co-Authored-By: Claude Opus 5 Claude-Session: https://claude.ai/code/session_01Fs1hTtmF4Hm9h617AG9Gse --- .github/workflows/ci.yml | 6 +- CHANGELOG.md | 26 +- README.md | 89 +++ ROADMAP.md | 30 +- pyisolate/runtime/cell_worker.py | 405 ++++++++++++ pyisolate/runtime/fabric.py | 963 ++++++++++++++++++++++++++++ pyisolate/runtime/subinterpreter.py | 21 + pyisolate/supervisor.py | 80 ++- tests/test_fabric.py | 499 ++++++++++++++ tests/test_supervisor.py | 61 +- 10 files changed, 2159 insertions(+), 21 deletions(-) create mode 100644 pyisolate/runtime/cell_worker.py create mode 100644 pyisolate/runtime/fabric.py create mode 100644 tests/test_fabric.py diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 682cedb..89129ab 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -102,7 +102,7 @@ jobs: run: pytest -q tests/test_matrix_hardening.py -m "race and no_gil" tests-subinterpreter: - name: sub-interpreter cells / py3.14t + name: cells + fabric / py3.14t runs-on: ubuntu-24.04 steps: - uses: actions/checkout@v4 @@ -132,6 +132,10 @@ jobs: PY - name: Run the sub-interpreter backend suite run: pytest -q tests/test_subinterpreter_backend.py tests/test_supervisor.py + - name: Run the fabric suite + # Spawns worker processes and kills them; the runaway cases run inline + # here because the fabric reclaims them, unlike the in-process backend. + run: pytest -q tests/test_fabric.py - name: Validate the no-GIL readiness axis on 3.14t run: pytest -q tests/test_nogil.py diff --git a/CHANGELOG.md b/CHANGELOG.md index b326869..d2a93d0 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -8,6 +8,16 @@ guarantees; **no release should be treated as a hardened security boundary**. ## [Unreleased] ### Added +- `backend="fabric"`: the multi-tenant mode. Sub-interpreter cells hosted in + worker processes the supervisor can kill, which is the only reclaim that + works on a running cell. A deadline that expires kills the worker and the + tenant keeps running on a fresh one; other tenants are untouched. Includes + tenant-aware placement (`tenant_isolation=True` by default, so a kill costs + one tenant), per-worker `RLIMIT_AS` caps, `prewarm` to move the ~160 ms + worker spawn off the request path, worker recycling with counters, and + `Supervisor.fabric_report()` for placement visibility. Needs CPython 3.14+ + and fails closed below it. + - `backend="subinterpreter"`: real CPython sub-interpreter cells on 3.14+, via `concurrent.interpreters`, with a pre-warmed `CellPool`. Each guest gets its own `sys.modules` and its own `builtins`, so the import allow-list is a @@ -63,11 +73,17 @@ guarantees; **no release should be treated as a hardened security boundary**. - The eBPF programs compile to loadable objects and are covered by ELF-level tests, but load/attach against a live verifier is still only exercised by the root-gated `PYISOLATE_LIVE_BPF_TESTS=1` tests, not by CI. -- A running sub-interpreter cell cannot be reclaimed: one that overruns its - wall-time deadline is abandoned, and its thread stays pinned until the - process exits. Cells enforce a wall-time deadline and no other quota; - `sys.getallocatedblocks()` is process-global, so per-cell memory - accounting needs a worker-process layer that does not exist yet. +- In `backend="subinterpreter"` a running cell still cannot be reclaimed: one + that overruns is abandoned and its thread stays pinned until the process + exits. Use `backend="fabric"`, where the worker is the kill domain. +- A fabric worker is an ordinary process with the supervisor's privileges: the + fabric bounds faults, not hostile Python. Kernel confinement of workers is + not implemented. +- The broker `request` op is surfaced by the fabric but, as with the other + backends, nothing executes it. +- Memory is capped per worker (`RLIMIT_AS`), not per cell: + `sys.getallocatedblocks()` is process-global on both free-threaded and GIL + builds, so there is no per-cell figure to limit. - Process-backed sandboxes are not attached to cgroups or watched by the resource watchdog (they get `rlimit` only). - `backend="microvm"` fails closed: the guest agent and vsock cell transport are diff --git a/README.md b/README.md index 63ffa5c..4f23103 100644 --- a/README.md +++ b/README.md @@ -263,6 +263,95 @@ domain. --- +## The fabric + +`backend="fabric"` is the multi-tenant mode: sub-interpreter cells hosted in +**worker processes the supervisor can kill**. + +``` + Supervisor + | + +-- Worker process <- the kill domain + | +-- cell cell cell (one CellPool, many cells) + | + +-- Worker process + +-- cell cell +``` + +It exists because of one limitation the in-process cell backend cannot fix: a +running sub-interpreter **cannot be reclaimed**. `close()` refuses while the +guest executes, there is no `kill`, and an async exception aimed at the thread +does not reach the interpreter running in it. In-process, a runaway guest is +permanent. Putting cells in a worker makes the process the unit of reclaim: + +```python +import pyisolate as iso + +sb = iso.spawn("report", backend="fabric", tenant="acme", wall_time_ms=500) +sb.exec("while True: pass") +# WallTimeExceeded: cell c1 exceeded 0.5s and did not return. A running +# sub-interpreter cannot be reclaimed, so the worker process hosting it was +# killed; cells sharing that worker were lost with it. +``` + +Each of the three levels is a different kind of boundary: + +| Level | Isolates | Reclaimable | +| --- | --- | --- | +| cell | `sys.modules`, `builtins`, globals | no | +| worker | the kill domain, and where a memory cap applies | **yes — SIGKILL** | +| fabric | decides which worker a tenant lands in | n/a | + +**Placement is the blast-radius decision.** With `tenant_isolation=True` (the +default) a worker only ever hosts one tenant's cells, so killing it for a +runaway costs that tenant and nobody else. Turning it off packs tenants +together for density and makes them share a fate. The fabric makes callers +state which they want rather than picking silently. + +```python +from pyisolate.runtime.fabric import WorkerPool + +pool = WorkerPool( + max_workers=8, + cells_per_worker=16, + tenant_isolation=True, # one tenant per worker + worker_mem_bytes=2 << 30, # RLIMIT_AS per worker +) +pool.prewarm("acme", 2) # pay the ~160 ms spawn before traffic arrives +``` + +`worker_mem_bytes` is where a memory limit can actually be enforced: +`sys.getallocatedblocks()` is process-global rather than per-interpreter on +both free-threaded and GIL builds, so there is no per-cell figure to cap. +Worker sizing is the control. + +### What it costs + +Measured on the same 4-core container as the figures above, free-threaded +CPython 3.14: + +| Operation | p50 | +| --- | --- | +| worker spawn (fresh interpreter + pyisolate import + pool) | 160.8 ms | +| first cell for a tenant (spawns its worker) | 190.3 ms | +| further cells in that worker | 29.8 ms | +| `exec` + `recv` round trip | 1.20 ms | + +So the kill domain costs about **0.4 ms per round trip** over an in-process +cell (1.20 ms against 0.82 ms) plus one worker spawn per tenant, which +`prewarm` moves off the request path. + +### What it is still not + +A guest that escapes its cell owns its worker, and a worker is an ordinary +process holding the supervisor's privileges. The fabric is **not** a boundary +against hostile Python. For untrusted code use `backend="process"` — one +confined process per sandbox — or a microVM. What the fabric buys is that +trusted-but-independent tenants cannot wedge each other, and that a tenant +which wedges itself is recoverable. + +--- + ## Canonical execution model A cell is intentionally limited to seven operations: `exec`, `call`, `post`, `recv`, `log`, `metric`, and `request`. diff --git a/ROADMAP.md b/ROADMAP.md index ef32644..716b4aa 100644 --- a/ROADMAP.md +++ b/ROADMAP.md @@ -10,7 +10,9 @@ normative statement. - **Backends** — `backend="thread"` (execution cell; a dedicated thread of the supervisor process, previously spelled `subinterpreter`), `backend="subinterpreter"` (real CPython sub-interpreter cells on 3.14+, with - a pre-warmed pool; an execution cell, not a boundary), and + a pre-warmed pool; an execution cell, not a boundary), `backend="fabric"` + (those cells hosted in worker processes, with tenant-aware placement, worker + memory caps, and kill-and-replace recovery — the multi-tenant mode), and `backend="process"` (the boundary mode): a real separate-process boundary with `no_new_privs` + a seccomp deny-list, Landlock filesystem rules, Landlock TCP-egress rules (ABI ≥ 4), a coarse per-cgroup eBPF/LSM deny-mask, and @@ -33,17 +35,21 @@ normative statement. ## Now / next -- **A kill domain for cells** — a running sub-interpreter cannot be reclaimed: - `close()` refuses while the guest is executing and there is no `kill`, so a - cell that overruns is abandoned and its thread is pinned for the life of the - process. The fix is not in-process. Add a layer of pre-forked worker - processes between the supervisor and the cells, size them by tenant, and make - the worker the unit that gets killed and replaced. This is what turns the cell - pool into something that survives a hostile-by-accident tenant. -- **Per-cell resource accounting** — there is none today. - `sys.getallocatedblocks()` is process-global on both free-threaded and GIL - builds, so memory has to be capped at the worker/cgroup level rather than per - cell. Cells currently enforce a wall-time deadline and nothing else. +- **Broker request execution for the fabric** — a fabric cell's `request` op + surfaces to the supervisor like the other backends', and still nothing + executes it. The fabric is where mediation matters most, because a worker + holds many tenants' cells: the handler has to be capability-scoped per cell, + not per worker. +- **Kernel confinement of fabric workers** — a worker is currently an ordinary + process with the supervisor's privileges. With `tenant_isolation=True` a + worker serves one tenant, so its policy is unambiguous and the process + backend's seccomp/Landlock/cgroup layers could be applied to it at spawn. + That would make the fabric a defence-in-depth boundary rather than only a + fault boundary. +- **Fabric admission from policy** — routing is explicit today + (`backend="fabric"`, `tenant=...`). Deriving placement and backend choice + from the policy's import list would let a tenant that needs numpy land on + `process` automatically instead of failing at import time. - **Broker request execution** — the `request` op currently surfaces a `BrokerRequest` to the host but nothing executes it or returns a result. Add a request/response round-trip and a pluggable, capability-scoped handler so the diff --git a/pyisolate/runtime/cell_worker.py b/pyisolate/runtime/cell_worker.py new file mode 100644 index 0000000..bfa49b3 --- /dev/null +++ b/pyisolate/runtime/cell_worker.py @@ -0,0 +1,405 @@ +"""Worker process that hosts sub-interpreter cells for the fabric. + +This module is the entry point executed in a *fresh* interpreter for every +fabric worker (``python -m pyisolate.runtime.cell_worker ``). A worker owns +a :class:`~pyisolate.runtime.subinterpreter.CellPool` and runs many cells inside +itself; the supervisor drives it over the same length-framed JSON protocol the +process backend uses. + +Why a separate process at all, when a cell is already isolated from its +siblings: **a running sub-interpreter cannot be reclaimed.** ``close()`` refuses +while the guest is executing, there is no ``kill``, and an async exception aimed +at the thread does not reach the interpreter running in it. The only reclaim +that works is killing the process. So the worker is the kill domain: when a cell +overruns, the supervisor SIGKILLs the whole worker and starts a replacement, +which costs the other cells in that worker and nothing outside it. Sizing a +worker -- ideally one tenant per worker -- is therefore how a deployment chooses +its blast radius. + +A worker is started with ``spawn`` semantics (a fresh ``sys.executable``), never +by forking the supervisor. Forking a process that already hosts sub-interpreters +and their threads segfaults, which is also why the backend's own runaway tests +shell out instead of calling ``os.fork``. + +Parent -> worker frames:: + + {"op": "hello", "warm_per_spec": N, "max_warm": N, "mem_bytes": N | null} + {"op": "open", "cell": "", "allowed_imports": [...], "preimport": [...]} + {"op": "exec", "cell": "", "seq": N, "source": "..."} + {"op": "call", "cell": "", "seq": N, "target": "mod.fn", + "args": [...], "kwargs": {...}} + {"op": "close", "cell": ""} + {"op": "stats"} + {"op": "stop"} + +Worker -> parent frames:: + + {"ev": "ready", "pid": N, "python": "...", "free_threaded": bool} + {"ev": "opened", "cell": "", "interp": N} + {"ev": "post" | "log" | "metric" | "request", "cell": "", ...} + {"ev": "done", "cell": "", "seq": N} + {"ev": "error", "cell": "", "seq": N, "exc_type": "...", "message": "..."} + {"ev": "closed", "cell": ""} + {"ev": "stats", "pool": {...}, "cells": N} +""" + +from __future__ import annotations + +import json +import resource +import socket +import struct +import sys +import threading +import traceback +from typing import Any, Optional + +from .subinterpreter import Cell, CellPool, CellSpec, require_available + +_LEN = struct.Struct("!I") + +#: How often the pump forwards whatever cells have posted. A cell writes to its +#: own queue whenever guest code calls ``post``; without a pump those messages +#: would only reach the supervisor when the operation returns, so a guest that +#: posts and then blocks would look silent. +_PUMP_INTERVAL = 0.005 + + +class _Channel: + """Length-framed JSON over the inherited socket, safe for many senders. + + Cell operations run on their own threads so one slow guest cannot stall the + worker's control loop, which means several threads frame concurrently. The + lock keeps a frame from being interleaved with another on the wire. + """ + + def __init__(self, sock: socket.socket) -> None: + self._sock = sock + self._lock = threading.Lock() + + def send(self, obj: dict[str, Any]) -> None: + try: + data = json.dumps(obj).encode("utf-8") + except TypeError: + # A cell posted something JSON cannot carry. The cell bootstrap + # already encodes guest payloads, so this is a bug on our side + # rather than guest input; report it rather than killing the loop. + data = json.dumps( + { + "ev": "error", + "cell": obj.get("cell"), + "seq": obj.get("seq"), + "exc_type": "TypeError", + "message": "worker produced a non-serialisable frame", + } + ).encode("utf-8") + with self._lock: + try: + self._sock.sendall(_LEN.pack(len(data)) + data) + except OSError: + # The supervisor is gone (it killed us, or it exited). Nothing + # to report to, and raising here would only spam tracebacks out + # of cell threads. + pass + + def recv(self) -> Optional[dict[str, Any]]: + header = self._recv_exact(_LEN.size) + if header is None: + return None + (size,) = _LEN.unpack(header) + body = self._recv_exact(size) + if body is None: + return None + return json.loads(body.decode("utf-8")) + + def _recv_exact(self, size: int) -> Optional[bytes]: + chunks: list[bytes] = [] + remaining = size + while remaining: + try: + chunk = self._sock.recv(remaining) + except OSError: + return None + if not chunk: + return None + chunks.append(chunk) + remaining -= len(chunk) + return b"".join(chunks) + + +class _Worker: + """Owns this process's cell pool and serves the supervisor's frames.""" + + def __init__(self, channel: _Channel) -> None: + self._channel = channel + self._pool: Optional[CellPool] = None + self._cells: dict[str, Cell] = {} + self._lock = threading.Lock() + self._stopping = threading.Event() + self._pump: Optional[threading.Thread] = None + + # --- lifecycle --- + + def hello(self, frame: dict[str, Any]) -> None: + mem_bytes = frame.get("mem_bytes") + if mem_bytes: + # The only memory cap available to a cell fabric. + # ``sys.getallocatedblocks()`` is process-global rather than + # per-interpreter on both free-threaded and GIL builds, so there is + # no per-cell accounting to enforce; capping the worker's address + # space is what makes a memory limit mean anything, and it is + # another reason worker sizing is how blast radius gets chosen. + try: + resource.setrlimit(resource.RLIMIT_AS, (mem_bytes, mem_bytes)) + except (ValueError, OSError) as exc: # pragma: no cover - host dependent + self._channel.send( + { + "ev": "error", + "cell": None, + "seq": None, + "exc_type": type(exc).__name__, + "message": f"could not apply worker memory cap: {exc}", + } + ) + self._pool = CellPool( + warm_per_spec=int(frame.get("warm_per_spec", 1)), + max_warm=int(frame.get("max_warm", 32)), + ) + self._pump = threading.Thread( + target=self._pump_loop, name="pyisolate-worker-pump", daemon=True + ) + self._pump.start() + self._channel.send( + { + "ev": "ready", + "pid": _getpid(), + "python": sys.version.split()[0], + "free_threaded": not getattr(sys, "_is_gil_enabled", lambda: True)(), + } + ) + + def open(self, frame: dict[str, Any]) -> None: + cell_id = frame["cell"] + assert self._pool is not None + spec = CellSpec.build( + frame.get("allowed_imports") or (), + frame.get("preimport"), + ) + cell = self._pool.acquire(spec) + with self._lock: + self._cells[cell_id] = cell + self._channel.send({"ev": "opened", "cell": cell_id, "interp": cell.id}) + + def close_cell(self, frame: dict[str, Any]) -> None: + cell_id = frame["cell"] + with self._lock: + cell = self._cells.pop(cell_id, None) + if cell is not None and self._pool is not None: + self._forward(cell_id, cell) + self._pool.release(cell) + self._channel.send({"ev": "closed", "cell": cell_id}) + + def stats(self) -> None: + with self._lock: + live = len(self._cells) + pool = self._pool.stats() if self._pool is not None else {} + self._channel.send({"ev": "stats", "pool": pool, "cells": live}) + + # --- running guest work --- + + def dispatch(self, frame: dict[str, Any]) -> None: + """Run one operation on a thread of its own. + + The worker must stay responsive while a guest runs: a cell that never + returns would otherwise block ``stats`` and ``close`` for every other + cell in this worker, and the supervisor would have no way to tell a + wedged worker from a busy one. + """ + cell_id = frame["cell"] + with self._lock: + cell = self._cells.get(cell_id) + if cell is None: + self._channel.send( + { + "ev": "error", + "cell": cell_id, + "seq": frame.get("seq"), + "exc_type": "SandboxError", + "message": f"no open cell {cell_id!r}", + } + ) + return + thread = threading.Thread( + target=self._run_op, + args=(cell_id, cell, frame), + name=f"pyisolate-cell-{cell_id}", + daemon=True, + ) + thread.start() + + def _run_op(self, cell_id: str, cell: Cell, frame: dict[str, Any]) -> None: + seq = frame.get("seq") + try: + if frame["op"] == "exec": + cell.exec(frame.get("source", "")) + else: + self._run_call(cell, frame) + except BaseException as exc: # noqa: BLE001 - every failure goes to the host + self._forward(cell_id, cell) + self._channel.send( + { + "ev": "error", + "cell": cell_id, + "seq": seq, + "exc_type": type(exc).__name__, + "message": _describe(exc), + } + ) + return + self._forward(cell_id, cell) + self._channel.send({"ev": "done", "cell": cell_id, "seq": seq}) + + @staticmethod + def _run_call(cell: Cell, frame: dict[str, Any]) -> None: + """Resolve a dotted name inside the cell and post its result. + + Resolution happens in the cell so it goes through that interpreter's + guarded ``__import__`` -- the allow-list applies to a ``call`` target + exactly as it does to an ``import`` in guest source -- and so the worker + never hands one of its own objects to the guest. + """ + target = frame.get("target", "") + module, _, attr = target.rpartition(".") + if not module or not attr: + raise ValueError(f"call target {target!r} must be a dotted name") + payload = json.dumps( + {"args": frame.get("args") or [], "kwargs": frame.get("kwargs") or {}} + ) + cell.exec( + f"_pyi_call = _pyi_json.loads({payload!r})\n" + f"_pyi_mod = __import__({module!r}, fromlist=[{attr!r}])\n" + f"post(getattr(_pyi_mod, {attr!r})" + "(*_pyi_call['args'], **_pyi_call['kwargs']))\n" + ) + + # --- message forwarding --- + + def _forward(self, cell_id: str, cell: Cell) -> None: + for kind, payload in cell.drain(): + self._channel.send({"ev": kind, "cell": cell_id, "payload": payload}) + + def _pump_loop(self) -> None: + while not self._stopping.wait(_PUMP_INTERVAL): + with self._lock: + items = list(self._cells.items()) + for cell_id, cell in items: + try: + self._forward(cell_id, cell) + except Exception: # pragma: no cover - a retired cell mid-drain + continue + + # --- shutdown --- + + def shutdown(self) -> None: + self._stopping.set() + with self._lock: + cells = list(self._cells.values()) + self._cells.clear() + pool = self._pool + if pool is None: + return + for cell in cells: + # A cell still running cannot be retired; the pool records that + # rather than blocking shutdown on a guest that never returns. + pool.release(cell) + pool.close() + + +def _describe(exc: BaseException) -> str: + text = str(exc).strip() + if text: + return text.splitlines()[-1] + return traceback.format_exception_only(type(exc), exc)[-1].strip() + + +def _getpid() -> int: + import os + + return os.getpid() + + +def _serve(sock: socket.socket) -> None: + channel = _Channel(sock) + worker = _Worker(channel) + try: + require_available() + except Exception as exc: # pragma: no cover - guarded by the supervisor too + channel.send( + { + "ev": "error", + "cell": None, + "seq": None, + "exc_type": type(exc).__name__, + "message": str(exc), + } + ) + return + + handlers = { + "hello": worker.hello, + "open": worker.open, + "close": worker.close_cell, + } + try: + while True: + frame = channel.recv() + if frame is None: + return + op = frame.get("op") + if op == "stop": + return + try: + if op in handlers: + handlers[op](frame) + elif op == "stats": + worker.stats() + elif op in ("exec", "call"): + worker.dispatch(frame) + else: + raise ValueError(f"unknown worker operation: {op!r}") + except BaseException as exc: # noqa: BLE001 - never drop the loop + channel.send( + { + "ev": "error", + "cell": frame.get("cell"), + "seq": frame.get("seq"), + "exc_type": type(exc).__name__, + "message": _describe(exc), + } + ) + finally: + worker.shutdown() + + +def main(argv: list[str]) -> int: + if len(argv) < 2: + return 2 + fd = int(argv[1]) + sock = socket.socket(socket.AF_UNIX, socket.SOCK_STREAM, fileno=fd) + try: + _serve(sock) + finally: + try: + sock.close() + except OSError: + pass + # A cell that never returned leaves its thread pinned inside a running + # interpreter, and a normal exit would block on it forever. The supervisor + # already has the result it is going to get, so leave immediately. + import os + + os._exit(0) + + +if __name__ == "__main__": + raise SystemExit(main(sys.argv)) diff --git a/pyisolate/runtime/fabric.py b/pyisolate/runtime/fabric.py new file mode 100644 index 0000000..5416f2c --- /dev/null +++ b/pyisolate/runtime/fabric.py @@ -0,0 +1,963 @@ +"""The fabric: sub-interpreter cells inside worker processes that can be killed. + +:mod:`pyisolate.runtime.subinterpreter` gives cheap, namespace-isolated cells, +but leaves one thing unsolved that a multi-tenant deployment cannot live +without: **a running cell cannot be reclaimed.** ``close()`` refuses while the +guest executes, there is no ``kill``, and an async exception aimed at the thread +does not reach the interpreter running in it. In-process, a runaway guest is +permanent -- the cell is abandoned and its thread stays pinned for the life of +the process. + +The fabric fixes that by putting the cells somewhere killable: + + Supervisor + | + +-- Worker process <- the kill domain + | +-- cell cell cell (one CellPool, many cells) + | + +-- Worker process + +-- cell cell + +Three levels, and each one is a different kind of boundary: + +* a **cell** separates namespaces -- its own ``sys.modules``, its own + ``builtins``, its own globals. Not a boundary against hostile Python. +* a **worker** is the unit that gets killed and replaced, so it bounds the + damage a runaway or wedged guest can do. It is also where a memory cap can + actually be enforced, because ``sys.getallocatedblocks()`` is process-global + rather than per-interpreter, so there is no per-cell figure to limit. +* the **fabric** decides which worker a tenant's cells land in, which is how a + deployment chooses its blast radius. + +Placement is the part worth getting right. With ``tenant_isolation=True`` (the +default) a worker only ever hosts one tenant's cells, so killing it for a +runaway costs that tenant and nobody else. Turning it off packs tenants together +for density and makes them share a fate; that is a real trade and the fabric +makes callers state which one they want rather than picking silently. + +This is still not a boundary against hostile Python -- a guest that escapes its +cell owns its worker, and the worker is an ordinary process with the +supervisor's privileges. For untrusted code use ``backend="process"``, one +sandbox per process, or a microVM. What the fabric buys is that trusted-but- +independent tenants cannot wedge each other, and that a tenant which does wedge +itself is recoverable. +""" + +from __future__ import annotations + +import errno +import itertools +import json +import logging +import socket +import struct +import subprocess +import sys +import threading +import time +from dataclasses import dataclass, field +from typing import Any, Optional + +from .. import errors +from . import subinterpreter +from .process_backend import build_child_env + +logger = logging.getLogger(__name__) + +_LEN = struct.Struct("!I") +_WORKER_MODULE = "pyisolate.runtime.cell_worker" + +#: How long to wait for a worker to report ``ready`` before treating the spawn +#: as failed. Generous: the worker imports pyisolate and builds a pool first. +_READY_TIMEOUT = 30.0 + +#: How long a SIGTERM gets before SIGKILL when retiring a worker cleanly. A +#: worker with a stranded cell never exits on its own, so this stays short. +_TERM_GRACE = 0.5 + + +class WorkerDied(errors.SandboxError): + """The worker hosting this cell went away, taking the cell with it.""" + + +@dataclass +class FabricStats: + """Counters that make placement and recycling legible in production.""" + + workers_started: int = 0 + workers_killed: int = 0 + workers_died: int = 0 + cells_opened: int = 0 + cells_closed: int = 0 + + def as_dict(self) -> dict[str, int]: + return { + "workers_started": self.workers_started, + "workers_killed": self.workers_killed, + "workers_died": self.workers_died, + "cells_opened": self.cells_opened, + "cells_closed": self.cells_closed, + } + + +@dataclass +class _Pending: + """One in-flight operation, waiting on its worker.""" + + done: threading.Event = field(default_factory=threading.Event) + error: Optional[tuple[str, str]] = None + messages: list[tuple[str, Any]] = field(default_factory=list) + + +class Worker: + """One worker process, its channel, and the cells currently open in it.""" + + _ids = itertools.count(1) + + def __init__( + self, + *, + tenant: Optional[str] = None, + warm_per_spec: int = 1, + max_warm: int = 32, + mem_bytes: Optional[int] = None, + env: Optional[dict[str, str]] = None, + ) -> None: + self.id = next(self._ids) + self.tenant = tenant + self.started_at = time.monotonic() + self.dead_reason: Optional[str] = None + self._lock = threading.Lock() + self._send_lock = threading.Lock() + self._closed = False + self._seq = itertools.count(1) + # Keyed by operation sequence number for exec/call, and by a string for + # the one-shot lifecycle replies ("opened:", "closed:", + # "stats") -- those have no sequence of their own because the worker + # answers them with a named event rather than a numbered completion. + self._pending: dict[object, _Pending] = {} + self._cell_inbox: dict[str, list[tuple[str, Any]]] = {} + self._ready = threading.Event() + self._ready_info: dict[str, Any] = {} + # Set the moment the channel reports EOF, before the process is reaped. + # ``Popen.poll()`` returns None when it cannot take the waitpid lock, so + # while the reader thread is inside ``wait()`` a dead worker would + # otherwise still look alive -- long enough for placement to hand it a + # new cell. + self._eof = threading.Event() + + parent_sock, child_sock = socket.socketpair(socket.AF_UNIX, socket.SOCK_STREAM) + try: + self._proc = subprocess.Popen( + [sys.executable, "-m", _WORKER_MODULE, str(child_sock.fileno())], + pass_fds=(child_sock.fileno(),), + close_fds=True, + env=build_child_env(env), + ) + except Exception: + parent_sock.close() + child_sock.close() + raise + # Drop our copy of the child's end so an unexpected worker exit shows up + # as EOF on the parent side rather than hanging the reader forever. + child_sock.close() + self._sock = parent_sock + + self._reader = threading.Thread( + target=self._read_loop, name=f"pyisolate-worker-{self.id}", daemon=True + ) + self._reader.start() + self._send( + { + "op": "hello", + "warm_per_spec": warm_per_spec, + "max_warm": max_warm, + "mem_bytes": mem_bytes, + } + ) + if not self._ready.wait(_READY_TIMEOUT): + self.kill("worker did not become ready") + raise errors.SandboxError( + f"fabric worker {self.id} did not report ready within " + f"{_READY_TIMEOUT}s" + ) + + # --- introspection --- + + @property + def pid(self) -> Optional[int]: + return self._proc.pid + + @property + def info(self) -> dict[str, Any]: + return dict(self._ready_info) + + def is_alive(self) -> bool: + if self._closed or self._eof.is_set(): + return False + return self._proc.poll() is None + + def cell_count(self) -> int: + with self._lock: + return len(self._cell_inbox) + + # --- transport --- + + def _send(self, obj: dict[str, Any]) -> None: + data = json.dumps(obj).encode("utf-8") + with self._send_lock: + if self._closed: + raise WorkerDied(f"fabric worker {self.id} is closed") + try: + self._sock.sendall(_LEN.pack(len(data)) + data) + except OSError as exc: + raise WorkerDied( + f"fabric worker {self.id} channel is gone: {exc}" + ) from exc + + def _recv_exact(self, size: int) -> Optional[bytes]: + chunks: list[bytes] = [] + remaining = size + while remaining: + try: + chunk = self._sock.recv(remaining) + except OSError: + return None + if not chunk: + return None + chunks.append(chunk) + remaining -= len(chunk) + return b"".join(chunks) + + def _read_loop(self) -> None: + while True: + header = self._recv_exact(_LEN.size) + if header is None: + break + (size,) = _LEN.unpack(header) + body = self._recv_exact(size) + if body is None: + break + try: + frame = json.loads(body.decode("utf-8")) + except ValueError: + logger.warning("worker %s sent an unparseable frame", self.id) + continue + try: + self._dispatch(frame) + except Exception: # pragma: no cover - never lose the reader + logger.exception("worker %s frame dispatch failed", self.id) + # EOF: the worker exited, was killed, or crashed. Mark it unusable + # first, so placement stops considering it before we block on the reap. + self._eof.set() + # Reap it here rather than waiting for the next placement decision to + # poll. A killed child stays a zombie until someone waits on it, and a + # fabric that recycles workers under load would accumulate one per + # recycle. The channel closing is the earliest reliable signal that the + # process is finished, so it is the right place to collect it. + try: + self._proc.wait(timeout=_TERM_GRACE) + except subprocess.TimeoutExpired: # pragma: no cover - EOF without exit + logger.warning("worker %s closed its channel but is still running", self.id) + # Everything waiting on it has to be told, or callers block forever on a + # process that is gone. + self._fail_all(self.dead_reason or "worker process exited") + + def _dispatch(self, frame: dict[str, Any]) -> None: + event = frame.get("ev") + if event == "ready": + self._ready_info = {k: v for k, v in frame.items() if k not in ("ev",)} + self._ready.set() + return + if event == "stats": + self._resolve_stats(frame) + return + + cell_id = frame.get("cell") + if event in ("post", "log", "metric", "request"): + if cell_id is not None: + with self._lock: + self._cell_inbox.setdefault(cell_id, []).append( + (event, frame.get("payload")) + ) + return + if event in ("opened", "closed"): + self._resolve_lifecycle(event, frame) + return + if event in ("done", "error"): + seq = frame.get("seq") + with self._lock: + pending = self._pending.pop(seq, None) if seq is not None else None + if pending is None: + # An error with no sequence is a worker-level failure (a failed + # memory cap, an unknown op). Nothing is waiting on it, so log + # it rather than dropping it silently. + if event == "error": + logger.warning( + "worker %s error: %s: %s", + self.id, + frame.get("exc_type"), + frame.get("message"), + ) + return + if event == "error": + pending.error = ( + str(frame.get("exc_type", "SandboxError")), + str(frame.get("message", "")), + ) + pending.done.set() + + def _resolve_lifecycle(self, event: str, frame: dict[str, Any]) -> None: + key = f"{event}:{frame.get('cell')}" + with self._lock: + pending = self._pending.pop(key, None) + if pending is not None: + pending.messages.append((event, frame)) + pending.done.set() + + def _resolve_stats(self, frame: dict[str, Any]) -> None: + with self._lock: + pending = self._pending.pop("stats", None) + if pending is not None: + pending.messages.append(("stats", frame)) + pending.done.set() + + def _fail_all(self, reason: str) -> None: + with self._lock: + pending = list(self._pending.values()) + self._pending.clear() + self._closed = True + for item in pending: + if item.error is None: + item.error = ("WorkerDied", reason) + item.done.set() + + def _await(self, key: Any, timeout: Optional[float]) -> _Pending: + pending = _Pending() + with self._lock: + if self._closed: + raise WorkerDied(f"fabric worker {self.id} is closed") + self._pending[key] = pending + return pending + + # --- operations --- + + def open_cell( + self, + cell_id: str, + *, + allowed_imports: Optional[list[str]] = None, + preimport: Optional[list[str]] = None, + timeout: Optional[float] = _READY_TIMEOUT, + ) -> dict[str, Any]: + pending = self._await(f"opened:{cell_id}", timeout) + self._send( + { + "op": "open", + "cell": cell_id, + "allowed_imports": list(allowed_imports or []), + "preimport": None if preimport is None else list(preimport), + } + ) + self._wait(pending, timeout, f"opening cell {cell_id}") + with self._lock: + self._cell_inbox.setdefault(cell_id, []) + return pending.messages[0][1] if pending.messages else {} + + def run( + self, + cell_id: str, + frame: dict[str, Any], + *, + timeout: Optional[float], + ) -> None: + """Dispatch one operation and wait for the worker to report on it. + + A timeout here is not "the call was slow" -- it means the guest never + came back, and nothing short of killing the worker will change that. + The caller (``WorkerPool.run``) is what turns that into a recycle. + """ + seq = next(self._seq) + pending = self._await(seq, timeout) + payload = dict(frame) + payload["cell"] = cell_id + payload["seq"] = seq + self._send(payload) + self._wait(pending, timeout, f"cell {cell_id} operation") + + def stats(self, timeout: float = 5.0) -> dict[str, Any]: + pending = self._await("stats", timeout) + self._send({"op": "stats"}) + self._wait(pending, timeout, "worker stats") + return pending.messages[0][1] if pending.messages else {} + + def close_cell(self, cell_id: str, timeout: float = 5.0) -> None: + try: + pending = self._await(f"closed:{cell_id}", timeout) + self._send({"op": "close", "cell": cell_id}) + self._wait(pending, timeout, f"closing cell {cell_id}") + except WorkerDied: + # Closing a cell in a worker that is already gone is a no-op, not + # an error: the kill freed it. + pass + finally: + with self._lock: + self._cell_inbox.pop(cell_id, None) + + def drain(self, cell_id: str) -> list[tuple[str, Any]]: + with self._lock: + messages = self._cell_inbox.get(cell_id) or [] + self._cell_inbox[cell_id] = [] + return messages + + def _wait(self, pending: _Pending, timeout: Optional[float], what: str) -> None: + if not pending.done.wait(timeout): + raise TimeoutError(f"{what} did not complete within {timeout}s") + if pending.error is not None: + exc_type, message = pending.error + if exc_type == "WorkerDied": + raise WorkerDied(message) + raise _rebuild(exc_type, message) + + # --- teardown --- + + def kill(self, reason: str) -> None: + """SIGKILL the worker. The only reclaim that works on a stranded cell.""" + self.dead_reason = reason + if self._proc.poll() is None: + try: + self._proc.kill() + except OSError as exc: # pragma: no cover - already reaped + if exc.errno != errno.ESRCH: + raise + self._reap() + + def stop(self, timeout: float = _TERM_GRACE) -> None: + """Ask the worker to exit, then kill it if it will not. + + A worker holding a stranded cell can never exit on its own, so the + escalation is not a fallback for slow shutdown -- it is the expected + path whenever a guest is still running. + """ + if self._proc.poll() is None: + try: + self._send({"op": "stop"}) + except (WorkerDied, OSError): + pass + try: + self._proc.wait(timeout) + except subprocess.TimeoutExpired: + self.kill("worker did not stop on request") + return + self._reap() + + def _reap(self) -> None: + try: + self._proc.wait(timeout=_TERM_GRACE) + except subprocess.TimeoutExpired: # pragma: no cover - SIGKILL is prompt + logger.warning("worker %s did not reap", self.id) + self._fail_all(self.dead_reason or "worker stopped") + try: + self._sock.close() + except OSError: + pass + + @property + def returncode(self) -> Optional[int]: + return self._proc.poll() + + +def _rebuild(exc_type: str, message: str) -> BaseException: + """Recreate a worker-side exception by name, without importing anything. + + Only the name crosses the boundary, so this maps the names the runtime is + known to produce and falls back to SandboxError. It never looks the name up + dynamically: that would let a worker name any class in the supervisor. + """ + known: dict[str, type[BaseException]] = { + "SandboxError": errors.SandboxError, + "PolicyError": errors.PolicyError, + "TimeoutError": errors.TimeoutError, + "WallTimeExceeded": errors.WallTimeExceeded, + "MemoryExceeded": errors.MemoryExceeded, + "ValueError": ValueError, + "TypeError": TypeError, + "ImportError": ImportError, + "NameError": NameError, + "AttributeError": AttributeError, + "ZeroDivisionError": ZeroDivisionError, + } + cls = known.get(exc_type, errors.SandboxError) + return cls(f"{exc_type}: {message}" if cls is errors.SandboxError else message) + + +class WorkerPool: + """Places cells into worker processes and recycles workers that wedge.""" + + def __init__( + self, + *, + max_workers: int = 4, + cells_per_worker: int = 8, + tenant_isolation: bool = True, + warm_per_spec: int = 1, + max_warm: int = 32, + worker_mem_bytes: Optional[int] = None, + env: Optional[dict[str, str]] = None, + ) -> None: + if max_workers < 1: + raise ValueError("max_workers must be >= 1") + if cells_per_worker < 1: + raise ValueError("cells_per_worker must be >= 1") + self.max_workers = max_workers + self.cells_per_worker = cells_per_worker + self.tenant_isolation = tenant_isolation + self._warm_per_spec = warm_per_spec + self._max_warm = max_warm + self._worker_mem_bytes = worker_mem_bytes + self._env = env + self._workers: list[Worker] = [] + self._lock = threading.Lock() + self._stats = FabricStats() + self._closed = False + self._cell_ids = itertools.count(1) + + # --- placement --- + + def placement_for(self, tenant: Optional[str]) -> Worker: + """Pick (or start) the worker a tenant's next cell belongs in.""" + with self._lock: + if self._closed: + raise errors.SandboxError("fabric worker pool is closed") + self._reap_dead_locked() + candidates = [ + w + for w in self._workers + if w.is_alive() + and w.cell_count() < self.cells_per_worker + and (not self.tenant_isolation or w.tenant == tenant) + ] + if candidates: + # Least-loaded, so cells spread rather than piling into the + # first worker and making one kill unusually expensive. + return min(candidates, key=lambda w: w.cell_count()) + if len(self._workers) >= self.max_workers: + raise errors.SandboxError( + f"fabric is at capacity: {len(self._workers)} workers x " + f"{self.cells_per_worker} cells. Raise max_workers or " + "cells_per_worker, or close some sandboxes." + ) + return self._start_worker(tenant) + + def prewarm(self, tenant: Optional[str] = None, count: int = 1) -> int: + """Start workers for *tenant* ahead of demand. Returns how many started. + + A worker costs ~160 ms to spawn: a fresh interpreter, the pyisolate + import, and a cell pool. Paying that on a tenant's first request is the + difference between a fabric that feels instant and one that does not, + and it is entirely avoidable -- the supervisor knows its tenants before + their traffic arrives. + """ + started = 0 + for _ in range(count): + with self._lock: + if self._closed or len(self._workers) >= self.max_workers: + break + try: + self._start_worker(tenant) + except errors.SandboxError: + break + started += 1 + return started + + def _start_worker(self, tenant: Optional[str]) -> Worker: + worker = Worker( + tenant=tenant, + warm_per_spec=self._warm_per_spec, + max_warm=self._max_warm, + mem_bytes=self._worker_mem_bytes, + env=self._env, + ) + with self._lock: + if self._closed: + worker.stop() + raise errors.SandboxError("fabric worker pool is closed") + self._workers.append(worker) + self._stats.workers_started += 1 + return worker + + def _reap_dead_locked(self) -> None: + alive = [] + for worker in self._workers: + if worker.is_alive(): + alive.append(worker) + else: + if worker.dead_reason is None: + # Nobody killed it; it crashed, was OOM-killed, or exited. + self._stats.workers_died += 1 + worker.dead_reason = ( + f"worker exited unexpectedly (rc={worker.returncode})" + ) + self._workers = alive + + # --- the kill domain --- + + def recycle(self, worker: Worker, reason: str) -> None: + """Kill *worker* and drop it. Everything it hosted goes with it. + + This is what the whole three-level design exists for: a cell that will + not stop is reclaimed by killing the process it lives in. The cost is + the other cells in that worker, which is why ``tenant_isolation`` + defaults to on -- so that cost lands on one tenant. + """ + logger.warning( + "recycling fabric worker %s (pid=%s, tenant=%r): %s", + worker.id, + worker.pid, + worker.tenant, + reason, + ) + worker.kill(reason) + with self._lock: + if worker in self._workers: + self._workers.remove(worker) + self._stats.workers_killed += 1 + + # --- cells --- + + def open_cell( + self, + *, + tenant: Optional[str] = None, + allowed_imports: Optional[list[str]] = None, + preimport: Optional[list[str]] = None, + ) -> tuple[Worker, str]: + worker = self.placement_for(tenant) + cell_id = f"c{next(self._cell_ids)}" + try: + worker.open_cell( + cell_id, allowed_imports=allowed_imports, preimport=preimport + ) + except WorkerDied: + # The worker died between placement and the open. Retire it and let + # the caller land on a fresh one rather than surfacing a race. + self.recycle(worker, "worker died during cell open") + worker = self.placement_for(tenant) + worker.open_cell( + cell_id, allowed_imports=allowed_imports, preimport=preimport + ) + with self._lock: + self._stats.cells_opened += 1 + return worker, cell_id + + def run( + self, + worker: Worker, + cell_id: str, + frame: dict[str, Any], + *, + timeout: Optional[float], + ) -> None: + """Run one operation, recycling the worker if the guest never returns.""" + try: + worker.run(cell_id, frame, timeout=timeout) + except TimeoutError as exc: + self.recycle( + worker, f"cell {cell_id} exceeded its deadline; killing the worker" + ) + raise errors.WallTimeExceeded( + f"cell {cell_id} exceeded {timeout}s and did not return. A " + "running sub-interpreter cannot be reclaimed, so the worker " + "process hosting it was killed; cells sharing that worker were " + "lost with it." + ) from exc + + def close_cell(self, worker: Worker, cell_id: str) -> None: + try: + worker.close_cell(cell_id) + finally: + with self._lock: + self._stats.cells_closed += 1 + + # --- observability --- + + def stats(self) -> dict[str, Any]: + with self._lock: + self._reap_dead_locked() + workers = list(self._workers) + counters = self._stats.as_dict() + counters["workers_live"] = len(workers) + counters["cells_live"] = sum(w.cell_count() for w in workers) + counters["tenant_isolation"] = self.tenant_isolation + return counters + + def worker_report(self) -> list[dict[str, Any]]: + """Per-worker placement view, for dashboards and admission checks.""" + with self._lock: + workers = list(self._workers) + return [ + { + "id": w.id, + "pid": w.pid, + "tenant": w.tenant, + "cells": w.cell_count(), + "alive": w.is_alive(), + "uptime_s": round(time.monotonic() - w.started_at, 3), + } + for w in workers + ] + + # --- teardown --- + + def close(self) -> None: + with self._lock: + self._closed = True + workers = list(self._workers) + self._workers.clear() + for worker in workers: + worker.stop() + + def __enter__(self) -> "WorkerPool": + return self + + def __exit__(self, *exc: object) -> None: + self.close() + + +class FabricSandbox: + """Sandbox handle over a cell hosted in a worker process.""" + + def __init__( + self, + name: str, + *, + pool: WorkerPool, + allowed_imports: Optional[list[str]] = None, + preimport: Optional[list[str]] = None, + tenant: Optional[str] = None, + wall_time_ms: Optional[int] = None, + ) -> None: + subinterpreter.require_available_for_fabric() + self.name = name + self._pool = pool + self._tenant = tenant + self._allowed_imports = list(allowed_imports or []) + self.wall_time_ms = wall_time_ms + self._lock = threading.Lock() + self._posted: list[Any] = [] + self._logs: list[Any] = [] + self._metrics: list[Any] = [] + self._requests: list[Any] = [] + self._backend = "fabric" + self._closed = False + # Handle-surface attributes the Sandbox wrapper reads directly. + self._cgroup_path: Optional[str] = None + self._quarantine_reason: Optional[str] = None + self.termination_reason: Optional[str] = None + self.quota_enforcement = "worker_kill" + self._worker, self._cell_id = pool.open_cell( + tenant=tenant, allowed_imports=self._allowed_imports, preimport=preimport + ) + + # --- the cell ABI --- + + def exec(self, src: str) -> None: + self._run({"op": "exec", "source": src}) + + def call( + self, + func: str, + *args: Any, + timeout: Optional[float] = None, + **kwargs: Any, + ) -> Any: + if not isinstance(func, str) or not func: + raise TypeError("func must be a dotted name") + if "." not in func: + raise ValueError(f"expected a dotted name, got {func!r}") + self._run( + { + "op": "call", + "target": func, + "args": list(args), + "kwargs": kwargs, + }, + timeout=timeout, + ) + with self._lock: + if not self._posted: + raise errors.SandboxError(f"call to {func!r} returned no result") + return self._posted.pop() + + def recv(self, timeout: Optional[float] = None) -> Any: + deadline = time.monotonic() + (timeout or 0.0) + while True: + self._collect() + with self._lock: + if self._posted: + return self._posted.pop(0) + if timeout is None or time.monotonic() >= deadline: + raise errors.TimeoutError(f"no message from sandbox '{self.name}'") + time.sleep(0.002) + + # --- supervisor surface --- + + def is_alive(self) -> bool: + return not self._closed and self._worker.is_alive() + + def kill(self, timeout: float = 0.2) -> bool: + """Kill the worker hosting this cell. Unlike a cell, this always works.""" + del timeout + if self._closed: + return True + self._pool.recycle(self._worker, f"sandbox {self.name} killed") + self._closed = True + return True + + def cancel(self, timeout: float = 0.2) -> bool: + return self.kill(timeout) + + def stop(self, timeout: float = 0.2) -> None: + self.close(timeout) + + def close(self, timeout: float = 0.2) -> None: + del timeout + if self._closed: + return + self._closed = True + self._collect() + self._pool.close_cell(self._worker, self._cell_id) + + def reap(self) -> bool: + self.close() + return True + + def quarantine(self, reason: str = "manual quarantine") -> None: + self._quarantine_reason = reason + self.kill() + + def stats(self) -> dict[str, Any]: + return { + "backend": self._backend, + "tenant": self._tenant, + "worker": self._worker.id, + "worker_pid": self._worker.pid, + "cell": self._cell_id, + "posted": len(self._posted), + "requests": len(self._requests), + "fabric": self._pool.stats(), + } + + def profile(self) -> dict[str, Any]: + return self.stats() + + def snapshot(self) -> dict[str, Any]: + """Configuration, not guest state: an interpreter cannot be captured.""" + return { + "name": self.name, + "backend": self._backend, + "tenant": self._tenant, + "allowed_imports": sorted(self._allowed_imports), + "wall_time_ms": self.wall_time_ms, + } + + def get_broker_requests(self) -> list[Any]: + self._collect() + with self._lock: + return list(self._requests) + + def reset_config(self) -> dict[str, Any]: + raise NotImplementedError( + "a sub-interpreter cannot be reset to a pristine state, so a cell " + "is never reused across tenants. Close this sandbox and spawn " + "another; the worker keeps a warm cell so that costs ~1ms." + ) + + def reset(self, *args: Any, **kwargs: Any) -> None: + self.reset_config() + + def enable_tracing(self) -> None: + raise NotImplementedError( + "syscall tracing is a process-backend feature; a fabric worker runs " + "many tenants' cells and has no per-cell syscall boundary to trace" + ) + + def get_syscall_log(self) -> list[str]: + return [] + + def get_denial_events(self) -> list[dict[str, str]]: + return [] + + def __enter__(self) -> "FabricSandbox": + return self + + def __exit__(self, *exc: object) -> None: + self.close() + + # --- internals --- + + def _run(self, frame: dict[str, Any], timeout: Optional[float] = None) -> None: + if self._closed: + raise errors.SandboxError(f"sandbox '{self.name}' is closed") + limit = timeout + if limit is None and self.wall_time_ms is not None: + limit = self.wall_time_ms / 1000.0 + try: + self._pool.run(self._worker, self._cell_id, frame, timeout=limit) + except (errors.WallTimeExceeded, WorkerDied): + self._closed = True + self._collect() + raise + finally: + self._collect() + + def _collect(self) -> None: + try: + messages = self._worker.drain(self._cell_id) + except Exception: # pragma: no cover - worker gone mid-drain + return + for kind, payload in messages: + with self._lock: + if kind == "post": + self._posted.append(payload) + elif kind == "log": + self._logs.append(payload) + elif kind == "metric": + self._metrics.append(payload) + elif kind == "request": + self._requests.append(payload) + + +#: Process-wide fabric, created lazily so importing pyisolate on a build without +#: sub-interpreters costs nothing. +_default_pool: Optional[WorkerPool] = None +_default_pool_lock = threading.Lock() + + +def default_pool() -> WorkerPool: + global _default_pool + with _default_pool_lock: + if _default_pool is None: + subinterpreter.require_available_for_fabric() + _default_pool = WorkerPool() + return _default_pool + + +def reset_default_pool() -> None: + """Drop the process-wide fabric. Used by tests and supervisor shutdown.""" + global _default_pool + with _default_pool_lock: + pool, _default_pool = _default_pool, None + if pool is not None: + pool.close() + + +__all__ = [ + "FabricSandbox", + "FabricStats", + "Worker", + "WorkerDied", + "WorkerPool", + "default_pool", + "reset_default_pool", +] diff --git a/pyisolate/runtime/subinterpreter.py b/pyisolate/runtime/subinterpreter.py index 5e2a425..7086b74 100644 --- a/pyisolate/runtime/subinterpreter.py +++ b/pyisolate/runtime/subinterpreter.py @@ -97,6 +97,26 @@ def require_available() -> None: ) +def require_available_for_fabric() -> None: + """Same requirement as a cell, phrased for ``backend="fabric"``. + + The fabric puts cells in worker processes so a runaway one can be killed; + it cannot invent cells on a build that has none, and degrading to threads + would give the caller a different isolation model under the same name. + """ + if is_available(): + return + running = f"{sys.version_info[0]}.{sys.version_info[1]}" + needed = f"{MIN_PYTHON[0]}.{MIN_PYTHON[1]}" + raise errors.SandboxError( + f"backend='fabric' needs CPython {needed}+ for concurrent.interpreters; " + f"this is {running}. The fabric hosts sub-interpreter cells in worker " + "processes, so it needs the same interpreter support a cell does. Use " + "backend='process' for one confined process per sandbox, which is a " + "real boundary and works on every supported Python." + ) + + # --- what a warm cell contains -------------------------------------------- @@ -928,5 +948,6 @@ def reset_default_pool() -> None: "default_pool", "is_available", "require_available", + "require_available_for_fabric", "reset_default_pool", ] diff --git a/pyisolate/supervisor.py b/pyisolate/supervisor.py index 013deb3..51a58ca 100644 --- a/pyisolate/supervisor.py +++ b/pyisolate/supervisor.py @@ -25,6 +25,7 @@ from .observability.alerts import AlertManager from .observability.trace import Tracer from .policy import resolve_policy +from .runtime import fabric as _fabric from .runtime import microvm as _microvm from .runtime import subinterpreter from .runtime.process_backend import ProcessSandbox @@ -44,17 +45,19 @@ DEFAULT_NAME_PATTERN = re.compile(r"^[A-Za-z0-9_-]+$") NAME_PATTERN = DEFAULT_NAME_PATTERN -BackendMode = Literal["thread", "subinterpreter", "process", "microvm"] +BackendMode = Literal["thread", "subinterpreter", "fabric", "process", "microvm"] DEFAULT_BACKEND: BackendMode = "thread" SUPPORTED_BACKENDS: tuple[BackendMode, ...] = ( "thread", "subinterpreter", + "fabric", "process", "microvm", ) IMPLEMENTED_BACKENDS: tuple[BackendMode, ...] = ( "thread", "subinterpreter", + "fabric", "process", ) @@ -126,7 +129,10 @@ def _require_implemented_backend(backend: BackendMode) -> None: #: classes that implement the same cell ABI rather than a shared base, so the #: union is the type. BackendSandbox = Union[ - "SandboxThread", "ProcessSandbox", "subinterpreter.SubinterpreterSandbox" + "SandboxThread", + "ProcessSandbox", + "subinterpreter.SubinterpreterSandbox", + "_fabric.FabricSandbox", ] @@ -274,6 +280,10 @@ def __init__( # Created on first use so that importing pyisolate on a build without # concurrent.interpreters costs nothing. self._pool: Optional["subinterpreter.CellPool"] = None + # Fabric cells live in worker processes and get their own registry for + # the same reason: they are neither threads nor ProcessSandbox. + self._fabric_sandboxes: Dict[str, "_fabric.FabricSandbox"] = {} + self._worker_pool: Optional["_fabric.WorkerPool"] = None self._lock = threading.Lock() self._alerts = AlertManager() self._tracer = Tracer() @@ -428,6 +438,15 @@ def spawn( wall_time_ms=wall_time_ms, ) + if backend == "fabric": + return self._spawn_fabric( + name, + allowed_imports=allowed_imports, + wall_time_ms=wall_time_ms, + tenant=tenant, + mem_bytes=mem_bytes, + ) + if backend == "process": return self._spawn_process( name, @@ -611,6 +630,51 @@ def _spawn_subinterpreter( self._cell_sandboxes[name] = cell_sandbox return Sandbox(cell_sandbox, self) + def _spawn_fabric( + self, + name: str, + *, + allowed_imports: Optional[list[str]] = None, + wall_time_ms: Optional[int] = None, + tenant: Optional[str] = None, + mem_bytes: Optional[int] = None, + ) -> Sandbox: + """Spawn a cell inside a worker process that the supervisor can kill. + + This is the boundary-less-but-recoverable mode: the cell isolates + namespaces and the worker is the kill domain, so a guest that will not + stop costs its worker rather than being unreclaimable. + """ + subinterpreter.require_available_for_fabric() + sandbox = _fabric.FabricSandbox( + name, + pool=self._fabric_pool(mem_bytes), + allowed_imports=allowed_imports, + tenant=tenant, + wall_time_ms=wall_time_ms, + ) + with self._lock: + existing = self._fabric_sandboxes.get(name) + if existing is not None and existing.is_alive(): + sandbox.close() + raise RuntimeError(f"sandbox '{name}' already exists") + self._fabric_sandboxes[name] = sandbox + return Sandbox(sandbox, self) + + def _fabric_pool(self, mem_bytes: Optional[int] = None) -> "_fabric.WorkerPool": + with self._lock: + if self._worker_pool is None: + self._worker_pool = _fabric.WorkerPool(worker_mem_bytes=mem_bytes) + return self._worker_pool + + def fabric_report(self) -> dict[str, Any]: + """Placement and recycling view of the fabric, for dashboards.""" + with self._lock: + pool = self._worker_pool + if pool is None: + return {"workers": [], "stats": {}} + return {"workers": pool.worker_report(), "stats": pool.stats()} + def _cell_pool(self) -> "subinterpreter.CellPool": with self._lock: if self._pool is None: @@ -798,15 +862,22 @@ def shutdown(self, cap: RootCapability = ROOT) -> None: procs = list(self._process_sandboxes.values()) cells = list(self._cell_sandboxes.values()) self._cell_sandboxes.clear() + fabric_cells = list(self._fabric_sandboxes.values()) + self._fabric_sandboxes.clear() pool, self._pool = self._pool, None + worker_pool, self._worker_pool = self._worker_pool, None for sb in sandboxes + warm: sb.stop() for proc in procs: proc.stop() for cell in cells: cell.stop() + for fabric_cell in fabric_cells: + fabric_cell.stop() if pool is not None: pool.close() + if worker_pool is not None: + worker_pool.close() self._cleanup() def quarantine(self, name: str, reason: str) -> None: @@ -886,6 +957,11 @@ def _cleanup(self) -> None: ] for n in dead_cells: del self._cell_sandboxes[n] + dead_fabric = [ + n for n, c in self._fabric_sandboxes.items() if not c.is_alive() + ] + for n in dead_fabric: + del self._fabric_sandboxes[n] _supervisor: Supervisor | None = None diff --git a/tests/test_fabric.py b/tests/test_fabric.py new file mode 100644 index 0000000..4c259a8 --- /dev/null +++ b/tests/test_fabric.py @@ -0,0 +1,499 @@ +"""Tests for the fabric: cells in worker processes that can actually be killed. + +Note what these tests do *not* need. The in-process cell backend's runaway +tests have to shell out to a fresh interpreter, because a stranded cell pins a +thread forever and would hang pytest at exit. Here the runaway tests run +inline: the fabric kills the worker, so the strand dies with it. That +difference is the whole point of this layer, and it shows up in the shape of +the test file. +""" + +import os +import sys +import time + +import pytest + +from pyisolate import errors +from pyisolate.runtime import fabric as F +from pyisolate.runtime import subinterpreter as S + +requires_interpreters = pytest.mark.skipif( + not S.is_available(), + reason=f"needs CPython {S.MIN_PYTHON[0]}.{S.MIN_PYTHON[1]}+ for cells", +) + + +@pytest.fixture +def pool(): + """A small fabric, always torn down so no worker outlives the test.""" + p = F.WorkerPool(max_workers=4, cells_per_worker=4, warm_per_spec=0) + try: + yield p + finally: + p.close() + + +def _wait_gone(pid, timeout=5.0): + """True once *pid* is no longer running. + + ``os.kill(pid, 0)`` still succeeds for a zombie -- a process that has + exited but has not been waited on yet -- so a signal probe alone would + report a killed worker as alive. Read the state out of procfs and treat + ``Z`` as gone, so the test measures "stopped running" rather than "already + collected", which is the fabric's actual guarantee. + """ + deadline = time.monotonic() + timeout + while time.monotonic() < deadline: + try: + with open(f"/proc/{pid}/stat", encoding="utf-8") as fh: + # "pid (comm) state ..." -- comm can contain spaces and parens, + # so split after the last ')'. + state = fh.read().rsplit(")", 1)[1].split()[0] + except (FileNotFoundError, ProcessLookupError, PermissionError, IndexError): + return True + if state == "Z": + return True + time.sleep(0.01) + return False + + +# --- availability --------------------------------------------------------- + + +def test_is_available_tracks_the_running_build(): + expected = sys.version_info[:2] >= S.MIN_PYTHON + assert S.is_available() is expected + + +@pytest.mark.skipif(S.is_available(), reason="this build has sub-interpreters") +def test_fabric_fails_closed_naming_the_process_backend(): + """A fabric cannot invent cells; it says what to use instead.""" + with pytest.raises(errors.SandboxError) as excinfo: + S.require_available_for_fabric() + message = str(excinfo.value) + assert "backend='fabric'" in message + assert "3.14" in message + assert "backend='process'" in message + + +# --- pool configuration (no worker needed) -------------------------------- + + +def test_pool_rejects_nonsense_sizes(): + with pytest.raises(ValueError, match="max_workers"): + F.WorkerPool(max_workers=0) + with pytest.raises(ValueError, match="cells_per_worker"): + F.WorkerPool(cells_per_worker=0) + + +def test_rebuild_never_resolves_an_arbitrary_name(): + """Only the exception's *name* crosses the boundary, so it is mapped. + + Looking the name up dynamically would let a worker -- which runs guest code + -- name any class in the supervisor and have it constructed here. + """ + assert isinstance(F._rebuild("ValueError", "x"), ValueError) + assert isinstance(F._rebuild("PolicyError", "x"), errors.PolicyError) + exc = F._rebuild("os.system", "rm -rf /") + assert isinstance(exc, errors.SandboxError) + assert "os.system" in str(exc) + + +# --- the cell ABI through a worker ---------------------------------------- + + +@requires_interpreters +def test_exec_and_recv_round_trip(pool): + sandbox = F.FabricSandbox("demo", pool=pool, allowed_imports=["math"]) + try: + sandbox.exec("from math import sqrt; post(sqrt(2))") + assert sandbox.recv(timeout=10) == pytest.approx(1.4142135623730951) + finally: + sandbox.close() + + +@requires_interpreters +def test_call_resolves_a_dotted_name_in_the_cell(pool): + sandbox = F.FabricSandbox("demo", pool=pool, allowed_imports=["math"]) + try: + assert sandbox.call("math.factorial", 5) == 120 + finally: + sandbox.close() + + +@requires_interpreters +def test_call_rejects_a_bare_name(pool): + sandbox = F.FabricSandbox("demo", pool=pool) + try: + with pytest.raises(ValueError, match="dotted name"): + sandbox.call("factorial") + with pytest.raises(TypeError, match="dotted name"): + sandbox.call(None) # type: ignore[arg-type] + finally: + sandbox.close() + + +@requires_interpreters +def test_the_import_allow_list_is_enforced_inside_the_worker(pool): + sandbox = F.FabricSandbox("demo", pool=pool, allowed_imports=["math"]) + try: + with pytest.raises(Exception, match="not permitted by policy"): + sandbox.exec("import os") + finally: + sandbox.close() + + +@requires_interpreters +def test_log_metric_and_request_reach_the_supervisor(pool): + sandbox = F.FabricSandbox("demo", pool=pool) + try: + sandbox.exec( + "log('info', 'hello', k=1)\n" + "metric('m', 3)\n" + "request('read_path', '/etc/hosts')\n" + ) + assert sandbox.get_broker_requests() == [ + {"capability": "read_path", "args": ["/etc/hosts"], "kwargs": {}} + ] + finally: + sandbox.close() + + +@requires_interpreters +def test_a_guest_error_surfaces_as_the_matching_exception(pool): + sandbox = F.FabricSandbox("demo", pool=pool) + try: + with pytest.raises(Exception) as excinfo: + sandbox.exec("raise ValueError('boom')") + assert "boom" in str(excinfo.value) + finally: + sandbox.close() + + +@requires_interpreters +def test_worker_reports_the_build_it_is_running(pool): + sandbox = F.FabricSandbox("demo", pool=pool) + try: + info = sandbox._worker.info + assert isinstance(info["pid"], int) + assert info["python"].startswith("3.") + assert isinstance(info["free_threaded"], bool) + finally: + sandbox.close() + + +# --- placement ------------------------------------------------------------ + + +@requires_interpreters +def test_tenants_do_not_share_a_worker_by_default(pool): + """Fate-sharing is the cost of packing tenants together, so it is opt-in.""" + a = F.FabricSandbox("a", pool=pool, tenant="acme") + b = F.FabricSandbox("b", pool=pool, tenant="globex") + try: + assert a._worker.id != b._worker.id + assert a._worker.tenant == "acme" + assert b._worker.tenant == "globex" + finally: + a.close() + b.close() + + +@requires_interpreters +def test_one_tenants_cells_share_a_worker_until_it_is_full(): + p = F.WorkerPool(max_workers=4, cells_per_worker=2, warm_per_spec=0) + try: + cells = [F.FabricSandbox(f"c{i}", pool=p, tenant="acme") for i in range(3)] + try: + workers = {c._worker.id for c in cells} + # Two fit in the first worker; the third opens a second one. + assert len(workers) == 2 + finally: + for c in cells: + c.close() + finally: + p.close() + + +@requires_interpreters +def test_tenant_isolation_off_packs_them_together(): + p = F.WorkerPool( + max_workers=4, cells_per_worker=4, tenant_isolation=False, warm_per_spec=0 + ) + try: + a = F.FabricSandbox("a", pool=p, tenant="acme") + b = F.FabricSandbox("b", pool=p, tenant="globex") + try: + assert a._worker.id == b._worker.id + assert p.stats()["tenant_isolation"] is False + finally: + a.close() + b.close() + finally: + p.close() + + +@requires_interpreters +def test_capacity_is_refused_with_an_actionable_message(): + p = F.WorkerPool(max_workers=1, cells_per_worker=1, warm_per_spec=0) + try: + first = F.FabricSandbox("a", pool=p, tenant="acme") + try: + with pytest.raises(errors.SandboxError, match="at capacity"): + F.FabricSandbox("b", pool=p, tenant="globex") + finally: + first.close() + finally: + p.close() + + +# --- the kill domain ------------------------------------------------------ + + +@requires_interpreters +def test_a_runaway_cell_is_reclaimed_by_killing_its_worker(pool): + """The thing that is impossible in-process. + + In ``backend="subinterpreter"`` this cell would be abandoned and its thread + pinned for the life of the process. Here the deadline kills the worker, and + the strand dies with it. + """ + sandbox = F.FabricSandbox("runaway", pool=pool, tenant="acme", wall_time_ms=400) + pid = sandbox._worker.pid + + with pytest.raises(errors.WallTimeExceeded) as excinfo: + sandbox.exec("while True:\n pass") + + message = str(excinfo.value) + assert "cannot be reclaimed" in message + assert "killed" in message + assert _wait_gone(pid), "the worker process outlived the recycle" + assert sandbox.is_alive() is False + assert pool.stats()["workers_killed"] == 1 + + +@requires_interpreters +def test_a_runaway_does_not_touch_another_tenant(pool): + victim = F.FabricSandbox("runaway", pool=pool, tenant="acme", wall_time_ms=400) + bystander = F.FabricSandbox( + "fine", pool=pool, tenant="globex", allowed_imports=["math"] + ) + try: + assert victim._worker.id != bystander._worker.id + with pytest.raises(errors.WallTimeExceeded): + victim.exec("while True:\n pass") + + # The bystander's worker was never touched. + assert bystander.is_alive() + bystander.exec("from math import sqrt; post(sqrt(9))") + assert bystander.recv(timeout=10) == 3.0 + finally: + bystander.close() + + +@requires_interpreters +def test_the_tenant_can_keep_working_after_its_worker_is_recycled(pool): + victim = F.FabricSandbox("runaway", pool=pool, tenant="acme", wall_time_ms=400) + with pytest.raises(errors.WallTimeExceeded): + victim.exec("while True:\n pass") + + replacement = F.FabricSandbox("after", pool=pool, tenant="acme") + try: + replacement.exec("post(1 + 1)") + assert replacement.recv(timeout=10) == 2 + finally: + replacement.close() + + +@requires_interpreters +def test_killing_a_sandbox_kills_its_worker(pool): + """Unlike a cell's kill(), this one can actually deliver.""" + sandbox = F.FabricSandbox("demo", pool=pool, tenant="acme") + pid = sandbox._worker.pid + assert sandbox.kill() is True + assert _wait_gone(pid) + assert sandbox.is_alive() is False + + +@requires_interpreters +def test_quarantine_kills_the_worker_and_records_the_reason(pool): + sandbox = F.FabricSandbox("demo", pool=pool, tenant="acme") + pid = sandbox._worker.pid + sandbox.quarantine("policy breach") + assert sandbox._quarantine_reason == "policy breach" + assert _wait_gone(pid) + + +@requires_interpreters +def test_a_worker_dying_under_us_surfaces_rather_than_hanging(pool): + """An OOM kill or a segfault must not leave a caller blocked forever.""" + sandbox = F.FabricSandbox("demo", pool=pool, tenant="acme") + pid = sandbox._worker.pid + os.kill(pid, 9) + assert _wait_gone(pid) + + with pytest.raises(errors.SandboxError): + sandbox.exec("post(1)") + assert sandbox.is_alive() is False + + +@requires_interpreters +def test_a_dead_worker_is_reaped_out_of_the_placement_set(pool): + sandbox = F.FabricSandbox("demo", pool=pool, tenant="acme") + pid = sandbox._worker.pid + os.kill(pid, 9) + assert _wait_gone(pid) + # The worker is marked unusable when its channel hits EOF, which the reader + # thread notices a moment after the process dies. + deadline = time.monotonic() + 5.0 + while sandbox._worker.is_alive() and time.monotonic() < deadline: + time.sleep(0.01) + + stats = pool.stats() + assert stats["workers_live"] == 0 + assert stats["workers_died"] == 1 + + # A new sandbox for the same tenant gets a fresh worker. + replacement = F.FabricSandbox("after", pool=pool, tenant="acme") + try: + assert replacement._worker.pid != pid + replacement.exec("post('alive')") + assert replacement.recv(timeout=10) == "alive" + finally: + replacement.close() + + +# --- observability -------------------------------------------------------- + + +@requires_interpreters +def test_worker_report_shows_placement(pool): + a = F.FabricSandbox("a", pool=pool, tenant="acme") + b = F.FabricSandbox("b", pool=pool, tenant="globex") + try: + report = {row["tenant"]: row for row in pool.worker_report()} + assert set(report) == {"acme", "globex"} + for row in report.values(): + assert row["cells"] == 1 + assert row["alive"] is True + assert row["uptime_s"] >= 0 + finally: + a.close() + b.close() + + +@requires_interpreters +def test_sandbox_stats_name_the_worker_hosting_the_cell(pool): + sandbox = F.FabricSandbox("demo", pool=pool, tenant="acme") + try: + stats = sandbox.stats() + assert stats["backend"] == "fabric" + assert stats["tenant"] == "acme" + assert stats["worker"] == sandbox._worker.id + assert stats["worker_pid"] == sandbox._worker.pid + assert stats["fabric"]["workers_live"] >= 1 + finally: + sandbox.close() + + +# --- the surface the Sandbox handle delegates to -------------------------- + + +@requires_interpreters +def test_snapshot_carries_configuration_not_guest_state(pool): + sandbox = F.FabricSandbox( + "demo", + pool=pool, + allowed_imports=["math"], + tenant="acme", + wall_time_ms=1000, + ) + try: + assert sandbox.snapshot() == { + "name": "demo", + "backend": "fabric", + "tenant": "acme", + "allowed_imports": ["math"], + "wall_time_ms": 1000, + } + finally: + sandbox.close() + + +@requires_interpreters +def test_reset_is_refused_with_the_reason(pool): + sandbox = F.FabricSandbox("demo", pool=pool) + try: + with pytest.raises(NotImplementedError, match="cannot be reset"): + sandbox.reset() + finally: + sandbox.close() + + +@requires_interpreters +def test_tracing_is_refused_rather_than_silently_empty(pool): + sandbox = F.FabricSandbox("demo", pool=pool) + try: + with pytest.raises(NotImplementedError, match="no per-cell syscall"): + sandbox.enable_tracing() + finally: + sandbox.close() + + +@requires_interpreters +def test_using_a_closed_sandbox_is_refused(pool): + sandbox = F.FabricSandbox("demo", pool=pool) + sandbox.close() + with pytest.raises(errors.SandboxError, match="closed"): + sandbox.exec("post(1)") + + +# --- pre-warming ---------------------------------------------------------- + + +@requires_interpreters +def test_prewarm_starts_workers_before_the_first_request(): + """A worker costs ~160ms to spawn; a tenant should not pay that inline.""" + p = F.WorkerPool(max_workers=3, cells_per_worker=4, warm_per_spec=0) + try: + assert p.prewarm("acme", 2) == 2 + assert p.stats()["workers_live"] == 2 + + started_before = p.stats()["workers_started"] + sandbox = F.FabricSandbox("demo", pool=p, tenant="acme") + try: + # The cell landed on an existing worker rather than spawning one. + assert p.stats()["workers_started"] == started_before + sandbox.exec("post('warm')") + assert sandbox.recv(timeout=10) == "warm" + finally: + sandbox.close() + finally: + p.close() + + +@requires_interpreters +def test_prewarm_respects_max_workers(): + p = F.WorkerPool(max_workers=1, cells_per_worker=4, warm_per_spec=0) + try: + assert p.prewarm("acme", 5) == 1 + assert p.stats()["workers_live"] == 1 + finally: + p.close() + + +@requires_interpreters +def test_prewarmed_workers_respect_tenant_isolation(): + """A warm worker for one tenant is not a warm worker for another.""" + p = F.WorkerPool(max_workers=3, cells_per_worker=4, warm_per_spec=0) + try: + p.prewarm("acme", 1) + sandbox = F.FabricSandbox("demo", pool=p, tenant="globex") + try: + assert sandbox._worker.tenant == "globex" + assert p.stats()["workers_live"] == 2 + finally: + sandbox.close() + finally: + p.close() diff --git a/tests/test_supervisor.py b/tests/test_supervisor.py index 5fea695..50e24ef 100644 --- a/tests/test_supervisor.py +++ b/tests/test_supervisor.py @@ -4,11 +4,13 @@ ROOT = Path(__file__).resolve().parents[1] sys.path.insert(0, str(ROOT)) +import os import warnings import pytest import pyisolate as iso +from pyisolate import supervisor as supervisor_mod from pyisolate.bpf.manager import BPFManager from pyisolate.runtime import subinterpreter @@ -191,10 +193,16 @@ def test_spawn_backend_is_explicit_thread(): assert iso.SUPPORTED_BACKENDS == ( "thread", "subinterpreter", + "fabric", "process", "microvm", ) - assert iso.IMPLEMENTED_BACKENDS == ("thread", "subinterpreter", "process") + assert iso.IMPLEMENTED_BACKENDS == ( + "thread", + "subinterpreter", + "fabric", + "process", + ) assert iso.DEFAULT_BACKEND == "thread" finally: sb.close() @@ -242,6 +250,57 @@ def test_subinterpreter_backend_fails_closed_below_314(): iso.spawn("backend-cell", backend="subinterpreter") +def test_fabric_is_an_implemented_backend(): + assert "fabric" in iso.SUPPORTED_BACKENDS + assert "fabric" in iso.IMPLEMENTED_BACKENDS + + +@pytest.mark.skipif( + not subinterpreter.is_available(), + reason="needs CPython 3.14+ for cells", +) +def test_fabric_backend_runs_a_cell_in_a_worker_process(): + sb = iso.spawn( + "backend-fabric", + backend="fabric", + allowed_imports=["math"], + tenant="acme", + ) + try: + assert sb.backend == "fabric" + sb.exec("from math import sqrt; post(sqrt(4))") + assert sb.recv(timeout=10) == 2.0 + stats = sb.profile() + # The cell runs in another process, which is what makes it reclaimable. + assert stats["worker_pid"] != os.getpid() + assert stats["tenant"] == "acme" + finally: + sb.close() + + +@pytest.mark.skipif( + not subinterpreter.is_available(), + reason="needs CPython 3.14+ for cells", +) +def test_supervisor_reports_fabric_placement(): + sb = iso.spawn("backend-fabric-report", backend="fabric", tenant="acme") + try: + report = supervisor_mod._get_supervisor().fabric_report() + assert report["stats"]["workers_live"] >= 1 + assert any(row["tenant"] == "acme" for row in report["workers"]) + finally: + sb.close() + + +@pytest.mark.skipif( + subinterpreter.is_available(), + reason="this build has sub-interpreters", +) +def test_fabric_backend_fails_closed_below_314(): + with pytest.raises(iso.SandboxError, match="backend='fabric'"): + iso.spawn("backend-fabric", backend="fabric") + + def test_unknown_backend_still_rejects_without_warning(): with warnings.catch_warnings(): warnings.simplefilter("error", DeprecationWarning)