From 4d37ca4c8081cb4f66066f5432fac68af3b69227 Mon Sep 17 00:00:00 2001 From: Josh Taillon Date: Fri, 9 Oct 2026 13:11:07 -0400 Subject: [PATCH 01/25] refactor: name the code-run display classes for either language The stylesheet's code and plot display classes were named for run_r. run_python renders the same display, so the classes become commons-run-display, commons-run-details, commons-run-code, and commons-run-plot, and run_r's HTML uses the new names. --- .../commons/www/commons-chat/commons-chat.css | 52 +++++++++---------- pkg-r/R/run-r.R | 8 +-- pkg-r/inst/www/commons-chat/commons-chat.css | 52 +++++++++---------- pkg-r/tests/testthat/test-run-r.R | 16 +++--- www/commons-chat/commons-chat.css | 52 +++++++++---------- 5 files changed, 90 insertions(+), 90 deletions(-) diff --git a/pkg-py/src/commons/www/commons-chat/commons-chat.css b/pkg-py/src/commons/www/commons-chat/commons-chat.css index 14af90dd..82b9df25 100644 --- a/pkg-py/src/commons/www/commons-chat/commons-chat.css +++ b/pkg-py/src/commons/www/commons-chat/commons-chat.css @@ -558,12 +558,12 @@ shiny-chat-container /* ---- run_r display ---------------------------------------------------- */ -.commons-run-r-display { +.commons-run-display { display: grid; gap: 0.6rem; } -.commons-run-r-details > summary { +.commons-run-details > summary { align-items: center; color: inherit; cursor: pointer; @@ -574,11 +574,11 @@ shiny-chat-container width: fit-content; } -.commons-run-r-details > summary::-webkit-details-marker { +.commons-run-details > summary::-webkit-details-marker { display: none; } -.commons-run-r-details > summary::before { +.commons-run-details > summary::before { border-bottom: 1px solid currentColor; border-right: 1px solid currentColor; content: ""; @@ -590,90 +590,90 @@ shiny-chat-container width: 0.45em; } -.commons-run-r-details[open] > summary::before { +.commons-run-details[open] > summary::before { transform: rotate(45deg); } -.commons-run-r-details > summary:focus-visible { +.commons-run-details > summary:focus-visible { border-radius: var(--bs-border-radius-sm, 0.25rem); outline: 2px solid var(--bs-primary, #007bc2); outline-offset: 2px; } -.commons-run-r-details[open] > .commons-run-r-code { - animation: commons-run-r-details-reveal 0.2s ease-out; +.commons-run-details[open] > .commons-run-code { + animation: commons-run-details-reveal 0.2s ease-out; } -@keyframes commons-run-r-details-reveal { +@keyframes commons-run-details-reveal { from { opacity: 0; } } @media (prefers-reduced-motion: reduce) { - .commons-run-r-details > summary::before { + .commons-run-details > summary::before { transition: none; } - .commons-run-r-details[open] > .commons-run-r-code { + .commons-run-details[open] > .commons-run-code { animation: none; } } -.commons-run-r-display .commons-run-r-code { +.commons-run-display .commons-run-code { margin: 0; overflow-x: auto; white-space: pre; } -.commons-run-r-code .hl.com { +.commons-run-code .hl.com { color: #a0a1a7; font-style: italic; } -.commons-run-r-code .hl.kwa { +.commons-run-code .hl.kwa { color: #a626a4; } -.commons-run-r-code .hl.kwc, -.commons-run-r-code .hl.num { +.commons-run-code .hl.kwc, +.commons-run-code .hl.num { color: #986801; } -.commons-run-r-code .hl.kwd { +.commons-run-code .hl.kwd { color: #4078f2; } -.commons-run-r-code .hl.sng { +.commons-run-code .hl.sng { color: #50a14f; } -[data-bs-theme="dark"] .commons-run-r-code .hl.com { +[data-bs-theme="dark"] .commons-run-code .hl.com { color: #5c6370; } -[data-bs-theme="dark"] .commons-run-r-code .hl.kwa { +[data-bs-theme="dark"] .commons-run-code .hl.kwa { color: #c678dd; } -[data-bs-theme="dark"] .commons-run-r-code .hl.kwc, -[data-bs-theme="dark"] .commons-run-r-code .hl.num { +[data-bs-theme="dark"] .commons-run-code .hl.kwc, +[data-bs-theme="dark"] .commons-run-code .hl.num { color: #d19a66; } -[data-bs-theme="dark"] .commons-run-r-code .hl.kwd { +[data-bs-theme="dark"] .commons-run-code .hl.kwd { color: #61aeee; } -[data-bs-theme="dark"] .commons-run-r-code .hl.sng { +[data-bs-theme="dark"] .commons-run-code .hl.sng { color: #98c379; } -.commons-run-r-details > .commons-run-r-code { +.commons-run-details > .commons-run-code { margin-top: 0.4rem; } -.commons-run-r-plot { +.commons-run-plot { border: 1px solid var(--bs-border-color, #dee2e6); border-radius: 0.5rem; height: auto; diff --git a/pkg-r/R/run-r.R b/pkg-r/R/run-r.R index 913e879d..20fc2553 100644 --- a/pkg-r/R/run-r.R +++ b/pkg-r/R/run-r.R @@ -226,7 +226,7 @@ run_r_html <- function(code, segments) { dims <- plot_dimensions() plot_html <- c(plot_html, sprintf( paste0( - "" ), plot_image_data(seg$path), @@ -241,18 +241,18 @@ run_r_html <- function(code, segments) { } } code_html <- sprintf( - "
%s
", + "
%s
", highlight_r_html(paste(c(code, output), collapse = "\n")) ) if (length(plot_html)) { code_html <- paste0( - "
Details", + "
Details", code_html, "
" ) } sprintf( - "
%s
", + "
%s
", paste(c(code_html, plot_html), collapse = "\n") ) } diff --git a/pkg-r/inst/www/commons-chat/commons-chat.css b/pkg-r/inst/www/commons-chat/commons-chat.css index 14af90dd..82b9df25 100644 --- a/pkg-r/inst/www/commons-chat/commons-chat.css +++ b/pkg-r/inst/www/commons-chat/commons-chat.css @@ -558,12 +558,12 @@ shiny-chat-container /* ---- run_r display ---------------------------------------------------- */ -.commons-run-r-display { +.commons-run-display { display: grid; gap: 0.6rem; } -.commons-run-r-details > summary { +.commons-run-details > summary { align-items: center; color: inherit; cursor: pointer; @@ -574,11 +574,11 @@ shiny-chat-container width: fit-content; } -.commons-run-r-details > summary::-webkit-details-marker { +.commons-run-details > summary::-webkit-details-marker { display: none; } -.commons-run-r-details > summary::before { +.commons-run-details > summary::before { border-bottom: 1px solid currentColor; border-right: 1px solid currentColor; content: ""; @@ -590,90 +590,90 @@ shiny-chat-container width: 0.45em; } -.commons-run-r-details[open] > summary::before { +.commons-run-details[open] > summary::before { transform: rotate(45deg); } -.commons-run-r-details > summary:focus-visible { +.commons-run-details > summary:focus-visible { border-radius: var(--bs-border-radius-sm, 0.25rem); outline: 2px solid var(--bs-primary, #007bc2); outline-offset: 2px; } -.commons-run-r-details[open] > .commons-run-r-code { - animation: commons-run-r-details-reveal 0.2s ease-out; +.commons-run-details[open] > .commons-run-code { + animation: commons-run-details-reveal 0.2s ease-out; } -@keyframes commons-run-r-details-reveal { +@keyframes commons-run-details-reveal { from { opacity: 0; } } @media (prefers-reduced-motion: reduce) { - .commons-run-r-details > summary::before { + .commons-run-details > summary::before { transition: none; } - .commons-run-r-details[open] > .commons-run-r-code { + .commons-run-details[open] > .commons-run-code { animation: none; } } -.commons-run-r-display .commons-run-r-code { +.commons-run-display .commons-run-code { margin: 0; overflow-x: auto; white-space: pre; } -.commons-run-r-code .hl.com { +.commons-run-code .hl.com { color: #a0a1a7; font-style: italic; } -.commons-run-r-code .hl.kwa { +.commons-run-code .hl.kwa { color: #a626a4; } -.commons-run-r-code .hl.kwc, -.commons-run-r-code .hl.num { +.commons-run-code .hl.kwc, +.commons-run-code .hl.num { color: #986801; } -.commons-run-r-code .hl.kwd { +.commons-run-code .hl.kwd { color: #4078f2; } -.commons-run-r-code .hl.sng { +.commons-run-code .hl.sng { color: #50a14f; } -[data-bs-theme="dark"] .commons-run-r-code .hl.com { +[data-bs-theme="dark"] .commons-run-code .hl.com { color: #5c6370; } -[data-bs-theme="dark"] .commons-run-r-code .hl.kwa { +[data-bs-theme="dark"] .commons-run-code .hl.kwa { color: #c678dd; } -[data-bs-theme="dark"] .commons-run-r-code .hl.kwc, -[data-bs-theme="dark"] .commons-run-r-code .hl.num { +[data-bs-theme="dark"] .commons-run-code .hl.kwc, +[data-bs-theme="dark"] .commons-run-code .hl.num { color: #d19a66; } -[data-bs-theme="dark"] .commons-run-r-code .hl.kwd { +[data-bs-theme="dark"] .commons-run-code .hl.kwd { color: #61aeee; } -[data-bs-theme="dark"] .commons-run-r-code .hl.sng { +[data-bs-theme="dark"] .commons-run-code .hl.sng { color: #98c379; } -.commons-run-r-details > .commons-run-r-code { +.commons-run-details > .commons-run-code { margin-top: 0.4rem; } -.commons-run-r-plot { +.commons-run-plot { border: 1px solid var(--bs-border-color, #dee2e6); border-radius: 0.5rem; height: auto; diff --git a/pkg-r/tests/testthat/test-run-r.R b/pkg-r/tests/testthat/test-run-r.R index f2640821..e084edba 100644 --- a/pkg-r/tests/testthat/test-run-r.R +++ b/pkg-r/tests/testthat/test-run-r.R @@ -31,7 +31,7 @@ test_that("run_r executes code against stored handles", { expect_false(res@extra$display$open) expect_match( res@extra$display$html, - '
',
+    '
',
     fixed = TRUE
   )
   expect_match(
@@ -145,8 +145,8 @@ test_that("run_r returns plots as images and opens the display", {
     "data:image/png;base64,",
     fixed = TRUE
   )
-  expect_match(res@extra$display$html, "commons-run-r-details")
-  expect_match(res@extra$display$html, "commons-run-r-code", fixed = TRUE)
+  expect_match(res@extra$display$html, "commons-run-details")
+  expect_match(res@extra$display$html, "commons-run-code", fixed = TRUE)
   expect_match(res@extra$display$html, "Details", fixed = TRUE)
 })
 
@@ -165,7 +165,7 @@ test_that("run_r collapses code and output above plots", {
   expect_match(res@value[[1]]@text, "private warning")
   expect_match(
     res@extra$display$html,
-    '
Details', + '
Details', fixed = TRUE ) expect_match(res@extra$display$html, "#> private text", fixed = TRUE) @@ -177,7 +177,7 @@ test_that("run_r collapses code and output above plots", { fixed = TRUE ) expect_lt( - as.integer(regexpr("commons-run-r-details", res@extra$display$html)), + as.integer(regexpr("commons-run-details", res@extra$display$html)), as.integer(regexpr( "data:image/png;base64,", res@extra$display$html, @@ -219,7 +219,7 @@ test_that("run_r surfaces errors from model code without failing the tool", { expect_match(res@value, "Error: boom") expect_false(res@extra$display$open) expect_match(res@extra$display$html, "#> boom", fixed = TRUE) - expect_match(res@extra$display$html, "commons-run-r-code", fixed = TRUE) + expect_match(res@extra$display$html, "commons-run-code", fixed = TRUE) expect_no_match(res@extra$display$html, "1', fixed = TRUE ) - expect_match(res@extra$display$html, "commons-run-r-code", fixed = TRUE) + expect_match(res@extra$display$html, "commons-run-code", fixed = TRUE) expect_no_match(res@extra$display$html, "#>", fixed = TRUE) expect_no_match(res@extra$display$html, " summary { +.commons-run-details > summary { align-items: center; color: inherit; cursor: pointer; @@ -574,11 +574,11 @@ shiny-chat-container width: fit-content; } -.commons-run-r-details > summary::-webkit-details-marker { +.commons-run-details > summary::-webkit-details-marker { display: none; } -.commons-run-r-details > summary::before { +.commons-run-details > summary::before { border-bottom: 1px solid currentColor; border-right: 1px solid currentColor; content: ""; @@ -590,90 +590,90 @@ shiny-chat-container width: 0.45em; } -.commons-run-r-details[open] > summary::before { +.commons-run-details[open] > summary::before { transform: rotate(45deg); } -.commons-run-r-details > summary:focus-visible { +.commons-run-details > summary:focus-visible { border-radius: var(--bs-border-radius-sm, 0.25rem); outline: 2px solid var(--bs-primary, #007bc2); outline-offset: 2px; } -.commons-run-r-details[open] > .commons-run-r-code { - animation: commons-run-r-details-reveal 0.2s ease-out; +.commons-run-details[open] > .commons-run-code { + animation: commons-run-details-reveal 0.2s ease-out; } -@keyframes commons-run-r-details-reveal { +@keyframes commons-run-details-reveal { from { opacity: 0; } } @media (prefers-reduced-motion: reduce) { - .commons-run-r-details > summary::before { + .commons-run-details > summary::before { transition: none; } - .commons-run-r-details[open] > .commons-run-r-code { + .commons-run-details[open] > .commons-run-code { animation: none; } } -.commons-run-r-display .commons-run-r-code { +.commons-run-display .commons-run-code { margin: 0; overflow-x: auto; white-space: pre; } -.commons-run-r-code .hl.com { +.commons-run-code .hl.com { color: #a0a1a7; font-style: italic; } -.commons-run-r-code .hl.kwa { +.commons-run-code .hl.kwa { color: #a626a4; } -.commons-run-r-code .hl.kwc, -.commons-run-r-code .hl.num { +.commons-run-code .hl.kwc, +.commons-run-code .hl.num { color: #986801; } -.commons-run-r-code .hl.kwd { +.commons-run-code .hl.kwd { color: #4078f2; } -.commons-run-r-code .hl.sng { +.commons-run-code .hl.sng { color: #50a14f; } -[data-bs-theme="dark"] .commons-run-r-code .hl.com { +[data-bs-theme="dark"] .commons-run-code .hl.com { color: #5c6370; } -[data-bs-theme="dark"] .commons-run-r-code .hl.kwa { +[data-bs-theme="dark"] .commons-run-code .hl.kwa { color: #c678dd; } -[data-bs-theme="dark"] .commons-run-r-code .hl.kwc, -[data-bs-theme="dark"] .commons-run-r-code .hl.num { +[data-bs-theme="dark"] .commons-run-code .hl.kwc, +[data-bs-theme="dark"] .commons-run-code .hl.num { color: #d19a66; } -[data-bs-theme="dark"] .commons-run-r-code .hl.kwd { +[data-bs-theme="dark"] .commons-run-code .hl.kwd { color: #61aeee; } -[data-bs-theme="dark"] .commons-run-r-code .hl.sng { +[data-bs-theme="dark"] .commons-run-code .hl.sng { color: #98c379; } -.commons-run-r-details > .commons-run-r-code { +.commons-run-details > .commons-run-code { margin-top: 0.4rem; } -.commons-run-r-plot { +.commons-run-plot { border: 1px solid var(--bs-border-color, #dee2e6); border-radius: 0.5rem; height: auto; From 32dd9ad040168e656de0e2e068e3966cbd175a3d Mon Sep 17 00:00:00 2001 From: Josh Taillon Date: Fri, 9 Oct 2026 13:19:58 -0400 Subject: [PATCH 02/25] feat(py): run a worker on an event loop of its own chatlas runs an async tool only from stream_async(), and runs a sync tool directly on the caller's event loop, where a call that takes a minute stalls every other task. WorkerThread keeps a Worker on a private loop in a daemon thread, started by the first call. An async caller awaits a call without blocking its loop, a sync caller blocks on the same session, and cancelling an async caller cancels the call. --- pkg-py/src/commons/_execution/_thread.py | 103 +++++++++++++++++++++++ pkg-py/tests/test_execution_thread.py | 84 ++++++++++++++++++ 2 files changed, 187 insertions(+) create mode 100644 pkg-py/src/commons/_execution/_thread.py create mode 100644 pkg-py/tests/test_execution_thread.py diff --git a/pkg-py/src/commons/_execution/_thread.py b/pkg-py/src/commons/_execution/_thread.py new file mode 100644 index 00000000..79ed7c42 --- /dev/null +++ b/pkg-py/src/commons/_execution/_thread.py @@ -0,0 +1,103 @@ +"""A worker on an event loop of its own, so any caller can reach one session. + +chatlas runs an async tool only from ``stream_async()``, and a sync tool +directly on the caller's event loop, where a long call would stall every +other task on it. Keeping the ``Worker`` on a private loop in a background +thread lets an async caller await a call without blocking its loop, and a +sync caller block on the same session, so variables persist whichever way +the agent is asked. +""" + +from __future__ import annotations + +import asyncio +import concurrent.futures +import contextlib +import threading + +from .._handles import HandleStore +from ._driver import Failure, Worker +from ._protocol import Error, Result + +__all__ = ["WorkerThread"] + +# How long close() waits for the worker's shutdown, which is itself bounded +# by its grace periods; the margin covers a loop slow to get to it. +CLOSE_TIMEOUT = 30.0 + +_CLOSED = Failure(message="the Python session is closed.") + + +class WorkerThread: + """Runs ``worker`` on a loop in a daemon thread, started by the first call.""" + + def __init__(self, worker: Worker) -> None: + self._worker = worker + self._loop: asyncio.AbstractEventLoop | None = None + self._thread: threading.Thread | None = None + self._start_lock = threading.Lock() + self._closed = False + + async def run( + self, code: str, handles: HandleStore | None = None + ) -> Result | Error | Failure: + """Run ``code`` without blocking the caller's event loop. + + Cancelling the caller cancels the call, which shuts the worker down + the way a cancelled ``Worker.run`` does. + """ + future = self._submit(code, handles) + if future is None: + return _CLOSED + return await asyncio.wrap_future(future) + + def run_sync( + self, code: str, handles: HandleStore | None = None + ) -> Result | Error | Failure: + """Run ``code``, blocking the calling thread until the reply arrives.""" + if threading.current_thread() is self._thread: + raise RuntimeError("run_sync() would deadlock on the worker's own loop") + future = self._submit(code, handles) + if future is None: + return _CLOSED + return future.result() + + def close(self) -> None: + """Close the worker, then stop the loop and its thread. Safe to repeat.""" + with self._start_lock: + if self._closed: + return + self._closed = True + loop, thread = self._loop, self._thread + if loop is None or thread is None: + return + closing = asyncio.run_coroutine_threadsafe(self._worker.aclose(), loop) + with contextlib.suppress(concurrent.futures.TimeoutError, RuntimeError): + closing.result(CLOSE_TIMEOUT) + loop.call_soon_threadsafe(loop.stop) + thread.join(CLOSE_TIMEOUT) + if not thread.is_alive(): + loop.close() + + def _submit( + self, code: str, handles: HandleStore | None + ) -> concurrent.futures.Future[Result | Error | Failure] | None: + loop = self._ensure_loop() + if loop is None: + return None + # The store is read from the worker's thread while the caller waits. + # A store only ever grows, and each read is a single dict operation. + return asyncio.run_coroutine_threadsafe(self._worker.run(code, handles), loop) + + def _ensure_loop(self) -> asyncio.AbstractEventLoop | None: + with self._start_lock: + if self._closed: + return None + if self._loop is None: + loop = asyncio.new_event_loop() + thread = threading.Thread( + target=loop.run_forever, name="commons-python-session", daemon=True + ) + thread.start() + self._loop, self._thread = loop, thread + return self._loop diff --git a/pkg-py/tests/test_execution_thread.py b/pkg-py/tests/test_execution_thread.py new file mode 100644 index 00000000..4dee70fc --- /dev/null +++ b/pkg-py/tests/test_execution_thread.py @@ -0,0 +1,84 @@ +"""The worker thread: one session, reachable from sync callers and any event loop.""" + +from __future__ import annotations + +import asyncio +import threading + +import pytest + +from commons._execution._driver import Failure, Worker +from commons._execution._protocol import Result +from commons._execution._thread import WorkerThread + + +@pytest.fixture +def runner(): + worker_thread = WorkerThread(Worker(call_timeout=10)) + yield worker_thread + worker_thread.close() + + +def test_sync_and_async_callers_share_one_session(runner: WorkerThread) -> None: + assert isinstance(runner.run_sync("x = 41"), Result) + + async def from_a_loop() -> Result: + reply = await runner.run("x + 1") + assert isinstance(reply, Result) + return reply + + assert asyncio.run(from_a_loop()).value == 42 + + +def test_the_session_runs_off_the_callers_loop(runner: WorkerThread) -> None: + async def caller() -> tuple[int, bool]: + ticks = 0 + stop = asyncio.Event() + + async def tick() -> None: + nonlocal ticks + while not stop.is_set(): + ticks += 1 + await asyncio.sleep(0.01) + + ticker = asyncio.ensure_future(tick()) + reply = await runner.run("import time; time.sleep(0.5)") + stop.set() + await ticker + return ticks, isinstance(reply, Result) + + ticks, ok = asyncio.run(caller()) + assert ok + # The caller's loop kept running while the call slept. + assert ticks > 10 + + +def test_no_thread_starts_before_the_first_call() -> None: + before = threading.active_count() + worker_thread = WorkerThread(Worker()) + assert threading.active_count() == before + worker_thread.close() + + +def test_a_cancelled_caller_cancels_the_call(runner: WorkerThread) -> None: + async def cancel_midway() -> None: + call = asyncio.ensure_future(runner.run("import time; time.sleep(30)")) + await asyncio.sleep(1) + call.cancel() + with pytest.raises(asyncio.CancelledError): + await call + + asyncio.run(cancel_midway()) + # The cancelled call's worker was shut down; the next one respawns. + reply = runner.run_sync("1 + 1") + assert isinstance(reply, Result) + assert reply.value == 2 + + +def test_close_ends_the_thread_and_later_calls_fail(runner: WorkerThread) -> None: + runner.run_sync("1") + runner.close() + runner.close() + reply = runner.run_sync("1") + assert isinstance(reply, Failure) + assert "closed" in reply.message From 911859a3141472f4ef8ed0a27ef15d57f70ffb5b Mon Sep 17 00:00:00 2001 From: Josh Taillon Date: Fri, 9 Oct 2026 13:20:05 -0400 Subject: [PATCH 03/25] feat(py): add run_python, which runs model code in the agent's sandboxed session Every agent registers run_python. Its description follows run_r's in Python's idiom: the sandbox framing, the session the user cannot reach, the handles preloaded from whichever registered tools store results, measure sources read with inspect.getsource(), and rules whose network line says whether pip can install packages. A result is tagged B and asks for citations; the model gets the text the call produced and each plot as an image, and the reader gets the highlighted code with its output and the plots at their size. Commons takes a network argument and builds its Worker during construction, so a host that cannot sandbox the session fails there, before any model asks to run code. The worker lives on a WorkerThread: the registered tool is async, and chat() swaps in a sync variant for its length, because chatlas refuses a synchronous chat while any async tool is registered. --- pkg-py/src/commons/_agent.py | 71 +++++-- pkg-py/src/commons/_display.py | 2 + pkg-py/src/commons/_run_python.py | 334 ++++++++++++++++++++++++++++++ pkg-py/tests/test_agent.py | 1 + pkg-py/tests/test_run_python.py | 282 +++++++++++++++++++++++++ 5 files changed, 676 insertions(+), 14 deletions(-) create mode 100644 pkg-py/src/commons/_run_python.py create mode 100644 pkg-py/tests/test_run_python.py diff --git a/pkg-py/src/commons/_agent.py b/pkg-py/src/commons/_agent.py index d852107a..d829e436 100644 --- a/pkg-py/src/commons/_agent.py +++ b/pkg-py/src/commons/_agent.py @@ -11,6 +11,7 @@ import copy import warnings +import weakref from collections.abc import AsyncGenerator, Mapping, Sequence from typing import Any, Literal, NoReturn @@ -29,6 +30,9 @@ from ._context_layer import ContextLayer, augment_context_layer from ._data_source import DataSource from ._definitions import Registry, build_registry +from ._execution._backend import Network +from ._execution._driver import Worker +from ._execution._thread import WorkerThread from ._handles import HandleStore from ._measures import SemanticLayer, resolve_injections, semantic_layer from ._prompt import ( @@ -40,6 +44,7 @@ ) from ._provenance import collect_appended_tags, derive_provenance_tag, provenance_aside from ._reminders import append_restored_conversation_reminder, append_turn_reminder +from ._run_python import run_python_description, run_python_tools from ._tools import FirstTouch, ToolContext, build_commons_tools __all__ = ["Commons"] @@ -84,11 +89,21 @@ class Commons(Chat[Any, Any]): `## Additional instructions` heading at the end of commons' built-in system prompt, as a string or the path to a text or Markdown file. + The agent runs the Python code its model writes in a sandboxed session, + which starts on the first call. `network` is whether that session can + reach the network: `"none"` (the default) or `"full"`. The session is + sandboxed on Linux and macOS. On any other host, local development can + opt in to best-effort guardrails by setting the + `COMMONS_ALLOW_UNSAFE_FALLBACK` environment variable; these guardrails + are not a security boundary. + Construction raises a TypeError if `client` is not a `chatlas.Chat`, if an entry of `data_sources` is not a `DataSource`, or if a layer is not the layer its argument claims; a ValueError if `data_sources` names no - source or a measure asks for an injection no named source can fill; and - a FileNotFoundError if `instructions` names a file that does not exist. + source, a measure asks for an injection no named source can fill, or + `network` is neither `"none"` nor `"full"`; a FileNotFoundError if + `instructions` names a file that does not exist; and a RuntimeError if + this host cannot sandbox the session and the opt-in is not set. """ def __init__( @@ -99,6 +114,7 @@ def __init__( context_layer: ContextLayer | None = None, *, instructions: str | None = None, + network: Network = "none", ) -> None: if not isinstance(client, Chat): raise TypeError( @@ -124,6 +140,12 @@ def __init__( f"{type(semantic_layer).__name__}." ) check_instructions(instructions) + # Built here so a host that cannot sandbox the session fails now, + # before any model asks to run code. The process starts on first use. + worker = Worker( + network=network, + measure_sources=list(semantic_layer.source_text.values()), + ) # Share the provider, which carries the chosen model; shallow-copy # the chat kwargs so later changes don't cross between the two. @@ -152,18 +174,31 @@ def __init__( ) self._restore_reminder_pending = False - tools = build_commons_tools( - ToolContext( - sources=sources, - measures=self._measures, - definitions=self._definitions, - context_layer=self._context_layer, - handles=self._handles, - citation_request=self._citation_request, - injections=self._injections, - first_touch=self._first_touch, - ) + context = ToolContext( + sources=sources, + measures=self._measures, + definitions=self._definitions, + context_layer=self._context_layer, + handles=self._handles, + citation_request=self._citation_request, + injections=self._injections, + first_touch=self._first_touch, + ) + tools = build_commons_tools(context) + self._python = WorkerThread(worker) + # The session's process and thread go with the agent. + weakref.finalize(self, self._python.close) + self._run_python, self._run_python_sync = run_python_tools( + self._python, + context, + run_python_description( + [tool.name for tool in tools], + has_measures=bool(self._measures), + network=network, + ), + network, ) + tools.append(self._run_python) self.set_tools(list(tools)) self.system_prompt = _system_prompt( sources, @@ -204,7 +239,15 @@ def chat( was_pending = self._restore_reminder_pending inputs = self._prepare_turn_inputs(args) self._citation_request.reset() - response = super().chat(*inputs, echo=echo, stream=stream, kwargs=kwargs) + # chatlas refuses a synchronous chat while an async tool is + # registered, so the sync run_python stands in for this call. + self.register_tool(self._run_python_sync, force=True) + try: + response = super().chat( + *inputs, echo=echo, stream=stream, kwargs=kwargs + ) + finally: + self.register_tool(self._run_python, force=True) self._consume_restore_reminder(was_pending) return response diff --git a/pkg-py/src/commons/_display.py b/pkg-py/src/commons/_display.py index d36a2706..220342e8 100644 --- a/pkg-py/src/commons/_display.py +++ b/pkg-py/src/commons/_display.py @@ -24,6 +24,7 @@ _URL = re.compile(r"https?://", re.IGNORECASE) __all__ = [ + "CODE_ANALYSIS", "CONTEXT_SEARCH", "DATA_RETRIEVAL", "DISPLAY_EXTRA_KEY", @@ -60,6 +61,7 @@ class Title: CONTEXT_SEARCH = Title("Searching context", "Searched context") TABLE_INSPECTION = Title("Inspecting a table", "Inspected a table") DATA_RETRIEVAL = Title("Retrieving data", "Retrieved data") +CODE_ANALYSIS = Title("Analyzing data", "Analyzed data") def tool_display( diff --git a/pkg-py/src/commons/_run_python.py b/pkg-py/src/commons/_run_python.py new file mode 100644 index 00000000..fb33f28f --- /dev/null +++ b/pkg-py/src/commons/_run_python.py @@ -0,0 +1,334 @@ +"""The `run_python` tool: model-written Python run in the agent's sandboxed session. + +`pkg-r/R/run-r.R` builds R's `run_r`, and the two tools describe themselves +and shape their results the same way, each in its own language's idiom. + +The agent registers the async tool, which awaits the session without blocking +the caller's event loop. chatlas refuses a synchronous `chat()` while any async +tool is registered, so `Commons.chat()` swaps in the sync tool for the call. +""" + +from __future__ import annotations + +import base64 +import html +import importlib.util +import io +import keyword +import tokenize +from collections.abc import Sequence +from typing import Any + +from chatlas import ContentToolResult, Tool +from chatlas.types import ContentImageInline, ContentText +from htmltools import HTML, Tag, div, tags + +from ._citations import tool_result +from ._display import CODE_ANALYSIS, visible_result_note +from ._execution._backend import Network +from ._execution._driver import Failure +from ._execution._protocol import Error, OpaqueValue, Plot, Result +from ._execution._thread import WorkerThread +from ._frames import describe_frame, is_frame +from ._prompt import EXECUTION_TOOL +from ._provenance import Tag as ProvenanceTag +from ._rows import frame_rows, rows_to_markdown +from ._tools import ToolContext + +__all__ = [ + "HANDLE_TOOLS", + "run_python_description", + "run_python_result", + "run_python_tools", +] + +# The tools whose results are stored as handles, in the order the +# description names them. +HANDLE_TOOLS = ("call_measure", "call_metrics", "run_sql") + +NO_OUTPUT = "(The code ran but produced no output.)" + + +def run_python_description( + tool_names: Sequence[str], + *, + has_measures: bool, + network: Network, + can_install: bool | None = None, +) -> str: + """What `run_python` tells the model it is for. + + ``tool_names`` are the agent's other registered tools; the preloaded + handles are named after whichever of them store results. ``can_install`` + says whether the session's interpreter has pip, and defaults to asking + this one, which is the interpreter the session runs. + """ + if can_install is None: + can_install = importlib.util.find_spec("pip") is not None + handle_tools = [name for name in HANDLE_TOOLS if name in tool_names] + parts = [ + ( + "Run Python code in your sandboxed Python session to analyze results or " + "render plots. Python code and textual output are visible only to you; " + "rendered plots are also shown to the user." + ), + ( + "The user cannot access or interact with this session. Never direct them " + "to run code or inspect its variables or files; perform follow-up " + "analysis yourself and report the result in your response." + ), + ( + "Your session persists across calls: variables you assign and modules " + "you import remain available." + ), + ] + if handle_tools: + parts.append( + f"Results from {_listed(handle_tools)} are preloaded as variables " + "(r1, r2, ...)." + ) + if has_measures: + parts.append( + "Measure definitions and their helper functions are predefined under " + "their own names: call inspect.getsource() on a measure to read its " + "source. These are source-only copies without their original " + "environment or database connections, so treat them as reference " + "material; to compute a measure, use call_measure." + ) + rules = [ + "Work incrementally: each call should do one small, well-defined task.", + ( + "Follow PEP 8: put separate statements on separate lines and wrap long " + "calls for readability." + ), + ( + "Create at most one matplotlib figure per call and leave it open rather " + "than saving it." + ), + ( + "Do not use this tool to talk to the user; explanations belong in your " + "reply." + ), + ( + "Return results by ending with an expression (`df`, not `print(df)`) " + "and prefer brief summaries (df.head(), df.describe()) over large " + "outputs." + ), + "The session can only write to its own temporary directory.", + ] + if network == "none": + rules.append("The session has no network access.") + elif can_install: + rules.append( + "The session has network access. To use a package that is not " + "installed, run `sys.executable -m pip install --target` with the " + "temporary directory through subprocess, then add that directory " + "to sys.path." + ) + else: + rules.append( + "The session has network access, but only packages that are " + "already installed can be imported." + ) + return " ".join(parts) + "\n\nRules:" + "".join(f"\n- {rule}" for rule in rules) + + +def _listed(names: Sequence[str]) -> str: + if len(names) <= 2: + return " and ".join(names) + return f"{', '.join(names[:-1])}, and {names[-1]}" + + +def run_python_tools( + runner: WorkerThread, context: ToolContext, description: str, network: Network +) -> tuple[Tool, Tool]: + """The async tool the agent registers, and the sync one `chat()` swaps in.""" + + def finish(code: str, reply: Result | Error | Failure) -> ContentToolResult: + result = run_python_result(code, reply) + if context.citation_request is None: + return result + return context.citation_request.add_request(result) + + async def run_python(code: str) -> ContentToolResult: + return finish(code, await runner.run(code, context.handles)) + + def run_python_sync(code: str) -> ContentToolResult: + return finish(code, runner.run_sync(code, context.handles)) + + parameters = { + "type": "object", + "properties": { + "code": {"type": "string", "description": "The Python code to run."} + }, + "required": ["code"], + "additionalProperties": False, + } + annotations: Any = { + "title": CODE_ANALYSIS.running, + "readOnlyHint": False, + "openWorldHint": network == "full", + } + + def build(func: Any) -> Tool: + return Tool( + func=func, + name=EXECUTION_TOOL, + description=description, + parameters=parameters, + annotations=annotations, + ) + + return build(run_python), build(run_python_sync) + + +def run_python_result(code: str, reply: Result | Error | Failure) -> ContentToolResult: + """The tool result for one call: the model's view and the reader's. + + The model gets the text the call produced and each plot as an image. + The reader gets the code with its output, and the plots at their size. + """ + plots: tuple[Plot, ...] = () + if isinstance(reply, Failure): + texts = [f"Error: {reply.message}"] + elif isinstance(reply, Error): + texts = [reply.traceback.strip() or f"Error: {reply.message}"] + plots = reply.plots + else: + texts = [ + text + for text in ( + reply.stdout.rstrip("\n"), + reply.stderr.rstrip("\n"), + _value_text(reply.value), + ) + if text + ] + plots = reply.plots + return tool_result( + _model_value(texts, plots), + ProvenanceTag.B, + title=CODE_ANALYSIS.settled, + html=_display_html(code, texts, plots), + open=bool(plots), + ) + + +def _value_text(value: Any) -> str: + """The value a call ended on, as a REPL would show it; empty for None.""" + if value is None: + return "" + if isinstance(value, OpaqueValue): + return value.text + if is_frame(value): + rows = frame_rows(value) + return rows_to_markdown(rows) if rows is not None else describe_frame(value) + return repr(value) + + +def _model_value(texts: list[str], plots: tuple[Plot, ...]) -> Any: + text = "\n".join(texts) + if not plots: + return text or NO_OUTPUT + parts: list[Any] = [ContentText(text=text)] if text else [] + parts.extend( + ContentImageInline( + image_content_type="image/png", + data=base64.b64encode(plot.png).decode("ascii"), + ) + for plot in plots + ) + parts.append(ContentText(text=visible_result_note("plot"))) + return parts + + +def _display_html(code: str, texts: list[str], plots: tuple[Plot, ...]) -> Tag: + output = [f"#> {line}" for text in texts for line in text.split("\n")] + block: Tag = tags.pre( + tags.code( + HTML(highlight_python("\n".join([code, *output]))), + class_="language-python", + ), + class_="commons-run-code", + ) + if plots: + block = tags.details( + tags.summary("Details"), block, class_="commons-run-details" + ) + images = [ + tags.img( + class_="commons-run-plot", + src="data:image/png;base64," + + base64.b64encode(plot.display_png).decode("ascii"), + alt="Plot produced by Python code", + width=str(plot.width), + height=str(plot.height), + ) + for plot in plots + ] + return div(block, *images, class_="commons-run-display") + + +# The highlight classes the shared stylesheet styles, by token kind. +_COMMENT, _KEYWORD, _CALL, _NUMBER, _STRING = "com", "kwa", "kwd", "num", "sng" +_STRING_TOKENS = { + name + for name in ("STRING", "FSTRING_START", "FSTRING_MIDDLE", "FSTRING_END") + if hasattr(tokenize, name) +} + + +def highlight_python(source: str) -> str: + """``source`` as escaped HTML, with its tokens wrapped for the stylesheet. + + Text that does not tokenize as Python is escaped and left plain. + """ + lines = source.splitlines(keepends=True) + starts = [0] + for line in lines: + starts.append(starts[-1] + len(line)) + + def offset(position: tuple[int, int]) -> int: + row, column = position + return starts[row - 1] + column if row <= len(lines) else len(source) + + spans: list[tuple[int, int, str]] = [] + try: + tokens = list(tokenize.generate_tokens(io.StringIO(source).readline)) + except (tokenize.TokenError, SyntaxError): + return html.escape(source) + for index, token in enumerate(tokens): + kind = _token_class( + token, tokens[index + 1] if index + 1 < len(tokens) else None + ) + if kind is not None: + spans.append((offset(token.start), offset(token.end), kind)) + + out: list[str] = [] + cursor = 0 + for start, end, kind in spans: + if start < cursor: + continue + out.append(html.escape(source[cursor:start])) + out.append(f'{html.escape(source[start:end])}') + cursor = end + out.append(html.escape(source[cursor:])) + return "".join(out) + + +def _token_class( + token: tokenize.TokenInfo, following: tokenize.TokenInfo | None +) -> str | None: + name = tokenize.tok_name[token.type] + if token.type == tokenize.COMMENT: + return _COMMENT + if token.type == tokenize.NUMBER: + return _NUMBER + if name in _STRING_TOKENS: + return _STRING + if token.type == tokenize.NAME: + if keyword.iskeyword(token.string): + return _KEYWORD + if following is not None and following.string == "(": + return _CALL + return None diff --git a/pkg-py/tests/test_agent.py b/pkg-py/tests/test_agent.py index 295fb317..deca85fc 100644 --- a/pkg-py/tests/test_agent.py +++ b/pkg-py/tests/test_agent.py @@ -345,6 +345,7 @@ def test_an_agent_registers_the_tools_its_composition_earns( "describe_table", "run_sql", "load_skill", + "run_python", ] diff --git a/pkg-py/tests/test_run_python.py b/pkg-py/tests/test_run_python.py new file mode 100644 index 00000000..6f553167 --- /dev/null +++ b/pkg-py/tests/test_run_python.py @@ -0,0 +1,282 @@ +"""run_python: its description, the result it builds, and calls through an agent.""" + +from __future__ import annotations + +import base64 +import importlib.util +from typing import Any + +import pandas as pd +import pytest +from chatlas import ContentToolRequest, ContentToolResult, Tool +from chatlas.types import ContentImageInline, ContentText + +from commons import Injected, data_source, measure, semantic_layer +from commons._agent import Commons +from commons._display import DISPLAY_EXTRA_KEY +from commons._execution._driver import Failure +from commons._execution._protocol import Error, OpaqueValue, Plot, Result +from commons._provenance import TAG_EXTRA_KEY, Tag +from commons._run_python import ( + NO_OUTPUT, + highlight_python, + run_python_description, + run_python_result, +) + +from ._provider import scripted_chat, text + + +@measure(description="Total revenue.") +def total_revenue(sales: Injected[Any]) -> float: + return float(sales.execute("SELECT SUM(revenue) FROM sales").fetchone()[0]) + + +def frame() -> pd.DataFrame: + return pd.DataFrame({"revenue": [500.0, 900.0, 300.0]}) + + +def png(width: int, height: int) -> bytes: + """A PNG signature and IHDR header of the given size, which is all Plot reads.""" + header = width.to_bytes(4, "big") + height.to_bytes(4, "big") + return ( + b"\x89PNG\r\n\x1a\n" + b"\x00\x00\x00\rIHDR" + header + b"\x08\x06\x00\x00\x00" + ) + + +def plot() -> Plot: + return Plot(png=png(400, 300), display_png=png(800, 600)) + + +def display(result: ContentToolResult) -> dict[str, Any]: + assert result.extra is not None + return result.extra[DISPLAY_EXTRA_KEY] + + +def tool_request(**arguments: Any) -> list[Any]: + return [ + ContentToolRequest(id="call-run_python", name="run_python", arguments=arguments) + ] + + +def run_python_tool(agent: Commons) -> Tool: + tool = {tool.name: tool for tool in agent.get_tools()}["run_python"] + assert isinstance(tool, Tool) + return tool + + +def tool_results(agent: Commons) -> list[ContentToolResult]: + return [ + content + for turn in agent.get_turns() + for content in turn.contents + if isinstance(content, ContentToolResult) + ] + + +# ---- the description -------------------------------------------------------- + + +def test_the_preloaded_handles_name_the_registered_tools_that_store_results() -> None: + description = run_python_description( + ["search_pool", "call_measure", "run_sql"], has_measures=True, network="none" + ) + assert "Results from call_measure and run_sql are preloaded" in description + description = run_python_description( + ["call_measure", "call_metrics", "run_sql"], has_measures=True, network="none" + ) + assert "Results from call_measure, call_metrics, and run_sql are preloaded" in ( + description + ) + description = run_python_description( + ["run_sql"], has_measures=False, network="none" + ) + assert "Results from run_sql are preloaded" in description + + +def test_measure_sources_are_described_only_when_there_are_measures() -> None: + with_measures = run_python_description( + ["run_sql"], has_measures=True, network="none" + ) + without = run_python_description(["run_sql"], has_measures=False, network="none") + assert "inspect.getsource()" in with_measures + assert "inspect.getsource()" not in without + + +@pytest.mark.parametrize( + ("network", "can_install", "rule"), + [ + ("none", True, "- The session has no network access."), + ("full", True, "`sys.executable -m pip install --target`"), + ("full", False, "only packages that are already installed can be imported"), + ], +) +def test_the_network_rule_says_what_the_session_can_reach( + network: Any, can_install: bool, rule: str +) -> None: + description = run_python_description( + ["run_sql"], has_measures=False, network=network, can_install=can_install + ) + assert rule in description.split("\n\nRules:")[1] + + +# ---- the result ------------------------------------------------------------- + + +def test_a_result_shows_what_was_printed_and_the_value() -> None: + result = run_python_result( + "print('hi')\n1 + 1", Result(id="c1", value=2, stdout="hi\n", stderr="warn\n") + ) + assert result.value == "hi\nwarn\n2" + assert result.extra is not None + assert result.extra[TAG_EXTRA_KEY] == Tag.B + assert display(result)["title"] == "Analyzed data" + assert display(result)["open"] is False + + +def test_a_call_that_shows_nothing_says_so() -> None: + assert run_python_result("x = 1", Result(id="c1")).value == NO_OUTPUT + + +def test_a_frame_value_is_a_table_and_an_opaque_value_is_its_repr() -> None: + table = run_python_result("df", Result(id="c1", value=frame())).value + assert "| revenue |" in table + opaque = OpaqueValue(type_name="Connection", text="") + assert run_python_result("con", Result(id="c1", value=opaque)).value == ( + "" + ) + + +def test_an_error_shows_the_traceback_and_a_failure_its_message() -> None: + error = Error( + id="c1", + message="NameError: name 'y' is not defined", + traceback="Traceback...\nNameError", + ) + assert run_python_result("y", error).value == "Traceback...\nNameError" + failure = run_python_result("1", Failure(message="the Python session crashed.")) + assert failure.value == "Error: the Python session crashed." + assert failure.extra is not None + assert failure.extra[TAG_EXTRA_KEY] == Tag.B + + +def test_plots_reach_the_model_as_images_and_the_reader_at_their_size() -> None: + result = run_python_result( + "fig", Result(id="c1", stdout="drawn\n", plots=(plot(),)) + ) + parts = result.value + assert isinstance(parts, list) + assert isinstance(parts[0], ContentText) and parts[0].text == "drawn" + image = parts[1] + assert isinstance(image, ContentImageInline) + assert base64.b64decode(image.data) == plot().png + assert "already visible to the user" in parts[-1].text + shown = display(result) + assert shown["open"] is True + html = str(shown["html"]) + assert 'class="commons-run-details"' in html + assert base64.b64encode(plot().display_png).decode() in html + assert 'width="400"' in html and 'height="300"' in html + + +def test_the_display_shows_the_code_and_its_output_escaped() -> None: + html = str( + display(run_python_result("''", Result(id="c1", value="")))["html"] + ) + assert 'class="commons-run-code"' in html + assert "#> '<b>'" in html or "#> '<b>'" in html + assert "" not in html + + +# ---- highlighting ----------------------------------------------------------- + + +def test_highlighting_marks_tokens_with_the_stylesheets_classes() -> None: + out = highlight_python("def f(x):\n return len('a') + 1 # note\n") + assert 'def' in out + assert 'len' in out + assert ''a'' in out + assert '1' in out + assert '# note' in out + + +def test_text_that_does_not_tokenize_is_escaped_plainly() -> None: + assert highlight_python("''' None: + agent = Commons( + scripted_chat(), + {"sales": data_source(sales=frame())}, + semantic_layer(total_revenue), + ) + description = run_python_tool(agent).schema["function"]["description"] + assert "inspect.getsource()" in description + assert "Results from call_measure and run_sql are preloaded" in description + + +def test_chat_runs_code_against_a_preloaded_handle() -> None: + agent = Commons( + scripted_chat( + [ + [ + ContentToolRequest( + id="q", name="run_sql", arguments={"sql": "SELECT * FROM sales"} + ) + ], + tool_request(code="sum(row['revenue'] for row in r1)"), + text("Done."), + ] + ), + data_source(sales=frame()), + ) + agent.chat("Total revenue?", echo="none") + run = tool_results(agent)[-1] + assert run.value.startswith("1700.0") + # The async tool is back in place for stream_async. + assert run_python_tool(agent)._is_async + + +async def test_stream_async_keeps_session_state_and_reads_measure_sources() -> None: + agent = Commons( + scripted_chat( + [ + tool_request(code="x = 41"), + tool_request( + code="import inspect\nprint(inspect.getsource(total_revenue))\nx + 1" + ), + text("Done."), + ] + ), + {"sales": data_source(sales=frame())}, + semantic_layer(total_revenue), + ) + stream = await agent.stream_async("Go.") + [chunk async for chunk in stream] + run = tool_results(agent)[-1] + assert "def total_revenue(sales" in run.value + assert "\n42" in run.value + + +@pytest.mark.skipif( + importlib.util.find_spec("matplotlib") is None, reason="needs matplotlib" +) +async def test_a_plot_reaches_the_model_as_an_image() -> None: + provider_chat = scripted_chat( + [ + tool_request(code="import matplotlib.pyplot as plt\nplt.plot([1, 2, 3])"), + text("Done."), + ] + ) + agent = Commons(provider_chat, data_source(sales=frame())) + stream = await agent.stream_async("Plot it.") + [chunk async for chunk in stream] + run = tool_results(agent)[-1] + assert any(isinstance(part, ContentImageInline) for part in run.value) + # chatlas moves the image out of the tool result for the provider. + last_request = agent.provider.requests[-1] # type: ignore[attr-defined] + sent = [content for turn in last_request for content in turn.contents] + assert any(isinstance(content, ContentImageInline) for content in sent) From 97f93dd90e0a6816ebe839a65a5e510a222888e3 Mon Sep 17 00:00:00 2001 From: Josh Taillon Date: Fri, 9 Oct 2026 13:24:53 -0400 Subject: [PATCH 04/25] fix(py): cancel a worker thread's calls when it closes, and refuse calls after close() waited for the worker's shutdown, which waits out a running call, so a call longer than the close timeout left the loop stopped under a shutdown still in progress. close() now cancels the calls in flight first, and stops the loop only once the shutdown finishes. Submitting a call holds the same lock close() takes, so a call is either cancelled by the close or refused as closed, and a caller whose call the close cancelled is told the session is closed. --- pkg-py/src/commons/_execution/_thread.py | 95 ++++++++++++++++-------- pkg-py/tests/test_execution_thread.py | 19 +++++ 2 files changed, 81 insertions(+), 33 deletions(-) diff --git a/pkg-py/src/commons/_execution/_thread.py b/pkg-py/src/commons/_execution/_thread.py index 79ed7c42..a9c046b8 100644 --- a/pkg-py/src/commons/_execution/_thread.py +++ b/pkg-py/src/commons/_execution/_thread.py @@ -22,11 +22,13 @@ __all__ = ["WorkerThread"] # How long close() waits for the worker's shutdown, which is itself bounded -# by its grace periods; the margin covers a loop slow to get to it. +# by its grace periods once the calls in flight are cancelled. CLOSE_TIMEOUT = 30.0 _CLOSED = Failure(message="the Python session is closed.") +_Reply = Result | Error | Failure + class WorkerThread: """Runs ``worker`` on a loop in a daemon thread, started by the first call.""" @@ -35,12 +37,13 @@ def __init__(self, worker: Worker) -> None: self._worker = worker self._loop: asyncio.AbstractEventLoop | None = None self._thread: threading.Thread | None = None - self._start_lock = threading.Lock() + # Guards the lifecycle: a call is either submitted before close() + # begins, and so cancelled by it, or refused as closed. + self._lock = threading.Lock() self._closed = False + self._calls: set[concurrent.futures.Future[_Reply]] = set() - async def run( - self, code: str, handles: HandleStore | None = None - ) -> Result | Error | Failure: + async def run(self, code: str, handles: HandleStore | None = None) -> _Reply: """Run ``code`` without blocking the caller's event loop. Cancelling the caller cancels the call, which shuts the worker down @@ -49,55 +52,81 @@ async def run( future = self._submit(code, handles) if future is None: return _CLOSED - return await asyncio.wrap_future(future) - - def run_sync( - self, code: str, handles: HandleStore | None = None - ) -> Result | Error | Failure: + try: + return await asyncio.wrap_future(future) + except asyncio.CancelledError: + task = asyncio.current_task() + if self._closed and future.cancelled() and not (task and task.cancelling()): + return _CLOSED + raise + + def run_sync(self, code: str, handles: HandleStore | None = None) -> _Reply: """Run ``code``, blocking the calling thread until the reply arrives.""" if threading.current_thread() is self._thread: raise RuntimeError("run_sync() would deadlock on the worker's own loop") future = self._submit(code, handles) if future is None: return _CLOSED - return future.result() + try: + return future.result() + except concurrent.futures.CancelledError: + if self._closed: + return _CLOSED + raise def close(self) -> None: - """Close the worker, then stop the loop and its thread. Safe to repeat.""" - with self._start_lock: + """Close the worker, then stop the loop and its thread. Safe to repeat. + + Calls in flight are cancelled first, so the worker shuts down within + its grace periods rather than after a call's timeout. A shutdown that + still overruns ``CLOSE_TIMEOUT`` finishes in the background, and the + loop stops only once it has. + """ + with self._lock: if self._closed: return self._closed = True loop, thread = self._loop, self._thread + calls = list(self._calls) if loop is None or thread is None: return + for call in calls: + call.cancel() closing = asyncio.run_coroutine_threadsafe(self._worker.aclose(), loop) - with contextlib.suppress(concurrent.futures.TimeoutError, RuntimeError): + closing.add_done_callback(lambda _: loop.call_soon_threadsafe(loop.stop)) + with contextlib.suppress(TimeoutError): closing.result(CLOSE_TIMEOUT) - loop.call_soon_threadsafe(loop.stop) thread.join(CLOSE_TIMEOUT) if not thread.is_alive(): loop.close() def _submit( self, code: str, handles: HandleStore | None - ) -> concurrent.futures.Future[Result | Error | Failure] | None: - loop = self._ensure_loop() - if loop is None: - return None - # The store is read from the worker's thread while the caller waits. - # A store only ever grows, and each read is a single dict operation. - return asyncio.run_coroutine_threadsafe(self._worker.run(code, handles), loop) - - def _ensure_loop(self) -> asyncio.AbstractEventLoop | None: - with self._start_lock: + ) -> concurrent.futures.Future[_Reply] | None: + with self._lock: if self._closed: return None - if self._loop is None: - loop = asyncio.new_event_loop() - thread = threading.Thread( - target=loop.run_forever, name="commons-python-session", daemon=True - ) - thread.start() - self._loop, self._thread = loop, thread - return self._loop + loop = self._ensure_loop() + # The store is read from the worker's thread while the caller + # waits. A store only grows, and each read is one dict operation. + future = asyncio.run_coroutine_threadsafe( + self._worker.run(code, handles), loop + ) + self._calls.add(future) + future.add_done_callback(self._forget) + return future + + def _forget(self, future: concurrent.futures.Future[_Reply]) -> None: + with self._lock: + self._calls.discard(future) + + def _ensure_loop(self) -> asyncio.AbstractEventLoop: + """The loop, started on first use. Called with the lock held.""" + if self._loop is None: + loop = asyncio.new_event_loop() + thread = threading.Thread( + target=loop.run_forever, name="commons-python-session", daemon=True + ) + thread.start() + self._loop, self._thread = loop, thread + return self._loop diff --git a/pkg-py/tests/test_execution_thread.py b/pkg-py/tests/test_execution_thread.py index 4dee70fc..2d1aeeeb 100644 --- a/pkg-py/tests/test_execution_thread.py +++ b/pkg-py/tests/test_execution_thread.py @@ -4,6 +4,7 @@ import asyncio import threading +import time import pytest @@ -82,3 +83,21 @@ def test_close_ends_the_thread_and_later_calls_fail(runner: WorkerThread) -> Non reply = runner.run_sync("1") assert isinstance(reply, Failure) assert "closed" in reply.message + + +def test_close_cancels_a_running_call_rather_than_waiting_it_out() -> None: + worker_thread = WorkerThread(Worker(call_timeout=60)) + worker_thread.run_sync("1") + replies: list[object] = [] + caller = threading.Thread( + target=lambda: replies.append( + worker_thread.run_sync("import time; time.sleep(60)") + ) + ) + caller.start() + time.sleep(1) + started = time.monotonic() + worker_thread.close() + caller.join(10) + assert time.monotonic() - started < 15 + assert replies == [Failure(message="the Python session is closed.")] From fb3c06cb32f58c0633e321017d277db9203063da Mon Sep 17 00:00:00 2001 From: Josh Taillon Date: Fri, 9 Oct 2026 13:24:54 -0400 Subject: [PATCH 05/25] fix(py): describe plotting and installs by what the -I session can import run_python's rules told every agent to draw matplotlib figures, though matplotlib is optional, and offered pip when the host found it only in the user site, which the session's -I leaves out. Both rules now follow session_can_import(), which ignores the user site. --- pkg-py/src/commons/_run_python.py | 32 +++++++++++++++++++++++++------ pkg-py/tests/test_run_python.py | 25 ++++++++++++++++++++++++ 2 files changed, 51 insertions(+), 6 deletions(-) diff --git a/pkg-py/src/commons/_run_python.py b/pkg-py/src/commons/_run_python.py index fb33f28f..735fef88 100644 --- a/pkg-py/src/commons/_run_python.py +++ b/pkg-py/src/commons/_run_python.py @@ -15,6 +15,8 @@ import importlib.util import io import keyword +import os +import site import tokenize from collections.abc import Sequence from typing import Any @@ -54,17 +56,20 @@ def run_python_description( *, has_measures: bool, network: Network, + can_plot: bool | None = None, can_install: bool | None = None, ) -> str: """What `run_python` tells the model it is for. ``tool_names`` are the agent's other registered tools; the preloaded - handles are named after whichever of them store results. ``can_install`` - says whether the session's interpreter has pip, and defaults to asking - this one, which is the interpreter the session runs. + handles are named after whichever of them store results. ``can_plot`` + and ``can_install`` say whether the session can import matplotlib and + pip, and default to asking the interpreter the session runs. """ + if can_plot is None: + can_plot = session_can_import("matplotlib") if can_install is None: - can_install = importlib.util.find_spec("pip") is not None + can_install = session_can_import("pip") handle_tools = [name for name in HANDLE_TOOLS if name in tool_names] parts = [ ( @@ -102,8 +107,10 @@ def run_python_description( "calls for readability." ), ( - "Create at most one matplotlib figure per call and leave it open rather " - "than saving it." + "Create at most one matplotlib figure per call and leave it open " + "rather than saving it." + if can_plot + else "matplotlib is not installed, so the session cannot draw plots." ), ( "Do not use this tool to talk to the user; explanations belong in your " @@ -133,6 +140,19 @@ def run_python_description( return " ".join(parts) + "\n\nRules:" + "".join(f"\n- {rule}" for rule in rules) +def session_can_import(module: str) -> bool: + """Whether the session's interpreter can import the top-level ``module``. + + The session runs this interpreter under ``-I``, which leaves out the user + site directory, so a module found only there does not count. + """ + spec = importlib.util.find_spec(module) + if spec is None or spec.origin is None: + return spec is not None + user_site = os.path.abspath(site.getusersitepackages()) + os.sep + return not os.path.abspath(spec.origin).startswith(user_site) + + def _listed(names: Sequence[str]) -> str: if len(names) <= 2: return " and ".join(names) diff --git a/pkg-py/tests/test_run_python.py b/pkg-py/tests/test_run_python.py index 6f553167..9f5fb534 100644 --- a/pkg-py/tests/test_run_python.py +++ b/pkg-py/tests/test_run_python.py @@ -4,6 +4,7 @@ import base64 import importlib.util +from pathlib import Path from typing import Any import pandas as pd @@ -22,6 +23,7 @@ highlight_python, run_python_description, run_python_result, + session_can_import, ) from ._provider import scripted_chat, text @@ -120,6 +122,29 @@ def test_the_network_rule_says_what_the_session_can_reach( assert rule in description.split("\n\nRules:")[1] +def test_the_plot_rule_says_whether_the_session_can_draw() -> None: + def rules(can_plot: bool) -> str: + return run_python_description( + ["run_sql"], has_measures=False, network="none", can_plot=can_plot + ).split("\n\nRules:")[1] + + assert "at most one matplotlib figure per call" in rules(True) + assert "the session cannot draw plots" in rules(False) + + +def test_a_module_only_in_the_user_site_is_not_importable_in_the_session( + monkeypatch: pytest.MonkeyPatch, +) -> None: + assert session_can_import("pytest") + assert not session_can_import("no_such_module_for_commons") + spec = importlib.util.find_spec("pytest") + assert spec is not None and spec.origin is not None + # The session's -I drops the user site; pretend pytest was found there. + site_dir = str(Path(spec.origin).parent.parent) + monkeypatch.setattr("site.getusersitepackages", lambda: site_dir) + assert not session_can_import("pytest") + + # ---- the result ------------------------------------------------------------- From cdb2658bb48155e7bca84c4a248e4d7471be865d Mon Sep 17 00:00:00 2001 From: Josh Taillon Date: Fri, 9 Oct 2026 13:28:11 -0400 Subject: [PATCH 06/25] fix(py): ask the session's own -I interpreter what it can import Filtering the user site out of the host's import path still counted modules found through PYTHONPATH or the current directory, which -I also leaves out. session_can_import() now asks the interpreter the session runs, under -I, once per module per process. --- pkg-py/src/commons/_run_python.py | 32 +++++++++++++++++++++---------- pkg-py/tests/test_run_python.py | 16 ++++++++-------- 2 files changed, 30 insertions(+), 18 deletions(-) diff --git a/pkg-py/src/commons/_run_python.py b/pkg-py/src/commons/_run_python.py index 735fef88..21efe831 100644 --- a/pkg-py/src/commons/_run_python.py +++ b/pkg-py/src/commons/_run_python.py @@ -11,12 +11,12 @@ from __future__ import annotations import base64 +import functools import html -import importlib.util import io import keyword -import os -import site +import subprocess +import sys import tokenize from collections.abc import Sequence from typing import Any @@ -50,6 +50,9 @@ NO_OUTPUT = "(The code ran but produced no output.)" +# How long to wait for the session's interpreter to say what it can import. +PROBE_TIMEOUT = 10.0 + def run_python_description( tool_names: Sequence[str], @@ -140,17 +143,26 @@ def run_python_description( return " ".join(parts) + "\n\nRules:" + "".join(f"\n- {rule}" for rule in rules) +@functools.cache def session_can_import(module: str) -> bool: - """Whether the session's interpreter can import the top-level ``module``. + """Whether the session's interpreter can find the top-level ``module``. The session runs this interpreter under ``-I``, which leaves out the user - site directory, so a module found only there does not count. + site directory, ``PYTHONPATH``, and the current directory, so the answer + comes from asking that interpreter the same way. It is cached, since the + interpreter's packages do not change while it runs. """ - spec = importlib.util.find_spec(module) - if spec is None or spec.origin is None: - return spec is not None - user_site = os.path.abspath(site.getusersitepackages()) + os.sep - return not os.path.abspath(spec.origin).startswith(user_site) + probe = "import importlib.util, sys; sys.exit(importlib.util.find_spec(sys.argv[1]) is None)" + try: + completed = subprocess.run( + [sys.executable, "-I", "-c", probe, module], + capture_output=True, + timeout=PROBE_TIMEOUT, + check=False, + ) + except (OSError, subprocess.SubprocessError): + return False + return completed.returncode == 0 def _listed(names: Sequence[str]) -> str: diff --git a/pkg-py/tests/test_run_python.py b/pkg-py/tests/test_run_python.py index 9f5fb534..9134933f 100644 --- a/pkg-py/tests/test_run_python.py +++ b/pkg-py/tests/test_run_python.py @@ -132,17 +132,17 @@ def rules(can_plot: bool) -> str: assert "the session cannot draw plots" in rules(False) -def test_a_module_only_in_the_user_site_is_not_importable_in_the_session( - monkeypatch: pytest.MonkeyPatch, +def test_a_module_the_isolated_session_cannot_see_is_not_importable( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch ) -> None: assert session_can_import("pytest") assert not session_can_import("no_such_module_for_commons") - spec = importlib.util.find_spec("pytest") - assert spec is not None and spec.origin is not None - # The session's -I drops the user site; pretend pytest was found there. - site_dir = str(Path(spec.origin).parent.parent) - monkeypatch.setattr("site.getusersitepackages", lambda: site_dir) - assert not session_can_import("pytest") + # On PYTHONPATH and so on this process's path, but -I ignores both. + (tmp_path / "only_on_pythonpath.py").write_text("") + monkeypatch.setenv("PYTHONPATH", str(tmp_path)) + monkeypatch.syspath_prepend(str(tmp_path)) + assert importlib.util.find_spec("only_on_pythonpath") is not None + assert not session_can_import("only_on_pythonpath") # ---- the result ------------------------------------------------------------- From ddac943838fd863f40195b0cd5a50d37ffe5aed3 Mon Sep 17 00:00:00 2001 From: Josh Taillon Date: Fri, 9 Oct 2026 13:28:12 -0400 Subject: [PATCH 07/25] test(py): cover an async call cancelled by its worker thread's close --- pkg-py/tests/test_execution_thread.py | 14 ++++++++++++++ 1 file changed, 14 insertions(+) diff --git a/pkg-py/tests/test_execution_thread.py b/pkg-py/tests/test_execution_thread.py index 2d1aeeeb..e519db25 100644 --- a/pkg-py/tests/test_execution_thread.py +++ b/pkg-py/tests/test_execution_thread.py @@ -101,3 +101,17 @@ def test_close_cancels_a_running_call_rather_than_waiting_it_out() -> None: caller.join(10) assert time.monotonic() - started < 15 assert replies == [Failure(message="the Python session is closed.")] + + +def test_an_async_call_running_when_the_thread_closes_is_told_it_closed() -> None: + worker_thread = WorkerThread(Worker(call_timeout=60)) + + async def caller() -> object: + await worker_thread.run("1") + call = asyncio.ensure_future(worker_thread.run("import time; time.sleep(60)")) + await asyncio.sleep(1) + # close() blocks, so it runs off this loop, which keeps serving the call. + await asyncio.to_thread(worker_thread.close) + return await asyncio.wait_for(call, 10) + + assert asyncio.run(caller()) == Failure(message="the Python session is closed.") From 5febd029081ad237e745d555d819da614010d96f Mon Sep 17 00:00:00 2001 From: Josh Taillon Date: Fri, 9 Oct 2026 13:29:54 -0400 Subject: [PATCH 08/25] fix(py): probe imports with the session's environment and scratch directory A sitecustomize or .pth hook still runs under -I and can move import paths by environment, so the probe now runs with worker_env() and an empty scratch directory as its working directory, as the session does. --- pkg-py/src/commons/_run_python.py | 20 +++++++++++++------- 1 file changed, 13 insertions(+), 7 deletions(-) diff --git a/pkg-py/src/commons/_run_python.py b/pkg-py/src/commons/_run_python.py index 21efe831..bf0af8e1 100644 --- a/pkg-py/src/commons/_run_python.py +++ b/pkg-py/src/commons/_run_python.py @@ -17,6 +17,7 @@ import keyword import subprocess import sys +import tempfile import tokenize from collections.abc import Sequence from typing import Any @@ -29,6 +30,7 @@ from ._display import CODE_ANALYSIS, visible_result_note from ._execution._backend import Network from ._execution._driver import Failure +from ._execution._env import worker_env from ._execution._protocol import Error, OpaqueValue, Plot, Result from ._execution._thread import WorkerThread from ._frames import describe_frame, is_frame @@ -149,17 +151,21 @@ def session_can_import(module: str) -> bool: The session runs this interpreter under ``-I``, which leaves out the user site directory, ``PYTHONPATH``, and the current directory, so the answer - comes from asking that interpreter the same way. It is cached, since the + comes from asking that interpreter the same way, with the session's + environment and an empty scratch directory. It is cached, since the interpreter's packages do not change while it runs. """ probe = "import importlib.util, sys; sys.exit(importlib.util.find_spec(sys.argv[1]) is None)" try: - completed = subprocess.run( - [sys.executable, "-I", "-c", probe, module], - capture_output=True, - timeout=PROBE_TIMEOUT, - check=False, - ) + with tempfile.TemporaryDirectory(prefix="commons-probe-") as scratch: + completed = subprocess.run( + [sys.executable, "-I", "-c", probe, module], + capture_output=True, + cwd=scratch, + env=worker_env(scratch), + timeout=PROBE_TIMEOUT, + check=False, + ) except (OSError, subprocess.SubprocessError): return False return completed.returncode == 0 From 23c92d0ee98e5715b317b70cc8899aedd6658631 Mon Sep 17 00:00:00 2001 From: Josh Taillon Date: Fri, 9 Oct 2026 13:32:28 -0400 Subject: [PATCH 09/25] fix(py): give the import probe the worker's single-thread environment LocalBackend sets the thread-count variables when the sandbox needs a single-threaded worker, so the probe builds its environment with the same needs_single_thread() answer. --- pkg-py/src/commons/_run_python.py | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/pkg-py/src/commons/_run_python.py b/pkg-py/src/commons/_run_python.py index bf0af8e1..b81ef544 100644 --- a/pkg-py/src/commons/_run_python.py +++ b/pkg-py/src/commons/_run_python.py @@ -32,6 +32,7 @@ from ._execution._driver import Failure from ._execution._env import worker_env from ._execution._protocol import Error, OpaqueValue, Plot, Result +from ._execution._sandbox import needs_single_thread from ._execution._thread import WorkerThread from ._frames import describe_frame, is_frame from ._prompt import EXECUTION_TOOL @@ -162,7 +163,7 @@ def session_can_import(module: str) -> bool: [sys.executable, "-I", "-c", probe, module], capture_output=True, cwd=scratch, - env=worker_env(scratch), + env=worker_env(scratch, single_thread=needs_single_thread()), timeout=PROBE_TIMEOUT, check=False, ) From 8cc2424a82b884b38779b97d1983d37c41126790 Mon Sep 17 00:00:00 2001 From: Josh Taillon Date: Fri, 9 Oct 2026 20:13:42 -0400 Subject: [PATCH 10/25] feat(py): build run_python's result from the reply's ordered output Replies now carry one ordered output of text and plots. run_python's result walks it: adjacent text joins into one run, each plot sits where it was drawn in what the model receives, and the value or the traceback starts a line after everything the call wrote. An error therefore shows what the call printed before it raised. --- pkg-py/src/commons/_run_python.py | 85 +++++++++++++++++++------------ pkg-py/tests/test_run_python.py | 40 +++++++++++---- 2 files changed, 82 insertions(+), 43 deletions(-) diff --git a/pkg-py/src/commons/_run_python.py b/pkg-py/src/commons/_run_python.py index b81ef544..49ffc4df 100644 --- a/pkg-py/src/commons/_run_python.py +++ b/pkg-py/src/commons/_run_python.py @@ -31,7 +31,7 @@ from ._execution._backend import Network from ._execution._driver import Failure from ._execution._env import worker_env -from ._execution._protocol import Error, OpaqueValue, Plot, Result +from ._execution._protocol import Error, OpaqueValue, Plot, Result, Text from ._execution._sandbox import needs_single_thread from ._execution._thread import WorkerThread from ._frames import describe_frame, is_frame @@ -224,35 +224,51 @@ def build(func: Any) -> Tool: def run_python_result(code: str, reply: Result | Error | Failure) -> ContentToolResult: """The tool result for one call: the model's view and the reader's. - The model gets the text the call produced and each plot as an image. - The reader gets the code with its output, and the plots at their size. + The model gets what the call wrote, with each plot as an image in the + place it was drawn. The reader gets the code with its output, and the + plots at their size. """ - plots: tuple[Plot, ...] = () - if isinstance(reply, Failure): - texts = [f"Error: {reply.message}"] - elif isinstance(reply, Error): - texts = [reply.traceback.strip() or f"Error: {reply.message}"] - plots = reply.plots - else: - texts = [ - text - for text in ( - reply.stdout.rstrip("\n"), - reply.stderr.rstrip("\n"), - _value_text(reply.value), - ) - if text - ] - plots = reply.plots + runs = _runs(reply) + plots = [run for run in runs if isinstance(run, Plot)] return tool_result( - _model_value(texts, plots), + _model_value(runs), ProvenanceTag.B, title=CODE_ANALYSIS.settled, - html=_display_html(code, texts, plots), + html=_display_html(code, runs), open=bool(plots), ) +def _runs(reply: Result | Error | Failure) -> list[str | Plot]: + """The reply's output in order, with adjacent text joined into one run. + + Streams are joined as written, so a line split across two writes stays + one line. The value a call ended on, or its error, starts a line of its + own after everything the call wrote. + """ + if isinstance(reply, Failure): + return [f"Error: {reply.message}"] + if isinstance(reply, Error): + last = reply.traceback.strip() or f"Error: {reply.message}" + else: + last = _value_text(reply.value) + runs: list[str | Plot] = [] + text = "" + for segment in reply.output: + if isinstance(segment, Text): + text += segment.text + continue + if text.strip("\n"): + runs.append(text.rstrip("\n")) + text = "" + runs.append(segment) + if last: + text += ("\n" if text and not text.endswith("\n") else "") + last + if text.strip("\n"): + runs.append(text.rstrip("\n")) + return runs + + def _value_text(value: Any) -> str: """The value a call ended on, as a REPL would show it; empty for None.""" if value is None: @@ -265,24 +281,27 @@ def _value_text(value: Any) -> str: return repr(value) -def _model_value(texts: list[str], plots: tuple[Plot, ...]) -> Any: - text = "\n".join(texts) - if not plots: - return text or NO_OUTPUT - parts: list[Any] = [ContentText(text=text)] if text else [] - parts.extend( +def _model_value(runs: list[str | Plot]) -> Any: + if not any(isinstance(run, Plot) for run in runs): + return "\n".join(run for run in runs if isinstance(run, str)) or NO_OUTPUT + parts: list[Any] = [ ContentImageInline( image_content_type="image/png", - data=base64.b64encode(plot.png).decode("ascii"), + data=base64.b64encode(run.png).decode("ascii"), ) - for plot in plots - ) + if isinstance(run, Plot) + else ContentText(text=run) + for run in runs + ] parts.append(ContentText(text=visible_result_note("plot"))) return parts -def _display_html(code: str, texts: list[str], plots: tuple[Plot, ...]) -> Tag: - output = [f"#> {line}" for text in texts for line in text.split("\n")] +def _display_html(code: str, runs: list[str | Plot]) -> Tag: + output = [ + f"#> {line}" for run in runs if isinstance(run, str) for line in run.split("\n") + ] + plots = [run for run in runs if isinstance(run, Plot)] block: Tag = tags.pre( tags.code( HTML(highlight_python("\n".join([code, *output]))), diff --git a/pkg-py/tests/test_run_python.py b/pkg-py/tests/test_run_python.py index 9134933f..cef064e0 100644 --- a/pkg-py/tests/test_run_python.py +++ b/pkg-py/tests/test_run_python.py @@ -16,7 +16,7 @@ from commons._agent import Commons from commons._display import DISPLAY_EXTRA_KEY from commons._execution._driver import Failure -from commons._execution._protocol import Error, OpaqueValue, Plot, Result +from commons._execution._protocol import Error, OpaqueValue, Plot, Result, Text from commons._provenance import TAG_EXTRA_KEY, Tag from commons._run_python import ( NO_OUTPUT, @@ -149,10 +149,17 @@ def test_a_module_the_isolated_session_cannot_see_is_not_importable( def test_a_result_shows_what_was_printed_and_the_value() -> None: + output = ( + Text(stream="stdout", text="hi\n"), + Text(stream="stderr", text="warn\n"), + Text(stream="stdout", text="partial"), + Text(stream="stdout", text=" line\n"), + ) result = run_python_result( - "print('hi')\n1 + 1", Result(id="c1", value=2, stdout="hi\n", stderr="warn\n") + "print('hi')\n1 + 1", Result(id="c1", value=2, output=output) ) - assert result.value == "hi\nwarn\n2" + # Streams interleave as written; the value starts its own line. + assert result.value == "hi\nwarn\npartial line\n2" assert result.extra is not None assert result.extra[TAG_EXTRA_KEY] == Tag.B assert display(result)["title"] == "Analyzed data" @@ -177,8 +184,12 @@ def test_an_error_shows_the_traceback_and_a_failure_its_message() -> None: id="c1", message="NameError: name 'y' is not defined", traceback="Traceback...\nNameError", + output=(Text(stream="stdout", text="got this far"),), + ) + # What ran before the error is kept, so the model sees how far it got. + assert run_python_result("y", error).value == ( + "got this far\nTraceback...\nNameError" ) - assert run_python_result("y", error).value == "Traceback...\nNameError" failure = run_python_result("1", Failure(message="the Python session crashed.")) assert failure.value == "Error: the Python session crashed." assert failure.extra is not None @@ -186,15 +197,24 @@ def test_an_error_shows_the_traceback_and_a_failure_its_message() -> None: def test_plots_reach_the_model_as_images_and_the_reader_at_their_size() -> None: - result = run_python_result( - "fig", Result(id="c1", stdout="drawn\n", plots=(plot(),)) + output = ( + Text(stream="stdout", text="before\n"), + plot(), + Text(stream="stdout", text="after\n"), ) + result = run_python_result("fig", Result(id="c1", value=3, output=output)) parts = result.value assert isinstance(parts, list) - assert isinstance(parts[0], ContentText) and parts[0].text == "drawn" - image = parts[1] - assert isinstance(image, ContentImageInline) - assert base64.b64decode(image.data) == plot().png + # The plot keeps its place between the text written before and after it. + assert [type(part) for part in parts] == [ + ContentText, + ContentImageInline, + ContentText, + ContentText, + ] + assert parts[0].text == "before" + assert base64.b64decode(parts[1].data) == plot().png + assert parts[2].text == "after\n3" assert "already visible to the user" in parts[-1].text shown = display(result) assert shown["open"] is True From 954f96aa3aad5794eb6ae606a5563f2b72cd4120 Mon Sep 17 00:00:00 2001 From: Josh Taillon Date: Fri, 9 Oct 2026 21:42:12 -0400 Subject: [PATCH 11/25] fix(py): tell the model that plt.show() places a plot and savefig hides it --- pkg-py/src/commons/_run_python.py | 6 ++++-- 1 file changed, 4 insertions(+), 2 deletions(-) diff --git a/pkg-py/src/commons/_run_python.py b/pkg-py/src/commons/_run_python.py index 49ffc4df..6d8e9aee 100644 --- a/pkg-py/src/commons/_run_python.py +++ b/pkg-py/src/commons/_run_python.py @@ -113,8 +113,10 @@ def run_python_description( "calls for readability." ), ( - "Create at most one matplotlib figure per call and leave it open " - "rather than saving it." + "Create at most one matplotlib figure per call. Draw it with pyplot " + "and do not call savefig, since a saved file reaches neither you nor " + "the user. The figure appears where you call plt.show(), or after " + "the call's text if you do not." if can_plot else "matplotlib is not installed, so the session cannot draw plots." ), From 04bc1170d287f23b4a141331e43b97139ac8f760 Mon Sep 17 00:00:00 2001 From: Josh Taillon Date: Fri, 9 Oct 2026 21:54:59 -0400 Subject: [PATCH 12/25] fix(py): tell the model that fig.show() places a plot too --- pkg-py/src/commons/_run_python.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/pkg-py/src/commons/_run_python.py b/pkg-py/src/commons/_run_python.py index 6d8e9aee..7973172e 100644 --- a/pkg-py/src/commons/_run_python.py +++ b/pkg-py/src/commons/_run_python.py @@ -115,8 +115,8 @@ def run_python_description( ( "Create at most one matplotlib figure per call. Draw it with pyplot " "and do not call savefig, since a saved file reaches neither you nor " - "the user. The figure appears where you call plt.show(), or after " - "the call's text if you do not." + "the user. The figure appears where you call plt.show() or " + "fig.show(), or after the call's text if you call neither." if can_plot else "matplotlib is not installed, so the session cannot draw plots." ), From fbe74d60cce30ec600d2e389b3ee8efae5a9b7ae Mon Sep 17 00:00:00 2001 From: Josh Taillon Date: Fri, 9 Oct 2026 22:42:33 -0400 Subject: [PATCH 13/25] fix(py): stop a call on Ctrl-C, and close safely from the session's thread Ctrl-C while a sync caller waits now cancels the call, as cancelling an async caller does, so the next call does not queue behind it until its timeout. close() called from the loop's own thread (a garbage collection there) starts the shutdown and returns instead of blocking the loop it waits on, and a shutdown that raises still joins the thread and closes the loop. The no-thread-before-first-call test checks the thread itself rather than the process's thread count, and the docstrings say what they mean in plain terms. --- pkg-py/src/commons/_execution/_thread.py | 56 ++++++++++++++---------- pkg-py/tests/test_execution_thread.py | 46 ++++++++++++++++++- 2 files changed, 78 insertions(+), 24 deletions(-) diff --git a/pkg-py/src/commons/_execution/_thread.py b/pkg-py/src/commons/_execution/_thread.py index a9c046b8..82724f09 100644 --- a/pkg-py/src/commons/_execution/_thread.py +++ b/pkg-py/src/commons/_execution/_thread.py @@ -1,18 +1,17 @@ -"""A worker on an event loop of its own, so any caller can reach one session. - -chatlas runs an async tool only from ``stream_async()``, and a sync tool -directly on the caller's event loop, where a long call would stall every -other task on it. Keeping the ``Worker`` on a private loop in a background -thread lets an async caller await a call without blocking its loop, and a -sync caller block on the same session, so variables persist whichever way -the agent is asked. +"""Runs the session's ``Worker`` in a background thread with its own event loop. + +chatlas runs an async tool only from ``stream_async()``, and runs a sync tool +on the caller's event loop, where a long call would stop every other task on +that loop. With the ``Worker`` in its own thread, an async caller can wait for +a call without blocking its loop, and a sync caller can block on the same +session. Variables therefore persist whether the agent is used through +``chat()`` or ``stream_async()``. """ from __future__ import annotations import asyncio import concurrent.futures -import contextlib import threading from .._handles import HandleStore @@ -31,7 +30,7 @@ class WorkerThread: - """Runs ``worker`` on a loop in a daemon thread, started by the first call.""" + """Runs ``worker`` in a background thread, which starts on the first call.""" def __init__(self, worker: Worker) -> None: self._worker = worker @@ -46,8 +45,8 @@ def __init__(self, worker: Worker) -> None: async def run(self, code: str, handles: HandleStore | None = None) -> _Reply: """Run ``code`` without blocking the caller's event loop. - Cancelling the caller cancels the call, which shuts the worker down - the way a cancelled ``Worker.run`` does. + If the caller is cancelled, the call is cancelled too, and the worker + shuts down as it does when ``Worker.run`` is cancelled. """ future = self._submit(code, handles) if future is None: @@ -73,14 +72,20 @@ def run_sync(self, code: str, handles: HandleStore | None = None) -> _Reply: if self._closed: return _CLOSED raise + except BaseException: + # Ctrl-C while waiting stops the call, as cancelling an async + # caller does, rather than leaving it to run out its timeout. + future.cancel() + raise def close(self) -> None: - """Close the worker, then stop the loop and its thread. Safe to repeat. + """Close the worker, then stop the loop and its thread. - Calls in flight are cancelled first, so the worker shuts down within - its grace periods rather than after a call's timeout. A shutdown that - still overruns ``CLOSE_TIMEOUT`` finishes in the background, and the - loop stops only once it has. + Calling it again does nothing. Running calls are cancelled first, so + the worker stops quickly instead of waiting for a call to time out. If + the shutdown takes longer than ``CLOSE_TIMEOUT``, it continues in the + background, and the loop stops when it ends. When called from the + loop's own thread, it starts the shutdown and returns at once. """ with self._lock: if self._closed: @@ -94,11 +99,18 @@ def close(self) -> None: call.cancel() closing = asyncio.run_coroutine_threadsafe(self._worker.aclose(), loop) closing.add_done_callback(lambda _: loop.call_soon_threadsafe(loop.stop)) - with contextlib.suppress(TimeoutError): + # Waiting here on the loop's own thread would block the shutdown. + if threading.current_thread() is thread: + return + try: closing.result(CLOSE_TIMEOUT) - thread.join(CLOSE_TIMEOUT) - if not thread.is_alive(): - loop.close() + except TimeoutError: + return + finally: + if closing.done(): + thread.join(CLOSE_TIMEOUT) + if not thread.is_alive(): + loop.close() def _submit( self, code: str, handles: HandleStore | None @@ -121,7 +133,7 @@ def _forget(self, future: concurrent.futures.Future[_Reply]) -> None: self._calls.discard(future) def _ensure_loop(self) -> asyncio.AbstractEventLoop: - """The loop, started on first use. Called with the lock held.""" + """Return the loop, starting it and its thread if needed. Needs the lock.""" if self._loop is None: loop = asyncio.new_event_loop() thread = threading.Thread( diff --git a/pkg-py/tests/test_execution_thread.py b/pkg-py/tests/test_execution_thread.py index e519db25..c74d9d25 100644 --- a/pkg-py/tests/test_execution_thread.py +++ b/pkg-py/tests/test_execution_thread.py @@ -3,6 +3,7 @@ from __future__ import annotations import asyncio +import signal import threading import time @@ -55,10 +56,51 @@ async def tick() -> None: def test_no_thread_starts_before_the_first_call() -> None: - before = threading.active_count() worker_thread = WorkerThread(Worker()) - assert threading.active_count() == before + assert worker_thread._thread is None worker_thread.close() + assert worker_thread._thread is None + + +def test_run_sync_refuses_the_workers_own_thread(runner: WorkerThread) -> None: + runner.run_sync("1") + assert runner._loop is not None + + async def from_the_loop() -> None: + runner.run_sync("1") + + future = asyncio.run_coroutine_threadsafe(from_the_loop(), runner._loop) + with pytest.raises(RuntimeError, match="deadlock"): + future.result(10) + + +def test_ctrl_c_while_waiting_stops_the_call() -> None: + worker_thread = WorkerThread(Worker(call_timeout=60)) + try: + worker_thread.run_sync("1") + main = threading.main_thread().ident + assert main is not None + threading.Timer(1, signal.pthread_kill, (main, signal.SIGINT)).start() + with pytest.raises(KeyboardInterrupt): + worker_thread.run_sync("import time; time.sleep(60)") + started = time.monotonic() + # Without the cancel, this call would queue behind the sleep. + reply = worker_thread.run_sync("1 + 1") + assert time.monotonic() - started < 15 + assert isinstance(reply, Result) + assert reply.value == 2 + finally: + worker_thread.close() + + +def test_close_from_the_loops_own_thread_returns_and_the_thread_stops() -> None: + worker_thread = WorkerThread(Worker()) + worker_thread.run_sync("1") + loop, thread = worker_thread._loop, worker_thread._thread + assert loop is not None and thread is not None + loop.call_soon_threadsafe(worker_thread.close) + thread.join(15) + assert not thread.is_alive() def test_a_cancelled_caller_cancels_the_call(runner: WorkerThread) -> None: From ac20001fb6289d8535b7602eceec3b593fab9c64 Mon Sep 17 00:00:00 2001 From: Josh Taillon Date: Fri, 9 Oct 2026 22:42:33 -0400 Subject: [PATCH 14/25] fix(py): tighten run_python's wiring, highlighting, and probe - Highlighting splits lines as tokenize does, so a carriage return or form feed in the output no longer shifts spans onto the wrong text. - pip is probed only when the session has network access. - Without matplotlib, the description no longer opens by offering plots. - chat() swaps in the sync tool only while run_python is registered, so a tool the user removed stays removed. - Commons' network argument is typed inline, since the alias is private. Tests now cover the network argument reaching the description and annotations, the async tool coming back after a failed turn, the citation request on a run_python result, and the finalizer stopping the session's thread. The docstrings are rewritten in plain terms. --- pkg-py/src/commons/_agent.py | 25 ++++++----- pkg-py/src/commons/_run_python.py | 73 +++++++++++++++++-------------- pkg-py/tests/test_run_python.py | 70 ++++++++++++++++++++++++++++- 3 files changed, 124 insertions(+), 44 deletions(-) diff --git a/pkg-py/src/commons/_agent.py b/pkg-py/src/commons/_agent.py index d829e436..577f5533 100644 --- a/pkg-py/src/commons/_agent.py +++ b/pkg-py/src/commons/_agent.py @@ -30,12 +30,12 @@ from ._context_layer import ContextLayer, augment_context_layer from ._data_source import DataSource from ._definitions import Registry, build_registry -from ._execution._backend import Network from ._execution._driver import Worker from ._execution._thread import WorkerThread from ._handles import HandleStore from ._measures import SemanticLayer, resolve_injections, semantic_layer from ._prompt import ( + EXECUTION_TOOL, check_instructions, read_instructions, render_system_prompt, @@ -89,13 +89,13 @@ class Commons(Chat[Any, Any]): `## Additional instructions` heading at the end of commons' built-in system prompt, as a string or the path to a text or Markdown file. - The agent runs the Python code its model writes in a sandboxed session, - which starts on the first call. `network` is whether that session can - reach the network: `"none"` (the default) or `"full"`. The session is - sandboxed on Linux and macOS. On any other host, local development can - opt in to best-effort guardrails by setting the - `COMMONS_ALLOW_UNSAFE_FALLBACK` environment variable; these guardrails - are not a security boundary. + The agent can run Python code that its model writes. The code runs in a + separate, sandboxed Python process, which starts the first time the model + runs code. `network` sets whether that process can reach the network: + `"none"` (the default) or `"full"`. The sandbox works on Linux and macOS. + For local development on another system, set the + `COMMONS_ALLOW_UNSAFE_FALLBACK` environment variable to run the code with + limited checks instead. These checks are not a security boundary. Construction raises a TypeError if `client` is not a `chatlas.Chat`, if an entry of `data_sources` is not a `DataSource`, or if a layer is not @@ -114,7 +114,7 @@ def __init__( context_layer: ContextLayer | None = None, *, instructions: str | None = None, - network: Network = "none", + network: Literal["none", "full"] = "none", ) -> None: if not isinstance(client, Chat): raise TypeError( @@ -241,13 +241,16 @@ def chat( self._citation_request.reset() # chatlas refuses a synchronous chat while an async tool is # registered, so the sync run_python stands in for this call. - self.register_tool(self._run_python_sync, force=True) + swap = any(tool.name == EXECUTION_TOOL for tool in self.get_tools()) + if swap: + self.register_tool(self._run_python_sync, force=True) try: response = super().chat( *inputs, echo=echo, stream=stream, kwargs=kwargs ) finally: - self.register_tool(self._run_python, force=True) + if swap: + self.register_tool(self._run_python, force=True) self._consume_restore_reminder(was_pending) return response diff --git a/pkg-py/src/commons/_run_python.py b/pkg-py/src/commons/_run_python.py index 7973172e..7b6b1217 100644 --- a/pkg-py/src/commons/_run_python.py +++ b/pkg-py/src/commons/_run_python.py @@ -1,11 +1,13 @@ -"""The `run_python` tool: model-written Python run in the agent's sandboxed session. +"""The `run_python` tool, which runs the model's code in a sandboxed session. -`pkg-r/R/run-r.R` builds R's `run_r`, and the two tools describe themselves -and shape their results the same way, each in its own language's idiom. +`pkg-r/R/run-r.R` builds R's `run_r`. The two tools follow one contract for +what they tell the model and what they return, each in its own language's +idiom. -The agent registers the async tool, which awaits the session without blocking -the caller's event loop. chatlas refuses a synchronous `chat()` while any async -tool is registered, so `Commons.chat()` swaps in the sync tool for the call. +The agent registers the async version of the tool, which waits for the session +without blocking the caller's event loop. chatlas refuses a synchronous +`chat()` while any async tool is registered, so `Commons.chat()` registers the +sync version for the length of the call. """ from __future__ import annotations @@ -65,16 +67,17 @@ def run_python_description( can_plot: bool | None = None, can_install: bool | None = None, ) -> str: - """What `run_python` tells the model it is for. + """The description that tells the model what `run_python` does. - ``tool_names`` are the agent's other registered tools; the preloaded - handles are named after whichever of them store results. ``can_plot`` - and ``can_install`` say whether the session can import matplotlib and - pip, and default to asking the interpreter the session runs. + ``tool_names`` are the agent's other registered tools. The description + lists the results of those that store them as preloaded variables. + ``can_plot`` and ``can_install`` say whether the session can import + matplotlib and pip. When omitted, they are found by asking the session's + interpreter; pip is checked only when the session has network access. """ if can_plot is None: can_plot = session_can_import("matplotlib") - if can_install is None: + if can_install is None and network != "none": can_install = session_can_import("pip") handle_tools = [name for name in HANDLE_TOOLS if name in tool_names] parts = [ @@ -82,6 +85,9 @@ def run_python_description( "Run Python code in your sandboxed Python session to analyze results or " "render plots. Python code and textual output are visible only to you; " "rendered plots are also shown to the user." + if can_plot + else "Run Python code in your sandboxed Python session to analyze " + "results. Python code and its output are visible only to you." ), ( "The user cannot access or interact with this session. Never direct them " @@ -150,13 +156,14 @@ def run_python_description( @functools.cache def session_can_import(module: str) -> bool: - """Whether the session's interpreter can find the top-level ``module``. - - The session runs this interpreter under ``-I``, which leaves out the user - site directory, ``PYTHONPATH``, and the current directory, so the answer - comes from asking that interpreter the same way, with the session's - environment and an empty scratch directory. It is cached, since the - interpreter's packages do not change while it runs. + """Whether the session's interpreter can import the top-level ``module``. + + The session starts Python in isolated mode (``-I``), which ignores the + user's site-packages directory, ``PYTHONPATH``, and the current directory. + The check therefore starts the same interpreter the same way, with the + session's environment variables and an empty temporary directory. The + answer is cached, because the installed packages do not change while the + process runs. """ probe = "import importlib.util, sys; sys.exit(importlib.util.find_spec(sys.argv[1]) is None)" try: @@ -183,7 +190,7 @@ def _listed(names: Sequence[str]) -> str: def run_python_tools( runner: WorkerThread, context: ToolContext, description: str, network: Network ) -> tuple[Tool, Tool]: - """The async tool the agent registers, and the sync one `chat()` swaps in.""" + """The async tool the agent registers, and the sync tool `chat()` uses instead.""" def finish(code: str, reply: Result | Error | Failure) -> ContentToolResult: result = run_python_result(code, reply) @@ -224,11 +231,11 @@ def build(func: Any) -> Tool: def run_python_result(code: str, reply: Result | Error | Failure) -> ContentToolResult: - """The tool result for one call: the model's view and the reader's. + """The tool result for one call: one view for the model, one for the user. - The model gets what the call wrote, with each plot as an image in the - place it was drawn. The reader gets the code with its output, and the - plots at their size. + The model gets the call's output, with each plot as an image at the point + where it was drawn. The user sees the code and its output, followed by the + plots. """ runs = _runs(reply) plots = [run for run in runs if isinstance(run, Plot)] @@ -242,11 +249,11 @@ def run_python_result(code: str, reply: Result | Error | Failure) -> ContentTool def _runs(reply: Result | Error | Failure) -> list[str | Plot]: - """The reply's output in order, with adjacent text joined into one run. + """The reply's output in order, with consecutive text joined into one string. - Streams are joined as written, so a line split across two writes stays - one line. The value a call ended on, or its error, starts a line of its - own after everything the call wrote. + Text written to stdout and stderr is joined in the order it was written, so + a line printed in two parts stays one line. The call's final value, or its + error, goes on a new line after all the other output. """ if isinstance(reply, Failure): return [f"Error: {reply.message}"] @@ -272,7 +279,7 @@ def _runs(reply: Result | Error | Failure) -> list[str | Plot]: def _value_text(value: Any) -> str: - """The value a call ended on, as a REPL would show it; empty for None.""" + """The call's final value as the Python REPL shows it; empty for None.""" if value is None: return "" if isinstance(value, OpaqueValue): @@ -339,11 +346,13 @@ def _display_html(code: str, runs: list[str | Plot]) -> Tag: def highlight_python(source: str) -> str: - """``source`` as escaped HTML, with its tokens wrapped for the stylesheet. + """``source`` as escaped HTML, with each token in a span the stylesheet colors. - Text that does not tokenize as Python is escaped and left plain. + Text that Python cannot tokenize is escaped without highlighting. """ - lines = source.splitlines(keepends=True) + # Split as tokenize reads, on newlines only, so offsets line up with its + # rows even when the text holds a carriage return or a form feed. + lines = io.StringIO(source).readlines() starts = [0] for line in lines: starts.append(starts[-1] + len(line)) diff --git a/pkg-py/tests/test_run_python.py b/pkg-py/tests/test_run_python.py index cef064e0..b134517b 100644 --- a/pkg-py/tests/test_run_python.py +++ b/pkg-py/tests/test_run_python.py @@ -3,9 +3,12 @@ from __future__ import annotations import base64 +import gc +import html import importlib.util +import re from pathlib import Path -from typing import Any +from typing import Any, NoReturn import pandas as pd import pytest @@ -249,6 +252,12 @@ def test_text_that_does_not_tokenize_is_escaped_plainly() -> None: assert highlight_python("''' None: + out = highlight_python('#> 50%\r100%\n# a\x0cb\ng("z")') + assert 'g("z")' in out + assert html.unescape(re.sub(r"<[^>]+>", "", out)) == '#> 50%\r100%\n# a\x0cb\ng("z")' + + # ---- through an agent ------------------------------------------------------- @@ -325,3 +334,62 @@ async def test_a_plot_reaches_the_model_as_an_image() -> None: last_request = agent.provider.requests[-1] # type: ignore[attr-defined] sent = [content for turn in last_request for content in turn.contents] assert any(isinstance(content, ContentImageInline) for content in sent) + + +def test_the_network_argument_reaches_the_description_and_annotations() -> None: + agent = Commons(scripted_chat(), data_source(sales=frame()), network="full") + tool = run_python_tool(agent) + rules = tool.schema["function"]["description"].split("\n\nRules:")[1] + assert "The session has network access" in rules + assert tool.annotations is not None + assert tool.annotations.get("openWorldHint") is True + + +def test_a_network_other_than_none_or_full_is_refused() -> None: + with pytest.raises(ValueError): + Commons(scripted_chat(), data_source(sales=frame()), network="some") # type: ignore[arg-type] + + +def test_chat_restores_the_async_tool_when_the_turn_fails( + monkeypatch: pytest.MonkeyPatch, +) -> None: + agent = Commons(scripted_chat(), data_source(sales=frame())) + + def unavailable(**_: Any) -> NoReturn: + raise ConnectionError("provider unavailable") + + monkeypatch.setattr(agent.provider, "chat_perform", unavailable) + with pytest.raises(ConnectionError): + agent.chat("Hello?", echo="none") + assert run_python_tool(agent)._is_async + + +def test_chat_leaves_a_removed_run_python_removed() -> None: + agent = Commons(scripted_chat([text("Hi.")]), data_source(sales=frame())) + agent.set_tools([tool for tool in agent.get_tools() if tool.name != "run_python"]) + agent.chat("Hello?", echo="none") + assert "run_python" not in {tool.name for tool in agent.get_tools()} + + +def test_the_first_run_python_result_of_a_turn_asks_for_citations() -> None: + agent = Commons( + scripted_chat([tool_request(code="1 + 1"), text("Done.")]), + data_source(sales=frame()), + ) + agent.chat("What is one plus one?", echo="none") + run = tool_results(agent)[-1] + assert run.value == f"2\n\n{agent._citation_request.reminder}" + + +def test_dropping_the_agent_stops_its_session_thread() -> None: + agent = Commons( + scripted_chat([tool_request(code="1"), text("Done.")]), + data_source(sales=frame()), + ) + agent.chat("Go.", echo="none") + thread = agent._python._thread + assert thread is not None and thread.is_alive() + del agent + gc.collect() + thread.join(15) + assert not thread.is_alive() From 9a0cd5664e3558bfbe4b074baee4e00ac501f1d2 Mon Sep 17 00:00:00 2001 From: Josh Taillon Date: Fri, 9 Oct 2026 22:42:34 -0400 Subject: [PATCH 15/25] docs: note the code-run class rename in NEWS and name both tools in the stylesheet --- pkg-py/src/commons/www/commons-chat/commons-chat.css | 2 +- pkg-r/NEWS.md | 2 ++ pkg-r/inst/www/commons-chat/commons-chat.css | 2 +- www/commons-chat/commons-chat.css | 2 +- 4 files changed, 5 insertions(+), 3 deletions(-) diff --git a/pkg-py/src/commons/www/commons-chat/commons-chat.css b/pkg-py/src/commons/www/commons-chat/commons-chat.css index 82b9df25..fb8f4647 100644 --- a/pkg-py/src/commons/www/commons-chat/commons-chat.css +++ b/pkg-py/src/commons/www/commons-chat/commons-chat.css @@ -556,7 +556,7 @@ shiny-chat-container white-space: normal; } -/* ---- run_r display ---------------------------------------------------- */ +/* ---- run_r and run_python display -------------------------------------- */ .commons-run-display { display: grid; diff --git a/pkg-r/NEWS.md b/pkg-r/NEWS.md index 0913d478..9482608a 100644 --- a/pkg-r/NEWS.md +++ b/pkg-r/NEWS.md @@ -2,6 +2,8 @@ * `commons()` gains a `mode` argument. With `mode = "trusted only"`, the agent answers only with trusted calculations. It never writes its own SQL or R, and it tells the user when no trusted calculation answers a question. These answers carry no provenance markers. +* The CSS classes on `run_r`'s display are now `commons-run-display`, `commons-run-details`, `commons-run-code`, and `commons-run-plot`, without the `-r`, because the Python package's code tool uses the same classes. App CSS that targets the old `commons-run-r-*` names needs updating. + # commons 0.1.1 * Fixes an issue with the R code sandbox where the generated policy would be diff --git a/pkg-r/inst/www/commons-chat/commons-chat.css b/pkg-r/inst/www/commons-chat/commons-chat.css index 82b9df25..fb8f4647 100644 --- a/pkg-r/inst/www/commons-chat/commons-chat.css +++ b/pkg-r/inst/www/commons-chat/commons-chat.css @@ -556,7 +556,7 @@ shiny-chat-container white-space: normal; } -/* ---- run_r display ---------------------------------------------------- */ +/* ---- run_r and run_python display -------------------------------------- */ .commons-run-display { display: grid; diff --git a/www/commons-chat/commons-chat.css b/www/commons-chat/commons-chat.css index 82b9df25..fb8f4647 100644 --- a/www/commons-chat/commons-chat.css +++ b/www/commons-chat/commons-chat.css @@ -556,7 +556,7 @@ shiny-chat-container white-space: normal; } -/* ---- run_r display ---------------------------------------------------- */ +/* ---- run_r and run_python display -------------------------------------- */ .commons-run-display { display: grid; From e238778873d32705ebcfeabd8ec5ebf62c5ed1b9 Mon Sep 17 00:00:00 2001 From: Josh Taillon Date: Fri, 9 Oct 2026 22:44:18 -0400 Subject: [PATCH 16/25] test(py): cover a worker shutdown that raises --- pkg-py/tests/test_execution_thread.py | 17 +++++++++++++++++ 1 file changed, 17 insertions(+) diff --git a/pkg-py/tests/test_execution_thread.py b/pkg-py/tests/test_execution_thread.py index c74d9d25..5c16ca69 100644 --- a/pkg-py/tests/test_execution_thread.py +++ b/pkg-py/tests/test_execution_thread.py @@ -157,3 +157,20 @@ async def caller() -> object: return await asyncio.wait_for(call, 10) assert asyncio.run(caller()) == Failure(message="the Python session is closed.") + + +class FailingClose(Worker): + async def aclose(self) -> None: + await super().aclose() + raise OSError("shutdown failed") + + +def test_a_shutdown_that_raises_still_stops_the_thread_and_closes_the_loop() -> None: + worker_thread = WorkerThread(FailingClose()) + worker_thread.run_sync("1") + loop, thread = worker_thread._loop, worker_thread._thread + assert loop is not None and thread is not None + with pytest.raises(OSError, match="shutdown failed"): + worker_thread.close() + assert not thread.is_alive() + assert loop.is_closed() From a3ba0c22a2d246c5bd44ea49545ee5856dfb0863 Mon Sep 17 00:00:00 2001 From: Josh Taillon Date: Fri, 9 Oct 2026 22:51:04 -0400 Subject: [PATCH 17/25] fix(py): close the session's loop from its own thread once it stops The loop's thread now closes the loop when run_forever() returns, so it is closed on every path, including a close() called from that thread, which returns without waiting. --- pkg-py/src/commons/_execution/_thread.py | 15 ++++++++++++--- pkg-py/tests/test_execution_thread.py | 1 + 2 files changed, 13 insertions(+), 3 deletions(-) diff --git a/pkg-py/src/commons/_execution/_thread.py b/pkg-py/src/commons/_execution/_thread.py index 82724f09..e4b8c9b1 100644 --- a/pkg-py/src/commons/_execution/_thread.py +++ b/pkg-py/src/commons/_execution/_thread.py @@ -109,8 +109,6 @@ def close(self) -> None: finally: if closing.done(): thread.join(CLOSE_TIMEOUT) - if not thread.is_alive(): - loop.close() def _submit( self, code: str, handles: HandleStore | None @@ -137,8 +135,19 @@ def _ensure_loop(self) -> asyncio.AbstractEventLoop: if self._loop is None: loop = asyncio.new_event_loop() thread = threading.Thread( - target=loop.run_forever, name="commons-python-session", daemon=True + target=_run_loop, + args=(loop,), + name="commons-python-session", + daemon=True, ) thread.start() self._loop, self._thread = loop, thread return self._loop + + +def _run_loop(loop: asyncio.AbstractEventLoop) -> None: + """Run ``loop`` until it is stopped, then close it, however close() was called.""" + try: + loop.run_forever() + finally: + loop.close() diff --git a/pkg-py/tests/test_execution_thread.py b/pkg-py/tests/test_execution_thread.py index 5c16ca69..e9b512f3 100644 --- a/pkg-py/tests/test_execution_thread.py +++ b/pkg-py/tests/test_execution_thread.py @@ -101,6 +101,7 @@ def test_close_from_the_loops_own_thread_returns_and_the_thread_stops() -> None: loop.call_soon_threadsafe(worker_thread.close) thread.join(15) assert not thread.is_alive() + assert loop.is_closed() def test_a_cancelled_caller_cancels_the_call(runner: WorkerThread) -> None: From f20e7fb4e21ed4ac0153d05be4b672ea1457f019 Mon Sep 17 00:00:00 2001 From: Josh Taillon Date: Sat, 10 Oct 2026 17:05:23 -0500 Subject: [PATCH 18/25] docs(py): clarify run_python wiring comments and docstrings --- pkg-py/src/commons/_agent.py | 14 +++++++++----- pkg-py/src/commons/_execution/_thread.py | 2 ++ pkg-py/src/commons/_run_python.py | 6 +++--- 3 files changed, 14 insertions(+), 8 deletions(-) diff --git a/pkg-py/src/commons/_agent.py b/pkg-py/src/commons/_agent.py index 577f5533..468eb274 100644 --- a/pkg-py/src/commons/_agent.py +++ b/pkg-py/src/commons/_agent.py @@ -95,7 +95,8 @@ class Commons(Chat[Any, Any]): `"none"` (the default) or `"full"`. The sandbox works on Linux and macOS. For local development on another system, set the `COMMONS_ALLOW_UNSAFE_FALLBACK` environment variable to run the code with - limited checks instead. These checks are not a security boundary. + limited checks instead. Warning! These checks are not intended to be a + security boundary. Construction raises a TypeError if `client` is not a `chatlas.Chat`, if an entry of `data_sources` is not a `DataSource`, or if a layer is not @@ -147,7 +148,7 @@ def __init__( measure_sources=list(semantic_layer.source_text.values()), ) - # Share the provider, which carries the chosen model; shallow-copy + # Share the provider, which knows the chosen model; shallow-copy # the chat kwargs so later changes don't cross between the two. super().__init__( provider=client.provider, kwargs_chat=copy.copy(client.kwargs_chat) @@ -186,7 +187,9 @@ def __init__( ) tools = build_commons_tools(context) self._python = WorkerThread(worker) - # The session's process and thread go with the agent. + # Close the session's process and thread when the agent is + # collected or the interpreter exits. The callback names only the + # WorkerThread, so the finalizer never keeps the agent alive. weakref.finalize(self, self._python.close) self._run_python, self._run_python_sync = run_python_tools( self._python, @@ -239,8 +242,9 @@ def chat( was_pending = self._restore_reminder_pending inputs = self._prepare_turn_inputs(args) self._citation_request.reset() - # chatlas refuses a synchronous chat while an async tool is - # registered, so the sync run_python stands in for this call. + # A synchronous chat cannot await an async tool, so chatlas + # refuses to run while the async run_python tool is registered. + # This workaround swaps in the sync variant of run_python for this call. swap = any(tool.name == EXECUTION_TOOL for tool in self.get_tools()) if swap: self.register_tool(self._run_python_sync, force=True) diff --git a/pkg-py/src/commons/_execution/_thread.py b/pkg-py/src/commons/_execution/_thread.py index e4b8c9b1..6b85c8cb 100644 --- a/pkg-py/src/commons/_execution/_thread.py +++ b/pkg-py/src/commons/_execution/_thread.py @@ -113,6 +113,7 @@ def close(self) -> None: def _submit( self, code: str, handles: HandleStore | None ) -> concurrent.futures.Future[_Reply] | None: + """Queue ``code`` on the worker's loop, or return ``None`` once closed.""" with self._lock: if self._closed: return None @@ -127,6 +128,7 @@ def _submit( return future def _forget(self, future: concurrent.futures.Future[_Reply]) -> None: + """Drop a finished call from the set that ``close()`` cancels.""" with self._lock: self._calls.discard(future) diff --git a/pkg-py/src/commons/_run_python.py b/pkg-py/src/commons/_run_python.py index 7b6b1217..4295596b 100644 --- a/pkg-py/src/commons/_run_python.py +++ b/pkg-py/src/commons/_run_python.py @@ -5,9 +5,9 @@ idiom. The agent registers the async version of the tool, which waits for the session -without blocking the caller's event loop. chatlas refuses a synchronous -`chat()` while any async tool is registered, so `Commons.chat()` registers the -sync version for the length of the call. +without blocking the caller's event loop. Because chatlas refuses a synchronous +`chat()` while any async tool is registered, `Commons.chat()` dynamically swaps the +sync version for the length of the call if required. """ from __future__ import annotations From 8447669a1dfc7edd3a4915d1671ef81a19a253aa Mon Sep 17 00:00:00 2001 From: Josh Taillon Date: Sat, 10 Oct 2026 17:12:50 -0500 Subject: [PATCH 19/25] fix(py): tell the model when a cancelled call restarted the session A cancelled call or Ctrl-C shuts the worker down, so the next call starts a fresh session. WorkerThread now records that, and take_restart() reports it once. The cancel and Ctrl-C tests use a 60-second call timeout with a time bound, so they fail if cancelling does nothing; the Ctrl-C test cancels its signal timer; and every test that builds its own worker closes it in a finally. --- pkg-py/src/commons/_execution/_thread.py | 11 ++++ pkg-py/tests/test_execution_thread.py | 70 ++++++++++++++++-------- 2 files changed, 57 insertions(+), 24 deletions(-) diff --git a/pkg-py/src/commons/_execution/_thread.py b/pkg-py/src/commons/_execution/_thread.py index 6b85c8cb..cc3f9e96 100644 --- a/pkg-py/src/commons/_execution/_thread.py +++ b/pkg-py/src/commons/_execution/_thread.py @@ -41,6 +41,9 @@ def __init__(self, worker: Worker) -> None: self._lock = threading.Lock() self._closed = False self._calls: set[concurrent.futures.Future[_Reply]] = set() + # A cancelled call shuts the worker down, so the next call starts a + # fresh session; take_restart() reports that once. + self._restarted = False async def run(self, code: str, handles: HandleStore | None = None) -> _Reply: """Run ``code`` without blocking the caller's event loop. @@ -57,6 +60,7 @@ async def run(self, code: str, handles: HandleStore | None = None) -> _Reply: task = asyncio.current_task() if self._closed and future.cancelled() and not (task and task.cancelling()): return _CLOSED + self._restarted = True raise def run_sync(self, code: str, handles: HandleStore | None = None) -> _Reply: @@ -71,13 +75,20 @@ def run_sync(self, code: str, handles: HandleStore | None = None) -> _Reply: except concurrent.futures.CancelledError: if self._closed: return _CLOSED + self._restarted = True raise except BaseException: # Ctrl-C while waiting stops the call, as cancelling an async # caller does, rather than leaving it to run out its timeout. future.cancel() + self._restarted = True raise + def take_restart(self) -> bool: + """Whether a cancelled call restarted the session since the last ask.""" + restarted, self._restarted = self._restarted, False + return restarted + def close(self) -> None: """Close the worker, then stop the loop and its thread. diff --git a/pkg-py/tests/test_execution_thread.py b/pkg-py/tests/test_execution_thread.py index e9b512f3..c08886e4 100644 --- a/pkg-py/tests/test_execution_thread.py +++ b/pkg-py/tests/test_execution_thread.py @@ -76,13 +76,16 @@ async def from_the_loop() -> None: def test_ctrl_c_while_waiting_stops_the_call() -> None: worker_thread = WorkerThread(Worker(call_timeout=60)) + main = threading.main_thread().ident + assert main is not None + interrupt = threading.Timer(1, signal.pthread_kill, (main, signal.SIGINT)) try: worker_thread.run_sync("1") - main = threading.main_thread().ident - assert main is not None - threading.Timer(1, signal.pthread_kill, (main, signal.SIGINT)).start() + interrupt.start() with pytest.raises(KeyboardInterrupt): worker_thread.run_sync("import time; time.sleep(60)") + assert worker_thread.take_restart() + assert not worker_thread.take_restart() started = time.monotonic() # Without the cancel, this call would queue behind the sleep. reply = worker_thread.run_sync("1 + 1") @@ -90,6 +93,8 @@ def test_ctrl_c_while_waiting_stops_the_call() -> None: assert isinstance(reply, Result) assert reply.value == 2 finally: + # A call that returned early must not leave the signal for a later test. + interrupt.cancel() worker_thread.close() @@ -104,19 +109,28 @@ def test_close_from_the_loops_own_thread_returns_and_the_thread_stops() -> None: assert loop.is_closed() -def test_a_cancelled_caller_cancels_the_call(runner: WorkerThread) -> None: +def test_a_cancelled_caller_cancels_the_call() -> None: + worker_thread = WorkerThread(Worker(call_timeout=60)) + async def cancel_midway() -> None: - call = asyncio.ensure_future(runner.run("import time; time.sleep(30)")) + call = asyncio.ensure_future(worker_thread.run("import time; time.sleep(60)")) await asyncio.sleep(1) call.cancel() with pytest.raises(asyncio.CancelledError): await call - asyncio.run(cancel_midway()) - # The cancelled call's worker was shut down; the next one respawns. - reply = runner.run_sync("1 + 1") - assert isinstance(reply, Result) - assert reply.value == 2 + try: + asyncio.run(cancel_midway()) + assert worker_thread.take_restart() + started = time.monotonic() + # The cancelled call's worker was shut down, so the next one respawns + # rather than queue behind the sleep. + reply = worker_thread.run_sync("1 + 1") + assert time.monotonic() - started < 15 + assert isinstance(reply, Result) + assert reply.value == 2 + finally: + worker_thread.close() def test_close_ends_the_thread_and_later_calls_fail(runner: WorkerThread) -> None: @@ -130,20 +144,23 @@ def test_close_ends_the_thread_and_later_calls_fail(runner: WorkerThread) -> Non def test_close_cancels_a_running_call_rather_than_waiting_it_out() -> None: worker_thread = WorkerThread(Worker(call_timeout=60)) - worker_thread.run_sync("1") - replies: list[object] = [] - caller = threading.Thread( - target=lambda: replies.append( - worker_thread.run_sync("import time; time.sleep(60)") + try: + worker_thread.run_sync("1") + replies: list[object] = [] + caller = threading.Thread( + target=lambda: replies.append( + worker_thread.run_sync("import time; time.sleep(60)") + ) ) - ) - caller.start() - time.sleep(1) - started = time.monotonic() - worker_thread.close() - caller.join(10) - assert time.monotonic() - started < 15 - assert replies == [Failure(message="the Python session is closed.")] + caller.start() + time.sleep(1) + started = time.monotonic() + worker_thread.close() + caller.join(10) + assert time.monotonic() - started < 15 + assert replies == [Failure(message="the Python session is closed.")] + finally: + worker_thread.close() def test_an_async_call_running_when_the_thread_closes_is_told_it_closed() -> None: @@ -157,7 +174,12 @@ async def caller() -> object: await asyncio.to_thread(worker_thread.close) return await asyncio.wait_for(call, 10) - assert asyncio.run(caller()) == Failure(message="the Python session is closed.") + try: + assert asyncio.run(caller()) == Failure( + message="the Python session is closed." + ) + finally: + worker_thread.close() class FailingClose(Worker): From c52a102ebf499726724d12642c863ee4d808543f Mon Sep 17 00:00:00 2001 From: Josh Taillon Date: Sat, 10 Oct 2026 17:12:50 -0500 Subject: [PATCH 20/25] fix(py): keep run_python's output when highlighting fails, and fix pip advice - The highlighter also falls back to plain escaping on UnicodeDecodeError, which Python 3.12+ tokenize raises for a carriage return before non-ASCII text. That error used to replace the call's whole result. - The network="full" rule tells the model to run pip in the session. The macOS sandbox aborts every child process and guardrails refuse them, so the subprocess route it gave before could not work. - A model that cancelled a call is told, on the next result, that the session restarted and its variables were reset. - A probe that fails to run is no longer cached, so one slow start does not leave every later agent saying matplotlib is missing. --- pkg-py/src/commons/_run_python.py | 50 +++++++++++++++++++++---------- pkg-py/tests/test_run_python.py | 25 +++++++++++++++- 2 files changed, 58 insertions(+), 17 deletions(-) diff --git a/pkg-py/src/commons/_run_python.py b/pkg-py/src/commons/_run_python.py index 4295596b..f9a1811d 100644 --- a/pkg-py/src/commons/_run_python.py +++ b/pkg-py/src/commons/_run_python.py @@ -13,7 +13,6 @@ from __future__ import annotations import base64 -import functools import html import io import keyword @@ -55,6 +54,11 @@ NO_OUTPUT = "(The code ran but produced no output.)" +RESTART_NOTE = ( + "(A previous call was cancelled, so the Python session was restarted " + "before this call. Session variables were reset.)" +) + # How long to wait for the session's interpreter to say what it can import. PROBE_TIMEOUT = 10.0 @@ -141,10 +145,11 @@ def run_python_description( rules.append("The session has no network access.") elif can_install: rules.append( - "The session has network access. To use a package that is not " - "installed, run `sys.executable -m pip install --target` with the " - "temporary directory through subprocess, then add that directory " - "to sys.path." + "The session has network access. The sandbox stops subprocesses, so " + "to use a package that is not installed, run pip in the session: " + "call `pip._internal.cli.main.main(['install', '--no-cache-dir', " + "'--target', path, name])` with a path in the temporary directory, " + "then add that path to sys.path." ) else: rules.append( @@ -154,17 +159,24 @@ def run_python_description( return " ".join(parts) + "\n\nRules:" + "".join(f"\n- {rule}" for rule in rules) -@functools.cache +# Definite answers from the probe, by module. A probe that failed to run is +# not kept, so a slow first start does not settle the answer for the process. +_IMPORTABLE: dict[str, bool] = {} + + def session_can_import(module: str) -> bool: """Whether the session's interpreter can import the top-level ``module``. The session starts Python in isolated mode (``-I``), which ignores the user's site-packages directory, ``PYTHONPATH``, and the current directory. - The check therefore starts the same interpreter the same way, with the - session's environment variables and an empty temporary directory. The - answer is cached, because the installed packages do not change while the - process runs. + The check therefore starts the same interpreter with the same ``-I`` + isolation, the session's environment variables, and an empty temporary + directory, though outside the sandbox. A definite answer is cached, + because the installed packages do not change while the process runs; a + probe that fails to run counts as no, and is asked again next time. """ + if module in _IMPORTABLE: + return _IMPORTABLE[module] probe = "import importlib.util, sys; sys.exit(importlib.util.find_spec(sys.argv[1]) is None)" try: with tempfile.TemporaryDirectory(prefix="commons-probe-") as scratch: @@ -178,7 +190,8 @@ def session_can_import(module: str) -> bool: ) except (OSError, subprocess.SubprocessError): return False - return completed.returncode == 0 + _IMPORTABLE[module] = completed.returncode == 0 + return _IMPORTABLE[module] def _listed(names: Sequence[str]) -> str: @@ -193,7 +206,7 @@ def run_python_tools( """The async tool the agent registers, and the sync tool `chat()` uses instead.""" def finish(code: str, reply: Result | Error | Failure) -> ContentToolResult: - result = run_python_result(code, reply) + result = run_python_result(code, reply, restarted=runner.take_restart()) if context.citation_request is None: return result return context.citation_request.add_request(result) @@ -230,17 +243,20 @@ def build(func: Any) -> Tool: return build(run_python), build(run_python_sync) -def run_python_result(code: str, reply: Result | Error | Failure) -> ContentToolResult: +def run_python_result( + code: str, reply: Result | Error | Failure, *, restarted: bool = False +) -> ContentToolResult: """The tool result for one call: one view for the model, one for the user. The model gets the call's output, with each plot as an image at the point where it was drawn. The user sees the code and its output, followed by the - plots. + plots. ``restarted`` says a cancelled call restarted the session before + this one, and the model is told so before the output. """ runs = _runs(reply) plots = [run for run in runs if isinstance(run, Plot)] return tool_result( - _model_value(runs), + _model_value([RESTART_NOTE, *runs] if restarted else runs), ProvenanceTag.B, title=CODE_ANALYSIS.settled, html=_display_html(code, runs), @@ -364,7 +380,9 @@ def offset(position: tuple[int, int]) -> int: spans: list[tuple[int, int, str]] = [] try: tokens = list(tokenize.generate_tokens(io.StringIO(source).readline)) - except (tokenize.TokenError, SyntaxError): + except (tokenize.TokenError, SyntaxError, UnicodeDecodeError): + # Python 3.12+ raises UnicodeDecodeError for a lone carriage return + # before non-ASCII text. return html.escape(source) for index, token in enumerate(tokens): kind = _token_class( diff --git a/pkg-py/tests/test_run_python.py b/pkg-py/tests/test_run_python.py index b134517b..f5ae0786 100644 --- a/pkg-py/tests/test_run_python.py +++ b/pkg-py/tests/test_run_python.py @@ -23,6 +23,7 @@ from commons._provenance import TAG_EXTRA_KEY, Tag from commons._run_python import ( NO_OUTPUT, + RESTART_NOTE, highlight_python, run_python_description, run_python_result, @@ -112,7 +113,7 @@ def test_measure_sources_are_described_only_when_there_are_measures() -> None: ("network", "can_install", "rule"), [ ("none", True, "- The session has no network access."), - ("full", True, "`sys.executable -m pip install --target`"), + ("full", True, "`pip._internal.cli.main.main(["), ("full", False, "only packages that are already installed can be imported"), ], ) @@ -148,6 +149,16 @@ def test_a_module_the_isolated_session_cannot_see_is_not_importable( assert not session_can_import("only_on_pythonpath") +def test_a_probe_that_fails_to_run_is_asked_again( + monkeypatch: pytest.MonkeyPatch, +) -> None: + # Too short for the interpreter to start, so the probe times out. + monkeypatch.setattr("commons._run_python.PROBE_TIMEOUT", 1e-6) + assert not session_can_import("email") + monkeypatch.undo() + assert session_can_import("email") + + # ---- the result ------------------------------------------------------------- @@ -169,6 +180,12 @@ def test_a_result_shows_what_was_printed_and_the_value() -> None: assert display(result)["open"] is False +def test_a_restarted_session_is_noted_for_the_model_only() -> None: + result = run_python_result("x", Result(id="c1", value=2), restarted=True) + assert result.value == f"{RESTART_NOTE}\n2" + assert "restarted" not in str(display(result)["html"]) + + def test_a_call_that_shows_nothing_says_so() -> None: assert run_python_result("x = 1", Result(id="c1")).value == NO_OUTPUT @@ -252,6 +269,11 @@ def test_text_that_does_not_tokenize_is_escaped_plainly() -> None: assert highlight_python("''' None: + # Python 3.12+ tokenize raises UnicodeDecodeError on this input. + assert highlight_python("#> 50%\r\u00e9t\u00e9") == "#> 50%\r\u00e9t\u00e9" + + def test_a_carriage_return_or_form_feed_does_not_shift_the_highlighting() -> None: out = highlight_python('#> 50%\r100%\n# a\x0cb\ng("z")') assert 'g("z")' in out @@ -343,6 +365,7 @@ def test_the_network_argument_reaches_the_description_and_annotations() -> None: assert "The session has network access" in rules assert tool.annotations is not None assert tool.annotations.get("openWorldHint") is True + assert agent._python._worker._network == "full" def test_a_network_other_than_none_or_full_is_refused() -> None: From 79356766e7d15ba7e11b0653f9b0c0e01da6d96f Mon Sep 17 00:00:00 2001 From: Josh Taillon Date: Sat, 10 Oct 2026 17:18:43 -0500 Subject: [PATCH 21/25] fix(py): report a restart only when a cancel shut the session down The restart flag moves from WorkerThread to Worker, where it is set in the one place a cancel shuts the session down: a call cancelled while it ran. A call cancelled while it waited for the lock leaves the session and its variables alone, and no longer makes the next result claim a reset. --- pkg-py/src/commons/_execution/_driver.py | 13 ++++++++++ pkg-py/src/commons/_execution/_thread.py | 9 +------ pkg-py/tests/test_execution_thread.py | 31 +++++++++++++++++++++--- 3 files changed, 42 insertions(+), 11 deletions(-) diff --git a/pkg-py/src/commons/_execution/_driver.py b/pkg-py/src/commons/_execution/_driver.py index 67e3d253..8a3df010 100644 --- a/pkg-py/src/commons/_execution/_driver.py +++ b/pkg-py/src/commons/_execution/_driver.py @@ -98,6 +98,9 @@ def __init__( self._last_used = 0.0 self._lock = asyncio.Lock() self._ids = itertools.count(1) + # Set when a cancelled call shut the session down; take_restart() + # reports it once, so the model learns its variables were reset. + self._restarted = False async def __aenter__(self) -> Self: return self @@ -136,6 +139,15 @@ async def run( self._pending -= 1 self._schedule_reap() + def take_restart(self) -> bool: + """Whether a cancelled call restarted the session since the last ask. + + Only a call cancelled while it ran counts; one cancelled while it + waited for the lock left the session as it was. + """ + restarted, self._restarted = self._restarted, False + return restarted + async def aclose(self) -> None: """Close the worker, cancelling the idle reap. @@ -203,6 +215,7 @@ async def _call( # still running the code, or about to answer into a channel the # next call would misread as its own reply. Shut it down before # releasing the lock. + self._restarted = True await self._shutdown() raise except Exception as exc: # noqa: BLE001 - any start or write failure fails the call diff --git a/pkg-py/src/commons/_execution/_thread.py b/pkg-py/src/commons/_execution/_thread.py index cc3f9e96..64b54911 100644 --- a/pkg-py/src/commons/_execution/_thread.py +++ b/pkg-py/src/commons/_execution/_thread.py @@ -41,9 +41,6 @@ def __init__(self, worker: Worker) -> None: self._lock = threading.Lock() self._closed = False self._calls: set[concurrent.futures.Future[_Reply]] = set() - # A cancelled call shuts the worker down, so the next call starts a - # fresh session; take_restart() reports that once. - self._restarted = False async def run(self, code: str, handles: HandleStore | None = None) -> _Reply: """Run ``code`` without blocking the caller's event loop. @@ -60,7 +57,6 @@ async def run(self, code: str, handles: HandleStore | None = None) -> _Reply: task = asyncio.current_task() if self._closed and future.cancelled() and not (task and task.cancelling()): return _CLOSED - self._restarted = True raise def run_sync(self, code: str, handles: HandleStore | None = None) -> _Reply: @@ -75,19 +71,16 @@ def run_sync(self, code: str, handles: HandleStore | None = None) -> _Reply: except concurrent.futures.CancelledError: if self._closed: return _CLOSED - self._restarted = True raise except BaseException: # Ctrl-C while waiting stops the call, as cancelling an async # caller does, rather than leaving it to run out its timeout. future.cancel() - self._restarted = True raise def take_restart(self) -> bool: """Whether a cancelled call restarted the session since the last ask.""" - restarted, self._restarted = self._restarted, False - return restarted + return self._worker.take_restart() def close(self) -> None: """Close the worker, then stop the loop and its thread. diff --git a/pkg-py/tests/test_execution_thread.py b/pkg-py/tests/test_execution_thread.py index c08886e4..58d9782d 100644 --- a/pkg-py/tests/test_execution_thread.py +++ b/pkg-py/tests/test_execution_thread.py @@ -84,14 +84,14 @@ def test_ctrl_c_while_waiting_stops_the_call() -> None: interrupt.start() with pytest.raises(KeyboardInterrupt): worker_thread.run_sync("import time; time.sleep(60)") - assert worker_thread.take_restart() - assert not worker_thread.take_restart() started = time.monotonic() # Without the cancel, this call would queue behind the sleep. reply = worker_thread.run_sync("1 + 1") assert time.monotonic() - started < 15 assert isinstance(reply, Result) assert reply.value == 2 + assert worker_thread.take_restart() + assert not worker_thread.take_restart() finally: # A call that returned early must not leave the signal for a later test. interrupt.cancel() @@ -121,7 +121,6 @@ async def cancel_midway() -> None: try: asyncio.run(cancel_midway()) - assert worker_thread.take_restart() started = time.monotonic() # The cancelled call's worker was shut down, so the next one respawns # rather than queue behind the sleep. @@ -129,6 +128,7 @@ async def cancel_midway() -> None: assert time.monotonic() - started < 15 assert isinstance(reply, Result) assert reply.value == 2 + assert worker_thread.take_restart() finally: worker_thread.close() @@ -197,3 +197,28 @@ def test_a_shutdown_that_raises_still_stops_the_thread_and_closes_the_loop() -> worker_thread.close() assert not thread.is_alive() assert loop.is_closed() + + +def test_a_call_cancelled_while_it_waits_leaves_the_session_alone() -> None: + worker_thread = WorkerThread(Worker(call_timeout=60)) + + async def cancel_the_queued_one() -> object: + await worker_thread.run("x = 41") + running = asyncio.ensure_future( + worker_thread.run("import time; time.sleep(2)\nx + 1") + ) + await asyncio.sleep(0.5) + queued = asyncio.ensure_future(worker_thread.run("1")) + await asyncio.sleep(0.5) + queued.cancel() + with pytest.raises(asyncio.CancelledError): + await queued + return await running + + try: + reply = asyncio.run(cancel_the_queued_one()) + assert isinstance(reply, Result) + assert reply.value == 42 + assert not worker_thread.take_restart() + finally: + worker_thread.close() From 49f25eebe4de06ce4d468147546a7430d75814fb Mon Sep 17 00:00:00 2001 From: Josh Taillon Date: Sat, 10 Oct 2026 17:19:16 -0500 Subject: [PATCH 22/25] fix(py): give the model a pip invocation that runs in a fresh session The install rule now names the import and where the target path comes from, since the session starts with neither pip nor sys bound. Verified in the macOS sandbox: the steps as written install and import tabulate. --- pkg-py/src/commons/_run_python.py | 7 ++++--- pkg-py/tests/test_run_python.py | 2 +- 2 files changed, 5 insertions(+), 4 deletions(-) diff --git a/pkg-py/src/commons/_run_python.py b/pkg-py/src/commons/_run_python.py index f9a1811d..d014904d 100644 --- a/pkg-py/src/commons/_run_python.py +++ b/pkg-py/src/commons/_run_python.py @@ -147,9 +147,10 @@ def run_python_description( rules.append( "The session has network access. The sandbox stops subprocesses, so " "to use a package that is not installed, run pip in the session: " - "call `pip._internal.cli.main.main(['install', '--no-cache-dir', " - "'--target', path, name])` with a path in the temporary directory, " - "then add that path to sys.path." + "`from pip._internal.cli.main import main`, then " + "`main(['install', '--no-cache-dir', '--target', path, name])` " + "with `path` a directory under `tempfile.gettempdir()`, then " + "`sys.path.insert(0, path)`." ) else: rules.append( diff --git a/pkg-py/tests/test_run_python.py b/pkg-py/tests/test_run_python.py index f5ae0786..b5f2118f 100644 --- a/pkg-py/tests/test_run_python.py +++ b/pkg-py/tests/test_run_python.py @@ -113,7 +113,7 @@ def test_measure_sources_are_described_only_when_there_are_measures() -> None: ("network", "can_install", "rule"), [ ("none", True, "- The session has no network access."), - ("full", True, "`pip._internal.cli.main.main(["), + ("full", True, "`from pip._internal.cli.main import main`"), ("full", False, "only packages that are already installed can be imported"), ], ) From 0c2fb0e11387317b0112eceb13d0810a02aef154 Mon Sep 17 00:00:00 2001 From: Josh Taillon Date: Sat, 10 Oct 2026 21:28:11 -0500 Subject: [PATCH 23/25] fix(py): spell out the in-session pip install as one complete snippet The rule now gives every import and the target directory, and the snippet, taken from the description verbatim, installs and imports tabulate in the macOS sandbox. --- pkg-py/src/commons/_run_python.py | 7 +++---- pkg-py/tests/test_run_python.py | 2 +- 2 files changed, 4 insertions(+), 5 deletions(-) diff --git a/pkg-py/src/commons/_run_python.py b/pkg-py/src/commons/_run_python.py index d014904d..90ea1b8c 100644 --- a/pkg-py/src/commons/_run_python.py +++ b/pkg-py/src/commons/_run_python.py @@ -147,10 +147,9 @@ def run_python_description( rules.append( "The session has network access. The sandbox stops subprocesses, so " "to use a package that is not installed, run pip in the session: " - "`from pip._internal.cli.main import main`, then " - "`main(['install', '--no-cache-dir', '--target', path, name])` " - "with `path` a directory under `tempfile.gettempdir()`, then " - "`sys.path.insert(0, path)`." + "`import sys, tempfile; from pip._internal.cli.main import main; " + "path = tempfile.mkdtemp(); main(['install', '--no-cache-dir', " + "'--target', path, name]); sys.path.insert(0, path)`." ) else: rules.append( diff --git a/pkg-py/tests/test_run_python.py b/pkg-py/tests/test_run_python.py index b5f2118f..c2155422 100644 --- a/pkg-py/tests/test_run_python.py +++ b/pkg-py/tests/test_run_python.py @@ -113,7 +113,7 @@ def test_measure_sources_are_described_only_when_there_are_measures() -> None: ("network", "can_install", "rule"), [ ("none", True, "- The session has no network access."), - ("full", True, "`from pip._internal.cli.main import main`"), + ("full", True, "from pip._internal.cli.main import main"), ("full", False, "only packages that are already installed can be imported"), ], ) From 24f6c21ebd0fd7bb7ee230991055b42dfc318081 Mon Sep 17 00:00:00 2001 From: Josh Taillon Date: Sat, 10 Oct 2026 21:29:11 -0500 Subject: [PATCH 24/25] fix(py): name the package in the pip snippet with a quoted placeholder --- pkg-py/src/commons/_run_python.py | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/pkg-py/src/commons/_run_python.py b/pkg-py/src/commons/_run_python.py index 90ea1b8c..716a4bc2 100644 --- a/pkg-py/src/commons/_run_python.py +++ b/pkg-py/src/commons/_run_python.py @@ -149,7 +149,8 @@ def run_python_description( "to use a package that is not installed, run pip in the session: " "`import sys, tempfile; from pip._internal.cli.main import main; " "path = tempfile.mkdtemp(); main(['install', '--no-cache-dir', " - "'--target', path, name]); sys.path.insert(0, path)`." + "'--target', path, 'PACKAGE']); sys.path.insert(0, path)`, with " + "PACKAGE replaced by the package's name." ) else: rules.append( From 688aff47f9407bad94ad1d5e35ce0e2d707dca61 Mon Sep 17 00:00:00 2001 From: Josh Taillon Date: Sat, 10 Oct 2026 21:55:01 -0500 Subject: [PATCH 25/25] test(py): check the carriage-return case by its text on every Python Python 3.11 tokenizes a carriage return before non-ASCII text and highlights it, while 3.12+ raises and falls back to plain escaping, so the test asserts what both share: the displayed text is the source. --- pkg-py/tests/test_run_python.py | 8 +++++--- 1 file changed, 5 insertions(+), 3 deletions(-) diff --git a/pkg-py/tests/test_run_python.py b/pkg-py/tests/test_run_python.py index c2155422..4a6cbe03 100644 --- a/pkg-py/tests/test_run_python.py +++ b/pkg-py/tests/test_run_python.py @@ -269,9 +269,11 @@ def test_text_that_does_not_tokenize_is_escaped_plainly() -> None: assert highlight_python("''' None: - # Python 3.12+ tokenize raises UnicodeDecodeError on this input. - assert highlight_python("#> 50%\r\u00e9t\u00e9") == "#> 50%\r\u00e9t\u00e9" +def test_a_carriage_return_before_non_ascii_text_keeps_its_text() -> None: + # Python 3.12+ tokenize raises UnicodeDecodeError on this input, and + # earlier versions highlight it, so only the text is the same everywhere. + source = "#> 50%\r\u00e9t\u00e9" + assert html.unescape(re.sub(r"<[^>]+>", "", highlight_python(source))) == source def test_a_carriage_return_or_form_feed_does_not_shift_the_highlighting() -> None: