Repository navigation
feat(py): add Worker, which runs model code in a persistent sandboxed Python session - #387
Merged
Merged
Conversation
Contributor
|
Preview root: https://posit-dev.github.io/commons/pr-387/ Python site preview: https://posit-dev.github.io/commons/pr-387/py/ Built from the latest commit on this branch. The R links in it point at the published R site, which no pull request rebuilds. |
jat255
added this pull request to stack #390
September 23, 2026 17:44
jat255
force-pushed
the
jat255/bgk1-worker-lifecycle
branch
from
September 24, 2026 02:00
71f01fd to
bc9ae64
Compare
jat255
force-pushed
the
jat255/bgk1-worker-lifecycle
branch
2 times, most recently
from
October 7, 2026 19:20
9d3e1bd to
20b6f3d
Compare
The driver owns everything around a call, mirroring run_r_tool() and worker_await() in pkg-r/R/run-r.R: no process exists until the first call, an asyncio.Lock runs one call at a time, a call past its time limit is interrupted with SIGINT first and killed only if the interrupt goes unanswered (an interrupted session keeps its variables; a restarted one has lost them, and the two messages say so), and a worker idle for ten minutes is reaped and respawned on the next call. A pending counter keeps the reaper from closing a worker with a call in flight. Handles sync per call from the store's high-water mark. Measure sources are exec'd as text at spawn, so a respawned worker is defined exactly the way the first one was. The harvest in _measures now keeps the definitions exec-ready — dedented, decorator lines dropped — and the driver compiles them with annotations deferred, since a measure's own imports do not exist in the worker's session. The worker entry point claims the protocol channel at the fd level: duplicates of fds 0 and 1 become the channel and the model-visible fds point at sinks for the process's whole life, so model code writing straight to the descriptor (os.write, a C extension, a thread still printing after its call returned) or reading stdin from a thread lands in the sink rather than on the channel. Startup applies the address-space limit and engages the sandbox before the ready announcement and any model-written code; fd 2 stays attached until then so a failed start can still reach the driver's diagnostics.
…se() The spawn failure path waited on every tracked shutdown while holding the call lock. aclose() tracks its own close task there, and that task waits for the same lock, so a close racing a failing spawn never finished. The failure path now waits only for its own stderr drain.
Resolving symlinked entries one level down turns /bin's pkill into /usr/bin/pgrep, a regular file. Landlock refuses a rule that names directory rights on a file with EINVAL, so the worker could not start on a stock Ubuntu host. A rule on a non-directory is now masked to the file rights.
On macOS, 3.14's asyncio takes waitid()'s report of a stopped child for an exit and blocks the event loop in waitpid(), so the SIGSTOP test hangs there; it is skipped on that combination until the driver stops relying on asyncio's child watcher. The crashed-children test checked the process group the instant the call returned, while the killed child was still finishing its exit; it now allows two seconds for the group to go.
…o parse textwrap.dedent finds no common indent when a nested measure holds a string whose lines sit at column 0, such as a SQL query, so ast.parse raised IndentationError and semantic_layer() failed. The indent is now removed from code lines only, and a line that begins inside a multi-line string keeps its text, so the string also keeps its value. Source that still does not parse is kept as the "source unavailable" comment, because the worker defines every harvested source at spawn and a fragment would fail every spawn.
The harvest keeps every module-level function, lambdas included, but the define shim guarded defaults on def statements only. A helper such as `scale = lambda x, k=FACTOR: x * k` raised NameError at define time, and a failed define fails the spawn, so every call failed to start. The shim now guards every def or lambda whose defaults are evaluated as the source is defined, and leaves function bodies alone.
The driver spawned its own subprocess and carried copies of what _backend.py already had: the stderr tail, the shielded shutdown set, and the process-group signalling. The seam's one-shot exec() could not host a persistent session, so nothing used it. ExecBackend now starts a worker session instead, and LocalBackend owns everything about where the process runs: the scratch directory, the environment, the protection mode, the spawn, the interrupt, and the shutdown. Worker keeps the lifecycle policy and the messages the model sees, so a container-hosted backend is a new implementation of the seam rather than an edit to the driver. exec() had no caller and is removed; the guarantees its tests pinned (the allowlisted environment, kill escalation, group kill, cancellation) are now tested through start(). A close that is cancelled still finishes, and the cancellation now reaches the caller instead of being absorbed.
…n Linux Worker passes its network access to the backend, which nothing checked; the worker's own argv shows it. The dead-leader test's child held the inherited stdio open, and asyncio's wait() does not return until every pipe closes, so on Linux the child finished before the close could reach it. The child now lets go of the pipes, as the real worker's fds are sinks.
The driver's SIGINT is meant for the call it timed out, but it can arrive as the call finishes. Raised while the worker encoded the reply, decoded the next line, or looped back to read, KeyboardInterrupt either wrote the reply twice or ended the worker after the driver had been told the session survived. The worker now installs its own SIGINT handler, which raises only while a call's code runs and ignores the signal otherwise. Installing it explicitly also means a worker whose parent ignores SIGINT, as a shell does for a background job, can still be interrupted.
Replaces the transport metaphors (ride along, land, arrive, hold) in the comments and one test name this branch adds with the literal verb.
Several tests passed whether or not the code they named worked. The timed-out-children test used subprocess.run(), which kills its own child on KeyboardInterrupt, so a SIGINT to the leader alone also passed; it now waits on Popen(). The aclose() cancellation test also cancelled the call, whose own shutdown closed the worker; it now cancels only the close. The respawn tests now check that a respawn happened, and the define failure test checks that the error's message is passed on. New tests cover the branches nothing reached: a reply that does not parse or is out of turn, a crash or out-of-turn answer during the interrupt grace, a start that never becomes ready, a define that ends the worker, and a bad line on the worker's stdin.
…says so Handles went with the call in one message. A restarted worker is sent the whole store again, so a large store outgrew the channel and every handle was shrunk to its repr and then counted as sent. Each handle is now its own message and counts as synced once the worker acknowledges it. The mark is also tied to the store it came from, and a call that passes no store leaves it alone: resetting it made the next call resend the store and overwrite any handle the model had reassigned. A call the worker cannot decode now gets an Error keyed to its id, where the driver waited out the call timeout and the interrupt grace and then reported an unresponsive worker. A write to a worker that exited unseen reports the reset instead of the bare OS error, and a call queued behind aclose() no longer starts a worker only for the close to kill it. The cancelled-aclose test now waits for the original close to finish before closing again, so it fails if that close was abandoned.
A sync reply that never comes now reports that the session stopped answering, rather than that it stopped reading. Garbage in answer to the interrupt reports an unparseable line rather than a crash. A process that cannot be created at all says the session failed to start, and a start that times out includes the worker's stderr. A call too large for the channel is refused before any of it is written, so it no longer restarts the session; a handle too large is skipped.
The harvest keeps every module-level function, and some cannot run once lifted out of their module: a lambda whose line names a module global, or a closure that rebinds a name its module defined. One such source failed every start, so the session was unusable. The worker now skips a source that raises; only a worker that stops answering fails the start. A default the harvest could not evaluate became a placeholder that was truthy and formatted as text, so a measure called with it could filter a query on the placeholder's repr and silently match nothing. Using the placeholder now raises a NameError naming the expression.
The flush-left case could not tell whether string lines were exempt from the dedent, since their indent is already zero. An indented string inside a nested measure loses its value without the exemption.
The group kill ran only when the session was closed, which for a worker that died between calls could be the idle window later, by which time an empty group's id may belong to another process. It now runs as the exit is observed. The same kill covers a child that ignores SIGTERM after the leader honoured it, which used to outlive the shutdown. An interrupt is no longer sent to a leader that has exited.
The reap tests slept a fixed two seconds against a reap due at one and a half; they now poll with a deadline, which is faster and leaves a slow runner room. The timed-out-children test passed vacuously if the timeout fired before the child started, so it now asserts the child did start. The check jobs get a timeout, so a deadlock fails a leg instead of holding it for six hours.
The interrupt docstring said children get a KeyboardInterrupt; they get SIGINT, and the worker raises only while model code runs. The define helper's docstring claimed it stood in for missing globals, where it guards only defaults. The R reference names the functions that run the lifecycle. Also drops contrast framing, em-dashes, and transport verbs.
A handle too large for the channel even as a repr is skipped and the call still runs; the earlier test reached the worker's decode error instead. The dead-worker test now checks for the child before the close, so it fails if the group kill waits for the close.
Guardrails mode applied only the address-space cap, so under network="none" nothing stopped model code from opening a socket, and nothing kept it inside its roots. An audit hook now checks, while model code runs, that reads stay inside the sandbox's read roots and the scratch directory, that writes stay inside the scratch directory, that no process is started, and that the network is reached only when allowed. Paths are resolved through symlinks first, a denial is a PermissionError naming the access and the path, and the worker's own work between calls is not checked. The policy is the one pkg-r/inst/worker/worker.R applies to R; like it, this provides no security boundary.
jat255
force-pushed
the
jat255/bgk1-worker-lifecycle
branch
from
October 9, 2026 02:47
9eacdf5 to
ec21831
Compare
jat255
marked this pull request as ready for review
October 9, 2026 02:47
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
This PR adds
Worker, which runs one persistent Python session in a sandboxed child process. The process starts on the first call and closes after 10 idle minutes. When a call times out, the worker interrupts it, and restarts the process only if the interrupt gets no answer.Workerdoes not start the Python process itself. It asks a backend for one. The only backend today,LocalBackend, runs the process on this machine: it sets up the process's scratch folder, environment, and sandbox, and stops the process when asked. A backend that runs the process in a container could replace it later, with no change toWorker. This PR also removes the backend'sexec()method, which was an errant addition in a previous PR that had no callers yet.On a host that cannot sandbox, a user who opts in with
COMMONS_ALLOW_UNSAFE_FALLBACKnow gets guardrails: best-effort checks on files, processes, and the network. They are not a security boundary.It also corrects a fault on
mainthat stopped the worker from starting on stock Ubuntu. The Landlock sandbox asked for directory rights on a file, and the kernel refused the rule.Details (agent-written)
KeyboardInterruptonly while model code runs, so an interrupt that comes as a call ends cannot lose the reply or the session. Every other failure message names its cause: code too large to send (the session stays), a worker that stopped answering, or a start that failed, with the worker's stderr.NameErroronly when the measure uses it. A source that still fails to define is skipped. The model reads a measure's source withinspect.getsource(name), because the worker registers each source's text withlinecache.LocalSessionkills the rest of its process group at once, so a child left by a crash, or one that ignores SIGTERM, does not outlive it.network="none", and raisePermissionError. The read roots are the Python sandbox's, which are wider than those of R'sworker_guardrails().run_r_tool()andworker_await()inpkg-r/R/run-r.R. A shared fixture for the timeouts and messages both drivers use is tracked as kata y154, for the wiring PR./bin/pkillresolves to the regular file/usr/bin/pgrep. The kernel accepts only file rights in a rule on a file.add_rootsinpkg-r/src/sandbox.cdoes not limit a file's rights either. This PR does not change it.waitid()report of a stopped child as an exit. The SIGSTOP test is skipped there. Tracked as kata d5sr.setsid()escapes the process-group kill (kata y06q), the worker can signal its parent (kata gs5e), and model code can reach the protocol descriptors (recorded on kata bgk1)./devroot (shared with R) lets model code open other terminals. A worker that exits between calls restarts with no message to the model, as in R.Lifecycle state diagram
stateDiagram-v2 state "No process" as none none : Worker._session is None state "Starting" as starting starting : Worker._ensure() starting : LocalBackend.start() spawns _runtime/_worker.py starting : waits for Ready, then Worker._define() per measure source state "Running a call" as running running : Worker._call() running : Worker._sync() sends each new handle in its own Call running : Worker._send(Call), then Worker._await_reply() state "Idle" as idle idle : Worker._schedule_reap() timer is running state "Interrupting" as interrupting interrupting : Worker._escalate() interrupting : LocalSession.interrupt() sends SIGINT to the process group state "Shutting down" as closing closing : Worker._shutdown() then LocalSession.close() closing : _terminate() sends SIGTERM, then SIGKILL closing : LocalSession._kill_group() kills any children left closing : scratch directory removed [*] --> none none --> starting: Worker.run() starting --> running: Ready, sources defined starting --> closing: SPAWN_TIMEOUT 30 s, raises RuntimeError idle --> running: Worker.run() running --> idle: Result or Error running --> interrupting: CALL_TIMEOUT 60 s interrupting --> idle: reply within INTERRUPT_GRACE 5 s, variables kept interrupting --> closing: no reply, Failure, variables lost running --> closing: crash or bad reply, Failure idle --> closing: IDLE_TIMEOUT 10 min, Worker._reap_when_idle() closing --> none