Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
25 commits
Select commit Hold shift + click to select a range
4d37ca4
refactor: name the code-run display classes for either language
jat255 Oct 9, 2026
32dd9ad
feat(py): run a worker on an event loop of its own
jat255 Oct 9, 2026
911859a
feat(py): add run_python, which runs model code in the agent's sandbo…
jat255 Oct 9, 2026
97f93dd
fix(py): cancel a worker thread's calls when it closes, and refuse ca…
jat255 Oct 9, 2026
fb3c06c
fix(py): describe plotting and installs by what the -I session can im…
jat255 Oct 9, 2026
cdb2658
fix(py): ask the session's own -I interpreter what it can import
jat255 Oct 9, 2026
ddac943
test(py): cover an async call cancelled by its worker thread's close
jat255 Oct 9, 2026
5febd02
fix(py): probe imports with the session's environment and scratch dir…
jat255 Oct 9, 2026
23c92d0
fix(py): give the import probe the worker's single-thread environment
jat255 Oct 9, 2026
8cc2424
feat(py): build run_python's result from the reply's ordered output
jat255 Oct 10, 2026
954f96a
fix(py): tell the model that plt.show() places a plot and savefig hid…
jat255 Oct 10, 2026
04bc117
fix(py): tell the model that fig.show() places a plot too
jat255 Oct 10, 2026
fbe74d6
fix(py): stop a call on Ctrl-C, and close safely from the session's t…
jat255 Oct 10, 2026
ac20001
fix(py): tighten run_python's wiring, highlighting, and probe
jat255 Oct 10, 2026
9a0cd56
docs: note the code-run class rename in NEWS and name both tools in t…
jat255 Oct 10, 2026
e238778
test(py): cover a worker shutdown that raises
jat255 Oct 10, 2026
a3ba0c2
fix(py): close the session's loop from its own thread once it stops
jat255 Oct 10, 2026
f20e7fb
docs(py): clarify run_python wiring comments and docstrings
jat255 Oct 10, 2026
8447669
fix(py): tell the model when a cancelled call restarted the session
jat255 Oct 10, 2026
c52a102
fix(py): keep run_python's output when highlighting fails, and fix pi…
jat255 Oct 10, 2026
7935676
fix(py): report a restart only when a cancel shut the session down
jat255 Oct 10, 2026
49f25ee
fix(py): give the model a pip invocation that runs in a fresh session
jat255 Oct 10, 2026
0c2fb0e
fix(py): spell out the in-session pip install as one complete snippet
jat255 Oct 11, 2026
24f6c21
fix(py): name the package in the pip snippet with a quoted placeholder
jat255 Oct 11, 2026
688aff4
test(py): check the carriage-return case by its text on every Python
jat255 Oct 11, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
80 changes: 65 additions & 15 deletions pkg-py/src/commons/_agent.py
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@

import copy
import warnings
import weakref
from collections.abc import AsyncGenerator, Mapping, Sequence
from typing import Any, Literal, NoReturn

Expand All @@ -29,9 +30,12 @@
from ._context_layer import ContextLayer, augment_context_layer
from ._data_source import DataSource
from ._definitions import Registry, build_registry
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,
Expand All @@ -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"]
Expand Down Expand Up @@ -84,11 +89,22 @@ 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 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. 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
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__(
Expand All @@ -99,6 +115,7 @@ def __init__(
context_layer: ContextLayer | None = None,
*,
instructions: str | None = None,
network: Literal["none", "full"] = "none",
) -> None:
if not isinstance(client, Chat):
raise TypeError(
Expand All @@ -124,8 +141,14 @@ 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
# 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)
Expand All @@ -152,18 +175,33 @@ 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)
# 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,
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,
Expand Down Expand Up @@ -204,7 +242,19 @@ 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)
# 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)
try:
response = super().chat(
*inputs, echo=echo, stream=stream, kwargs=kwargs
)
finally:
if swap:
self.register_tool(self._run_python, force=True)
self._consume_restore_reminder(was_pending)
return response

Expand Down
2 changes: 2 additions & 0 deletions pkg-py/src/commons/_display.py
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,7 @@
_URL = re.compile(r"https?://", re.IGNORECASE)

__all__ = [
"CODE_ANALYSIS",
"CONTEXT_SEARCH",
"DATA_RETRIEVAL",
"DISPLAY_EXTRA_KEY",
Expand Down Expand Up @@ -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(
Expand Down
13 changes: 13 additions & 0 deletions pkg-py/src/commons/_execution/_driver.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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.

Expand Down Expand Up @@ -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
Expand Down
159 changes: 159 additions & 0 deletions pkg-py/src/commons/_execution/_thread.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,159 @@
"""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 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 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`` in a background thread, which starts on the first call."""

def __init__(self, worker: Worker) -> None:
self._worker = worker
self._loop: asyncio.AbstractEventLoop | None = None
self._thread: threading.Thread | None = None
# 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) -> _Reply:
"""Run ``code`` without blocking the caller's event loop.

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:
return _CLOSED
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
try:
return future.result()
except concurrent.futures.CancelledError:
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 take_restart(self) -> bool:
"""Whether a cancelled call restarted the session since the last ask."""
return self._worker.take_restart()

def close(self) -> None:
"""Close the worker, then stop the loop and its thread.

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:
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)
closing.add_done_callback(lambda _: loop.call_soon_threadsafe(loop.stop))
# Waiting here on the loop's own thread would block the shutdown.
if threading.current_thread() is thread:
return
try:
closing.result(CLOSE_TIMEOUT)
except TimeoutError:
return
finally:
if closing.done():
thread.join(CLOSE_TIMEOUT)

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
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:
"""Drop a finished call from the set that ``close()`` cancels."""
with self._lock:
self._calls.discard(future)

def _ensure_loop(self) -> asyncio.AbstractEventLoop:
"""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(
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()
Loading
Loading