This is an automated email from the ASF dual-hosted git repository.
shahar1 pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/airflow.git
The following commit(s) were added to refs/heads/main by this push:
new 0566d6cd5ef Add OpenSandbox backend for sandbox tools (#71676)
0566d6cd5ef is described below
commit 0566d6cd5efe0e317de1c113142cd4e5b6e2a8c8
Author: Yossi Eliaz <[email protected]>
AuthorDate: Mon Sep 28 19:03:20 2026 +0300
Add OpenSandbox backend for sandbox tools (#71676)
Rebase the reviewed OpenSandbox backend onto current Airflow main,
preserving current dependency state and dropping the expired temporary SDK
quarantine override.
---
docs/spelling_wordlist.txt | 3 +
providers/common/ai/README.rst | 47 +-
providers/common/ai/docs/index.rst | 47 +-
providers/common/ai/docs/installation.rst | 7 +-
providers/common/ai/docs/sandbox/backends.rst | 89 ++-
providers/common/ai/docs/sandbox/index.rst | 40 +-
providers/common/ai/provider.yaml | 3 +
providers/common/ai/pyproject.toml | 2 +
.../providers/common/ai/get_provider_info.py | 5 +
.../providers/common/ai/sandbox/__init__.py | 2 +
.../providers/common/ai/sandbox/opensandbox.py | 535 +++++++++++++++
.../src/airflow/providers/common/ai/sandbox/sbx.py | 9 +-
.../common/ai/tests/system/common/ai/README.md | 37 +-
.../ai/example_sandbox_toolset_opensandbox.py | 142 ++++
.../unit/common/ai/sandbox/test_opensandbox.py | 725 +++++++++++++++++++++
uv.lock | 35 +-
16 files changed, 1646 insertions(+), 82 deletions(-)
diff --git a/docs/spelling_wordlist.txt b/docs/spelling_wordlist.txt
index c31da3d2400..04f847f4906 100644
--- a/docs/spelling_wordlist.txt
+++ b/docs/spelling_wordlist.txt
@@ -632,6 +632,7 @@ ExaConnection
Exasol
exasol
exc
+execd
executables
exitcode
expanduser
@@ -958,6 +959,7 @@ jwt
Kafka
kafka
Kamil
+Kata
KEDA
keepalive
keepalives
@@ -1213,6 +1215,7 @@ openfaas
OpenID
openjdk
openlineage
+opensandbox
OpenSearch
opensearch
OpenTelemetry
diff --git a/providers/common/ai/README.rst b/providers/common/ai/README.rst
index 48d868189a1..b4b86e5f091 100644
--- a/providers/common/ai/README.rst
+++ b/providers/common/ai/README.rst
@@ -82,29 +82,30 @@ Dependent package
Optional dependencies
----------------------
-==============
=======================================================================================================================================
-Extra Dependencies
-==============
=======================================================================================================================================
-``anthropic`` ``pydantic-ai-slim[anthropic]>=2.33.0``, ``anthropic>=1.0.0``
-``bedrock`` ``pydantic-ai-slim[bedrock]>=2.33.0``
-``google`` ``pydantic-ai-slim[google]>=2.33.0``
-``openai`` ``pydantic-ai-slim[openai]>=2.33.0``, ``openai>=2.47.0``
-``typesafe`` ``typesafe-sdk>=0.6.0``
-``mcp`` ``pydantic-ai-slim[mcp]>=2.33.0``
-``modal`` ``modal>=1.5.0``
-``code-mode`` ``pydantic-ai-harness[codemode]>=0.3.0``
-``shields`` ``pydantic-ai-shields>=0.3.4``
-``skills`` ``apache-airflow-providers-git>=0.4.0``,
``pydantic-ai-skills>=1.2.0``
-``avro`` ``fastavro>=1.10.0; python_version < "3.14"``,
``fastavro>=1.12.1; python_version >= "3.14"``
-``parquet`` ``pyarrow>=18.0.0; python_version < '3.14'``,
``pyarrow>=22.0.0; python_version >= '3.14'``
-``sql`` ``apache-airflow-providers-common-sql>=1.33.0``,
``sqlglot>=30.0.0``
-``common.sql`` ``apache-airflow-providers-common-sql>=1.33.0``
-``langchain`` ``langchain>=1.0.0``
-``llamaindex`` ``dataclasses-json>=0.6.7``, ``llama-index-core>=0.14.5``,
``llama-index-embeddings-openai>=0.6.0``, ``llama-index-llms-openai>=0.6.8``
-``pdf`` ``pypdf>=4.0.0``
-``docx`` ``python-docx>=1.0.0``
-``git`` ``apache-airflow-providers-git``
-==============
=======================================================================================================================================
+===============
=======================================================================================================================================
+Extra Dependencies
+===============
=======================================================================================================================================
+``anthropic`` ``pydantic-ai-slim[anthropic]>=2.33.0``, ``anthropic>=1.0.0``
+``bedrock`` ``pydantic-ai-slim[bedrock]>=2.33.0``
+``google`` ``pydantic-ai-slim[google]>=2.33.0``
+``openai`` ``pydantic-ai-slim[openai]>=2.33.0``, ``openai>=2.47.0``
+``typesafe`` ``typesafe-sdk>=0.6.0``
+``mcp`` ``pydantic-ai-slim[mcp]>=2.33.0``
+``modal`` ``modal>=1.5.0``
+``opensandbox`` ``opensandbox>=1.1.0``
+``code-mode`` ``pydantic-ai-harness[codemode]>=0.3.0``
+``shields`` ``pydantic-ai-shields>=0.3.4``
+``skills`` ``apache-airflow-providers-git>=0.4.0``,
``pydantic-ai-skills>=1.2.0``
+``avro`` ``fastavro>=1.10.0; python_version < "3.14"``,
``fastavro>=1.12.1; python_version >= "3.14"``
+``parquet`` ``pyarrow>=18.0.0; python_version < '3.14'``,
``pyarrow>=22.0.0; python_version >= '3.14'``
+``sql`` ``apache-airflow-providers-common-sql>=1.33.0``,
``sqlglot>=30.0.0``
+``common.sql`` ``apache-airflow-providers-common-sql>=1.33.0``
+``langchain`` ``langchain>=1.0.0``
+``llamaindex`` ``dataclasses-json>=0.6.7``, ``llama-index-core>=0.14.5``,
``llama-index-embeddings-openai>=0.6.0``, ``llama-index-llms-openai>=0.6.8``
+``pdf`` ``pypdf>=4.0.0``
+``docx`` ``python-docx>=1.0.0``
+``git`` ``apache-airflow-providers-git``
+===============
=======================================================================================================================================
The changelog for the provider package can be found in the
`changelog
<https://airflow.apache.org/docs/apache-airflow-providers-common-ai/0.10.0/changelog.html>`_.
diff --git a/providers/common/ai/docs/index.rst
b/providers/common/ai/docs/index.rst
index d986d937455..dd7eccfe2aa 100644
--- a/providers/common/ai/docs/index.rst
+++ b/providers/common/ai/docs/index.rst
@@ -212,29 +212,30 @@ Install them when installing from PyPI. For example:
pip install apache-airflow-providers-common-ai[anthropic]
-==============
=======================================================================================================================================
-Extra Dependencies
-==============
=======================================================================================================================================
-``anthropic`` ``pydantic-ai-slim[anthropic]>=2.33.0``, ``anthropic>=1.0.0``
-``bedrock`` ``pydantic-ai-slim[bedrock]>=2.33.0``
-``google`` ``pydantic-ai-slim[google]>=2.33.0``
-``openai`` ``pydantic-ai-slim[openai]>=2.33.0``, ``openai>=2.47.0``
-``typesafe`` ``typesafe-sdk>=0.6.0``
-``mcp`` ``pydantic-ai-slim[mcp]>=2.33.0``
-``modal`` ``modal>=1.5.0``
-``code-mode`` ``pydantic-ai-harness[codemode]>=0.3.0``
-``shields`` ``pydantic-ai-shields>=0.3.4``
-``skills`` ``apache-airflow-providers-git>=0.4.0``,
``pydantic-ai-skills>=1.2.0``
-``avro`` ``fastavro>=1.10.0; python_version < "3.14"``,
``fastavro>=1.12.1; python_version >= "3.14"``
-``parquet`` ``pyarrow>=18.0.0; python_version < '3.14'``,
``pyarrow>=22.0.0; python_version >= '3.14'``
-``sql`` ``apache-airflow-providers-common-sql>=1.33.0``,
``sqlglot>=30.0.0``
-``common.sql`` ``apache-airflow-providers-common-sql>=1.33.0``
-``langchain`` ``langchain>=1.0.0``
-``llamaindex`` ``dataclasses-json>=0.6.7``, ``llama-index-core>=0.14.5``,
``llama-index-embeddings-openai>=0.6.0``, ``llama-index-llms-openai>=0.6.8``
-``pdf`` ``pypdf>=4.0.0``
-``docx`` ``python-docx>=1.0.0``
-``git`` ``apache-airflow-providers-git``
-==============
=======================================================================================================================================
+===============
=======================================================================================================================================
+Extra Dependencies
+===============
=======================================================================================================================================
+``anthropic`` ``pydantic-ai-slim[anthropic]>=2.33.0``, ``anthropic>=1.0.0``
+``bedrock`` ``pydantic-ai-slim[bedrock]>=2.33.0``
+``google`` ``pydantic-ai-slim[google]>=2.33.0``
+``openai`` ``pydantic-ai-slim[openai]>=2.33.0``, ``openai>=2.47.0``
+``typesafe`` ``typesafe-sdk>=0.6.0``
+``mcp`` ``pydantic-ai-slim[mcp]>=2.33.0``
+``modal`` ``modal>=1.5.0``
+``opensandbox`` ``opensandbox>=1.1.0``
+``code-mode`` ``pydantic-ai-harness[codemode]>=0.3.0``
+``shields`` ``pydantic-ai-shields>=0.3.4``
+``skills`` ``apache-airflow-providers-git>=0.4.0``,
``pydantic-ai-skills>=1.2.0``
+``avro`` ``fastavro>=1.10.0; python_version < "3.14"``,
``fastavro>=1.12.1; python_version >= "3.14"``
+``parquet`` ``pyarrow>=18.0.0; python_version < '3.14'``,
``pyarrow>=22.0.0; python_version >= '3.14'``
+``sql`` ``apache-airflow-providers-common-sql>=1.33.0``,
``sqlglot>=30.0.0``
+``common.sql`` ``apache-airflow-providers-common-sql>=1.33.0``
+``langchain`` ``langchain>=1.0.0``
+``llamaindex`` ``dataclasses-json>=0.6.7``, ``llama-index-core>=0.14.5``,
``llama-index-embeddings-openai>=0.6.0``, ``llama-index-llms-openai>=0.6.8``
+``pdf`` ``pypdf>=4.0.0``
+``docx`` ``python-docx>=1.0.0``
+``git`` ``apache-airflow-providers-git``
+===============
=======================================================================================================================================
Downloading official packages
-----------------------------
diff --git a/providers/common/ai/docs/installation.rst
b/providers/common/ai/docs/installation.rst
index 5380e724fc2..43e12ca82c8 100644
--- a/providers/common/ai/docs/installation.rst
+++ b/providers/common/ai/docs/installation.rst
@@ -46,9 +46,10 @@ The provider's extras split into a few groups:
the built-in adapter talks to; pydantic-ai supports more model providers
than these, each under its own extra name, so check the
`pydantic-ai install docs <https://ai.pydantic.dev/install/#slim-install>`__
for the full list.
-* **Agent tooling** (``mcp``, ``skills``, ``code-mode``, ``shields``,
``modal``): MCP servers,
- Agent Skills, code-mode tool execution, shield capabilities (input/output
guards, tool
- guards, cost tracking), and the hosted Modal backend for :doc:`sandboxed
execution <sandbox/index>`.
+* **Agent tooling** (``mcp``, ``skills``, ``code-mode``, ``shields``,
``modal``,
+ ``opensandbox``): MCP servers, Agent Skills, code-mode tool execution, shield
+ capabilities (input/output guards, tool guards, cost tracking), and the
hosted Modal and
+ self-hosted OpenSandbox backends for :doc:`sandboxed execution
<sandbox/index>`.
* **Document loading** (``pdf``, ``docx``, ``avro``, ``parquet``): file
formats for
document pipelines.
* **Retrieval / SQL** (``sql``, ``common.sql``, ``langchain``,
``llamaindex``): RAG and
diff --git a/providers/common/ai/docs/sandbox/backends.rst
b/providers/common/ai/docs/sandbox/backends.rst
index dcfec0989dd..ea72b5d57df 100644
--- a/providers/common/ai/docs/sandbox/backends.rst
+++ b/providers/common/ai/docs/sandbox/backends.rst
@@ -25,7 +25,8 @@ Modal (hosted)
:class:`~airflow.providers.common.ai.sandbox.modal.ModalSandboxBackend` runs
each
sandbox in Modal, provisioned over the API. Of the backends that ship with the
-provider, **this is the one to use in production**, and the only one that runs
on
+provider, **this is the managed one to use in production**, and with
+:ref:`OpenSandbox <sandbox-backend-opensandbox>` one of the two that run on
Kubernetes: nothing has to be installed on the worker, model-written code never
executes on the worker host, and Modal reclaims a sandbox at its own lifetime
whether or not the worker survives. It needs the ``modal`` extra and ambient
@@ -159,6 +160,69 @@ them.
that is a symlink replaces the link with a regular file and leaves the original
target untouched, where a shell redirect would follow the link.
+.. _sandbox-backend-opensandbox:
+
+OpenSandbox (self-hosted remote)
+--------------------------------
+
+:class:`~airflow.providers.common.ai.sandbox.opensandbox.OpenSandboxBackend`
+runs sandboxes through an `OpenSandbox <https://open-sandbox.ai/>`__ server.
+The server may use Docker or Kubernetes; Airflow workers only use its HTTP API
+and do not need access to the container runtime.
+
+Install the SDK extra:
+
+.. code-block:: bash
+
+ pip install "apache-airflow-providers-common-ai[opensandbox]"
+
+Use a generic Airflow connection, resolved lazily on first use:
+
+.. code-block:: python
+
+ from airflow.providers.common.ai.sandbox import OpenSandboxBackend
+ from airflow.providers.common.ai.toolsets import SandboxToolset
+
+
SandboxToolset(OpenSandboxBackend(opensandbox_conn_id="opensandbox_default"))
+
+The connection ``host`` is required; ``port`` is optional, ``schema`` defaults
to
+``http``, and ``password`` carries the API key when required. Extras may set
+``request_timeout`` (default 30 seconds) and ``use_server_proxy`` (default
+``true``). Set ``opensandbox_conn_id=None`` to let the SDK read
+``OPEN_SANDBOX_DOMAIN`` and ``OPEN_SANDBOX_API_KEY``.
+
+``SandboxSpec.env`` is sent at creation. A default spec sends a deny-all
+network policy; ``allow_egress_to`` becomes explicit allow rules. With
+``block_network=False``, the backend omits network policy entirely so a
+deployment without the egress sidecar can still run an intentionally open
+sandbox. Deny/allowlist policy requires the sidecar, so the backend reads the
+enforced policy back after creation and destroys the sandbox if it does not
+match the requested spec. ``allow_egress_to_cidrs`` is refused: OpenSandbox
only
+enforces CIDR targets in ``dns+nft`` mode, and the Python SDK does not expose
+that enforcement mode on policy read-back, so this backend cannot prove the
+address-layer restriction is active.
+
+Every sandbox carries ``created-by: airflow`` metadata and an
+``airflow-sandbox-*`` name for attribution and cleanup. The server enforces a
+sandbox lifetime (default 3600 seconds). If the SDK event stream stalls, the
+worker abandons the call after the command budget plus a grace period, destroys
+the sandbox, and reports ``sandbox_terminated`` so the toolset provisions a
fresh
+one. Output is bounded per stream after the SDK yields it; a single
newline-free
+line is the SDK-level exception, because the SDK assembles that line before the
+backend sees it.
+
+Constructor parameters:
+
+- ``image``: image used by the server. Default ``"python:3.12-slim"``.
+- ``cpu`` and ``memory``: resource limits. Defaults ``"1"`` and ``"2Gi"``.
+- ``sandbox_timeout``: server-side lifetime in seconds. Default ``3600``.
+- ``ready_timeout``: provisioning/reconnect timeout. Default ``120``.
+- ``use_server_proxy``: override the connection extra for file and command
calls.
+
+The runtime remains a deployment choice. The default Docker runtime shares the
+host kernel; choose a stronger runtime such as Kata when your threat model
needs
+a VM boundary.
+
sbx (Docker Sandboxes, local)
-----------------------------
@@ -202,23 +266,28 @@ Constructor parameters:
``sbx policy init deny-all``, or ``"allow-all"`` to state that egress is open
and pass ``SandboxSpec(block_network=False)`` to match.
-What differs between the two
-----------------------------
+What differs between the backends
+---------------------------------
Swapping the backend is one constructor argument, and tool names, spec and
prompt
do not change. Four behaviours do, so read them before assuming the same Dag
-behaves identically in both places:
+behaves identically everywhere:
- **CPU.** ``sbx`` gives a sandbox every host CPU; Modal defaults to a request
of
- 0.125 of one, so set ``cpu``.
+ 0.125 of one, so set ``cpu``; OpenSandbox takes ``cpu`` as a limit the
server enforces.
- **Egress allowlists.** ``sbx`` enforces ``allow_egress_to`` at the host
policy
layer; Modal matches TLS handshake names, which is weaker and has to be opted
- into. ``allow_egress_to_cidrs`` is enforced at the address layer on Modal and
- refused on ``sbx``, which has no per-sandbox address rule.
-- **Command timeouts.** A timeout destroys an ``sbx`` sandbox and its files; a
- Modal sandbox survives with its files intact.
+ into; OpenSandbox enforces it in an egress sidecar, and the backend reads the
+ enforced policy back rather than trusting the create request.
``allow_egress_to_cidrs``
+ is enforced at the address layer on Modal, refused on ``sbx``, and refused by
+ OpenSandbox because its SDK cannot prove that the sidecar is running in the
+ ``dns+nft`` mode required for CIDR enforcement.
+- **Command timeouts.** A timeout destroys an ``sbx`` sandbox and its files;
+ Modal and a server-enforced OpenSandbox timeout preserve the sandbox and
files.
+ OpenSandbox destroys it only if the command event stream itself stalls past
the
+ client-side grace period.
- **Symlinks.** ``write_file`` through a symlink follows the link on ``sbx``
and
- replaces it on Modal.
+ replaces it on Modal and OpenSandbox.
Bringing your own backend
-------------------------
diff --git a/providers/common/ai/docs/sandbox/index.rst
b/providers/common/ai/docs/sandbox/index.rst
index d7877372ef4..0e54281f69e 100644
--- a/providers/common/ai/docs/sandbox/index.rst
+++ b/providers/common/ai/docs/sandbox/index.rst
@@ -50,10 +50,11 @@ model a disposable workspace for that code instead. It
exposes four tools:
The sandbox is provisioned by a
:class:`~airflow.providers.common.ai.sandbox.SandboxBackend` on the model's
first
-tool call and torn down when the agent run ends. Two backends ship: a hosted
one on
-`Modal <https://modal.com/docs/guide/sandbox>`__ for production and Kubernetes,
-and a local microVM one on `Docker Sandboxes
<https://docs.docker.com/ai/sandboxes/>`__
-for development. The four tool names and shapes match pydantic-ai's own sandbox
+tool call and torn down when the agent run ends. Three backends ship: a hosted
one on
+`Modal <https://modal.com/docs/guide/sandbox>`__ and a self-hosted one on
+`OpenSandbox <https://open-sandbox.ai/>`__ for production and Kubernetes, and
a local
+microVM one on `Docker Sandboxes <https://docs.docker.com/ai/sandboxes/>`__ for
+development. The four tool names and shapes match pydantic-ai's own sandbox
capabilities, so a model that has seen one already knows this one.
**Adding this toolset gives the agent shell and file operations in a separate
@@ -275,11 +276,16 @@ sandbox, and every other toolset stays on the worker.
- Modal's container runtime, which Modal documents as gVisor.
- Modal's infrastructure, off the worker.
- Ended by Modal at ``sandbox_timeout``, or ``idle_timeout`` if set.
+ * - ``OpenSandboxBackend``
+ - The container runtime your OpenSandbox deployment configures, on Docker
or
+ Kubernetes.
+ - Your OpenSandbox server's Docker host or Kubernetes cluster, off the
worker.
+ - Ended by the OpenSandbox server at ``sandbox_timeout``.
When a run ends normally, the task calls the backend's ``destroy``. ``sbx``
runs its
-removal command and waits up to two minutes for it; Modal sends a termination
-request and returns without waiting for the sandbox to stop. Either can return
with
-the sandbox still present, and neither case fails the task. A SIGKILL, an
+removal command and waits up to two minutes for it; Modal and OpenSandbox each
send a
+termination request and return without waiting for the sandbox to stop. Any of
them
+can return with the sandbox still present, and none of those cases fails the
task. A SIGKILL, an
out-of-memory kill or a lost node skips that teardown entirely, and then only
the
last column applies.
@@ -357,16 +363,19 @@ within reach whether or not code mode is on. See
:ref:`code-mode` and
**What it cannot do**
-- Only one of its two backends runs on Kubernetes. ``SbxSandboxBackend`` drives
+- One of its three backends does not run on Kubernetes. ``SbxSandboxBackend``
drives
Docker Sandboxes on the worker host, and its own documentation says to use it
for local development: it wants the ``sbx`` binary on the host, an
authenticated Docker account, a one-time ``sbx policy init``, and on Linux
KVM
or nested virtualization, which an unprivileged container cannot provide.
- Production and Kubernetes use
- :class:`~airflow.providers.common.ai.sandbox.modal.ModalSandboxBackend`, a
- hosted backend behind the ``modal`` extra that installs nothing on the worker
- and reclaims a sandbox at its own lifetime if the worker dies. Both implement
- :class:`~airflow.providers.common.ai.sandbox.SandboxBackend`, and a third
+ Production and Kubernetes use a remote backend instead, either
+ :class:`~airflow.providers.common.ai.sandbox.modal.ModalSandboxBackend`
behind
+ the ``modal`` extra for a managed service, or
+ :class:`~airflow.providers.common.ai.sandbox.opensandbox.OpenSandboxBackend`
+ behind ``opensandbox`` for a self-hosted one. Neither installs anything on
the
+ worker, and each reclaims a sandbox at its own server-side lifetime if the
+ worker dies. All three implement
+ :class:`~airflow.providers.common.ai.sandbox.SandboxBackend`, and another
vendor can too.
- It does not contain the agent. Only what these tools do runs in the sandbox;
the agent loop, the model calls, and every other toolset on the same agent
stay
@@ -384,7 +393,10 @@ within reach whether or not code mode is on. See
:ref:`code-mode` and
under the default ``host_network_policy="unknown"``. On Modal the same
default
maps onto the sandbox's own ``block_network`` and is enforced exactly; a
hostname allowlist there is matched on the TLS handshake name and has to be
- opted into, for the reasons set out on :doc:`backends`.
+ opted into, for the reasons set out on :doc:`backends`. OpenSandbox enforces
+ hostname allowlists with its egress sidecar, but refuses
+ ``allow_egress_to_cidrs`` because the SDK cannot prove the sidecar is in the
+ ``dns+nft`` mode required for CIDR enforcement.
- Reclamation depends on the backend. A failed teardown is logged as a warning
rather than raised, deliberately, so that a teardown blip cannot fail a
finished run. On ``sbx`` nothing else picks up the slack: there is no
diff --git a/providers/common/ai/provider.yaml
b/providers/common/ai/provider.yaml
index fb4a01d613d..cc025d58d53 100644
--- a/providers/common/ai/provider.yaml
+++ b/providers/common/ai/provider.yaml
@@ -73,6 +73,9 @@ integrations:
- integration-name: Modal
external-doc-url: https://modal.com/docs/guide/sandbox
tags: [service]
+ - integration-name: OpenSandbox
+ external-doc-url: https://open-sandbox.ai/
+ tags: [software]
hooks:
- integration-name: Pydantic AI
diff --git a/providers/common/ai/pyproject.toml
b/providers/common/ai/pyproject.toml
index 1cdec81dece..1fa52b03c89 100644
--- a/providers/common/ai/pyproject.toml
+++ b/providers/common/ai/pyproject.toml
@@ -103,6 +103,7 @@ dependencies = [
# cannot run microVMs -- Kubernetes, most notably. Needs no host-level
installation,
# only Modal credentials.
"modal" = ["modal>=1.5.0"]
+"opensandbox" = ["opensandbox>=1.1.0"]
# Code mode: collapse tool calls into a single `run_code` tool that the model
# drives by writing Python, executed in the Monty sandbox (pydantic-monty).
# Enables AgentOperator(code_mode=True). Monty is pre-1.0; pinned here as an
@@ -173,6 +174,7 @@ dev = [
# needs them so the adapter tests, which build real SDK models, run
instead of skipping.
"openai>=2.47.0",
"anthropic>=1.0.0",
+ "opensandbox>=1.1.0",
]
# To build docs:
diff --git
a/providers/common/ai/src/airflow/providers/common/ai/get_provider_info.py
b/providers/common/ai/src/airflow/providers/common/ai/get_provider_info.py
index b70a7f8382b..b5ca27bb775 100644
--- a/providers/common/ai/src/airflow/providers/common/ai/get_provider_info.py
+++ b/providers/common/ai/src/airflow/providers/common/ai/get_provider_info.py
@@ -76,6 +76,11 @@ def get_provider_info():
"external-doc-url": "https://modal.com/docs/guide/sandbox",
"tags": ["service"],
},
+ {
+ "integration-name": "OpenSandbox",
+ "external-doc-url": "https://open-sandbox.ai/",
+ "tags": ["software"],
+ },
],
"hooks": [
{
diff --git
a/providers/common/ai/src/airflow/providers/common/ai/sandbox/__init__.py
b/providers/common/ai/src/airflow/providers/common/ai/sandbox/__init__.py
index 3d4ae8dccd4..4f77ad0619d 100644
--- a/providers/common/ai/src/airflow/providers/common/ai/sandbox/__init__.py
+++ b/providers/common/ai/src/airflow/providers/common/ai/sandbox/__init__.py
@@ -26,9 +26,11 @@ from airflow.providers.common.ai.sandbox.base import (
SandboxSpec,
SandboxTerminalError,
)
+from airflow.providers.common.ai.sandbox.opensandbox import OpenSandboxBackend
from airflow.providers.common.ai.sandbox.sbx import SbxSandboxBackend
__all__ = [
+ "OpenSandboxBackend",
"ModalSandboxBackend",
"SandboxBackend",
"SandboxError",
diff --git
a/providers/common/ai/src/airflow/providers/common/ai/sandbox/opensandbox.py
b/providers/common/ai/src/airflow/providers/common/ai/sandbox/opensandbox.py
new file mode 100644
index 00000000000..e2523627cef
--- /dev/null
+++ b/providers/common/ai/src/airflow/providers/common/ai/sandbox/opensandbox.py
@@ -0,0 +1,535 @@
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements. See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership. The ASF licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License. You may obtain a copy of the License at
+#
+# http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing,
+# software distributed under the License is distributed on an
+# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+# KIND, either express or implied. See the License for the
+# specific language governing permissions and limitations
+# under the License.
+"""OpenSandbox backend for
:class:`~airflow.providers.common.ai.toolsets.sandbox.SandboxToolset`."""
+
+from __future__ import annotations
+
+import logging
+import posixpath
+import threading
+import time
+from contextlib import contextmanager, suppress
+from datetime import timedelta
+from typing import TYPE_CHECKING, Any
+
+from airflow.providers.common.ai.sandbox.base import (
+ SandboxBackend,
+ SandboxError,
+ SandboxExecResult,
+ SandboxFileTooLargeError,
+ SandboxTerminalError,
+ _new_sandbox_name,
+ _validate_positive_finite,
+)
+from airflow.providers.common.compat.sdk import BaseHook
+
+if TYPE_CHECKING:
+ from collections.abc import Callable, Iterator
+
+ from opensandbox import SandboxSync
+ from opensandbox.config import ConnectionConfigSync
+ from opensandbox.models.sandboxes import NetworkPolicy
+
+ from airflow.providers.common.ai.sandbox.base import SandboxSpec
+
+log = logging.getLogger(__name__)
+
+# Wall-clock allowance past the per-command budget before a streaming call is
+# treated as hung. execd enforces the budget itself; this covers the case where
+# its events stop arriving at all, which the SDK's read-timeout-free SSE client
+# would otherwise wait on forever.
+_EXEC_GRACE = 30.0
+# A listing is metadata, never bulk content. Past this many entries the model
is
+# told to narrow the path instead of the worker holding the whole array.
+_LIST_DIRECTORY_MAX_ENTRIES = 10_000
+# Stamped on every sandbox so an operator on a shared server can find Airflow's
+# through ``SandboxFilter(metadata=...)``; the sbx backend uses the sandbox
name
+# for the same purpose.
+_CREATED_BY_METADATA = {"created-by": "airflow"}
+
+
+def _get_status_code(error: Exception) -> int | None:
+ status_code = getattr(error, "status_code", None)
+ return status_code if isinstance(status_code, int) else None
+
+
+@contextmanager
+def _translate_opensandbox_errors(
+ operation: str, *, recoverable_statuses: frozenset[int] = frozenset()
+) -> Iterator[None]:
+ try:
+ yield
+ except SandboxError:
+ raise
+ except Exception as e:
+ try:
+ from opensandbox.exceptions import SandboxApiException
+ except ImportError:
+ raise SandboxTerminalError(
+ "The OpenSandbox SDK is not installed. Install "
+ '"apache-airflow-providers-common-ai[opensandbox]".'
+ ) from e
+ status_code = _get_status_code(e) if isinstance(e,
SandboxApiException) else None
+ status = f" (HTTP {status_code})" if status_code is not None else ""
+ message = f"OpenSandbox could not {operation}{status}."
+ if status_code in recoverable_statuses:
+ raise SandboxError(message) from e
+ raise SandboxTerminalError(message) from e
+
+
+class _BoundedTail:
+ def __init__(self, max_bytes: int) -> None:
+ self._max_bytes = max_bytes
+ self._data = bytearray()
+ self.truncated = False
+
+ def add_text(self, text: str) -> None:
+ self._data.extend(text.encode("utf-8"))
+ if len(self._data) > self._max_bytes:
+ del self._data[: len(self._data) - self._max_bytes]
+ self.truncated = True
+
+ def add_message(self, message: Any) -> None:
+ # execd streams one message per output line with the delimiter
stripped,
+ # so the newline has to be put back or every line runs together. A
blank
+ # line already arrives as "\n", hence the guard.
+ text = message.text
+ self.add_text(text if text.endswith("\n") else text + "\n")
+
+ def get_text(self) -> str:
+ return bytes(self._data).decode("utf-8", errors="ignore")
+
+
+class _CallStillRunning(Exception):
+ """The call handed to :func:`_call_with_deadline` outlived its deadline."""
+
+
+def _call_with_deadline(fn: Callable[[], Any], deadline: float) -> Any:
+ """
+ Run ``fn`` on a daemon thread and wait at most ``deadline`` seconds for it.
+
+ A daemon thread so a call that never returns cannot hold up interpreter
+ exit; the caller is expected to make it return by destroying the sandbox.
+ """
+ outcome: dict[str, Any] = {}
+
+ def target() -> None:
+ try:
+ outcome["result"] = fn()
+ except BaseException as e:
+ outcome["error"] = e
+
+ thread = threading.Thread(target=target, name="opensandbox-command",
daemon=True)
+ thread.start()
+ thread.join(deadline)
+ if thread.is_alive():
+ raise _CallStillRunning
+ if "error" in outcome:
+ raise outcome["error"]
+ return outcome["result"]
+
+
+def _parse_bool(value: Any, name: str) -> bool:
+ if isinstance(value, bool):
+ return value
+ if isinstance(value, str):
+ normalized = value.strip().lower()
+ if normalized in {"true", "1", "yes"}:
+ return True
+ if normalized in {"false", "0", "no"}:
+ return False
+ raise SandboxTerminalError(f"The OpenSandbox connection extra {name} must
be a boolean.")
+
+
+class OpenSandboxBackend(SandboxBackend):
+ """
+ Run sandbox tools through an OpenSandbox server.
+
+ OpenSandbox supports Docker and Kubernetes runtimes behind the same API.
+ Airflow workers need only network access to that API; the OpenSandbox
+ deployment owns container provisioning and isolation.
+
+ A generic Airflow connection supplies the server configuration. ``host``
+ and ``port`` identify the lifecycle API, ``schema`` selects ``http`` or
+ ``https``, and ``password`` carries the optional API key. Connection extras
+ may set ``request_timeout`` and ``use_server_proxy``.
+
+ For a deny-by-default spec, the create API accepts a network policy whether
+ or not the server runs the egress sidecar that enforces it, so after
creating
+ the sandbox the backend reads the enforced policy back and destroys the
+ sandbox if it differs from what
+ :class:`~airflow.providers.common.ai.sandbox.SandboxSpec` asked for. The
+ fail-closed contract is this backend's to keep, not the server's.
+
+ Command deadlines are enforced by execd. If its event stream stalls, the
+ call is abandoned ``_EXEC_GRACE`` seconds past the budget, the sandbox is
+ destroyed to end it, and the result reports ``timed_out`` with
+ ``sandbox_terminated`` so the toolset provisions a fresh one. Output is
+ streamed and each stream is kept to ``max_output_bytes`` on the worker,
with
+ one caveat: the SDK reassembles a whole output line before handing it over,
+ so a single line with no newline in it is resident in full first.
+
+ :param opensandbox_conn_id: Generic Airflow connection ID. ``None`` lets
the
+ SDK resolve ``OPEN_SANDBOX_DOMAIN`` and ``OPEN_SANDBOX_API_KEY``.
+ :param image: Container image used for each sandbox.
+ :param cpu: OpenSandbox CPU resource limit.
+ :param memory: OpenSandbox memory resource limit.
+ :param sandbox_timeout: Server-side sandbox lifetime in seconds.
+ :param ready_timeout: Seconds to wait for a newly created sandbox to
become healthy.
+ :param use_server_proxy: Route sandbox service calls through the lifecycle
+ server. ``None`` reads the connection extra and otherwise defaults to
``True``.
+ """
+
+ name = "opensandbox"
+
+ def __init__(
+ self,
+ opensandbox_conn_id: str | None = "opensandbox_default",
+ *,
+ image: str = "python:3.12-slim",
+ cpu: str = "1",
+ memory: str = "2Gi",
+ sandbox_timeout: float = 3600.0,
+ ready_timeout: float = 120.0,
+ use_server_proxy: bool | None = None,
+ ) -> None:
+ if not image:
+ raise ValueError("image must not be empty.")
+ if not cpu:
+ raise ValueError("cpu must not be empty.")
+ if not memory:
+ raise ValueError("memory must not be empty.")
+ _validate_positive_finite(sandbox_timeout, "sandbox_timeout")
+ _validate_positive_finite(ready_timeout, "ready_timeout")
+ self._opensandbox_conn_id = opensandbox_conn_id
+ self._image = image
+ self._resource = {"cpu": cpu, "memory": memory}
+ self._sandbox_timeout = sandbox_timeout
+ self._ready_timeout = ready_timeout
+ self._use_server_proxy = use_server_proxy
+ self._connection_config: ConnectionConfigSync | None = None
+ self._sandboxes: dict[str, SandboxSync] = {}
+
+ def _get_connection_config(self) -> ConnectionConfigSync:
+ if self._connection_config is not None:
+ return self._connection_config
+ with _translate_opensandbox_errors("initialize its client"):
+ from opensandbox.config import ConnectionConfigSync
+
+ if self._opensandbox_conn_id is None:
+ self._connection_config = ConnectionConfigSync(
+ use_server_proxy=True if self._use_server_proxy is None
else self._use_server_proxy
+ )
+ return self._connection_config
+
+ conn = BaseHook.get_connection(self._opensandbox_conn_id)
+ if not conn.host:
+ # The SDK would otherwise fall back to localhost:8080 and point
+ # the worker at itself.
+ raise SandboxTerminalError(
+ f"Connection {self._opensandbox_conn_id!r} has no host;
set it to the OpenSandbox "
+ "server address, or pass opensandbox_conn_id=None to use
OPEN_SANDBOX_DOMAIN."
+ )
+ extra = conn.extra_dejson
+ request_timeout = extra.get("request_timeout", 30)
+ try:
+ request_timeout = float(request_timeout)
+ _validate_positive_finite(request_timeout, "connection extra
request_timeout")
+ except (TypeError, ValueError) as e:
+ raise SandboxTerminalError(
+ "The OpenSandbox connection extra request_timeout must be
a positive finite number."
+ ) from e
+
+ use_server_proxy = self._use_server_proxy
+ if use_server_proxy is None:
+ value = extra.get("use_server_proxy", True)
+ use_server_proxy = _parse_bool(value, "use_server_proxy")
+
+ domain = conn.host
+ if conn.port:
+ domain = f"{domain}:{conn.port}"
+ self._connection_config = ConnectionConfigSync(
+ api_key=conn.password or None,
+ domain=domain,
+ protocol=conn.schema or "http",
+ request_timeout=timedelta(seconds=request_timeout),
+ use_server_proxy=use_server_proxy,
+ )
+ return self._connection_config
+
+ @staticmethod
+ def _get_network_policy(spec: SandboxSpec | None) -> NetworkPolicy | None:
+ if spec is None:
+ return None
+ if spec.allow_egress_to_cidrs:
+ raise SandboxTerminalError(
+ "SandboxSpec names allow_egress_to_cidrs, which this backend
cannot safely enforce: "
+ "OpenSandbox CIDR rules require the egress sidecar's dns+nft
mode, but the SDK "
+ "does not expose that enforcement mode on policy read-back.
Use allow_egress_to "
+ "with hostnames or a backend with verifiable address-layer
enforcement."
+ )
+ if not spec.block_network:
+ if spec.allow_egress_to:
+ raise SandboxTerminalError(
+ "SandboxSpec.allow_egress_to only narrows a
deny-by-default policy; "
+ "set block_network=True or remove the allowlist."
+ )
+ # No policy means intentionally open egress. Sending an explicit
allow policy
+ # unnecessarily requires the OpenSandbox egress sidecar.
+ return None
+
+ from opensandbox.models.sandboxes import NetworkPolicy, NetworkRule
+
+ rules = [NetworkRule(action="allow", target=target) for target in
spec.allow_egress_to or ()]
+ # default_action is declared under its wire alias. populate_by_name
means both
+ # spellings work at runtime, but only the alias is in the typed
signature.
+ return NetworkPolicy(defaultAction="deny", egress=rules or None)
+
+ @staticmethod
+ def _verify_network_policy(sandbox: SandboxSync, requested: NetworkPolicy)
-> None:
+ wanted = {rule.target for rule in requested.egress or ()}
+ try:
+ enforced = sandbox.get_egress_policy()
+ allowed = {rule.target for rule in enforced.egress or () if
rule.action == "allow"}
+ matches = enforced.default_action == "deny" and allowed == wanted
+ detail = f"enforced policy is default {enforced.default_action!r}
with allow {sorted(allowed)}"
+ except Exception as e:
+ matches = False
+ detail = f"the policy could not be read back ({type(e).__name__})"
+ if matches:
+ return
+ with suppress(Exception):
+ sandbox.destroy()
+ raise SandboxTerminalError(
+ f"OpenSandbox did not enforce the requested network policy
({detail}), so the sandbox was "
+ "destroyed. The server may be running without its egress sidecar."
+ )
+
+ def create(self, *, spec: SandboxSpec | None = None) -> str:
+ with _translate_opensandbox_errors("create a sandbox"):
+ from opensandbox import SandboxSync
+
+ network_policy = self._get_network_policy(spec)
+ sandbox = SandboxSync.create(
+ self._image,
+ timeout=timedelta(seconds=self._sandbox_timeout),
+ ready_timeout=timedelta(seconds=self._ready_timeout),
+ env=dict(spec.env) if spec is not None and spec.env else None,
+ metadata={**_CREATED_BY_METADATA, "name": _new_sandbox_name()},
+ resource=dict(self._resource),
+ network_policy=network_policy,
+ connection_config=self._get_connection_config(),
+ )
+ if network_policy is not None and network_policy.default_action ==
"deny":
+ self._verify_network_policy(sandbox, network_policy)
+ self._sandboxes[sandbox.id] = sandbox
+ return sandbox.id
+
+ def _get_sandbox(self, sandbox_id: str) -> SandboxSync:
+ if sandbox := self._sandboxes.get(sandbox_id):
+ return sandbox
+ with _translate_opensandbox_errors("connect to a sandbox"):
+ from opensandbox import SandboxSync
+
+ sandbox = SandboxSync.connect(
+ sandbox_id,
+ connection_config=self._get_connection_config(),
+ connect_timeout=timedelta(seconds=self._ready_timeout),
+ )
+ self._sandboxes[sandbox_id] = sandbox
+ return sandbox
+
+ def run_command(
+ self, sandbox: str, command: str, *, timeout: float, max_output_bytes:
int
+ ) -> SandboxExecResult:
+ _validate_positive_finite(timeout, "timeout")
+ _validate_positive_finite(max_output_bytes, "max_output_bytes")
+ # Reconnecting an uncached handle polls for readiness, which is not the
+ # command's time to spend against ``timeout``.
+ sandbox_client = self._get_sandbox(sandbox)
+ stdout = _BoundedTail(max_output_bytes)
+ stderr = _BoundedTail(max_output_bytes)
+ started = time.monotonic()
+ with _translate_opensandbox_errors("run a sandbox command"):
+ from opensandbox.models.execd import RunCommandOpts
+ from opensandbox.models.execd_sync import ExecutionHandlersSync
+
+ opts = RunCommandOpts(timeout=timedelta(seconds=timeout))
+ handlers = ExecutionHandlersSync(
+ on_stdout=stdout.add_message,
+ on_stderr=stderr.add_message,
+ skip_accumulation=True,
+ )
+ try:
+ execution = _call_with_deadline(
+ lambda: sandbox_client.commands.run(command, opts=opts,
handlers=handlers),
+ timeout + _EXEC_GRACE,
+ )
+ except _CallStillRunning:
+ return self._abandon_command(sandbox, stdout, stderr)
+ elapsed = time.monotonic() - started
+
+ exit_code = execution.exit_code
+ if exit_code is None:
+ # execd reports no exit status for a foreground run. The SDK parses
+ # one out of the free-text ``error.value`` and yields None when
that
+ # text is prose, so None means "unknown", never "sandbox unusable".
+ exit_code = 1 if execution.error is not None else -1
+ if execution.error is not None:
+ if not stderr.get_text():
+ stderr.add_text("\n".join(execution.error.traceback) or
execution.error.value)
+ elif execution.exit_code is None:
+ stderr.add_text("The command ended without reporting an exit
status.")
+ return SandboxExecResult(
+ exit_code=exit_code,
+ stdout=stdout.get_text(),
+ stderr=stderr.get_text(),
+ timed_out=exit_code != 0 and elapsed >= timeout,
+ stdout_truncated=stdout.truncated,
+ stderr_truncated=stderr.truncated,
+ )
+
+ def _abandon_command(self, sandbox: str, stdout: _BoundedTail, stderr:
_BoundedTail) -> SandboxExecResult:
+ # Destroying the sandbox is what ends the stalled stream. If that fails
+ # the server-side lifetime reclaims it; either way this sandbox is not
+ # one to reuse, and the toolset is told so.
+ try:
+ self.destroy(sandbox)
+ except SandboxError:
+ log.warning(
+ "Timed out running a command in OpenSandbox sandbox %s and
could not destroy it; "
+ "its server-side lifetime will reclaim it",
+ sandbox,
+ exc_info=True,
+ )
+ return SandboxExecResult(
+ exit_code=-1,
+ stdout=stdout.get_text(),
+ stderr=stderr.get_text(),
+ timed_out=True,
+ stdout_truncated=stdout.truncated,
+ stderr_truncated=stderr.truncated,
+ sandbox_terminated=True,
+ )
+
+ @staticmethod
+ def _confirm_sandbox_exists(sandbox: SandboxSync) -> None:
+ with _translate_opensandbox_errors("confirm that a sandbox still
exists"):
+ sandbox.get_info()
+
+ @staticmethod
+ def _get_file_size(sandbox: SandboxSync, path: str, *, at_least: int) ->
int:
+ """Return the file's size for the too-large message, never less than
what was already read."""
+ try:
+ info = sandbox.files.get_file_info([path])
+ entry = info.get(path) or next(iter(info.values()))
+ except Exception:
+ return at_least
+ return max(entry.size, at_least)
+
+ def read_file(self, sandbox: str, path: str, *, max_bytes: int) -> bytes:
+ _validate_positive_finite(max_bytes, "max_bytes")
+ sandbox_client = self._get_sandbox(sandbox)
+ chunks = None
+ data = bytearray()
+ try:
+ chunks = sandbox_client.files.read_bytes_stream(
+ path,
+ chunk_size=min(65536, max_bytes + 1),
+ range_header=f"bytes=0-{max_bytes}",
+ )
+ for chunk in chunks:
+ data.extend(chunk[: max_bytes + 1 - len(data)])
+ if len(data) > max_bytes:
+ size = self._get_file_size(sandbox_client, path,
at_least=len(data))
+ raise SandboxFileTooLargeError(path, size, max_bytes)
+ except SandboxFileTooLargeError:
+ raise
+ except Exception as e:
+ if _get_status_code(e) == 404:
+ self._confirm_sandbox_exists(sandbox_client)
+ raise SandboxError(f"{path!r} does not exist in the sandbox,
or is not readable.") from e
+ with _translate_opensandbox_errors("read a sandbox file",
recoverable_statuses=frozenset({400})):
+ raise
+ finally:
+ close = getattr(chunks, "close", None)
+ if close is not None:
+ with suppress(Exception):
+ close()
+ return bytes(data)
+
+ def write_file(self, sandbox: str, path: str, content: bytes) -> None:
+ sandbox_client = self._get_sandbox(sandbox)
+ try:
+ from opensandbox.models.filesystem import WriteEntry
+
+ parent = posixpath.dirname(path)
+ if parent and parent != "/":
+
sandbox_client.files.create_directories([WriteEntry(path=parent, mode=755)])
+ sandbox_client.files.write_file(path, content, mode=644)
+ except Exception as e:
+ if _get_status_code(e) == 404:
+ self._confirm_sandbox_exists(sandbox_client)
+ raise SandboxError(f"Could not write {path!r} in the
sandbox.") from e
+ with _translate_opensandbox_errors("write a sandbox file",
recoverable_statuses=frozenset({400})):
+ raise
+
+ def list_directory(self, sandbox: str, path: str) -> list[tuple[str,
bool]]:
+ sandbox_client = self._get_sandbox(sandbox)
+ try:
+ from opensandbox.models.filesystem import DirectoryListEntry
+
+ entries =
sandbox_client.files.list_directory(DirectoryListEntry(path=path, depth=1))
+ except Exception as e:
+ if _get_status_code(e) == 404:
+ self._confirm_sandbox_exists(sandbox_client)
+ raise SandboxError(f"{path!r} does not exist in the sandbox,
or is not readable.") from e
+ with _translate_opensandbox_errors(
+ "list a sandbox directory",
recoverable_statuses=frozenset({400})
+ ):
+ raise
+ # The SDK has parsed the whole listing by now; the cap bounds what goes
+ # any further and gives the model something to do about it.
+ if len(entries) > _LIST_DIRECTORY_MAX_ENTRIES:
+ raise SandboxError(
+ f"{path!r} has more than {_LIST_DIRECTORY_MAX_ENTRIES}
entries; list a subdirectory "
+ "instead, or use a shell command such as `ls | head`."
+ )
+ return [
+ (posixpath.basename(entry.path.rstrip("/")), entry.entry_type ==
"directory") for entry in entries
+ ]
+
+ def destroy(self, sandbox: str) -> None:
+ sandbox_client = self._sandboxes.pop(sandbox, None)
+ try:
+ if sandbox_client is None:
+ from opensandbox import SandboxSync
+
+ # No readiness poll: a paused or unhealthy sandbox is still one
+ # that must be destroyable, and waiting on it only delays that.
+ sandbox_client = SandboxSync.connect(
+ sandbox,
+ connection_config=self._get_connection_config(),
+ connect_timeout=timedelta(seconds=self._ready_timeout),
+ skip_health_check=True,
+ )
+ sandbox_client.destroy()
+ except Exception as e:
+ if _get_status_code(e) == 404:
+ return
+ with _translate_opensandbox_errors("destroy a sandbox"):
+ raise
diff --git a/providers/common/ai/src/airflow/providers/common/ai/sandbox/sbx.py
b/providers/common/ai/src/airflow/providers/common/ai/sandbox/sbx.py
index d9dbdaf7454..ea307885fe4 100644
--- a/providers/common/ai/src/airflow/providers/common/ai/sandbox/sbx.py
+++ b/providers/common/ai/src/airflow/providers/common/ai/sandbox/sbx.py
@@ -76,9 +76,12 @@ class SbxSandboxBackend(SandboxBackend):
driving it from an Airflow worker is off-label use. A production worker
would
need the ``sbx`` binary on the host, an authenticated Docker account
(``sbx login``), a one-time ``sbx policy init``, and on Linux, KVM or
nested
- virtualization -- which an unprivileged container cannot provide. No hosted
- backend ships with the provider yet; add one behind :class:`SandboxBackend`
- if you need Kubernetes.
+ virtualization -- which an unprivileged container cannot provide. If you
need
+ Kubernetes, use a remote backend --
+ :class:`~airflow.providers.common.ai.sandbox.ModalSandboxBackend` for a
+ managed service or
+ :class:`~airflow.providers.common.ai.sandbox.OpenSandboxBackend` for a
+ self-hosted one -- or add your own behind :class:`SandboxBackend`.
**Network policy is layered on a host-level setting, not independent of
one.** ``sbx`` governs egress through a host-level ``sbx policy``.
diff --git a/providers/common/ai/tests/system/common/ai/README.md
b/providers/common/ai/tests/system/common/ai/README.md
index 84fd62a61f1..9314e3ed5c1 100644
--- a/providers/common/ai/tests/system/common/ai/README.md
+++ b/providers/common/ai/tests/system/common/ai/README.md
@@ -17,16 +17,18 @@
under the License.
-->
-# SandboxToolset system test
+# SandboxToolset system tests
-The `example_sandbox_toolset_sbx.py` system test exercises the complete local
boundary:
+Each test exercises the complete toolset boundary with a deterministic
pydantic-ai model:
1. an Airflow task runs a deterministic pydantic-ai agent;
-2. the agent calls `SandboxToolset` twice;
-3. `SbxSandboxBackend` creates a Docker Sandbox microVM;
-4. the second Python call reads the file written by the first call, proving
per-run persistence; and
+2. the agent calls every sandbox tool;
+3. the selected backend creates a sandbox;
+4. later calls read the file written by the first call, proving per-run
persistence; and
5. the agent run tears the sandbox down.
+## Docker Sandboxes
+
Install the `sbx` CLI on every worker that can run this Dag and initialize its
network policy once:
```console
@@ -39,3 +41,28 @@ Airflow system-test environment whose task process has
access to the host `sbx`
```console
pytest --system
providers/common/ai/tests/system/common/ai/example_sandbox_toolset_sbx.py
```
+
+## Modal
+
+Install the Modal extra and authenticate, then:
+
+```console
+pip install "apache-airflow-providers-common-ai[modal]"
+modal token new
+pytest --system
providers/common/ai/tests/system/common/ai/example_sandbox_toolset_modal.py
+```
+
+## OpenSandbox
+
+Install the OpenSandbox extra, point the SDK at a running server, and run:
+
+```console
+pip install "apache-airflow-providers-common-ai[opensandbox]"
+export OPEN_SANDBOX_DOMAIN="opensandbox.example.com"
+export OPEN_SANDBOX_API_KEY="..."
+pytest --system
providers/common/ai/tests/system/common/ai/example_sandbox_toolset_opensandbox.py
+```
+
+The test requests the default deny-all egress policy, so the OpenSandbox server
+must have its egress sidecar configured. It also applies a 15-minute
server-side
+sandbox lifetime.
diff --git
a/providers/common/ai/tests/system/common/ai/example_sandbox_toolset_opensandbox.py
b/providers/common/ai/tests/system/common/ai/example_sandbox_toolset_opensandbox.py
new file mode 100644
index 00000000000..87d70c00f75
--- /dev/null
+++
b/providers/common/ai/tests/system/common/ai/example_sandbox_toolset_opensandbox.py
@@ -0,0 +1,142 @@
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements. See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership. The ASF licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License. You may obtain a copy of the License at
+#
+# http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing,
+# software distributed under the License is distributed on an
+# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+# KIND, either express or implied. See the License for the
+# specific language governing permissions and limitations
+# under the License.
+"""End-to-end system test for SandboxToolset with OpenSandbox."""
+
+from __future__ import annotations
+
+import os
+from datetime import datetime, timezone
+
+from airflow.providers.common.compat.sdk import dag as airflow_dag, task
+
+ENV_ID = os.environ.get("SYSTEM_TESTS_ENV_ID")
+DAG_ID = (
+ f"common_ai_sandbox_toolset_opensandbox_{ENV_ID}" if ENV_ID else
"common_ai_sandbox_toolset_opensandbox"
+)
+
+MARKER = "boundary-ok"
+STATE_PATH = "/tmp/airflow_sandbox_e2e"
+
+
+@airflow_dag(
+ dag_id=DAG_ID,
+ schedule="@once",
+ start_date=datetime(2024, 1, 1, tzinfo=timezone.utc),
+ catchup=False,
+ tags=["common.ai", "sandbox", "opensandbox", "system_test"],
+)
+def example_sandbox_toolset_opensandbox():
+ @task
+ def run_sandbox_agent() -> str:
+ from pydantic_ai import Agent
+ from pydantic_ai.messages import ModelMessage, ModelResponse,
TextPart, ToolCallPart
+ from pydantic_ai.models.function import AgentInfo, FunctionModel
+
+ from airflow.providers.common.ai.sandbox import OpenSandboxBackend
+ from airflow.providers.common.ai.toolsets import SandboxToolset
+
+ def model_function(messages: list[ModelMessage], _info: AgentInfo) ->
ModelResponse:
+ returns = [
+ part.content
+ for message in messages
+ for part in message.parts
+ if part.part_kind == "tool-return"
+ ]
+
+ if not returns:
+ return ModelResponse(
+ parts=[
+ ToolCallPart(
+ tool_name="write_file",
+ args={"path": STATE_PATH, "content": MARKER},
+ tool_call_id="write",
+ )
+ ]
+ )
+ if "Wrote" not in str(returns[0]):
+ raise RuntimeError(f"Unexpected write_file result:
{returns[0]!r}")
+
+ if len(returns) == 1:
+ return ModelResponse(
+ parts=[
+ ToolCallPart(
+ tool_name="run_command",
+ args={"command": f"cat {STATE_PATH} && python3 -c
'print(6 * 7)'"},
+ tool_call_id="shell",
+ )
+ ]
+ )
+ shell_out = str(returns[1])
+ if MARKER not in shell_out or "42" not in shell_out:
+ raise RuntimeError(f"Unexpected run_command result:
{shell_out!r}")
+
+ if len(returns) == 2:
+ return ModelResponse(
+ parts=[
+ ToolCallPart(tool_name="read_file", args={"path":
STATE_PATH}, tool_call_id="read")
+ ]
+ )
+ if MARKER not in str(returns[2]):
+ raise RuntimeError(f"Unexpected read_file result:
{returns[2]!r}")
+
+ if len(returns) == 3:
+ return ModelResponse(
+ parts=[
+ ToolCallPart(
+ tool_name="run_command",
+ args={"command": "echo to-stderr >&2; exit 3"},
+ tool_call_id="fail",
+ )
+ ]
+ )
+ failed = str(returns[3])
+ if "[exit code: 3]" not in failed or "to-stderr" not in failed:
+ raise RuntimeError(f"Unexpected failure result: {failed!r}")
+
+ if len(returns) == 4:
+ return ModelResponse(
+ parts=[ToolCallPart(tool_name="list_directory",
args={"path": "/tmp"}, tool_call_id="ls")]
+ )
+ if "airflow_sandbox_e2e" not in str(returns[4]):
+ raise RuntimeError(f"Unexpected listing: {returns[4]!r}")
+
+ return ModelResponse(parts=[TextPart(content="sandbox boundary e2e
passed")])
+
+ agent = Agent(
+ FunctionModel(model_function),
+ instructions="Use the sandbox tools as requested.",
+ toolsets=[
+ SandboxToolset(
+ OpenSandboxBackend(opensandbox_conn_id=None,
sandbox_timeout=900),
+ default_command_timeout=30.0,
+ max_command_timeout=30.0,
+ )
+ ],
+ )
+ result = agent.run_sync("Run the sandbox boundary system test.")
+ if result.output != "sandbox boundary e2e passed":
+ raise RuntimeError(f"Unexpected agent output: {result.output!r}")
+ return result.output
+
+ run_sandbox_agent()
+
+
+dag = example_sandbox_toolset_opensandbox()
+
+from tests_common.test_utils.system_tests import get_test_run # noqa: E402
+
+test_run = get_test_run(dag)
diff --git
a/providers/common/ai/tests/unit/common/ai/sandbox/test_opensandbox.py
b/providers/common/ai/tests/unit/common/ai/sandbox/test_opensandbox.py
new file mode 100644
index 00000000000..d1006d9519f
--- /dev/null
+++ b/providers/common/ai/tests/unit/common/ai/sandbox/test_opensandbox.py
@@ -0,0 +1,725 @@
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements. See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership. The ASF licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License. You may obtain a copy of the License at
+#
+# http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing,
+# software distributed under the License is distributed on an
+# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+# KIND, either express or implied. See the License for the
+# specific language governing permissions and limitations
+# under the License.
+from __future__ import annotations
+
+import builtins
+import threading
+import time
+from types import SimpleNamespace
+from unittest import mock
+
+import pytest
+
+pytest.importorskip("opensandbox")
+
+from opensandbox.exceptions import SandboxApiException
+from opensandbox.models.sandboxes import NetworkPolicy, NetworkRule
+
+from airflow.providers.common.ai.sandbox.base import (
+ SandboxError,
+ SandboxFileTooLargeError,
+ SandboxSpec,
+ SandboxTerminalError,
+)
+from airflow.providers.common.ai.sandbox.opensandbox import OpenSandboxBackend
+
+_BASE_HOOK_PATH = "airflow.providers.common.ai.sandbox.opensandbox.BaseHook"
+_MONOTONIC_PATH =
"airflow.providers.common.ai.sandbox.opensandbox.time.monotonic"
+
+
+def _api_error(status_code: int) -> SandboxApiException:
+ return SandboxApiException(status_code=status_code)
+
+
+def _connection(
+ *,
+ password: str | None = "secret",
+ host: str | None = "sandbox.example",
+ port: int | None = 443,
+ schema: str | None = "https",
+ extra: dict | None = None,
+):
+ return SimpleNamespace(
+ password=password,
+ host=host,
+ port=port,
+ schema=schema,
+ extra_dejson=extra or {},
+ )
+
+
+def _error(value: str, *, traceback: list[str] | None = None):
+ return SimpleNamespace(name="CommandError", value=value,
traceback=traceback or [])
+
+
+def _execution(*, error=None, complete=None):
+ """Shape a foreground result as the SDK does: an int parsed from
error.value, 0 on complete, else None."""
+ if error is not None:
+ try:
+ exit_code = int(error.value)
+ except ValueError:
+ exit_code = None
+ else:
+ exit_code = 0 if complete is not None else None
+ return SimpleNamespace(exit_code=exit_code, error=error, complete=complete)
+
+
+def _completed():
+ return _execution(complete=SimpleNamespace())
+
+
+def _deny_policy(*targets: str) -> NetworkPolicy:
+ return NetworkPolicy(
+ defaultAction="deny",
+ egress=[NetworkRule(action="allow", target=target) for target in
targets] or None,
+ )
+
+
+def _created(policy: NetworkPolicy | Exception | None = None):
+ sandbox = mock.MagicMock(spec=["id", "get_egress_policy", "destroy"])
+ sandbox.id = "created"
+ if isinstance(policy, Exception):
+ sandbox.get_egress_policy.side_effect = policy
+ else:
+ sandbox.get_egress_policy.return_value = policy
+ return sandbox
+
+
+def _backend_with_sandbox(**kwargs) -> tuple[OpenSandboxBackend,
mock.MagicMock]:
+ backend = OpenSandboxBackend(**kwargs)
+ backend._connection_config = mock.sentinel.connection_config
+ sandbox = mock.MagicMock(spec=["id", "commands", "files", "get_info",
"destroy"])
+ sandbox.id = "box-1"
+ sandbox.commands = mock.MagicMock(spec=["run"])
+ sandbox.files = mock.MagicMock(
+ spec=["read_bytes_stream", "create_directories", "write_file",
"list_directory", "get_file_info"]
+ )
+ backend._sandboxes[sandbox.id] = sandbox
+ return backend, sandbox
+
+
+def test_missing_sdk_error_is_actionable():
+ real_import = builtins.__import__
+
+ def blocked_import(name, *args, **kwargs):
+ if name.startswith("opensandbox"):
+ raise ImportError("blocked for test")
+ return real_import(name, *args, **kwargs)
+
+ backend = OpenSandboxBackend(opensandbox_conn_id=None)
+ with mock.patch("builtins.__import__", side_effect=blocked_import):
+ with pytest.raises(SandboxTerminalError, match="opensandbox"):
+ backend.create()
+
+
[email protected](
+ ("kwargs", "message"),
+ [
+ ({"image": ""}, "image"),
+ ({"cpu": ""}, "cpu"),
+ ({"memory": ""}, "memory"),
+ ({"sandbox_timeout": 0}, "sandbox_timeout"),
+ ({"ready_timeout": 0}, "ready_timeout"),
+ ],
+)
+def test_constructor_rejects_invalid_values(kwargs, message):
+ with pytest.raises(ValueError, match=message):
+ OpenSandboxBackend(**kwargs)
+
+
+class TestConnection:
+ @mock.patch("opensandbox.config.ConnectionConfigSync", autospec=True)
+ @mock.patch(_BASE_HOOK_PATH, autospec=True)
+ def
test_airflow_connection_fields_and_allowlisted_extras_are_forwarded(self, hook,
config):
+ hook.get_connection.return_value = _connection(
+ extra={"request_timeout": "12.5", "use_server_proxy": "false",
"ignored": "value"}
+ )
+ backend = OpenSandboxBackend(opensandbox_conn_id="my_opensandbox")
+
+ backend._get_connection_config()
+
+ hook.get_connection.assert_called_once_with("my_opensandbox")
+ kwargs = config.call_args.kwargs
+ assert kwargs["api_key"] == "secret"
+ assert kwargs["domain"] == "sandbox.example:443"
+ assert kwargs["protocol"] == "https"
+ assert kwargs["request_timeout"].total_seconds() == 12.5
+ assert kwargs["use_server_proxy"] is False
+ assert "ignored" not in kwargs
+
+ @mock.patch("opensandbox.config.ConnectionConfigSync", autospec=True)
+ @mock.patch(_BASE_HOOK_PATH, autospec=True)
+ def test_connection_is_resolved_once_and_cached(self, hook, _config):
+ hook.get_connection.return_value = _connection()
+ backend = OpenSandboxBackend()
+
+ backend._get_connection_config()
+ backend._get_connection_config()
+
+ hook.get_connection.assert_called_once_with("opensandbox_default")
+
+ @mock.patch("opensandbox.config.ConnectionConfigSync", autospec=True)
+ def test_none_connection_id_defers_to_sdk_environment(self, config):
+ OpenSandboxBackend(opensandbox_conn_id=None)._get_connection_config()
+
+ config.assert_called_once_with(use_server_proxy=True)
+
+ @pytest.mark.parametrize(
+ ("extra", "message"),
+ [
+ ({"request_timeout": "never"}, "request_timeout"),
+ ({"use_server_proxy": "sometimes"}, "use_server_proxy"),
+ ],
+ )
+ @mock.patch(_BASE_HOOK_PATH, autospec=True)
+ def test_invalid_connection_extra_is_terminal(self, hook, extra, message):
+ hook.get_connection.return_value = _connection(extra=extra)
+
+ with pytest.raises(SandboxTerminalError, match=message):
+ OpenSandboxBackend()._get_connection_config()
+
+ @pytest.mark.parametrize("host", [None, ""])
+ @mock.patch("opensandbox.config.ConnectionConfigSync", autospec=True)
+ @mock.patch(_BASE_HOOK_PATH, autospec=True)
+ def test_connection_without_host_is_terminal_rather_than_localhost(self,
hook, config, host):
+ hook.get_connection.return_value = _connection(host=host)
+
+ with pytest.raises(SandboxTerminalError, match="has no host"):
+ OpenSandboxBackend()._get_connection_config()
+
+ config.assert_not_called()
+
+
+class TestCreate:
+ @mock.patch("opensandbox.SandboxSync.create", autospec=True)
+ def test_spec_resources_and_timeouts_are_forwarded(self, create):
+ create.return_value = _created(_deny_policy("pypi.org"))
+ backend = OpenSandboxBackend(
+ image="python:3.13-slim",
+ cpu="2",
+ memory="4Gi",
+ sandbox_timeout=300,
+ ready_timeout=45,
+ )
+ backend._connection_config = mock.sentinel.connection_config
+
+ sandbox_id = backend.create(spec=SandboxSpec(env={"TOKEN": "value"},
allow_egress_to=["pypi.org"]))
+
+ assert sandbox_id == "created"
+ kwargs = create.call_args.kwargs
+ assert create.call_args.args == ("python:3.13-slim",)
+ assert kwargs["env"] == {"TOKEN": "value"}
+ assert kwargs["resource"] == {"cpu": "2", "memory": "4Gi"}
+ assert kwargs["timeout"].total_seconds() == 300
+ assert kwargs["ready_timeout"].total_seconds() == 45
+ assert kwargs["connection_config"] is mock.sentinel.connection_config
+ assert kwargs["network_policy"].default_action == "deny"
+ assert [(rule.action, rule.target) for rule in
kwargs["network_policy"].egress] == [
+ ("allow", "pypi.org")
+ ]
+
+ @mock.patch("opensandbox.SandboxSync.create", autospec=True)
+ def test_sandboxes_are_tagged_as_airflows(self, create):
+ create.return_value = _created(_deny_policy())
+ backend = OpenSandboxBackend()
+ backend._connection_config = mock.sentinel.connection_config
+
+ backend.create(spec=SandboxSpec())
+
+ metadata = create.call_args.kwargs["metadata"]
+ assert metadata["created-by"] == "airflow"
+ assert metadata["name"].startswith("airflow-sandbox-")
+
+ @pytest.mark.parametrize(
+ ("spec", "expected"),
+ [
+ (SandboxSpec(), "deny"),
+ (SandboxSpec(block_network=False), None),
+ ],
+ )
+ @mock.patch("opensandbox.SandboxSync.create", autospec=True)
+ def test_network_default_is_mapped(self, create, spec, expected):
+ create.return_value = _created(_deny_policy())
+ backend = OpenSandboxBackend()
+ backend._connection_config = mock.sentinel.connection_config
+
+ backend.create(spec=spec)
+
+ policy = create.call_args.kwargs["network_policy"]
+ assert (policy.default_action if policy is not None else None) ==
expected
+
+ @mock.patch("opensandbox.SandboxSync.create", autospec=True)
+ def test_none_spec_states_no_network_requirement(self, create):
+ create.return_value = _created()
+ backend = OpenSandboxBackend()
+ backend._connection_config = mock.sentinel.connection_config
+
+ backend.create()
+
+ assert create.call_args.kwargs["network_policy"] is None
+ create.return_value.get_egress_policy.assert_not_called()
+
+ def test_open_network_with_allowlist_is_refused(self):
+ backend = OpenSandboxBackend()
+
+ with pytest.raises(SandboxTerminalError, match="block_network=True"):
+ backend.create(spec=SandboxSpec(block_network=False,
allow_egress_to=["example.com"]))
+
+ def test_cidr_allowlist_is_refused_fail_closed(self):
+ backend = OpenSandboxBackend()
+
+ with pytest.raises(SandboxTerminalError,
match="allow_egress_to_cidrs"):
+
backend.create(spec=SandboxSpec(allow_egress_to_cidrs=["203.0.113.0/24"]))
+
+ @mock.patch("opensandbox.SandboxSync.create", autospec=True)
+ def
test_enforced_deny_policy_is_read_back_before_the_sandbox_is_handed_out(self,
create):
+ create.return_value = _created(_deny_policy("pypi.org"))
+ backend = OpenSandboxBackend()
+ backend._connection_config = mock.sentinel.connection_config
+
+ assert backend.create(spec=SandboxSpec(allow_egress_to=["pypi.org"]))
== "created"
+
+ create.return_value.get_egress_policy.assert_called_once_with()
+ create.return_value.destroy.assert_not_called()
+
+ @pytest.mark.parametrize(
+ "enforced",
+ [
+ NetworkPolicy(defaultAction="allow", egress=None),
+ _deny_policy("pypi.org", "example.com"),
+ _deny_policy(),
+ _api_error(404),
+ ],
+ ids=["open-egress", "wider-allowlist", "missing-rule",
"no-sidecar-endpoint"],
+ )
+ @mock.patch("opensandbox.SandboxSync.create", autospec=True)
+ def test_unenforced_policy_destroys_the_sandbox_and_is_terminal(self,
create, enforced):
+ create.return_value = _created(enforced)
+ backend = OpenSandboxBackend()
+ backend._connection_config = mock.sentinel.connection_config
+
+ with pytest.raises(SandboxTerminalError, match="did not enforce the
requested network policy"):
+ backend.create(spec=SandboxSpec(allow_egress_to=["pypi.org"]))
+
+ create.return_value.destroy.assert_called_once_with()
+ assert backend._sandboxes == {}
+
+ @mock.patch("opensandbox.SandboxSync.create", autospec=True)
+ def test_open_network_needs_no_read_back(self, create):
+ create.return_value = _created()
+ backend = OpenSandboxBackend()
+ backend._connection_config = mock.sentinel.connection_config
+
+ backend.create(spec=SandboxSpec(block_network=False))
+
+ assert create.call_args.kwargs["network_policy"] is None
+ create.return_value.get_egress_policy.assert_not_called()
+
+ @mock.patch("opensandbox.SandboxSync.create", autospec=True)
+ def test_api_failure_is_terminal(self, create):
+ create.side_effect = _api_error(503)
+ backend = OpenSandboxBackend()
+ backend._connection_config = mock.sentinel.connection_config
+
+ with pytest.raises(SandboxTerminalError, match="HTTP 503"):
+ backend.create(spec=SandboxSpec())
+
+
+class TestRunCommand:
+ def test_streams_output_into_byte_bounded_tails(self):
+ backend, sandbox = _backend_with_sandbox()
+
+ def run(_command, *, opts, handlers):
+ assert opts.timeout.total_seconds() == 9
+ assert handlers.skip_accumulation
+ handlers.on_stdout(SimpleNamespace(text="prefix-"))
+ handlers.on_stdout(SimpleNamespace(text="ééé"))
+ handlers.on_stderr(SimpleNamespace(text="stderr-tail"))
+ return _execution(error=_error("3"))
+
+ sandbox.commands.run.side_effect = run
+
+ result = backend.run_command("box-1", "echo hi", timeout=9,
max_output_bytes=6)
+
+ assert result.exit_code == 3
+ assert result.stdout == "éé\n"
+ assert result.stderr == "-tail\n"
+ assert result.stdout_truncated
+ assert result.stderr_truncated
+ assert not result.sandbox_terminated
+
+ def test_line_delimiters_are_restored(self):
+ """execd strips the delimiter from each streamed line; the tail must
put it back."""
+ backend, sandbox = _backend_with_sandbox()
+
+ def run(_command, *, opts, handlers):
+ for text in ("first", "second", "\n", "third"):
+ handlers.on_stdout(SimpleNamespace(text=text))
+ handlers.on_stderr(SimpleNamespace(text="err one"))
+ handlers.on_stderr(SimpleNamespace(text="err two"))
+ return _completed()
+
+ sandbox.commands.run.side_effect = run
+
+ result = backend.run_command("box-1", "printf ...", timeout=9,
max_output_bytes=4096)
+
+ assert result.exit_code == 0
+ assert result.stdout == "first\nsecond\n\nthird\n"
+ assert result.stderr == "err one\nerr two\n"
+
+ def test_execution_error_is_returned_on_stderr(self):
+ backend, sandbox = _backend_with_sandbox()
+ sandbox.commands.run.return_value = _execution(
+ error=_error("failed", traceback=["line one", "line two"])
+ )
+
+ result = backend.run_command("box-1", "bad", timeout=9,
max_output_bytes=100)
+
+ assert result.exit_code == 1
+ assert result.stderr == "line one\nline two"
+
+ @pytest.mark.parametrize("value", ["exit status 1", "signal: killed"])
+ def test_prose_error_value_is_a_failed_command_not_a_terminal_error(self,
value):
+ """The SDK cannot parse an exit code out of prose; that is an unknown
status, not a dead sandbox."""
+ backend, sandbox = _backend_with_sandbox()
+ sandbox.commands.run.return_value = _execution(error=_error(value))
+
+ result = backend.run_command("box-1", "bad", timeout=9,
max_output_bytes=100)
+
+ assert result.exit_code == 1
+ assert result.stderr == value
+ assert not result.timed_out
+
+ def test_no_terminal_event_is_reported_not_terminal(self):
+ backend, sandbox = _backend_with_sandbox()
+ sandbox.commands.run.return_value = _execution()
+
+ result = backend.run_command("box-1", "echo hi", timeout=5,
max_output_bytes=100)
+
+ assert result.exit_code == -1
+ assert "without reporting an exit status" in result.stderr
+
+ @pytest.mark.parametrize("value", ["-9", "124", "signal: killed"])
+ @mock.patch(_MONOTONIC_PATH, side_effect=[0.0, 5.0])
+ def test_nonzero_exit_at_the_deadline_is_a_timeout(self, _monotonic,
value):
+ backend, sandbox = _backend_with_sandbox()
+ sandbox.commands.run.return_value = _execution(error=_error(value))
+
+ result = backend.run_command("box-1", "sleep 60", timeout=5,
max_output_bytes=100)
+
+ assert result.timed_out
+ assert not result.sandbox_terminated
+
+ @mock.patch(_MONOTONIC_PATH, side_effect=[0.0, 5.0])
+ def test_clean_exit_at_the_deadline_is_not_a_timeout(self, _monotonic):
+ backend, sandbox = _backend_with_sandbox()
+ sandbox.commands.run.return_value = _completed()
+
+ result = backend.run_command("box-1", "sleep 5", timeout=5,
max_output_bytes=100)
+
+ assert not result.timed_out
+
+ @mock.patch(_MONOTONIC_PATH, side_effect=[0.0, 0.2])
+ def test_negative_exit_well_inside_the_deadline_is_not_a_timeout(self,
_monotonic):
+ """A signal kill is only a timeout when the call also outlived the
deadline."""
+ backend, sandbox = _backend_with_sandbox()
+ sandbox.commands.run.return_value = _execution(error=_error("-9"))
+
+ result = backend.run_command("box-1", "kill -9 $$", timeout=30,
max_output_bytes=100)
+
+ assert not result.timed_out
+
+ @mock.patch("opensandbox.SandboxSync.connect", autospec=True)
+ def test_reconnect_time_is_not_charged_to_the_command(self, connect):
+ remote = mock.MagicMock(spec=["commands"])
+ remote.commands = mock.MagicMock(spec=["run"])
+ remote.commands.run.return_value = _execution(error=_error("-9"))
+
+ def slow_connect(*_args, **_kwargs):
+ time.sleep(0.3)
+ return remote
+
+ connect.side_effect = slow_connect
+ backend = OpenSandboxBackend()
+ backend._connection_config = mock.sentinel.connection_config
+
+ result = backend.run_command("remote", "kill -9 $$", timeout=0.25,
max_output_bytes=100)
+
+ assert not result.timed_out
+
+ def test_error_details_do_not_overwrite_streamed_stderr(self):
+ backend, sandbox = _backend_with_sandbox()
+
+ def run(_command, *, opts, handlers):
+ handlers.on_stderr(SimpleNamespace(text="real stderr"))
+ return _execution(error=_error("1", traceback=["exit status 1"]))
+
+ sandbox.commands.run.side_effect = run
+
+ result = backend.run_command("box-1", "bad", timeout=9,
max_output_bytes=100)
+
+ assert result.stderr == "real stderr\n"
+
+ @mock.patch("airflow.providers.common.ai.sandbox.opensandbox._EXEC_GRACE",
0.0)
+ def test_stalled_stream_destroys_the_sandbox_and_reports_a_timeout(self):
+ backend, sandbox = _backend_with_sandbox()
+ release = threading.Event()
+
+ def hang(_command, *, opts, handlers):
+ handlers.on_stdout(SimpleNamespace(text="partial"))
+ release.wait(5)
+ return _completed()
+
+ sandbox.commands.run.side_effect = hang
+ try:
+ result = backend.run_command("box-1", "yes", timeout=0.1,
max_output_bytes=100)
+ finally:
+ release.set()
+
+ assert result.timed_out
+ assert result.sandbox_terminated
+ assert result.exit_code == -1
+ assert result.stdout == "partial\n"
+ sandbox.destroy.assert_called_once_with()
+ assert "box-1" not in backend._sandboxes
+
+ @mock.patch("airflow.providers.common.ai.sandbox.opensandbox._EXEC_GRACE",
0.0)
+ def
test_stalled_stream_whose_sandbox_cannot_be_destroyed_still_reports_a_timeout(self):
+ backend, sandbox = _backend_with_sandbox()
+ release = threading.Event()
+ sandbox.commands.run.side_effect = lambda *_a, **_k: release.wait(5)
+ sandbox.destroy.side_effect = _api_error(503)
+ try:
+ result = backend.run_command("box-1", "yes", timeout=0.1,
max_output_bytes=100)
+ finally:
+ release.set()
+
+ assert result.timed_out
+ assert result.sandbox_terminated
+
+ def test_api_failure_is_terminal(self):
+ backend, sandbox = _backend_with_sandbox()
+ sandbox.commands.run.side_effect = _api_error(401)
+
+ with pytest.raises(SandboxTerminalError, match="HTTP 401"):
+ backend.run_command("box-1", "echo hi", timeout=5,
max_output_bytes=100)
+
+ @pytest.mark.parametrize(
+ ("timeout", "max_bytes", "message"),
+ [(0, 1, "timeout"), (1, 0, "max_output_bytes")],
+ )
+ def test_rejects_invalid_budgets(self, timeout, max_bytes, message):
+ backend, _ = _backend_with_sandbox()
+
+ with pytest.raises(ValueError, match=message):
+ backend.run_command("box-1", "echo hi", timeout=timeout,
max_output_bytes=max_bytes)
+
+
+class TestFileOperations:
+ def test_read_uses_range_and_native_streaming(self):
+ backend, sandbox = _backend_with_sandbox()
+ stream = mock.MagicMock(spec=["__iter__", "close"])
+ stream.__iter__.return_value = iter([b"he", b"llo"])
+ sandbox.files.read_bytes_stream.return_value = stream
+
+ assert backend.read_file("box-1", "/w/a", max_bytes=10) == b"hello"
+ sandbox.files.read_bytes_stream.assert_called_once_with(
+ "/w/a", chunk_size=11, range_header="bytes=0-10"
+ )
+ stream.close.assert_called_once()
+
+ def
test_oversized_read_stops_at_the_sentinel_byte_and_reports_the_real_size(self):
+ backend, sandbox = _backend_with_sandbox()
+ stream = mock.MagicMock(spec=["__iter__", "close"])
+ stream.__iter__.return_value = iter([b"x" * 11, b"must-not-be-read"])
+ sandbox.files.read_bytes_stream.return_value = stream
+ sandbox.files.get_file_info.return_value = {"/w/a":
SimpleNamespace(size=1_000_000)}
+
+ with pytest.raises(SandboxFileTooLargeError) as error:
+ backend.read_file("box-1", "/w/a", max_bytes=10)
+
+ assert error.value.size_bytes == 1_000_000
+ sandbox.files.get_file_info.assert_called_once_with(["/w/a"])
+ stream.close.assert_called_once()
+
+ @pytest.mark.parametrize(
+ "lookup",
+ [{"/w/a": SimpleNamespace(size=0)}, _api_error(500)],
+ ids=["streaming-source", "lookup-failed"],
+ )
+ def test_oversized_read_never_reports_less_than_was_read(self, lookup):
+ backend, sandbox = _backend_with_sandbox()
+ stream = mock.MagicMock(spec=["__iter__", "close"])
+ stream.__iter__.return_value = iter([b"x" * 11])
+ sandbox.files.read_bytes_stream.return_value = stream
+ if isinstance(lookup, Exception):
+ sandbox.files.get_file_info.side_effect = lookup
+ else:
+ sandbox.files.get_file_info.return_value = lookup
+
+ with pytest.raises(SandboxFileTooLargeError) as error:
+ backend.read_file("box-1", "/w/a", max_bytes=10)
+
+ assert error.value.size_bytes == 11
+
+ def test_missing_file_is_recoverable_when_sandbox_exists(self):
+ backend, sandbox = _backend_with_sandbox()
+ sandbox.files.read_bytes_stream.side_effect = _api_error(404)
+
+ with pytest.raises(SandboxError, match="does not exist") as error:
+ backend.read_file("box-1", "/w/missing", max_bytes=10)
+
+ assert not isinstance(error.value, SandboxTerminalError)
+ sandbox.get_info.assert_called_once()
+
+ def test_missing_sandbox_is_terminal(self):
+ backend, sandbox = _backend_with_sandbox()
+ sandbox.files.read_bytes_stream.side_effect = _api_error(404)
+ sandbox.get_info.side_effect = _api_error(404)
+
+ with pytest.raises(SandboxTerminalError, match="confirm that a sandbox
still exists"):
+ backend.read_file("box-1", "/w/a", max_bytes=10)
+
+ def test_bad_read_request_is_recoverable(self):
+ backend, sandbox = _backend_with_sandbox()
+ sandbox.files.read_bytes_stream.side_effect = _api_error(400)
+
+ with pytest.raises(SandboxError) as error:
+ backend.read_file("box-1", "bad", max_bytes=10)
+
+ assert not isinstance(error.value, SandboxTerminalError)
+
+ def test_write_creates_parent_and_uses_native_file_api(self):
+ backend, sandbox = _backend_with_sandbox()
+
+ backend.write_file("box-1", "/w/sub/a", b"data")
+
+ entry = sandbox.files.create_directories.call_args.args[0][0]
+ assert (entry.path, entry.mode) == ("/w/sub", 755)
+ sandbox.files.write_file.assert_called_once_with("/w/sub/a", b"data",
mode=644)
+
+ def test_missing_write_path_is_recoverable_when_sandbox_exists(self):
+ backend, sandbox = _backend_with_sandbox()
+ sandbox.files.write_file.side_effect = _api_error(404)
+
+ with pytest.raises(SandboxError, match="Could not write") as error:
+ backend.write_file("box-1", "a", b"data")
+
+ assert not isinstance(error.value, SandboxTerminalError)
+ sandbox.get_info.assert_called_once()
+
+ def test_list_returns_direct_children_and_marks_directories(self):
+ backend, sandbox = _backend_with_sandbox()
+ sandbox.files.list_directory.return_value = [
+ SimpleNamespace(path="/w/a.txt", entry_type="file"),
+ SimpleNamespace(path="/w/sub/", entry_type="directory"),
+ ]
+
+ assert backend.list_directory("box-1", "/w") == [("a.txt", False),
("sub", True)]
+ entry = sandbox.files.list_directory.call_args.args[0]
+ assert (entry.path, entry.depth) == ("/w", 1)
+
+
@mock.patch("airflow.providers.common.ai.sandbox.opensandbox._LIST_DIRECTORY_MAX_ENTRIES",
2)
+ def test_oversized_listing_is_refused_with_guidance(self):
+ backend, sandbox = _backend_with_sandbox()
+ sandbox.files.list_directory.return_value = [
+ SimpleNamespace(path=f"/w/{index}", entry_type="file") for index
in range(3)
+ ]
+
+ with pytest.raises(SandboxError, match="more than 2 entries") as error:
+ backend.list_directory("box-1", "/w")
+
+ assert not isinstance(error.value, SandboxTerminalError)
+
+ def test_missing_directory_is_recoverable_when_sandbox_exists(self):
+ backend, sandbox = _backend_with_sandbox()
+ sandbox.files.list_directory.side_effect = _api_error(404)
+
+ with pytest.raises(SandboxError, match="does not exist") as error:
+ backend.list_directory("box-1", "/missing")
+
+ assert not isinstance(error.value, SandboxTerminalError)
+
+ def test_read_rejects_invalid_budget(self):
+ backend, _ = _backend_with_sandbox()
+
+ with pytest.raises(ValueError, match="max_bytes"):
+ backend.read_file("box-1", "/w/a", max_bytes=0)
+
+
+class TestGetSandbox:
+ @mock.patch("opensandbox.SandboxSync.connect", autospec=True)
+ def test_uncached_sandbox_is_reconnected_and_then_cached(self, connect):
+ remote = mock.MagicMock(spec=["commands"])
+ remote.commands = mock.MagicMock(spec=["run"])
+ remote.commands.run.return_value = _completed()
+ connect.return_value = remote
+ backend = OpenSandboxBackend()
+ backend._connection_config = mock.sentinel.connection_config
+
+ backend.run_command("remote", "echo hi", timeout=5,
max_output_bytes=100)
+ backend.run_command("remote", "echo hi again", timeout=5,
max_output_bytes=100)
+
+ connect.assert_called_once()
+ assert connect.call_args.kwargs.get("skip_health_check", False) is
False
+ assert backend._sandboxes["remote"] is remote
+
+ @mock.patch("opensandbox.SandboxSync.connect", autospec=True)
+ def test_reconnect_failure_is_terminal(self, connect):
+ connect.side_effect = _api_error(503)
+ backend = OpenSandboxBackend()
+ backend._connection_config = mock.sentinel.connection_config
+
+ with pytest.raises(SandboxTerminalError, match="HTTP 503"):
+ backend.run_command("remote", "echo hi", timeout=5,
max_output_bytes=100)
+
+
+class TestDestroy:
+ def test_cached_sandbox_is_destroyed_and_evicted(self):
+ backend, sandbox = _backend_with_sandbox()
+
+ backend.destroy("box-1")
+
+ sandbox.destroy.assert_called_once()
+ assert "box-1" not in backend._sandboxes
+
+ @mock.patch("opensandbox.SandboxSync.connect", autospec=True)
+ def test_uncached_sandbox_is_reconnected_without_a_readiness_poll(self,
connect):
+ remote = mock.MagicMock(spec=["destroy"])
+ connect.return_value = remote
+ backend = OpenSandboxBackend()
+ backend._connection_config = mock.sentinel.connection_config
+
+ backend.destroy("remote")
+
+ remote.destroy.assert_called_once()
+ assert connect.call_args.kwargs["skip_health_check"] is True
+
+ @mock.patch("opensandbox.SandboxSync.connect", autospec=True)
+ def test_already_gone_sandbox_is_idempotent(self, connect):
+ connect.side_effect = _api_error(404)
+ backend = OpenSandboxBackend()
+ backend._connection_config = mock.sentinel.connection_config
+
+ backend.destroy("gone")
+
+ def test_destroy_failure_is_terminal(self):
+ backend, sandbox = _backend_with_sandbox()
+ sandbox.destroy.side_effect = _api_error(503)
+
+ with pytest.raises(SandboxTerminalError, match="HTTP 503"):
+ backend.destroy("box-1")
diff --git a/uv.lock b/uv.lock
index ca75240a384..52020147fa8 100644
--- a/uv.lock
+++ b/uv.lock
@@ -4675,6 +4675,9 @@ parquet = [
pdf = [
{ name = "pypdf" },
]
+opensandbox = [
+ { name = "opensandbox" },
+]
shields = [
{ name = "pydantic-ai-shields" },
]
@@ -4705,6 +4708,7 @@ dev = [
{ name = "llama-index-embeddings-openai" },
{ name = "llama-index-llms-openai" },
{ name = "openai" },
+ { name = "opensandbox" },
{ name = "pydantic-ai-skills" },
{ name = "pydantic-ai-slim", extra = ["mcp"] },
{ name = "sqlglot" },
@@ -4732,6 +4736,7 @@ requires-dist = [
{ name = "llama-index-llms-openai", marker = "extra == 'llamaindex'",
specifier = ">=0.6.8" },
{ name = "modal", marker = "extra == 'modal'", specifier = ">=1.5.0" },
{ name = "openai", marker = "extra == 'openai'", specifier = ">=2.47.0" },
+ { name = "opensandbox", marker = "extra == 'opensandbox'", specifier =
">=1.1.0" },
{ name = "pyarrow", marker = "python_full_version >= '3.14' and extra ==
'parquet'", specifier = ">=22.0.0" },
{ name = "pyarrow", marker = "python_full_version < '3.14' and extra ==
'parquet'", specifier = ">=18.0.0" },
{ name = "pydantic-ai-harness", extras = ["codemode"], marker = "extra ==
'code-mode'", specifier = ">=0.3.0" },
@@ -4748,7 +4753,7 @@ requires-dist = [
{ name = "sqlglot", marker = "extra == 'sql'", specifier = ">=30.0.0" },
{ name = "typesafe-sdk", marker = "extra == 'typesafe'", specifier =
">=0.6.0" },
]
-provides-extras = ["anthropic", "bedrock", "google", "openai", "typesafe",
"mcp", "modal", "code-mode", "shields", "skills", "avro", "parquet", "sql",
"common-sql", "langchain", "llamaindex", "pdf", "docx", "git"]
+provides-extras = ["anthropic", "bedrock", "google", "openai", "typesafe",
"mcp", "modal", "opensandbox", "code-mode", "shields", "skills", "avro",
"parquet", "sql", "common-sql", "langchain", "llamaindex", "pdf", "docx", "git"]
[package.metadata.requires-dev]
dev = [
@@ -4766,6 +4771,7 @@ dev = [
{ name = "llama-index-embeddings-openai", specifier = ">=0.6.0" },
{ name = "llama-index-llms-openai", specifier = ">=0.6.8" },
{ name = "openai", specifier = ">=2.47.0" },
+ { name = "opensandbox", specifier = ">=1.1.0" },
{ name = "pydantic-ai-skills", specifier = ">=1.2.0" },
{ name = "pydantic-ai-slim", extras = ["mcp"], specifier = ">=2.33.0" },
{ name = "sqlglot", specifier = ">=30.0.0" },
@@ -15204,6 +15210,15 @@ http2 = [
{ name = "h2" },
]
+[[package]]
+name = "httpx-sse"
+version = "0.4.3"
+source = { registry = "https://pypi.org/simple" }
+sdist = { url =
"https://files.pythonhosted.org/packages/0f/4c/751061ffa58615a32c31b2d82e8482be8dd4a89154f003147acee90f2be9/httpx_sse-0.4.3.tar.gz",
hash =
"sha256:9b1ed0127459a66014aec3c56bebd93da3c1bc8bb6618c8082039a44889a755d", size
= 15943, upload-time = "2025-10-10T21:48:22.271Z" }
+wheels = [
+ { url =
"https://files.pythonhosted.org/packages/d2/fd/6668e5aec43ab844de6fc74927e155a3b37bf40d7c3790e49fc0406b6578/httpx_sse-0.4.3-py3-none-any.whl",
hash =
"sha256:0ac1c9fe3c0afad2e0ebb25a934a59f4c7823b60792691f779fad2c5568830fc", size
= 8960, upload-time = "2025-10-10T21:48:21.158Z" },
+]
+
[[package]]
name = "httpx2"
version = "2.13.1"
@@ -19080,6 +19095,24 @@ wheels = [
{ url =
"https://files.pythonhosted.org/packages/c0/da/977ded879c29cbd04de313843e76868e6e13408a94ed6b987245dc7c8506/openpyxl-3.1.5-py2.py3-none-any.whl",
hash =
"sha256:5282c12b107bffeef825f4617dc029afaf41d0ea60823bbb665ef3079dc79de2", size
= 250910, upload-time = "2024-06-28T14:03:41.161Z" },
]
+[[package]]
+name = "opensandbox"
+version = "1.1.0"
+source = { registry = "https://pypi.org/simple" }
+dependencies = [
+ { name = "attrs" },
+ { name = "httpcore" },
+ { name = "httpx" },
+ { name = "httpx-sse" },
+ { name = "opentelemetry-api" },
+ { name = "pydantic" },
+ { name = "python-dateutil" },
+]
+sdist = { url =
"https://files.pythonhosted.org/packages/ba/81/2a8b1e2f94ce8802a425c1e004cb678bb5f1962fad13b7884be10833246d/opensandbox-1.1.0.tar.gz",
hash =
"sha256:b7bf74313a99cd719337f665db9f15265bb8b5b64032fedd24069ad9b031cc95", size
= 286824, upload-time = "2026-09-21T06:34:54.256Z" }
+wheels = [
+ { url =
"https://files.pythonhosted.org/packages/11/09/01540f419dba266477ce6cdec22831667e1f2898ae070906c3e973163398/opensandbox-1.1.0-py3-none-any.whl",
hash =
"sha256:aae471beb3e37ab206b134ff1ea590c1ca649dea0c5ec0391a7cfb416317f697", size
= 642050, upload-time = "2026-09-21T06:34:52.708Z" },
+]
+
[[package]]
name = "opensearch-protobufs"
version = "1.2.0"