kaxil commented on code in PR #73961:
URL: https://github.com/apache/airflow/pull/73961#discussion_r4211012485


##########
providers/common/ai/src/airflow/providers/common/ai/sandbox/openshell.py:
##########
@@ -0,0 +1,995 @@
+# 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.
+"""NVIDIA OpenShell backend for 
:class:`~airflow.providers.common.ai.toolsets.sandbox.SandboxToolset`."""
+
+from __future__ import annotations
+
+import logging
+import math
+import shlex
+import threading
+import time
+import uuid
+from contextlib import contextmanager, suppress
+from dataclasses import dataclass
+from typing import TYPE_CHECKING, Any
+
+from airflow.providers.common.ai.sandbox.base import (
+    _FILE_OP_OUTPUT_CAP,
+    _FILE_OP_TIMEOUT,
+    SandboxBackend,
+    SandboxError,
+    SandboxExecResult,
+    SandboxFileTooLargeError,
+    SandboxTerminalError,
+    _validate_positive_finite,
+)
+
+if TYPE_CHECKING:
+    from collections.abc import Iterator, Mapping, Sequence
+
+    from openshell import SandboxClient
+
+    from airflow.providers.common.ai.sandbox.base import SandboxSpec
+
+log = logging.getLogger(__name__)
+
+_CLIENT_BUILD_LOCK = threading.Lock()
+
+# Wall-clock allowance past the per-command budget before the exec stream is
+# treated as hung. The guest wrapper enforces the budget itself; this covers a
+# gateway or supervisor that stops delivering events at all.
+_EXEC_GRACE = 30.0
+# How long a command or policy read waits for a restarting gateway to come back
+# before the failure is reported. A restart measured at 6-11 s end to end.
+_GATEWAY_RECOVERY = 60.0
+_POLL_INTERVAL = 0.25
+# The gateway decodes at most 1 MiB per request, and the command or a file 
chunk
+# travels as stdin inside that request, so both stay under it with headroom.
+_MAX_STDIN_BYTES = 1_000_000
+_WRITE_CHUNK_BYTES = 768 * 1024
+# The gateway rejects a sandbox name longer than 19 characters, which
+# _new_sandbox_name's ``airflow-sandbox-<12 hex>`` is, so the name is shorter
+# here. Labels carry the attribution instead, and the creation time so an
+# operator can reap by age: OpenShell has no server-side lifetime.
+_NAME_PREFIX = "airflow-"
+_CREATED_BY_LABEL = ("created-by", "airflow")
+_CREATED_AT_LABEL = "airflow-created-at"
+_EGRESS_RULE = "airflow-egress"
+_HTTPS_PORT = 443
+_ANY_BINARY = "/**"
+# OpenShell's restrictive default filesystem policy, plus /dev/shm so Python's
+# multiprocessing and anything else using POSIX shared memory works.
+_READ_ONLY_PATHS = ("/bin", "/usr", "/lib", "/proc", "/dev/urandom", "/etc", 
"/var/log")
+_READ_WRITE_PATHS = ("/tmp", "/dev/null", "/dev/shm")
+# The sandbox supervisor removes the proxy variables from every command it runs
+# and overwrites the CA-bundle variables with its own TLS-terminating CA
+# (openshell-sandbox process.rs PROXY_ENV_VARS and child_env.rs tls_env_vars in
+# 0.1.2), in both cases without an error. A spec naming them would be silently
+# changed, so it is refused instead.
+_SUPERVISOR_OWNED_ENV = frozenset(
+    {
+        "ALL_PROXY",
+        "HTTP_PROXY",
+        "HTTPS_PROXY",
+        "NO_PROXY",
+        "all_proxy",
+        "http_proxy",
+        "https_proxy",
+        "no_proxy",
+        "grpc_proxy",
+        "NODE_USE_ENV_PROXY",
+        "SSL_CERT_FILE",
+        "REQUESTS_CA_BUNDLE",
+        "CURL_CA_BUNDLE",
+        "GIT_SSL_CAINFO",
+        "NODE_EXTRA_CA_CERTS",
+        "DENO_CERT",
+    }
+)
+_PROPOSAL_APPROVAL_MODE = "proposal_approval_mode"
+_AGENT_POLICY_PROPOSALS = "agent_policy_proposals_enabled"
+# Wrapper-private exit status for "the budget ran out". The wrapper only uses 
it
+# after its own timer fired, and the caller also checks the elapsed time, so a
+# command that exits 124 on its own inside the budget is not a timeout.
+_TIMEOUT_STATUS = 124
+# Wrapper-private exit status and stderr line for "the command could not be
+# staged in /tmp, so it did not run". The status alone is not enough, because a
+# command can exit 125 itself, as docker run, env and nohup do on their own 
errors.
+_STAGING_STATUS = 125
+_STAGING_FAILED = "airflow-exec: could not stage the command in /tmp"
+_STAGING_LINE = f"{_STAGING_FAILED}\n".encode()
+_SYSTEM_PATH = 
"PATH=/usr/local/sbin:/usr/local/bin:/usr/sbin:/usr/bin:/sbin:/bin"
+
+# Runs one command for run_command. OpenShell's own exec timeout returns a
+# synthetic 124 and leaves the process running, and a backgrounded child that
+# inherits stdout holds the call open until the supervisor gives up on it, so
+# the wrapper owns both problems:
+#
+# * The command arrives on stdin, which avoids the gateway's 32 KiB 
per-argument
+#   cap, and runs in its own session via setsid, with stdout and stderr going 
to
+#   files, so a background child holds a file rather than the exec stream.
+# * A timer reading the monotonic /proc/uptime signals the wrapper when the
+#   budget is spent; counting its own one-second sleeps instead drifts late on 
a
+#   CPU-starved sandbox, past _EXEC_GRACE. The wrapper then SIGKILLs the
+#   processes in the command's session, rescanning /proc until a pass kills
+#   nothing, at most 50 times, so children forked mid-sweep are usually caught
+#   too. A negative-pid kill of the group is blocked by the sandbox's seccomp
+#   filter. The sweep shares the CPU with the processes it is killing, so it
+#   reads /proc with builtins, joining the lines of a stat file whose process
+#   name holds a newline, and strips the name with the cheap shortest match,
+#   taking the longest only for a name that contains ")": a `cat` per entry let
+#   one sweep of 128 busy processes on one CPU run past _EXEC_GRACE.
+# * The exit status is the wrapper's own process status, which the supervisor
+#   reports out of band, so nothing the command prints can change it; a command
+#   that signals the wrapper only ends its own run early. Output is sent last 
as
+#   one byte more than the caller's cap, so an over-cap stream is detected by
+#   counting what arrives rather than trusting a trailer.
+#
+# A command that starts its own session (setsid, a daemonizing server) leaves
+# this one and is not killed on timeout, and one that keeps forking faster than
+# the sweep can outlast it.
+_RUN_WRAPPER = rf"""t=$1 c=$2 o=$PATH
+{_SYSTEM_PATH}
+d=$(mktemp -d /tmp/.airflow-exec.XXXXXX) && cat >"$d/c" || {{ rm -rf "$d"; 
echo "{_STAGING_FAILED}" >&2; exit {_STAGING_STATUS}; }}
+f=0
+trap f=1 ALRM
+trap f=2 HUP INT TERM
+z=$(command -v setsid) || {{ echo "setsid is not installed in the sandbox 
image" >&2; rm -rf "$d"; exit 125; }}

Review Comment:
   Without `setsid` in the image, every command comes back as `[exit code: 
125]` with this line on stderr, which the model will read as its own command 
failing and keep retrying. It can't install one with no network, so this looks 
like a `SandboxTerminalError` case per the base class contract. Could it exit 
with a private marker the way the staging failure does, and have `run_command` 
raise terminal on it with a message naming the image requirement?



##########
providers/common/ai/src/airflow/providers/common/ai/sandbox/openshell.py:
##########
@@ -0,0 +1,995 @@
+# 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.
+"""NVIDIA OpenShell backend for 
:class:`~airflow.providers.common.ai.toolsets.sandbox.SandboxToolset`."""
+
+from __future__ import annotations
+
+import logging
+import math
+import shlex
+import threading
+import time
+import uuid
+from contextlib import contextmanager, suppress
+from dataclasses import dataclass
+from typing import TYPE_CHECKING, Any
+
+from airflow.providers.common.ai.sandbox.base import (
+    _FILE_OP_OUTPUT_CAP,
+    _FILE_OP_TIMEOUT,
+    SandboxBackend,
+    SandboxError,
+    SandboxExecResult,
+    SandboxFileTooLargeError,
+    SandboxTerminalError,
+    _validate_positive_finite,
+)
+
+if TYPE_CHECKING:
+    from collections.abc import Iterator, Mapping, Sequence
+
+    from openshell import SandboxClient
+
+    from airflow.providers.common.ai.sandbox.base import SandboxSpec
+
+log = logging.getLogger(__name__)
+
+_CLIENT_BUILD_LOCK = threading.Lock()
+
+# Wall-clock allowance past the per-command budget before the exec stream is
+# treated as hung. The guest wrapper enforces the budget itself; this covers a
+# gateway or supervisor that stops delivering events at all.
+_EXEC_GRACE = 30.0
+# How long a command or policy read waits for a restarting gateway to come back
+# before the failure is reported. A restart measured at 6-11 s end to end.
+_GATEWAY_RECOVERY = 60.0
+_POLL_INTERVAL = 0.25
+# The gateway decodes at most 1 MiB per request, and the command or a file 
chunk
+# travels as stdin inside that request, so both stay under it with headroom.
+_MAX_STDIN_BYTES = 1_000_000
+_WRITE_CHUNK_BYTES = 768 * 1024
+# The gateway rejects a sandbox name longer than 19 characters, which
+# _new_sandbox_name's ``airflow-sandbox-<12 hex>`` is, so the name is shorter
+# here. Labels carry the attribution instead, and the creation time so an
+# operator can reap by age: OpenShell has no server-side lifetime.
+_NAME_PREFIX = "airflow-"
+_CREATED_BY_LABEL = ("created-by", "airflow")
+_CREATED_AT_LABEL = "airflow-created-at"
+_EGRESS_RULE = "airflow-egress"
+_HTTPS_PORT = 443
+_ANY_BINARY = "/**"
+# OpenShell's restrictive default filesystem policy, plus /dev/shm so Python's
+# multiprocessing and anything else using POSIX shared memory works.
+_READ_ONLY_PATHS = ("/bin", "/usr", "/lib", "/proc", "/dev/urandom", "/etc", 
"/var/log")
+_READ_WRITE_PATHS = ("/tmp", "/dev/null", "/dev/shm")
+# The sandbox supervisor removes the proxy variables from every command it runs
+# and overwrites the CA-bundle variables with its own TLS-terminating CA
+# (openshell-sandbox process.rs PROXY_ENV_VARS and child_env.rs tls_env_vars in
+# 0.1.2), in both cases without an error. A spec naming them would be silently
+# changed, so it is refused instead.
+_SUPERVISOR_OWNED_ENV = frozenset(
+    {
+        "ALL_PROXY",
+        "HTTP_PROXY",
+        "HTTPS_PROXY",
+        "NO_PROXY",
+        "all_proxy",
+        "http_proxy",
+        "https_proxy",
+        "no_proxy",
+        "grpc_proxy",
+        "NODE_USE_ENV_PROXY",
+        "SSL_CERT_FILE",
+        "REQUESTS_CA_BUNDLE",
+        "CURL_CA_BUNDLE",
+        "GIT_SSL_CAINFO",
+        "NODE_EXTRA_CA_CERTS",
+        "DENO_CERT",
+    }
+)
+_PROPOSAL_APPROVAL_MODE = "proposal_approval_mode"
+_AGENT_POLICY_PROPOSALS = "agent_policy_proposals_enabled"
+# Wrapper-private exit status for "the budget ran out". The wrapper only uses 
it
+# after its own timer fired, and the caller also checks the elapsed time, so a
+# command that exits 124 on its own inside the budget is not a timeout.
+_TIMEOUT_STATUS = 124
+# Wrapper-private exit status and stderr line for "the command could not be
+# staged in /tmp, so it did not run". The status alone is not enough, because a
+# command can exit 125 itself, as docker run, env and nohup do on their own 
errors.
+_STAGING_STATUS = 125
+_STAGING_FAILED = "airflow-exec: could not stage the command in /tmp"
+_STAGING_LINE = f"{_STAGING_FAILED}\n".encode()
+_SYSTEM_PATH = 
"PATH=/usr/local/sbin:/usr/local/bin:/usr/sbin:/usr/bin:/sbin:/bin"
+
+# Runs one command for run_command. OpenShell's own exec timeout returns a
+# synthetic 124 and leaves the process running, and a backgrounded child that
+# inherits stdout holds the call open until the supervisor gives up on it, so
+# the wrapper owns both problems:
+#
+# * The command arrives on stdin, which avoids the gateway's 32 KiB 
per-argument
+#   cap, and runs in its own session via setsid, with stdout and stderr going 
to
+#   files, so a background child holds a file rather than the exec stream.
+# * A timer reading the monotonic /proc/uptime signals the wrapper when the
+#   budget is spent; counting its own one-second sleeps instead drifts late on 
a
+#   CPU-starved sandbox, past _EXEC_GRACE. The wrapper then SIGKILLs the
+#   processes in the command's session, rescanning /proc until a pass kills
+#   nothing, at most 50 times, so children forked mid-sweep are usually caught
+#   too. A negative-pid kill of the group is blocked by the sandbox's seccomp
+#   filter. The sweep shares the CPU with the processes it is killing, so it
+#   reads /proc with builtins, joining the lines of a stat file whose process
+#   name holds a newline, and strips the name with the cheap shortest match,
+#   taking the longest only for a name that contains ")": a `cat` per entry let
+#   one sweep of 128 busy processes on one CPU run past _EXEC_GRACE.
+# * The exit status is the wrapper's own process status, which the supervisor
+#   reports out of band, so nothing the command prints can change it; a command
+#   that signals the wrapper only ends its own run early. Output is sent last 
as
+#   one byte more than the caller's cap, so an over-cap stream is detected by
+#   counting what arrives rather than trusting a trailer.
+#
+# A command that starts its own session (setsid, a daemonizing server) leaves
+# this one and is not killed on timeout, and one that keeps forking faster than
+# the sweep can outlast it.
+_RUN_WRAPPER = rf"""t=$1 c=$2 o=$PATH
+{_SYSTEM_PATH}
+d=$(mktemp -d /tmp/.airflow-exec.XXXXXX) && cat >"$d/c" || {{ rm -rf "$d"; 
echo "{_STAGING_FAILED}" >&2; exit {_STAGING_STATUS}; }}
+f=0
+trap f=1 ALRM
+trap f=2 HUP INT TERM
+z=$(command -v setsid) || {{ echo "setsid is not installed in the sandbox 
image" >&2; rm -rf "$d"; exit 125; }}
+PATH=$o "$z" /bin/sh "$d/c" </dev/null >"$d/o" 2>"$d/e" &
+p=$!
+(
+  read u _ </proc/uptime
+  e=$((${{u%.*}}${{u#*.}} + t * 100))
+  while [ "${{u%.*}}${{u#*.}}" -lt "$e" ]; do
+    sleep 1
+    [ -d "/proc/$$" ] || exit 0
+    read u _ </proc/uptime
+  done
+  kill -ALRM $$
+) </dev/null >/dev/null 2>&1 &
+w=$!
+wait "$p"
+r=$?
+[ "$f" = 0 ] && trap '' ALRM
+if [ "$f" != 0 ]; then
+  n=0
+  while [ "$n" -lt 50 ]; do
+    k=0
+    for x in /proc/[0-9]*/stat; do
+      s=
+      while read -r l; do s="$s $l"; done 2>/dev/null <"$x" || continue
+      s=${{s#*) }}
+      case $s in *")"*) s=${{s##*) }} ;; esac
+      set -- $s
+      [ "$4" = "$p" ] || continue
+      x=${{x#/proc/}}
+      kill -KILL "${{x%/stat}}" 2>/dev/null && k=1
+    done
+    [ "$k" = 0 ] && break
+    n=$((n + 1))
+  done
+  wait "$p" 2>/dev/null
+fi
+kill -KILL "$w" 2>/dev/null
+tail -c "$((c + 1))" "$d/o"
+tail -c "$((c + 1))" "$d/e" >&2
+rm -rf "$d"
+[ "$f" = 1 ] && exit {_TIMEOUT_STATUS}
+[ "$f" = 2 ] && exit 143
+exit "$r"
+"""
+# read_file: the base class's script with the content sent raw instead of as
+# base64, so a binary file arrives intact and the transfer is half the size.
+_READ_SCRIPT = f"""{_SYSTEM_PATH}
+sz=$(stat -Lc %s -- "$1" 2>/dev/null) || exit 
{SandboxBackend._MISSING_PATH_STATUS}
+[ -d "$1" ] && exit {SandboxBackend._IS_DIRECTORY_STATUS}
+printf '%s\\n' "$sz"
+exec head -c "$2" -- "$1"
+"""
+_WRITE_FIRST_SCRIPT = f'{_SYSTEM_PATH}\nmkdir -p -- "$(dirname -- "$1")" && 
cat >"$1"\n'
+_WRITE_NEXT_SCRIPT = f'{_SYSTEM_PATH}\ncat >>"$1"\n'
+
+
+def _new_openshell_name() -> str:
+    return f"{_NAME_PREFIX}{uuid.uuid4().hex[:11]}"
+
+
+def _status_name(error: BaseException) -> str | None:
+    code = getattr(error, "code", None)
+    if not callable(code):
+        return None
+    try:
+        status = code()
+    except Exception:
+        return None
+    return getattr(status, "name", None)
+
+
+def _describe_error(error: BaseException) -> str:
+    status = _status_name(error)
+    if status is None:
+        return f"{type(error).__name__}: {error}"
+    details = getattr(error, "details", None)
+    text = details() if callable(details) else ""
+    return f"{status}: {text}" if text else status
+
+
+@contextmanager
+def _translate_openshell_errors(
+    operation: str, *, recoverable_statuses: frozenset[str] = frozenset()
+) -> Iterator[None]:
+    try:
+        yield
+    except SandboxError:
+        raise
+    except ImportError as e:
+        raise SandboxTerminalError(
+            'The OpenShell SDK is not installed. Install 
"apache-airflow-providers-common-ai[openshell]".'
+        ) from e
+    except Exception as e:
+        message = f"OpenShell could not {operation} ({_describe_error(e)})."
+        if _status_name(e) in recoverable_statuses:
+            raise SandboxError(message) from e
+        raise SandboxTerminalError(message) from e
+
+
+def _is_hostname(value: object) -> bool:
+    """Whether ``value`` is a hostname OpenShell's policy can match, 
optionally with one leading ``*.``."""
+    if not isinstance(value, str) or not value or len(value) > 253:
+        return False
+    labels = value.removeprefix("*.").split(".")
+    # A single label ('localhost', or '*.com' once the wildcard is removed) is
+    # refused by the gateway for wildcards and cannot name a public endpoint.
+    if len(labels) < 2 or labels[-1].isdigit():
+        return False
+    return all(
+        0 < len(label) <= 63
+        and label[0] != "-"
+        and label[-1] != "-"
+        and all(ch.isascii() and (ch.isalnum() or ch == "-") for ch in label)
+        for label in labels
+    )
+
+
+class _BoundedTail:
+    """
+    The last ``max_bytes`` of a stream, and how many bytes the stream carried 
in total.
+
+    ``min_window`` keeps that many trailing bytes even under a smaller 
``max_bytes``,
+    for :meth:`ends_with` only; the text and the truncation flag still follow 
``max_bytes``.
+    """
+
+    def __init__(self, max_bytes: int, *, min_window: int = 0) -> None:
+        self._max_bytes = max_bytes
+        self._window = max(max_bytes, min_window)
+        self._data = bytearray()
+        self.received = 0
+
+    def add(self, chunk: bytes) -> None:
+        self.received += len(chunk)
+        self._data.extend(chunk)
+        if len(self._data) > self._window:
+            del self._data[: len(self._data) - self._window]
+
+    @property
+    def truncated(self) -> bool:
+        return self.received > self._max_bytes
+
+    def ends_with(self, suffix: bytes) -> bool:
+        return self._data.endswith(suffix)
+
+    def get_text(self) -> str:
+        data = bytes(self._data[max(0, len(self._data) - self._max_bytes) :])
+        if self.truncated:
+            # Drop the leading partial line so no fragment reads as a whole 
record, unless that
+            # would throw away most of the window: one line longer than the 
cap has its newline
+            # at the very end, which would leave "(no output)" for a command 
that wrote megabytes.
+            newline = data.find(b"\n")
+            if newline != -1 and len(data) - (newline + 1) >= self._max_bytes 
// 2:
+                data = data[newline + 1 :]
+        return data.decode("utf-8", errors="replace")
+
+
+class _BoundedHead:
+    """The first ``max_bytes`` of a stream; anything past that is dropped."""
+
+    def __init__(self, max_bytes: int) -> None:
+        self._max_bytes = max_bytes
+        self._data = bytearray()
+
+    def add(self, chunk: bytes) -> None:
+        room = self._max_bytes - len(self._data)
+        if room > 0:
+            self._data.extend(chunk[:room])
+
+    def get_bytes(self) -> bytes:
+        return bytes(self._data)
+
+
+@dataclass(frozen=True)
+class _ExecOutcome:
+    exit_code: int
+    stdout: _BoundedTail | _BoundedHead
+    stderr: _BoundedTail
+    elapsed: float
+
+
+class _ExecHung(Exception):
+    """The exec stream outlived its deadline without reporting an exit 
status."""
+
+
+class _GatewayDropped(Exception):
+    """The gateway went away after the exec request was sent; the command's 
outcome is unknown."""
+
+
+@dataclass(frozen=True)
+class _EgressPolicy:
+    """What an OpenShell policy lets a sandbox reach, in a form two policies 
can be compared by."""
+
+    destinations: frozenset[tuple[str, int, frozenset[str]]]
+    problems: tuple[str, ...]
+
+    @classmethod
+    def requested(cls, hosts: Sequence[str]) -> _EgressPolicy:
+        return cls(
+            destinations=frozenset((host, _HTTPS_PORT, 
frozenset({_ANY_BINARY})) for host in hosts),
+            problems=(),
+        )
+
+    @classmethod
+    def effective(cls, config: Any) -> _EgressPolicy:
+        """Read the egress a ``GetSandboxConfigResponse`` says is in force."""
+        from openshell._proto import sandbox_pb2
+
+        destinations: set[tuple[str, int, frozenset[str]]] = set()
+        problems: list[str] = []
+        if config.policy_source != sandbox_pb2.POLICY_SOURCE_SANDBOX:
+            problems.append(
+                "a gateway-wide policy is in force instead of the sandbox's 
own "
+                f"(policy_source 
{sandbox_pb2.PolicySource.Name(config.policy_source)})"
+            )
+        if not config.configuration_admitted:
+            problems.append(
+                f"the gateway has not admitted the sandbox's configuration 
({config.configuration_error or 'no reason given'})"
+            )
+        if config.policy.landlock.compatibility != "hard_requirement":
+            problems.append(
+                f"Landlock enforcement is 
{config.policy.landlock.compatibility or 'unset'!r}, not 'hard_requirement'"
+            )
+        if config.policy.network_middlewares:
+            problems.append(
+                f"network middlewares are configured: 
{sorted(config.policy.network_middlewares)}"
+            )
+        for rule in config.policy.network_policies.values():
+            binaries = frozenset(binary.path for binary in rule.binaries)
+            for endpoint in rule.endpoints:
+                if endpoint.allowed_ips:
+                    problems.append(
+                        f"{endpoint.host or 'a rule with no host'} admits 
addresses {list(endpoint.allowed_ips)}"
+                    )
+                ports = set(endpoint.ports)
+                if endpoint.port:
+                    ports.add(endpoint.port)
+                destinations.update((endpoint.host.lower(), port, binaries) 
for port in ports)
+        if _PROPOSAL_APPROVAL_MODE in config.settings:
+            setting = config.settings[_PROPOSAL_APPROVAL_MODE]
+            if setting.value.string_value == "auto":
+                problems.append(
+                    "proposal_approval_mode is 'auto' "
+                    f"({sandbox_pb2.SettingScope.Name(setting.scope)}), so a 
connection the sandbox "
+                    "is denied can become an allow rule without review"
+                )
+        if _AGENT_POLICY_PROPOSALS in config.settings:
+            if config.settings[_AGENT_POLICY_PROPOSALS].value.bool_value:
+                problems.append(
+                    "agent_policy_proposals_enabled is on, so the agent can 
propose its own egress"
+                )
+        return cls(destinations=frozenset(destinations), 
problems=tuple(problems))
+
+    def describe(self) -> list[str]:
+        return sorted(f"{host}:{port}" for host, port, _ in self.destinations)
+
+    def mismatch(self, wanted: _EgressPolicy) -> str | None:
+        """Why this policy is not ``wanted``, or ``None`` when it is exactly 
that."""
+        problems = list(self.problems)
+        if self.destinations != wanted.destinations:
+            problems.append(
+                f"it admits {self.describe() or 'nothing'} where 
{wanted.describe() or 'nothing'} was asked for"
+            )
+        return "; ".join(problems) or None
+
+
+class OpenShellSandboxBackend(SandboxBackend):
+    """
+    Run sandbox tools through an NVIDIA OpenShell gateway.
+
+    .. note::
+
+        Experimental: this can change or be removed in a minor release of this 
provider.
+        See :ref:`howto/stability`.
+
+    `OpenShell <https://github.com/NVIDIA/OpenShell>`__ is a self-hosted 
gateway
+    that runs each sandbox on Docker, Podman or Kubernetes with Landlock and
+    seccomp inside the container and a per-sandbox supervisor that mediates
+    every outbound connection. Airflow workers only need gRPC access to the
+    gateway, with the ``openshell`` extra installed.
+
+    **Credentials are ambient.** The backend reads the gateway registration the
+    ``openshell`` CLI keeps under 
``$XDG_CONFIG_HOME/openshell/gateways/<name>/``:
+    its endpoint and either mTLS material or an OIDC token. There is no Airflow
+    connection type for it; see the backend page for how to provision a worker.
+
+    **Egress is deny-all unless listed, and verified.** A spec's hostnames 
become
+    one policy rule admitting each host on port 443 only, and a default spec
+    sends no rule at all, which OpenShell enforces as no egress; name 
resolution
+    is answered by the supervisor without leaving the sandbox. After create the
+    backend reads the effective policy back and destroys the sandbox unless it
+    is exactly what was asked for, comes from the sandbox rather than a
+    gateway-wide override, and cannot be widened by auto-approved proposals.
+    The same check runs before and after every command, because a gateway
+    administrator, an approved draft or a global policy can widen a running
+    sandbox; a change fails the task. What it cannot see is the network the
+    gateway runs on: on Kubernetes the NetworkPolicy that keeps a sandbox pod
+    behind its supervisor is the cluster's to enforce.
+
+    **Commands run through a guest wrapper.** OpenShell's own exec timeout
+    leaves the command running, so the wrapper enforces the budget itself and
+    kills the processes in the command's session when it runs out. Output is
+    spooled to the sandbox's ``/tmp`` and each stream is capped to
+    ``max_output_bytes`` before it leaves the sandbox. 
``allow_egress_to_cidrs``,
+    an open network and ``SandboxSpec.owner`` are refused. The image needs
+    ``setsid`` and GNU coreutils and findutils; ``python:*-slim`` has them.
+
+    :param gateway: Name of the OpenShell CLI gateway registration to use.
+        ``None`` uses ``$OPENSHELL_GATEWAY``, then the CLI's active gateway.
+    :param workspace: OpenShell workspace the sandboxes are created in.
+    :param image: Container image used for each sandbox.
+    :param cpu: CPU limit, as a Kubernetes quantity.
+    :param memory: Memory limit, as a Kubernetes quantity.
+    :param ready_timeout: Seconds to wait for a new sandbox to become ready.
+    :param request_timeout: Seconds allowed for each gateway call other than a 
command.
+    """
+
+    name = "openshell"
+
+    def __init__(
+        self,
+        *,
+        gateway: str | None = None,
+        workspace: str = "default",
+        image: str = "python:3.12-slim",
+        cpu: str = "1",
+        memory: str = "2Gi",
+        ready_timeout: float = 120.0,
+        request_timeout: float = 30.0,
+    ) -> None:
+        if gateway is not None and not gateway:
+            raise ValueError("gateway must not be empty; pass None to use the 
CLI's active gateway.")
+        if not workspace:
+            raise ValueError("workspace must not be empty.")
+        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(ready_timeout, "ready_timeout")
+        _validate_positive_finite(request_timeout, "request_timeout")
+        self._gateway = gateway
+        self._workspace = workspace
+        self._image = image
+        self._resources = {"limits": {"cpu": cpu, "memory": memory}}
+        self._ready_timeout = ready_timeout
+        self._request_timeout = request_timeout
+        self._client: SandboxClient | None = None
+        # The egress each sandbox was created with, compared against the 
effective
+        # policy before and after every command.
+        self._egress: dict[str, _EgressPolicy] = {}
+
+    def _get_client(self) -> SandboxClient:
+        with _CLIENT_BUILD_LOCK:
+            if self._client is None:
+                with _translate_openshell_errors("load the gateway 
registration"):
+                    from openshell import SandboxClient
+
+                    self._client = SandboxClient.from_active_cluster(

Review Comment:
   `from_active_cluster` keeps the SDK defaults `auto_refresh=True, 
write_back=True`, so once an OIDC gateway's cached access token goes stale the 
SDK refreshes it and writes `oidc_token.json` back via 
`tempfile.mkstemp(dir=gateway_dir)`. On a worker whose registration is a 
Kubernetes Secret or ConfigMap mount (the provisioning the credentials section 
describes), that write raises `PermissionError` out of the auth interceptor and 
`_translate_openshell_errors` turns it into `SandboxTerminalError`, so every 
task fails from then on, retries included.
   
   `write_back=False` avoids the crash but breaks under refresh-token rotation 
once several workers share the same token. 
`client_credentials=ClientCredentialsAuth(...)` keeps tokens in memory, though 
it needs a secret from somewhere. If neither fits this PR, could the docs say 
OIDC needs a writable per-worker registration and that mTLS is the supported 
worker setup for now?



##########
providers/common/ai/tests/system/common/ai/example_sandbox_toolset_openshell.py:
##########
@@ -0,0 +1,114 @@
+# 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 NVIDIA OpenShell."""
+
+from __future__ import annotations
+
+import os
+from datetime import UTC, datetime
+
+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_openshell_{ENV_ID}" if ENV_ID else 
"common_ai_sandbox_toolset_openshell"
+
+MARKER = "boundary-ok"
+STATE_PATH = "/tmp/airflow_sandbox_e2e"
+# Counts processes whose command line is exactly "sleep 300", without needing 
ps in the image.
+COUNT_SLEEPERS = (
+    "n=0; for f in /proc/[0-9]*/cmdline; do "
+    "[ \"$(tr '\\0' ' ' < \"$f\" 2>/dev/null)\" = 'sleep 300 ' ] && n=$((n + 
1)); done; echo sleepers=$n"
+)
+CONNECT = (
+    'python3 -c "import socket\n'
+    "try:\n    socket.create_connection(('1.1.1.1', 443), 5)\n    
print('egress=open')\n"
+    "except OSError as e:\n    print('egress=denied', e.errno)\""
+)
+
+
+@airflow_dag(
+    dag_id=DAG_ID,
+    schedule="@once",
+    start_date=datetime(2024, 1, 1, tzinfo=UTC),
+    catchup=False,
+    tags=["common.ai", "sandbox", "openshell", "system_test"],
+)
+def example_sandbox_toolset_openshell():
+    @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 OpenShellSandboxBackend
+        from airflow.providers.common.ai.toolsets import SandboxToolset
+
+        # Each step is a tool call and a check of what the previous one 
returned.
+        steps: list[tuple[str, dict, str]] = [
+            ("write_file", {"path": STATE_PATH, "content": MARKER}, "Wrote"),
+            ("run_command", {"command": f"cat {STATE_PATH} && python3 -c 
'print(6 * 7)'"}, "42"),
+            ("read_file", {"path": STATE_PATH}, MARKER),
+            ("run_command", {"command": "echo to-stderr >&2; exit 3"}, "[exit 
code: 3]"),

Review Comment:
   This only checks `[exit code: 3]`, so lost stderr wouldn't fail it, and 
stderr takes the spool-then-tail path only this backend has. The OpenSandbox 
example also checks for `to-stderr` here, and for the marker in the `6 * 7` 
step.



##########
providers/common/ai/src/airflow/providers/common/ai/sandbox/openshell.py:
##########
@@ -0,0 +1,995 @@
+# 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.
+"""NVIDIA OpenShell backend for 
:class:`~airflow.providers.common.ai.toolsets.sandbox.SandboxToolset`."""
+
+from __future__ import annotations
+
+import logging
+import math
+import shlex
+import threading
+import time
+import uuid
+from contextlib import contextmanager, suppress
+from dataclasses import dataclass
+from typing import TYPE_CHECKING, Any
+
+from airflow.providers.common.ai.sandbox.base import (
+    _FILE_OP_OUTPUT_CAP,
+    _FILE_OP_TIMEOUT,
+    SandboxBackend,
+    SandboxError,
+    SandboxExecResult,
+    SandboxFileTooLargeError,
+    SandboxTerminalError,
+    _validate_positive_finite,
+)
+
+if TYPE_CHECKING:
+    from collections.abc import Iterator, Mapping, Sequence
+
+    from openshell import SandboxClient
+
+    from airflow.providers.common.ai.sandbox.base import SandboxSpec
+
+log = logging.getLogger(__name__)
+
+_CLIENT_BUILD_LOCK = threading.Lock()
+
+# Wall-clock allowance past the per-command budget before the exec stream is
+# treated as hung. The guest wrapper enforces the budget itself; this covers a
+# gateway or supervisor that stops delivering events at all.
+_EXEC_GRACE = 30.0
+# How long a command or policy read waits for a restarting gateway to come back
+# before the failure is reported. A restart measured at 6-11 s end to end.
+_GATEWAY_RECOVERY = 60.0
+_POLL_INTERVAL = 0.25
+# The gateway decodes at most 1 MiB per request, and the command or a file 
chunk
+# travels as stdin inside that request, so both stay under it with headroom.
+_MAX_STDIN_BYTES = 1_000_000
+_WRITE_CHUNK_BYTES = 768 * 1024
+# The gateway rejects a sandbox name longer than 19 characters, which
+# _new_sandbox_name's ``airflow-sandbox-<12 hex>`` is, so the name is shorter
+# here. Labels carry the attribution instead, and the creation time so an
+# operator can reap by age: OpenShell has no server-side lifetime.
+_NAME_PREFIX = "airflow-"
+_CREATED_BY_LABEL = ("created-by", "airflow")
+_CREATED_AT_LABEL = "airflow-created-at"
+_EGRESS_RULE = "airflow-egress"
+_HTTPS_PORT = 443
+_ANY_BINARY = "/**"
+# OpenShell's restrictive default filesystem policy, plus /dev/shm so Python's
+# multiprocessing and anything else using POSIX shared memory works.
+_READ_ONLY_PATHS = ("/bin", "/usr", "/lib", "/proc", "/dev/urandom", "/etc", 
"/var/log")
+_READ_WRITE_PATHS = ("/tmp", "/dev/null", "/dev/shm")
+# The sandbox supervisor removes the proxy variables from every command it runs
+# and overwrites the CA-bundle variables with its own TLS-terminating CA
+# (openshell-sandbox process.rs PROXY_ENV_VARS and child_env.rs tls_env_vars in
+# 0.1.2), in both cases without an error. A spec naming them would be silently
+# changed, so it is refused instead.
+_SUPERVISOR_OWNED_ENV = frozenset(
+    {
+        "ALL_PROXY",
+        "HTTP_PROXY",
+        "HTTPS_PROXY",
+        "NO_PROXY",
+        "all_proxy",
+        "http_proxy",
+        "https_proxy",
+        "no_proxy",
+        "grpc_proxy",
+        "NODE_USE_ENV_PROXY",
+        "SSL_CERT_FILE",
+        "REQUESTS_CA_BUNDLE",
+        "CURL_CA_BUNDLE",
+        "GIT_SSL_CAINFO",
+        "NODE_EXTRA_CA_CERTS",
+        "DENO_CERT",
+    }
+)
+_PROPOSAL_APPROVAL_MODE = "proposal_approval_mode"
+_AGENT_POLICY_PROPOSALS = "agent_policy_proposals_enabled"
+# Wrapper-private exit status for "the budget ran out". The wrapper only uses 
it
+# after its own timer fired, and the caller also checks the elapsed time, so a
+# command that exits 124 on its own inside the budget is not a timeout.
+_TIMEOUT_STATUS = 124
+# Wrapper-private exit status and stderr line for "the command could not be
+# staged in /tmp, so it did not run". The status alone is not enough, because a
+# command can exit 125 itself, as docker run, env and nohup do on their own 
errors.
+_STAGING_STATUS = 125
+_STAGING_FAILED = "airflow-exec: could not stage the command in /tmp"
+_STAGING_LINE = f"{_STAGING_FAILED}\n".encode()
+_SYSTEM_PATH = 
"PATH=/usr/local/sbin:/usr/local/bin:/usr/sbin:/usr/bin:/sbin:/bin"
+
+# Runs one command for run_command. OpenShell's own exec timeout returns a
+# synthetic 124 and leaves the process running, and a backgrounded child that
+# inherits stdout holds the call open until the supervisor gives up on it, so
+# the wrapper owns both problems:
+#
+# * The command arrives on stdin, which avoids the gateway's 32 KiB 
per-argument
+#   cap, and runs in its own session via setsid, with stdout and stderr going 
to
+#   files, so a background child holds a file rather than the exec stream.
+# * A timer reading the monotonic /proc/uptime signals the wrapper when the
+#   budget is spent; counting its own one-second sleeps instead drifts late on 
a
+#   CPU-starved sandbox, past _EXEC_GRACE. The wrapper then SIGKILLs the
+#   processes in the command's session, rescanning /proc until a pass kills
+#   nothing, at most 50 times, so children forked mid-sweep are usually caught
+#   too. A negative-pid kill of the group is blocked by the sandbox's seccomp
+#   filter. The sweep shares the CPU with the processes it is killing, so it
+#   reads /proc with builtins, joining the lines of a stat file whose process
+#   name holds a newline, and strips the name with the cheap shortest match,
+#   taking the longest only for a name that contains ")": a `cat` per entry let
+#   one sweep of 128 busy processes on one CPU run past _EXEC_GRACE.
+# * The exit status is the wrapper's own process status, which the supervisor
+#   reports out of band, so nothing the command prints can change it; a command
+#   that signals the wrapper only ends its own run early. Output is sent last 
as
+#   one byte more than the caller's cap, so an over-cap stream is detected by
+#   counting what arrives rather than trusting a trailer.
+#
+# A command that starts its own session (setsid, a daemonizing server) leaves
+# this one and is not killed on timeout, and one that keeps forking faster than
+# the sweep can outlast it.
+_RUN_WRAPPER = rf"""t=$1 c=$2 o=$PATH
+{_SYSTEM_PATH}
+d=$(mktemp -d /tmp/.airflow-exec.XXXXXX) && cat >"$d/c" || {{ rm -rf "$d"; 
echo "{_STAGING_FAILED}" >&2; exit {_STAGING_STATUS}; }}
+f=0
+trap f=1 ALRM
+trap f=2 HUP INT TERM
+z=$(command -v setsid) || {{ echo "setsid is not installed in the sandbox 
image" >&2; rm -rf "$d"; exit 125; }}
+PATH=$o "$z" /bin/sh "$d/c" </dev/null >"$d/o" 2>"$d/e" &
+p=$!
+(
+  read u _ </proc/uptime
+  e=$((${{u%.*}}${{u#*.}} + t * 100))
+  while [ "${{u%.*}}${{u#*.}}" -lt "$e" ]; do
+    sleep 1
+    [ -d "/proc/$$" ] || exit 0
+    read u _ </proc/uptime
+  done
+  kill -ALRM $$
+) </dev/null >/dev/null 2>&1 &
+w=$!
+wait "$p"
+r=$?
+[ "$f" = 0 ] && trap '' ALRM
+if [ "$f" != 0 ]; then
+  n=0
+  while [ "$n" -lt 50 ]; do
+    k=0
+    for x in /proc/[0-9]*/stat; do
+      s=
+      while read -r l; do s="$s $l"; done 2>/dev/null <"$x" || continue
+      s=${{s#*) }}
+      case $s in *")"*) s=${{s##*) }} ;; esac
+      set -- $s
+      [ "$4" = "$p" ] || continue
+      x=${{x#/proc/}}
+      kill -KILL "${{x%/stat}}" 2>/dev/null && k=1
+    done
+    [ "$k" = 0 ] && break
+    n=$((n + 1))
+  done
+  wait "$p" 2>/dev/null
+fi
+kill -KILL "$w" 2>/dev/null
+tail -c "$((c + 1))" "$d/o"
+tail -c "$((c + 1))" "$d/e" >&2
+rm -rf "$d"
+[ "$f" = 1 ] && exit {_TIMEOUT_STATUS}
+[ "$f" = 2 ] && exit 143
+exit "$r"
+"""
+# read_file: the base class's script with the content sent raw instead of as
+# base64, so a binary file arrives intact and the transfer is half the size.
+_READ_SCRIPT = f"""{_SYSTEM_PATH}
+sz=$(stat -Lc %s -- "$1" 2>/dev/null) || exit 
{SandboxBackend._MISSING_PATH_STATUS}
+[ -d "$1" ] && exit {SandboxBackend._IS_DIRECTORY_STATUS}
+printf '%s\\n' "$sz"
+exec head -c "$2" -- "$1"
+"""
+_WRITE_FIRST_SCRIPT = f'{_SYSTEM_PATH}\nmkdir -p -- "$(dirname -- "$1")" && 
cat >"$1"\n'
+_WRITE_NEXT_SCRIPT = f'{_SYSTEM_PATH}\ncat >>"$1"\n'
+
+
+def _new_openshell_name() -> str:
+    return f"{_NAME_PREFIX}{uuid.uuid4().hex[:11]}"
+
+
+def _status_name(error: BaseException) -> str | None:
+    code = getattr(error, "code", None)
+    if not callable(code):
+        return None
+    try:
+        status = code()
+    except Exception:
+        return None
+    return getattr(status, "name", None)
+
+
+def _describe_error(error: BaseException) -> str:
+    status = _status_name(error)
+    if status is None:
+        return f"{type(error).__name__}: {error}"
+    details = getattr(error, "details", None)
+    text = details() if callable(details) else ""
+    return f"{status}: {text}" if text else status
+
+
+@contextmanager
+def _translate_openshell_errors(
+    operation: str, *, recoverable_statuses: frozenset[str] = frozenset()
+) -> Iterator[None]:
+    try:
+        yield
+    except SandboxError:
+        raise
+    except ImportError as e:
+        raise SandboxTerminalError(
+            'The OpenShell SDK is not installed. Install 
"apache-airflow-providers-common-ai[openshell]".'
+        ) from e
+    except Exception as e:
+        message = f"OpenShell could not {operation} ({_describe_error(e)})."
+        if _status_name(e) in recoverable_statuses:
+            raise SandboxError(message) from e
+        raise SandboxTerminalError(message) from e
+
+
+def _is_hostname(value: object) -> bool:
+    """Whether ``value`` is a hostname OpenShell's policy can match, 
optionally with one leading ``*.``."""
+    if not isinstance(value, str) or not value or len(value) > 253:
+        return False
+    labels = value.removeprefix("*.").split(".")
+    # A single label ('localhost', or '*.com' once the wildcard is removed) is
+    # refused by the gateway for wildcards and cannot name a public endpoint.
+    if len(labels) < 2 or labels[-1].isdigit():
+        return False
+    return all(
+        0 < len(label) <= 63
+        and label[0] != "-"
+        and label[-1] != "-"
+        and all(ch.isascii() and (ch.isalnum() or ch == "-") for ch in label)
+        for label in labels
+    )
+
+
+class _BoundedTail:
+    """
+    The last ``max_bytes`` of a stream, and how many bytes the stream carried 
in total.
+
+    ``min_window`` keeps that many trailing bytes even under a smaller 
``max_bytes``,
+    for :meth:`ends_with` only; the text and the truncation flag still follow 
``max_bytes``.
+    """
+
+    def __init__(self, max_bytes: int, *, min_window: int = 0) -> None:
+        self._max_bytes = max_bytes
+        self._window = max(max_bytes, min_window)
+        self._data = bytearray()
+        self.received = 0
+
+    def add(self, chunk: bytes) -> None:
+        self.received += len(chunk)
+        self._data.extend(chunk)
+        if len(self._data) > self._window:
+            del self._data[: len(self._data) - self._window]
+
+    @property
+    def truncated(self) -> bool:
+        return self.received > self._max_bytes
+
+    def ends_with(self, suffix: bytes) -> bool:
+        return self._data.endswith(suffix)
+
+    def get_text(self) -> str:
+        data = bytes(self._data[max(0, len(self._data) - self._max_bytes) :])
+        if self.truncated:
+            # Drop the leading partial line so no fragment reads as a whole 
record, unless that
+            # would throw away most of the window: one line longer than the 
cap has its newline
+            # at the very end, which would leave "(no output)" for a command 
that wrote megabytes.
+            newline = data.find(b"\n")
+            if newline != -1 and len(data) - (newline + 1) >= self._max_bytes 
// 2:
+                data = data[newline + 1 :]
+        return data.decode("utf-8", errors="replace")
+
+
+class _BoundedHead:
+    """The first ``max_bytes`` of a stream; anything past that is dropped."""
+
+    def __init__(self, max_bytes: int) -> None:
+        self._max_bytes = max_bytes
+        self._data = bytearray()
+
+    def add(self, chunk: bytes) -> None:
+        room = self._max_bytes - len(self._data)
+        if room > 0:
+            self._data.extend(chunk[:room])
+
+    def get_bytes(self) -> bytes:
+        return bytes(self._data)
+
+
+@dataclass(frozen=True)
+class _ExecOutcome:
+    exit_code: int
+    stdout: _BoundedTail | _BoundedHead
+    stderr: _BoundedTail
+    elapsed: float
+
+
+class _ExecHung(Exception):
+    """The exec stream outlived its deadline without reporting an exit 
status."""
+
+
+class _GatewayDropped(Exception):
+    """The gateway went away after the exec request was sent; the command's 
outcome is unknown."""
+
+
+@dataclass(frozen=True)
+class _EgressPolicy:
+    """What an OpenShell policy lets a sandbox reach, in a form two policies 
can be compared by."""
+
+    destinations: frozenset[tuple[str, int, frozenset[str]]]
+    problems: tuple[str, ...]
+
+    @classmethod
+    def requested(cls, hosts: Sequence[str]) -> _EgressPolicy:
+        return cls(
+            destinations=frozenset((host, _HTTPS_PORT, 
frozenset({_ANY_BINARY})) for host in hosts),
+            problems=(),
+        )
+
+    @classmethod
+    def effective(cls, config: Any) -> _EgressPolicy:

Review Comment:
   The SDK ships `.pyi` stubs for `openshell._proto`, so `config` could be 
typed as `sandbox_pb2.GetSandboxConfigResponse` under `TYPE_CHECKING`. The 
fail-closed check reads its fields, and with `Any` mypy won't flag a renamed or 
misspelled one. `_build_spec`, `_get_sandbox` and `_workspace_scope` take `Any` 
too.



##########
providers/common/ai/tests/unit/common/ai/sandbox/test_openshell.py:
##########
@@ -0,0 +1,1002 @@
+# 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 copy
+import os
+import shlex
+import shutil
+import signal
+import subprocess
+import sys
+import time
+from unittest import mock
+
+import pytest
+
+pytest.importorskip("openshell")
+
+import grpc
+from openshell import SandboxClient
+from openshell._proto import openshell_pb2, sandbox_pb2
+
+from airflow.providers.common.ai.sandbox.base import (
+    SandboxError,
+    SandboxFileTooLargeError,
+    SandboxSpec,
+    SandboxTerminalError,
+)
+from airflow.providers.common.ai.sandbox.openshell import (
+    _EXEC_GRACE,
+    _READ_WRITE_PATHS,
+    _RUN_WRAPPER,
+    _STAGING_FAILED,
+    OpenShellSandboxBackend,
+)
+
+_MODULE = "airflow.providers.common.ai.sandbox.openshell"
+_MONOTONIC_PATH = f"{_MODULE}.time.monotonic"
+
+
+class _RpcError(grpc.RpcError):
+    def __init__(self, code: grpc.StatusCode, details: str = "") -> None:
+        super().__init__(details)
+        self._code = code
+        self._details = details
+
+    def code(self) -> grpc.StatusCode:
+        return self._code
+
+    def details(self) -> str:
+        return self._details
+
+
+class _Stream:
+    """An ExecSandbox response stream: events, optionally ending in an 
error."""
+
+    def __init__(self, events, error: Exception | None = None) -> None:
+        self._events = iter(events)
+        self._error = error
+        self.cancel = mock.Mock()
+
+    def __iter__(self):
+        return self
+
+    def __next__(self):
+        try:
+            return next(self._events)
+        except StopIteration:
+            if self._error is not None:
+                raise self._error from None
+            raise
+
+
+def _stdout(data: bytes):
+    return 
openshell_pb2.ExecSandboxEvent(stdout=openshell_pb2.ExecSandboxStdout(data=data))
+
+
+def _stderr(data: bytes):
+    return 
openshell_pb2.ExecSandboxEvent(stderr=openshell_pb2.ExecSandboxStderr(data=data))
+
+
+def _exit(code: int):
+    return 
openshell_pb2.ExecSandboxEvent(exit=openshell_pb2.ExecSandboxExit(exit_code=code))
+
+
+def _result(code: int = 0, out: bytes = b"", err: bytes = b"") -> _Stream:
+    events = []
+    if out:
+        events.append(_stdout(out))
+    if err:
+        events.append(_stderr(err))
+    events.append(_exit(code))
+    return _Stream(events)
+
+
+def _sandbox(phase: openshell_pb2.SandboxPhase = 
openshell_pb2.SANDBOX_PHASE_READY, *conditions):
+    response = openshell_pb2.SandboxResponse()
+    response.sandbox.status.phase = phase
+    for reason, message in conditions:
+        response.sandbox.status.conditions.add(type="Ready", status="False", 
reason=reason, message=message)
+    return response
+
+
+def _config(
+    hosts=(),
+    *,
+    source: sandbox_pb2.PolicySource = sandbox_pb2.POLICY_SOURCE_SANDBOX,
+    admitted: bool = True,
+    landlock: str = "hard_requirement",
+    approval_mode: str | None = None,
+    proposals: bool = False,
+    extra_rule: tuple[str, int] | None = None,
+    allowed_ips=(),
+    binaries=("/**",),
+    middlewares=(),
+):
+    config = sandbox_pb2.GetSandboxConfigResponse(
+        version=1, policy_source=source, configuration_admitted=admitted
+    )
+    config.policy.landlock.compatibility = landlock
+    for name in middlewares:
+        # Reading a missing key of a message map adds it.
+        config.policy.network_middlewares[name]
+    if hosts or allowed_ips:
+        rule = config.policy.network_policies["airflow-egress"]
+        rule.name = "airflow-egress"
+        for host in hosts:
+            rule.endpoints.add(host=host, ports=[443], 
allowed_ips=list(allowed_ips))
+        if not hosts:
+            rule.endpoints.add(ports=[443], allowed_ips=list(allowed_ips))
+        rule.binaries.extend(sandbox_pb2.NetworkBinary(path=path) for path in 
binaries)
+    if extra_rule is not None:
+        rule = config.policy.network_policies["allow_extra"]
+        rule.endpoints.add(host=extra_rule[0], port=extra_rule[1])
+        rule.binaries.add(path="/usr/local/bin/python3.12")
+    if approval_mode is not None:
+        config.settings["proposal_approval_mode"].value.string_value = 
approval_mode
+        config.settings["proposal_approval_mode"].scope = 
sandbox_pb2.SETTING_SCOPE_GLOBAL
+    if proposals:
+        config.settings["agent_policy_proposals_enabled"].value.bool_value = 
True
+    return config
+
+
+def _backend(**kwargs) -> tuple[OpenShellSandboxBackend, mock.MagicMock]:
+    backend = OpenShellSandboxBackend(gateway="test-gateway", **kwargs)
+    client = mock.MagicMock(spec=SandboxClient)
+    client._stub = mock.MagicMock(spec=["GetSandbox", "GetSandboxConfig", 
"ExecSandbox"])
+    client._stub.GetSandbox.return_value = _sandbox()
+    client._stub.GetSandboxConfig.return_value = _config()
+    client._stub.ExecSandbox.return_value = _result()
+    backend._client = client
+    return backend, client
+
+
+def _exec_request(client, call: int = -1):
+    return client._stub.ExecSandbox.call_args_list[call].args[0]
+
+
[email protected](autouse=True)
+def _no_sleep():
+    with mock.patch(f"{_MODULE}.time.sleep", autospec=True):
+        yield
+
+
+def test_missing_sdk_error_is_actionable():
+    real_import = builtins.__import__
+
+    def blocked_import(name, *args, **kwargs):
+        if name.startswith("openshell"):
+            raise ImportError("blocked for test")
+        return real_import(name, *args, **kwargs)
+
+    backend = OpenShellSandboxBackend()
+    with mock.patch("builtins.__import__", side_effect=blocked_import):
+        with pytest.raises(SandboxTerminalError, match=r"\[openshell\]"):
+            backend.create()
+
+
[email protected](
+    ("kwargs", "message"),
+    [
+        ({"gateway": ""}, "gateway"),
+        ({"workspace": ""}, "workspace"),
+        ({"image": ""}, "image"),
+        ({"cpu": ""}, "cpu"),
+        ({"memory": ""}, "memory"),
+        ({"ready_timeout": 0}, "ready_timeout"),
+        ({"request_timeout": float("inf")}, "request_timeout"),
+    ],
+)
+def test_constructor_rejects_invalid_values(kwargs, message):
+    with pytest.raises(ValueError, match=message):
+        OpenShellSandboxBackend(**kwargs)
+
+
+class TestClient:
+    @mock.patch("openshell.SandboxClient.from_active_cluster", autospec=True)
+    def test_construction_reads_no_gateway_registration(self, 
from_active_cluster):
+        OpenShellSandboxBackend(gateway="prod")
+
+        from_active_cluster.assert_not_called()
+
+    def test_the_backend_can_be_deep_copied(self):
+        backend = OpenShellSandboxBackend(gateway="prod")
+
+        assert copy.deepcopy(backend)._gateway == "prod"
+
+    @mock.patch("openshell.SandboxClient.from_active_cluster", autospec=True)
+    def test_the_cli_gateway_registration_is_loaded_once(self, 
from_active_cluster):
+        backend = OpenShellSandboxBackend(gateway="prod", request_timeout=12)
+
+        assert backend._get_client() is backend._get_client()
+
+        from_active_cluster.assert_called_once_with(cluster="prod", timeout=12)
+
+    @mock.patch("openshell.SandboxClient.from_active_cluster", autospec=True)
+    def test_a_missing_registration_is_terminal(self, from_active_cluster):
+        from openshell import SandboxError as SdkError
+
+        from_active_cluster.side_effect = SdkError("gateway 'prod' not found")
+
+        with pytest.raises(SandboxTerminalError, match="gateway 'prod' not 
found"):
+            OpenShellSandboxBackend(gateway="prod")._get_client()
+
+
+class TestSpecRefusals:
+    @pytest.mark.parametrize(
+        ("spec", "message"),
+        [
+            (SandboxSpec(owner="dag/run"), "owner"),
+            (SandboxSpec(allow_egress_to_cidrs=["203.0.113.0/24"]), 
"allow_egress_to_cidrs"),
+            (SandboxSpec(block_network=False), "open network"),
+            (SandboxSpec(block_network=False, allow_egress_to=["pypi.org"]), 
"open network"),
+            (SandboxSpec(allow_egress_to="pypi.org"), "not one string"),
+            (SandboxSpec(env={"FOO": 1}), "strings to strings"),  # type: 
ignore[dict-item]
+        ],
+    )
+    def test_unenforceable_specs_are_refused_before_the_gateway_is_asked(self, 
spec, message):
+        backend, client = _backend()
+
+        with pytest.raises(SandboxTerminalError, match=message):
+            backend.create(spec=spec)
+
+        client.create.assert_not_called()
+
+    @pytest.mark.parametrize(
+        "host",
+        ["https://pypi.org";, "pypi.org:443", "203.0.113.7", "localhost", 
"*.com", "*", "pypi.org/simple", ""],
+    )
+    def test_entries_that_are_not_hostnames_are_refused(self, host):
+        backend, client = _backend()
+
+        with pytest.raises(SandboxTerminalError, match="bare hostnames"):
+            backend.create(spec=SandboxSpec(allow_egress_to=[host]))
+
+        client.create.assert_not_called()
+
+    @pytest.mark.parametrize(
+        "key", ["HTTPS_PROXY", "no_proxy", "SSL_CERT_FILE", 
"REQUESTS_CA_BUNDLE", "OPENSHELL_SANDBOX"]
+    )
+    def test_env_the_supervisor_would_silently_change_is_refused(self, key):
+        backend, client = _backend()
+
+        with pytest.raises(SandboxTerminalError, match=f"sets {key}"):
+            backend.create(spec=SandboxSpec(env={key: "value"}))
+
+        client.create.assert_not_called()
+
+
+class TestCreate:
+    def test_default_spec_is_deny_all_under_hard_landlock(self):
+        backend, client = _backend(image="python:3.13-slim", cpu="2", 
memory="4Gi")
+
+        with mock.patch(f"{_MODULE}.time.time", return_value=1_700_000_000.4):
+            name = backend.create(spec=SandboxSpec(env={"TOKEN": "value"}))
+
+        kwargs = client.create.call_args.kwargs
+        assert kwargs["name"] == name
+        assert name.startswith("airflow-")
+        assert len(name) <= 19
+        assert kwargs["workspace"] == "default"
+        assert kwargs["labels"] == {"created-by": "airflow", 
"airflow-created-at": "1700000000"}
+        spec = kwargs["spec"]
+        assert dict(spec.environment) == {"TOKEN": "value"}
+        assert spec.template.image == "python:3.13-slim"
+        assert dict(spec.template.resources)["limits"] == {"cpu": "2", 
"memory": "4Gi"}
+        assert len(spec.policy.network_policies) == 0
+        assert spec.policy.landlock.compatibility == "hard_requirement"
+        assert tuple(spec.policy.filesystem.read_write) == _READ_WRITE_PATHS

Review Comment:
   This compares `_READ_WRITE_PATHS` with itself, and `read_only` and 
`include_workdir` aren't asserted anywhere, so widening the filesystem policy 
(adding `/` to `_READ_ONLY_PATHS`, say) keeps the suite green. Spelling the 
expected paths out as literals here would pin it.



##########
providers/common/ai/docs/sandbox/backends.rst:
##########
@@ -248,6 +249,156 @@ 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.
 
+.. _sandbox-backend-openshell:
+
+OpenShell (self-hosted remote)
+------------------------------
+
+:class:`~airflow.providers.common.ai.sandbox.openshell.OpenShellSandboxBackend`
+runs sandboxes through an `NVIDIA OpenShell 
<https://github.com/NVIDIA/OpenShell>`__
+gateway, which runs each one on Docker, Podman or Kubernetes. Inside the
+container the workload is confined by Landlock and seccomp and has no network
+interface of its own; a per-sandbox supervisor opens every outbound connection
+on its behalf, against a policy the backend writes and reads back. Airflow 
workers
+only need gRPC access to the gateway.
+
+Install the SDK extra:
+
+.. code-block:: bash
+
+    pip install "apache-airflow-providers-common-ai[openshell]"
+
+**Credentials are ambient; there is no Airflow connection for this backend.**
+The gateway's endpoint, mTLS material and OIDC token live in the gateway
+registration that the ``openshell`` CLI keeps under
+``$XDG_CONFIG_HOME/openshell/gateways/<name>/`` (``~/.config`` by default), and
+the OpenShell SDK reads them from there itself: ``metadata.json`` with the
+endpoint, and ``mtls/ca.crt``, ``mtls/tls.crt`` and ``mtls/tls.key`` for mTLS,
+or the CLI's cached token for an OIDC gateway. The backend only names the
+registration. A worker without the CLI needs the same files, provisioned by the
+Deployment Manager:
+
+.. code-block:: json
+
+    {"name": "prod", "gateway_endpoint": 
"https://openshell.example.com:17670";, "auth_mode": "mtls"}
+
+.. code-block:: python
+
+    from airflow.providers.common.ai.sandbox import OpenShellSandboxBackend
+    from airflow.providers.common.ai.toolsets import SandboxToolset
+
+    SandboxToolset(OpenShellSandboxBackend(gateway="prod"))
+
+Whoever holds that credential is inside the trust boundary. On a gateway
+without OIDC, an mTLS client is a gateway-wide administrator: it can change the
+policy and settings of every sandbox. ``openshell-gateway generate-certs`` 
issues
+a single client identity whose certificate does not expire for practical
+purposes, so use your own PKI, with expiry and rotation, for a worker 
credential.
+
+**Network policy.** A default ``SandboxSpec()`` sends a policy with no egress
+rule, which OpenShell enforces as no egress: a connection fails with
+``EACCES``, and a name lookup is answered by the supervisor with a synthetic
+``198.18.0.0/15`` address rather than failing, so no query leaves the sandbox
+although the tool description's "including DNS" wording reads stricter than the
+lookup behaves. ``allow_egress_to`` becomes one rule admitting each listed host
+on port 443 for any program in the sandbox; ports are part of every rule, so
+plain HTTP and other ports stay closed. HTTPS to a listed host is terminated by
+the supervisor, which injects its CA through ``SSL_CERT_FILE``,
+``REQUESTS_CA_BUNDLE`` and similar variables; a client with its own trust store
+(a JVM, a statically linked binary) has to be pointed at it. A leading ``*.``
+label is accepted on a name of three labels or more.
+
+The backend verifies the policy rather than trusting the request. After create
+it reads the effective policy back and destroys the sandbox unless it admits
+exactly the requested hosts, comes from the sandbox rather than a gateway-wide
+policy, keeps Landlock at ``hard_requirement``, and has neither
+``proposal_approval_mode=auto`` nor agent policy proposals enabled at any 
scope.
+Auto-approval is refused because OpenShell turns a denied connection into a
+proposed allow rule and, in that mode, approves it without review. The same
+check runs before and after every command, because an administrator, an 
approved
+draft or a gateway-wide policy can widen a running sandbox; a change destroys 
the
+sandbox and fails the task. That detects a widening, it does not prevent one: a
+command already running when the policy changes can use it. The check sees
+OpenShell's policy, not the network under the gateway: on Kubernetes, a sandbox
+pod is kept behind its supervisor by a NetworkPolicy, which only a CNI that
+enforces NetworkPolicy upholds. Every denied connection also becomes a draft
+proposal in the gateway's approval inbox, even in manual mode, so operators see
+which destinations an agent tried.
+
+Refused at create: ``block_network=False`` (every OpenShell rule names a host 
and
+its ports, so an open network cannot be expressed), ``allow_egress_to_cidrs``
+(address rules admit TCP on listed ports only), ``SandboxSpec.owner``, and
+``SandboxSpec.env`` keys the supervisor owns: it removes ``HTTP_PROXY`` and the
+other proxy variables from every command and replaces ``SSL_CERT_FILE``,
+``REQUESTS_CA_BUNDLE``, ``CURL_CA_BUNDLE``, ``GIT_SSL_CAINFO``,
+``NODE_EXTRA_CA_CERTS`` and ``DENO_CERT`` with its own, without an error, and
+reserves ``OPENSHELL_*``. Everything else in ``SandboxSpec.env``, ``PATH``
+included, reaches commands as given, since they run without a login shell.
+
+**Commands.** OpenShell's own exec timeout reports exit 124 and leaves the
+command running, so the backend runs each command through a small shell wrapper
+in the sandbox. The command arrives on stdin and runs in a session of its own,
+with the command and its output spooled to ``/tmp``; one that cannot be written
+there, because ``/tmp`` is full for example, is not run and is reported as an
+error. When the budget runs out the wrapper kills the processes in that 
session,
+scanning again until a pass finds none to kill, at most 50 times, and returns 
the
+tail of each stream. A process the command left in the background keeps running
+after a command that finishes in time, and one that started a session of its 
own
+(``setsid``, a daemonizing server) escapes the kill on timeout, as can a 
command
+that keeps forking faster than the sweep. If the gateway stops relaying the
+command for longer than its budget plus 30 seconds, the sandbox is destroyed 
and

Review Comment:
   `_abandon_command` only logs when the delete fails, and the delete goes to 
the same gateway that just stopped relaying, so it's the call most likely to 
fail at this point. Maybe "the backend asks the gateway to delete the sandbox 
and logs a failure for the reaper below"? The portability note at line 472 says 
the same.



##########
providers/common/ai/src/airflow/providers/common/ai/sandbox/openshell.py:
##########
@@ -0,0 +1,995 @@
+# 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.
+"""NVIDIA OpenShell backend for 
:class:`~airflow.providers.common.ai.toolsets.sandbox.SandboxToolset`."""
+
+from __future__ import annotations
+
+import logging
+import math
+import shlex
+import threading
+import time
+import uuid
+from contextlib import contextmanager, suppress
+from dataclasses import dataclass
+from typing import TYPE_CHECKING, Any
+
+from airflow.providers.common.ai.sandbox.base import (
+    _FILE_OP_OUTPUT_CAP,
+    _FILE_OP_TIMEOUT,
+    SandboxBackend,
+    SandboxError,
+    SandboxExecResult,
+    SandboxFileTooLargeError,
+    SandboxTerminalError,
+    _validate_positive_finite,
+)
+
+if TYPE_CHECKING:
+    from collections.abc import Iterator, Mapping, Sequence
+
+    from openshell import SandboxClient
+
+    from airflow.providers.common.ai.sandbox.base import SandboxSpec
+
+log = logging.getLogger(__name__)
+
+_CLIENT_BUILD_LOCK = threading.Lock()
+
+# Wall-clock allowance past the per-command budget before the exec stream is
+# treated as hung. The guest wrapper enforces the budget itself; this covers a
+# gateway or supervisor that stops delivering events at all.
+_EXEC_GRACE = 30.0
+# How long a command or policy read waits for a restarting gateway to come back
+# before the failure is reported. A restart measured at 6-11 s end to end.
+_GATEWAY_RECOVERY = 60.0
+_POLL_INTERVAL = 0.25
+# The gateway decodes at most 1 MiB per request, and the command or a file 
chunk
+# travels as stdin inside that request, so both stay under it with headroom.
+_MAX_STDIN_BYTES = 1_000_000
+_WRITE_CHUNK_BYTES = 768 * 1024
+# The gateway rejects a sandbox name longer than 19 characters, which
+# _new_sandbox_name's ``airflow-sandbox-<12 hex>`` is, so the name is shorter
+# here. Labels carry the attribution instead, and the creation time so an
+# operator can reap by age: OpenShell has no server-side lifetime.
+_NAME_PREFIX = "airflow-"
+_CREATED_BY_LABEL = ("created-by", "airflow")
+_CREATED_AT_LABEL = "airflow-created-at"
+_EGRESS_RULE = "airflow-egress"
+_HTTPS_PORT = 443
+_ANY_BINARY = "/**"
+# OpenShell's restrictive default filesystem policy, plus /dev/shm so Python's
+# multiprocessing and anything else using POSIX shared memory works.
+_READ_ONLY_PATHS = ("/bin", "/usr", "/lib", "/proc", "/dev/urandom", "/etc", 
"/var/log")
+_READ_WRITE_PATHS = ("/tmp", "/dev/null", "/dev/shm")
+# The sandbox supervisor removes the proxy variables from every command it runs
+# and overwrites the CA-bundle variables with its own TLS-terminating CA
+# (openshell-sandbox process.rs PROXY_ENV_VARS and child_env.rs tls_env_vars in
+# 0.1.2), in both cases without an error. A spec naming them would be silently
+# changed, so it is refused instead.
+_SUPERVISOR_OWNED_ENV = frozenset(
+    {
+        "ALL_PROXY",
+        "HTTP_PROXY",
+        "HTTPS_PROXY",
+        "NO_PROXY",
+        "all_proxy",
+        "http_proxy",
+        "https_proxy",
+        "no_proxy",
+        "grpc_proxy",
+        "NODE_USE_ENV_PROXY",
+        "SSL_CERT_FILE",
+        "REQUESTS_CA_BUNDLE",
+        "CURL_CA_BUNDLE",
+        "GIT_SSL_CAINFO",
+        "NODE_EXTRA_CA_CERTS",
+        "DENO_CERT",
+    }
+)
+_PROPOSAL_APPROVAL_MODE = "proposal_approval_mode"
+_AGENT_POLICY_PROPOSALS = "agent_policy_proposals_enabled"
+# Wrapper-private exit status for "the budget ran out". The wrapper only uses 
it
+# after its own timer fired, and the caller also checks the elapsed time, so a
+# command that exits 124 on its own inside the budget is not a timeout.
+_TIMEOUT_STATUS = 124
+# Wrapper-private exit status and stderr line for "the command could not be
+# staged in /tmp, so it did not run". The status alone is not enough, because a
+# command can exit 125 itself, as docker run, env and nohup do on their own 
errors.
+_STAGING_STATUS = 125
+_STAGING_FAILED = "airflow-exec: could not stage the command in /tmp"
+_STAGING_LINE = f"{_STAGING_FAILED}\n".encode()
+_SYSTEM_PATH = 
"PATH=/usr/local/sbin:/usr/local/bin:/usr/sbin:/usr/bin:/sbin:/bin"
+
+# Runs one command for run_command. OpenShell's own exec timeout returns a
+# synthetic 124 and leaves the process running, and a backgrounded child that
+# inherits stdout holds the call open until the supervisor gives up on it, so
+# the wrapper owns both problems:
+#
+# * The command arrives on stdin, which avoids the gateway's 32 KiB 
per-argument
+#   cap, and runs in its own session via setsid, with stdout and stderr going 
to
+#   files, so a background child holds a file rather than the exec stream.
+# * A timer reading the monotonic /proc/uptime signals the wrapper when the
+#   budget is spent; counting its own one-second sleeps instead drifts late on 
a
+#   CPU-starved sandbox, past _EXEC_GRACE. The wrapper then SIGKILLs the
+#   processes in the command's session, rescanning /proc until a pass kills
+#   nothing, at most 50 times, so children forked mid-sweep are usually caught
+#   too. A negative-pid kill of the group is blocked by the sandbox's seccomp
+#   filter. The sweep shares the CPU with the processes it is killing, so it
+#   reads /proc with builtins, joining the lines of a stat file whose process
+#   name holds a newline, and strips the name with the cheap shortest match,
+#   taking the longest only for a name that contains ")": a `cat` per entry let
+#   one sweep of 128 busy processes on one CPU run past _EXEC_GRACE.
+# * The exit status is the wrapper's own process status, which the supervisor
+#   reports out of band, so nothing the command prints can change it; a command
+#   that signals the wrapper only ends its own run early. Output is sent last 
as
+#   one byte more than the caller's cap, so an over-cap stream is detected by
+#   counting what arrives rather than trusting a trailer.
+#
+# A command that starts its own session (setsid, a daemonizing server) leaves
+# this one and is not killed on timeout, and one that keeps forking faster than
+# the sweep can outlast it.
+_RUN_WRAPPER = rf"""t=$1 c=$2 o=$PATH
+{_SYSTEM_PATH}
+d=$(mktemp -d /tmp/.airflow-exec.XXXXXX) && cat >"$d/c" || {{ rm -rf "$d"; 
echo "{_STAGING_FAILED}" >&2; exit {_STAGING_STATUS}; }}
+f=0
+trap f=1 ALRM
+trap f=2 HUP INT TERM
+z=$(command -v setsid) || {{ echo "setsid is not installed in the sandbox 
image" >&2; rm -rf "$d"; exit 125; }}
+PATH=$o "$z" /bin/sh "$d/c" </dev/null >"$d/o" 2>"$d/e" &
+p=$!
+(
+  read u _ </proc/uptime
+  e=$((${{u%.*}}${{u#*.}} + t * 100))
+  while [ "${{u%.*}}${{u#*.}}" -lt "$e" ]; do
+    sleep 1
+    [ -d "/proc/$$" ] || exit 0
+    read u _ </proc/uptime
+  done
+  kill -ALRM $$
+) </dev/null >/dev/null 2>&1 &
+w=$!
+wait "$p"
+r=$?
+[ "$f" = 0 ] && trap '' ALRM
+if [ "$f" != 0 ]; then
+  n=0
+  while [ "$n" -lt 50 ]; do
+    k=0
+    for x in /proc/[0-9]*/stat; do
+      s=
+      while read -r l; do s="$s $l"; done 2>/dev/null <"$x" || continue
+      s=${{s#*) }}
+      case $s in *")"*) s=${{s##*) }} ;; esac
+      set -- $s
+      [ "$4" = "$p" ] || continue
+      x=${{x#/proc/}}
+      kill -KILL "${{x%/stat}}" 2>/dev/null && k=1
+    done
+    [ "$k" = 0 ] && break
+    n=$((n + 1))
+  done
+  wait "$p" 2>/dev/null
+fi
+kill -KILL "$w" 2>/dev/null
+tail -c "$((c + 1))" "$d/o"
+tail -c "$((c + 1))" "$d/e" >&2
+rm -rf "$d"
+[ "$f" = 1 ] && exit {_TIMEOUT_STATUS}
+[ "$f" = 2 ] && exit 143
+exit "$r"
+"""
+# read_file: the base class's script with the content sent raw instead of as
+# base64, so a binary file arrives intact and the transfer is half the size.
+_READ_SCRIPT = f"""{_SYSTEM_PATH}
+sz=$(stat -Lc %s -- "$1" 2>/dev/null) || exit 
{SandboxBackend._MISSING_PATH_STATUS}
+[ -d "$1" ] && exit {SandboxBackend._IS_DIRECTORY_STATUS}
+printf '%s\\n' "$sz"
+exec head -c "$2" -- "$1"
+"""
+_WRITE_FIRST_SCRIPT = f'{_SYSTEM_PATH}\nmkdir -p -- "$(dirname -- "$1")" && 
cat >"$1"\n'
+_WRITE_NEXT_SCRIPT = f'{_SYSTEM_PATH}\ncat >>"$1"\n'
+
+
+def _new_openshell_name() -> str:
+    return f"{_NAME_PREFIX}{uuid.uuid4().hex[:11]}"
+
+
+def _status_name(error: BaseException) -> str | None:
+    code = getattr(error, "code", None)
+    if not callable(code):
+        return None
+    try:
+        status = code()
+    except Exception:
+        return None
+    return getattr(status, "name", None)
+
+
+def _describe_error(error: BaseException) -> str:
+    status = _status_name(error)
+    if status is None:
+        return f"{type(error).__name__}: {error}"
+    details = getattr(error, "details", None)
+    text = details() if callable(details) else ""
+    return f"{status}: {text}" if text else status
+
+
+@contextmanager
+def _translate_openshell_errors(
+    operation: str, *, recoverable_statuses: frozenset[str] = frozenset()
+) -> Iterator[None]:
+    try:
+        yield
+    except SandboxError:
+        raise
+    except ImportError as e:
+        raise SandboxTerminalError(
+            'The OpenShell SDK is not installed. Install 
"apache-airflow-providers-common-ai[openshell]".'
+        ) from e
+    except Exception as e:
+        message = f"OpenShell could not {operation} ({_describe_error(e)})."
+        if _status_name(e) in recoverable_statuses:
+            raise SandboxError(message) from e
+        raise SandboxTerminalError(message) from e
+
+
+def _is_hostname(value: object) -> bool:
+    """Whether ``value`` is a hostname OpenShell's policy can match, 
optionally with one leading ``*.``."""
+    if not isinstance(value, str) or not value or len(value) > 253:
+        return False
+    labels = value.removeprefix("*.").split(".")
+    # A single label ('localhost', or '*.com' once the wildcard is removed) is
+    # refused by the gateway for wildcards and cannot name a public endpoint.
+    if len(labels) < 2 or labels[-1].isdigit():
+        return False
+    return all(
+        0 < len(label) <= 63
+        and label[0] != "-"
+        and label[-1] != "-"
+        and all(ch.isascii() and (ch.isalnum() or ch == "-") for ch in label)
+        for label in labels
+    )
+
+
+class _BoundedTail:
+    """
+    The last ``max_bytes`` of a stream, and how many bytes the stream carried 
in total.
+
+    ``min_window`` keeps that many trailing bytes even under a smaller 
``max_bytes``,
+    for :meth:`ends_with` only; the text and the truncation flag still follow 
``max_bytes``.
+    """
+
+    def __init__(self, max_bytes: int, *, min_window: int = 0) -> None:
+        self._max_bytes = max_bytes
+        self._window = max(max_bytes, min_window)
+        self._data = bytearray()
+        self.received = 0
+
+    def add(self, chunk: bytes) -> None:
+        self.received += len(chunk)
+        self._data.extend(chunk)
+        if len(self._data) > self._window:
+            del self._data[: len(self._data) - self._window]
+
+    @property
+    def truncated(self) -> bool:
+        return self.received > self._max_bytes
+
+    def ends_with(self, suffix: bytes) -> bool:
+        return self._data.endswith(suffix)
+
+    def get_text(self) -> str:
+        data = bytes(self._data[max(0, len(self._data) - self._max_bytes) :])
+        if self.truncated:
+            # Drop the leading partial line so no fragment reads as a whole 
record, unless that
+            # would throw away most of the window: one line longer than the 
cap has its newline
+            # at the very end, which would leave "(no output)" for a command 
that wrote megabytes.
+            newline = data.find(b"\n")
+            if newline != -1 and len(data) - (newline + 1) >= self._max_bytes 
// 2:
+                data = data[newline + 1 :]
+        return data.decode("utf-8", errors="replace")
+
+
+class _BoundedHead:
+    """The first ``max_bytes`` of a stream; anything past that is dropped."""
+
+    def __init__(self, max_bytes: int) -> None:
+        self._max_bytes = max_bytes
+        self._data = bytearray()
+
+    def add(self, chunk: bytes) -> None:
+        room = self._max_bytes - len(self._data)
+        if room > 0:
+            self._data.extend(chunk[:room])
+
+    def get_bytes(self) -> bytes:
+        return bytes(self._data)
+
+
+@dataclass(frozen=True)
+class _ExecOutcome:
+    exit_code: int
+    stdout: _BoundedTail | _BoundedHead
+    stderr: _BoundedTail
+    elapsed: float
+
+
+class _ExecHung(Exception):
+    """The exec stream outlived its deadline without reporting an exit 
status."""
+
+
+class _GatewayDropped(Exception):
+    """The gateway went away after the exec request was sent; the command's 
outcome is unknown."""
+
+
+@dataclass(frozen=True)
+class _EgressPolicy:
+    """What an OpenShell policy lets a sandbox reach, in a form two policies 
can be compared by."""
+
+    destinations: frozenset[tuple[str, int, frozenset[str]]]
+    problems: tuple[str, ...]
+
+    @classmethod
+    def requested(cls, hosts: Sequence[str]) -> _EgressPolicy:
+        return cls(
+            destinations=frozenset((host, _HTTPS_PORT, 
frozenset({_ANY_BINARY})) for host in hosts),
+            problems=(),
+        )
+
+    @classmethod
+    def effective(cls, config: Any) -> _EgressPolicy:
+        """Read the egress a ``GetSandboxConfigResponse`` says is in force."""
+        from openshell._proto import sandbox_pb2
+
+        destinations: set[tuple[str, int, frozenset[str]]] = set()
+        problems: list[str] = []
+        if config.policy_source != sandbox_pb2.POLICY_SOURCE_SANDBOX:
+            problems.append(
+                "a gateway-wide policy is in force instead of the sandbox's 
own "
+                f"(policy_source 
{sandbox_pb2.PolicySource.Name(config.policy_source)})"
+            )
+        if not config.configuration_admitted:
+            problems.append(
+                f"the gateway has not admitted the sandbox's configuration 
({config.configuration_error or 'no reason given'})"
+            )
+        if config.policy.landlock.compatibility != "hard_requirement":
+            problems.append(
+                f"Landlock enforcement is 
{config.policy.landlock.compatibility or 'unset'!r}, not 'hard_requirement'"
+            )
+        if config.policy.network_middlewares:
+            problems.append(
+                f"network middlewares are configured: 
{sorted(config.policy.network_middlewares)}"
+            )
+        for rule in config.policy.network_policies.values():
+            binaries = frozenset(binary.path for binary in rule.binaries)
+            for endpoint in rule.endpoints:
+                if endpoint.allowed_ips:
+                    problems.append(
+                        f"{endpoint.host or 'a rule with no host'} admits 
addresses {list(endpoint.allowed_ips)}"
+                    )
+                ports = set(endpoint.ports)
+                if endpoint.port:
+                    ports.add(endpoint.port)
+                destinations.update((endpoint.host.lower(), port, binaries) 
for port in ports)
+        if _PROPOSAL_APPROVAL_MODE in config.settings:
+            setting = config.settings[_PROPOSAL_APPROVAL_MODE]
+            if setting.value.string_value == "auto":
+                problems.append(
+                    "proposal_approval_mode is 'auto' "
+                    f"({sandbox_pb2.SettingScope.Name(setting.scope)}), so a 
connection the sandbox "
+                    "is denied can become an allow rule without review"
+                )
+        if _AGENT_POLICY_PROPOSALS in config.settings:
+            if config.settings[_AGENT_POLICY_PROPOSALS].value.bool_value:
+                problems.append(
+                    "agent_policy_proposals_enabled is on, so the agent can 
propose its own egress"
+                )
+        return cls(destinations=frozenset(destinations), 
problems=tuple(problems))
+
+    def describe(self) -> list[str]:
+        return sorted(f"{host}:{port}" for host, port, _ in self.destinations)
+
+    def mismatch(self, wanted: _EgressPolicy) -> str | None:
+        """Why this policy is not ``wanted``, or ``None`` when it is exactly 
that."""
+        problems = list(self.problems)
+        if self.destinations != wanted.destinations:
+            problems.append(
+                f"it admits {self.describe() or 'nothing'} where 
{wanted.describe() or 'nothing'} was asked for"
+            )
+        return "; ".join(problems) or None
+
+
+class OpenShellSandboxBackend(SandboxBackend):
+    """
+    Run sandbox tools through an NVIDIA OpenShell gateway.
+
+    .. note::
+
+        Experimental: this can change or be removed in a minor release of this 
provider.
+        See :ref:`howto/stability`.
+
+    `OpenShell <https://github.com/NVIDIA/OpenShell>`__ is a self-hosted 
gateway
+    that runs each sandbox on Docker, Podman or Kubernetes with Landlock and
+    seccomp inside the container and a per-sandbox supervisor that mediates
+    every outbound connection. Airflow workers only need gRPC access to the
+    gateway, with the ``openshell`` extra installed.
+
+    **Credentials are ambient.** The backend reads the gateway registration the
+    ``openshell`` CLI keeps under 
``$XDG_CONFIG_HOME/openshell/gateways/<name>/``:
+    its endpoint and either mTLS material or an OIDC token. There is no Airflow
+    connection type for it; see the backend page for how to provision a worker.
+
+    **Egress is deny-all unless listed, and verified.** A spec's hostnames 
become
+    one policy rule admitting each host on port 443 only, and a default spec
+    sends no rule at all, which OpenShell enforces as no egress; name 
resolution
+    is answered by the supervisor without leaving the sandbox. After create the
+    backend reads the effective policy back and destroys the sandbox unless it
+    is exactly what was asked for, comes from the sandbox rather than a
+    gateway-wide override, and cannot be widened by auto-approved proposals.
+    The same check runs before and after every command, because a gateway
+    administrator, an approved draft or a global policy can widen a running
+    sandbox; a change fails the task. What it cannot see is the network the
+    gateway runs on: on Kubernetes the NetworkPolicy that keeps a sandbox pod
+    behind its supervisor is the cluster's to enforce.
+
+    **Commands run through a guest wrapper.** OpenShell's own exec timeout
+    leaves the command running, so the wrapper enforces the budget itself and
+    kills the processes in the command's session when it runs out. Output is
+    spooled to the sandbox's ``/tmp`` and each stream is capped to
+    ``max_output_bytes`` before it leaves the sandbox. 
``allow_egress_to_cidrs``,
+    an open network and ``SandboxSpec.owner`` are refused. The image needs
+    ``setsid`` and GNU coreutils and findutils; ``python:*-slim`` has them.
+
+    :param gateway: Name of the OpenShell CLI gateway registration to use.
+        ``None`` uses ``$OPENSHELL_GATEWAY``, then the CLI's active gateway.
+    :param workspace: OpenShell workspace the sandboxes are created in.
+    :param image: Container image used for each sandbox.
+    :param cpu: CPU limit, as a Kubernetes quantity.
+    :param memory: Memory limit, as a Kubernetes quantity.
+    :param ready_timeout: Seconds to wait for a new sandbox to become ready.
+    :param request_timeout: Seconds allowed for each gateway call other than a 
command.
+    """
+
+    name = "openshell"
+
+    def __init__(
+        self,
+        *,
+        gateway: str | None = None,
+        workspace: str = "default",
+        image: str = "python:3.12-slim",
+        cpu: str = "1",
+        memory: str = "2Gi",
+        ready_timeout: float = 120.0,
+        request_timeout: float = 30.0,
+    ) -> None:
+        if gateway is not None and not gateway:
+            raise ValueError("gateway must not be empty; pass None to use the 
CLI's active gateway.")
+        if not workspace:
+            raise ValueError("workspace must not be empty.")
+        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(ready_timeout, "ready_timeout")
+        _validate_positive_finite(request_timeout, "request_timeout")
+        self._gateway = gateway
+        self._workspace = workspace
+        self._image = image
+        self._resources = {"limits": {"cpu": cpu, "memory": memory}}
+        self._ready_timeout = ready_timeout
+        self._request_timeout = request_timeout
+        self._client: SandboxClient | None = None
+        # The egress each sandbox was created with, compared against the 
effective
+        # policy before and after every command.
+        self._egress: dict[str, _EgressPolicy] = {}
+
+    def _get_client(self) -> SandboxClient:
+        with _CLIENT_BUILD_LOCK:
+            if self._client is None:
+                with _translate_openshell_errors("load the gateway 
registration"):
+                    from openshell import SandboxClient
+
+                    self._client = SandboxClient.from_active_cluster(
+                        cluster=self._gateway, timeout=self._request_timeout
+                    )
+            return self._client
+
+    def _workspace_scope(self) -> Any:
+        from openshell._proto import datamodel_pb2
+
+        return datamodel_pb2.WorkspaceSelector(workspace=self._workspace)
+
+    @staticmethod
+    def _check_spec(spec: SandboxSpec | None) -> tuple[list[str], dict[str, 
str]]:
+        """
+        Return the hosts and environment to provision, or refuse a spec 
OpenShell cannot enforce.
+
+        ``None`` gets the same deny-all policy as a default spec: OpenShell 
has no
+        policy-free mode that is not the image's own, which may be permissive.
+        """
+        if spec is None:
+            return [], {}
+        if spec.owner is not None:
+            # An owner exists so that a later task can attach to the sandbox. 
OpenShell
+            # labels are fixed at create and its annotations change only 
through an
+            # admin-scoped config mutation, so this version does not attach.
+            raise SandboxTerminalError(
+                "SandboxSpec names an owner, but this backend does not 
implement attaching, so a "
+                "sandbox created here cannot be attached to from another task. 
Drop owner, or "
+                "provision the sandbox on a backend that supports attaching, 
such as ModalSandboxBackend."
+            )
+        if spec.allow_egress_to_cidrs:
+            raise SandboxTerminalError(
+                "SandboxSpec names allow_egress_to_cidrs, which this backend 
cannot enforce as "
+                "specified: OpenShell address rules admit TCP on listed ports 
only, never UDP or any "
+                "port. Use allow_egress_to with hostnames, or a backend with 
an address-layer allowlist."
+            )
+        if not spec.block_network:
+            raise SandboxTerminalError(
+                "SandboxSpec asks for an open network, which OpenShell cannot 
provide: every egress "
+                "rule names a host and its ports. Set block_network=True and 
list the hosts the "
+                "sandbox needs in allow_egress_to."
+            )
+        hosts: list[str] = []
+        if spec.allow_egress_to:
+            if isinstance(spec.allow_egress_to, str):
+                raise SandboxTerminalError(
+                    "SandboxSpec.allow_egress_to must be a sequence of 
hostnames, not one string: "
+                    f"{spec.allow_egress_to!r} would be read a character at a 
time. Wrap it in a list."
+                )
+            rejected = [host for host in spec.allow_egress_to if not 
_is_hostname(host)]
+            if rejected:
+                raise SandboxTerminalError(
+                    "SandboxSpec.allow_egress_to must contain bare hostnames, 
optionally with a "
+                    f"leading '*.' label, and these entries are not: 
{rejected}. Write 'pypi.org' or "
+                    "'*.pythonhosted.org', not a URL, a host:port, an address, 
or a single label."
+                )
+            hosts = sorted({host.lower() for host in spec.allow_egress_to})
+        env: dict[str, str] = {}
+        for key, value in (spec.env or {}).items():
+            if not isinstance(key, str) or not isinstance(value, str):
+                raise SandboxTerminalError(
+                    f"SandboxSpec.env must map strings to strings; {key!r} is 
set to {type(value).__name__}."
+                )
+            if key in _SUPERVISOR_OWNED_ENV or key.startswith("OPENSHELL_"):
+                raise SandboxTerminalError(
+                    f"SandboxSpec.env sets {key}, which the OpenShell 
supervisor owns: it removes the "
+                    "proxy variables and replaces the CA bundle variables with 
its own TLS-terminating "
+                    "CA in every command, and reserves OPENSHELL_*. Remove it 
from the spec."
+                )
+            env[key] = value
+        return hosts, env
+
+    def _build_spec(self, hosts: Sequence[str], env: Mapping[str, str]) -> Any:
+        from google.protobuf import struct_pb2
+        from openshell._proto import openshell_pb2, sandbox_pb2
+
+        policy = sandbox_pb2.SandboxPolicy(
+            version=1,
+            filesystem=sandbox_pb2.FilesystemPolicy(
+                include_workdir=True, read_only=list(_READ_ONLY_PATHS), 
read_write=list(_READ_WRITE_PATHS)
+            ),
+            # Refuse to start at all rather than run without the filesystem 
confinement.
+            
landlock=sandbox_pb2.LandlockPolicy(compatibility="hard_requirement"),
+        )
+        if hosts:
+            policy.network_policies[_EGRESS_RULE].CopyFrom(
+                sandbox_pb2.NetworkPolicyRule(
+                    name=_EGRESS_RULE,
+                    endpoints=[sandbox_pb2.NetworkEndpoint(host=host, 
ports=[_HTTPS_PORT]) for host in hosts],
+                    binaries=[sandbox_pb2.NetworkBinary(path=_ANY_BINARY)],
+                )
+            )
+        spec = openshell_pb2.SandboxSpec(environment=dict(env), policy=policy)
+        spec.template.image = self._image
+        resources = struct_pb2.Struct()
+        resources.update(self._resources)
+        spec.template.resources.CopyFrom(resources)
+        return spec
+
+    def create(self, *, spec: SandboxSpec | None = None) -> str:
+        hosts, env = self._check_spec(spec)
+        client = self._get_client()
+        with _translate_openshell_errors("create a sandbox"):
+            request = self._build_spec(hosts, env)
+        name = _new_openshell_name()
+        try:
+            with _translate_openshell_errors("create a sandbox"):
+                client.create(
+                    workspace=self._workspace,
+                    spec=request,
+                    name=name,
+                    labels={
+                        _CREATED_BY_LABEL[0]: _CREATED_BY_LABEL[1],
+                        _CREATED_AT_LABEL: str(int(time.time())),
+                    },
+                )
+            self._wait_until_ready(name, deadline=time.monotonic() + 
self._ready_timeout)
+            wanted = _EgressPolicy.requested(hosts)
+            mismatch = self._read_egress(name).mismatch(wanted)
+            if mismatch is not None:
+                raise SandboxTerminalError(
+                    f"OpenShell did not apply the requested network policy to 
sandbox {name} ({mismatch}), "
+                    "so the sandbox was destroyed."
+                )
+        except BaseException:
+            # Bound before the call, so a create that fails half way -- after 
the
+            # gateway accepted it, or at the policy check -- still gets 
deleted.
+            with suppress(Exception):
+                self.destroy(name)
+            raise
+        self._egress[name] = wanted
+        return name
+
+    def _get_sandbox(self, name: str) -> Any:
+        from openshell._proto import openshell_pb2
+
+        return (
+            self._get_client()
+            ._stub.GetSandbox(
+                
openshell_pb2.GetSandboxRequest(workspace_scope=self._workspace_scope(), 
name=name),
+                timeout=self._request_timeout,
+            )
+            .sandbox
+        )
+
+    def _wait_until_ready(self, name: str, *, deadline: float) -> None:
+        """
+        Poll until the sandbox is ready, riding out a gateway that is 
restarting.
+
+        Raises :class:`SandboxTerminalError` for a sandbox that is gone, has 
failed or
+        has stopped, or that is not ready by ``deadline``.
+        """
+        from openshell._proto import openshell_pb2
+
+        last = "it did not report a phase"
+        while True:
+            try:
+                sandbox = self._get_sandbox(name)
+            except Exception as e:
+                if _status_name(e) not in {"UNAVAILABLE", "DEADLINE_EXCEEDED"}:
+                    with _translate_openshell_errors(f"check the status of 
sandbox {name}"):
+                        raise
+                last = f"the gateway did not answer ({_describe_error(e)})"
+            else:
+                phase = sandbox.status.phase
+                if phase == openshell_pb2.SANDBOX_PHASE_READY:
+                    return
+                reasons = "; ".join(
+                    f"{condition.reason}: {condition.message}"
+                    for condition in sandbox.status.conditions
+                    if condition.status != "True" and (condition.reason or 
condition.message)
+                )
+                last = f"it is {openshell_pb2.SandboxPhase.Name(phase)}" + (
+                    f" ({reasons})" if reasons else ""
+                )
+                if phase in {
+                    openshell_pb2.SANDBOX_PHASE_ERROR,
+                    openshell_pb2.SANDBOX_PHASE_STOPPING,
+                    openshell_pb2.SANDBOX_PHASE_STOPPED,
+                    openshell_pb2.SANDBOX_PHASE_COMPLETED,
+                    openshell_pb2.SANDBOX_PHASE_DELETING,
+                }:
+                    raise SandboxTerminalError(f"OpenShell sandbox {name} 
cannot run commands: {last}.")
+            if time.monotonic() >= deadline:
+                raise SandboxTerminalError(f"OpenShell sandbox {name} is not 
ready: {last}.")
+            time.sleep(_POLL_INTERVAL)
+
+    def _read_egress(self, name: str) -> _EgressPolicy:
+        """Read the sandbox's effective egress, waiting out a restarting 
gateway for a bounded time."""
+        from openshell._proto import sandbox_pb2
+
+        request = 
sandbox_pb2.GetSandboxConfigRequest(workspace_scope=self._workspace_scope(), 
name=name)
+        deadline = time.monotonic() + _GATEWAY_RECOVERY
+        while True:
+            try:
+                config = self._get_client()._stub.GetSandboxConfig(request, 
timeout=self._request_timeout)
+            except Exception as e:
+                if (
+                    _status_name(e) not in {"UNAVAILABLE", "DEADLINE_EXCEEDED"}
+                    or time.monotonic() >= deadline
+                ):
+                    with _translate_openshell_errors(f"read the network policy 
of sandbox {name}"):
+                        raise
+                time.sleep(_POLL_INTERVAL)
+                continue
+            return _EgressPolicy.effective(config)
+
+    def _check_egress(self, name: str) -> None:
+        """Destroy the sandbox and fail the task if its egress is no longer 
what it was created with."""
+        effective = self._read_egress(name)
+        # A handle this instance did not create is held to what it admits the 
first
+        # time it is seen, on top of the same invariants.
+        wanted = self._egress.setdefault(name, 
_EgressPolicy(effective.destinations, ()))
+        mismatch = effective.mismatch(wanted)
+        if mismatch is None:
+            return
+        with suppress(SandboxError):
+            self.destroy(name)
+        raise SandboxTerminalError(
+            f"The network policy of OpenShell sandbox {name} changed after it 
was created ({mismatch}), "
+            "so the sandbox was destroyed. A gateway administrator, an 
approved policy draft or a "
+            "gateway-wide policy can widen a running sandbox; commands already 
run may have used it."
+        )
+
+    def _exec(
+        self,
+        sandbox: str,
+        argv: list[str],
+        *,
+        stdin: bytes,
+        deadline: float,
+        stdout: _BoundedTail | _BoundedHead,
+        stderr: _BoundedTail,
+    ) -> _ExecOutcome:
+        """
+        Run ``argv`` once, streaming both outputs into the given bounded 
buffers.
+
+        Streams the gateway's events directly: the SDK's ``exec_stream`` keeps 
every
+        chunk in memory until the command exits. Raises :class:`_ExecHung` past
+        ``deadline`` and :class:`_GatewayDropped` when the gateway goes away 
after the
+        request left, where the command may or may not have run.
+        """
+        from openshell._proto import openshell_pb2
+
+        client = self._get_client()
+        request = openshell_pb2.ExecSandboxRequest(
+            workspace_scope=self._workspace_scope(),
+            sandbox=sandbox,
+            command=argv,
+            stdin=stdin,
+            # A login shell would source profiles that can override 
SandboxSpec.env.
+            no_login_shell=True,
+            request_id=str(uuid.uuid4()),
+        )
+        exit_code: int | None = None
+        started = time.monotonic()
+        stream = None
+        try:
+            stream = client._stub.ExecSandbox(request, timeout=deadline)
+            for event in stream:
+                kind = event.WhichOneof("payload")
+                if kind == "stdout":
+                    stdout.add(event.stdout.data)
+                elif kind == "stderr":
+                    stderr.add(event.stderr.data)
+                elif kind == "exit":
+                    exit_code = event.exit.exit_code
+        except Exception as e:
+            status = _status_name(e)
+            if status == "DEADLINE_EXCEEDED":
+                raise _ExecHung from e
+            if status == "UNAVAILABLE":
+                raise _GatewayDropped(_describe_error(e)) from e
+            with _translate_openshell_errors(
+                f"run a command in sandbox {sandbox}",
+                recoverable_statuses=frozenset({"OUT_OF_RANGE", 
"RESOURCE_EXHAUSTED"}),
+            ):
+                raise
+        finally:
+            if stream is not None:
+                with suppress(Exception):
+                    stream.cancel()
+        if exit_code is None:
+            raise _GatewayDropped("the exec stream ended without an exit 
status")
+        return _ExecOutcome(exit_code, stdout, stderr, time.monotonic() - 
started)
+
+    def _exec_with_recovery(
+        self,
+        sandbox: str,
+        argv: list[str],
+        *,
+        stdin: bytes,
+        deadline: float,
+        stdout: _BoundedTail | _BoundedHead,
+        stderr: _BoundedTail,
+    ) -> _ExecOutcome:
+        """
+        :meth:`_exec`, with a gateway restart turned into an honest outcome.
+
+        A sandbox that is not ready yet, which is how it looks for a few 
seconds after a
+        gateway restart, has not launched anything, so the command is sent 
again once it
+        is. A gateway that drops an exec in flight has killed the command with 
every
+        other process in the sandbox, and whether it had already done its work 
is not
+        knowable, so that is reported to the model rather than retried or 
guessed.
+        """
+        try:
+            try:
+                return self._exec(sandbox, argv, stdin=stdin, 
deadline=deadline, stdout=stdout, stderr=stderr)
+            except SandboxTerminalError as e:
+                if _status_name(e.__cause__ or e) != "FAILED_PRECONDITION":
+                    raise
+                self._wait_until_ready(sandbox, deadline=time.monotonic() + 
self._ready_timeout)
+                return self._exec(sandbox, argv, stdin=stdin, 
deadline=deadline, stdout=stdout, stderr=stderr)
+        except _GatewayDropped as e:
+            self._wait_until_ready(sandbox, deadline=time.monotonic() + 
_GATEWAY_RECOVERY)
+            raise SandboxError(
+                f"The OpenShell gateway dropped the command before it reported 
an exit status ({e}); the "
+                "gateway may have restarted. The command may or may not have 
run, and a gateway restart "
+                "also stops processes that earlier commands left running. 
Files are kept. Check what "
+                "state you need and run the command again."
+            ) from e
+
+    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")
+        # Whole seconds, at least one: the wrapper takes its budget as an 
integer.
+        seconds = max(1, math.ceil(timeout))
+        payload = command.encode("utf-8")
+        if len(payload) > _MAX_STDIN_BYTES:
+            raise SandboxError(
+                f"The command is {len(payload)} bytes, over the 
{_MAX_STDIN_BYTES} bytes OpenShell carries "
+                "in one request. Write the script to a file with write_file 
and run the file."
+            )
+        self._check_egress(sandbox)
+        stdout = _BoundedTail(int(max_output_bytes))
+        stderr = _BoundedTail(int(max_output_bytes), 
min_window=len(_STAGING_LINE))
+        try:
+            outcome = self._exec_with_recovery(
+                sandbox,
+                ["/bin/sh", "-c", _RUN_WRAPPER, "airflow-exec", str(seconds), 
str(int(max_output_bytes))],
+                stdin=payload,
+                deadline=seconds + _EXEC_GRACE,
+                stdout=stdout,
+                stderr=stderr,
+            )
+        except _ExecHung:
+            return self._abandon_command(sandbox, stdout, stderr, 
seconds=seconds)
+        self._check_egress(sandbox)
+        if outcome.exit_code == _STAGING_STATUS and 
stderr.ends_with(_STAGING_LINE):
+            cause = stderr.get_text().rpartition(_STAGING_FAILED)[0].strip()
+            raise SandboxError(
+                "The command was not run: the sandbox could not write it to 
/tmp"
+                + (f" ({cause})." if cause else ".")
+            )
+        return SandboxExecResult(
+            exit_code=outcome.exit_code,
+            stdout=stdout.get_text(),
+            stderr=stderr.get_text(),
+            timed_out=outcome.exit_code == _TIMEOUT_STATUS and outcome.elapsed 
>= seconds,
+            stdout_truncated=stdout.truncated,
+            stderr_truncated=stderr.truncated,
+            applied_timeout=float(seconds),
+        )
+
+    def _abandon_command(
+        self, sandbox: str, stdout: _BoundedTail, stderr: _BoundedTail, *, 
seconds: int
+    ) -> SandboxExecResult:
+        # The wrapper reports within its budget unless the command disabled it 
or the
+        # gateway stopped relaying; either way the command may still be 
running, and
+        # OpenShell does not stop an exec whose caller went away, so only 
destroying
+        # the sandbox ends it. There is no server-side lifetime to fall back 
on.
+        try:
+            self.destroy(sandbox)
+        except SandboxError:
+            log.warning(
+                "Timed out running a command in OpenShell sandbox %s and could 
not destroy it; it has no "
+                "server-side lifetime and will need manual cleanup",
+                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,
+            applied_timeout=float(seconds),
+        )
+
+    def _run_file_op(
+        self, sandbox: str, script: str, *args: str, stdin: bytes = b"", 
stdout_cap: int = _FILE_OP_OUTPUT_CAP
+    ) -> tuple[_ExecOutcome, bytes]:
+        head = _BoundedHead(stdout_cap)
+        try:
+            outcome = self._exec_with_recovery(
+                sandbox,
+                ["/bin/sh", "-c", script, "airflow-file", *args],

Review Comment:
   The path reaches the gateway as an argv element, and at v0.1.2 the gateway 
rejects any argument containing a NUL or longer than 32 KiB with 
`INVALID_ARGUMENT` 
([validation.rs](https://github.com/NVIDIA/OpenShell/blob/v0.1.2/crates/openshell-server/src/grpc/validation.rs#L56-L64)).
 `_exec` maps that to `SandboxTerminalError`, so a malformed path from the 
model fails the task instead of coming back as something it can correct. 
Checking for NUL and length before the exec and raising `SandboxError` would 
keep it recoverable for `read_file` and `write_file`.



##########
providers/common/ai/tests/unit/common/ai/sandbox/test_openshell.py:
##########
@@ -0,0 +1,1002 @@
+# 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 copy
+import os
+import shlex
+import shutil
+import signal
+import subprocess
+import sys
+import time
+from unittest import mock
+
+import pytest
+
+pytest.importorskip("openshell")
+
+import grpc
+from openshell import SandboxClient
+from openshell._proto import openshell_pb2, sandbox_pb2
+
+from airflow.providers.common.ai.sandbox.base import (
+    SandboxError,
+    SandboxFileTooLargeError,
+    SandboxSpec,
+    SandboxTerminalError,
+)
+from airflow.providers.common.ai.sandbox.openshell import (
+    _EXEC_GRACE,
+    _READ_WRITE_PATHS,
+    _RUN_WRAPPER,
+    _STAGING_FAILED,
+    OpenShellSandboxBackend,
+)
+
+_MODULE = "airflow.providers.common.ai.sandbox.openshell"
+_MONOTONIC_PATH = f"{_MODULE}.time.monotonic"
+
+
+class _RpcError(grpc.RpcError):
+    def __init__(self, code: grpc.StatusCode, details: str = "") -> None:
+        super().__init__(details)
+        self._code = code
+        self._details = details
+
+    def code(self) -> grpc.StatusCode:
+        return self._code
+
+    def details(self) -> str:
+        return self._details
+
+
+class _Stream:
+    """An ExecSandbox response stream: events, optionally ending in an 
error."""
+
+    def __init__(self, events, error: Exception | None = None) -> None:
+        self._events = iter(events)
+        self._error = error
+        self.cancel = mock.Mock()
+
+    def __iter__(self):
+        return self
+
+    def __next__(self):
+        try:
+            return next(self._events)
+        except StopIteration:
+            if self._error is not None:
+                raise self._error from None
+            raise
+
+
+def _stdout(data: bytes):
+    return 
openshell_pb2.ExecSandboxEvent(stdout=openshell_pb2.ExecSandboxStdout(data=data))
+
+
+def _stderr(data: bytes):
+    return 
openshell_pb2.ExecSandboxEvent(stderr=openshell_pb2.ExecSandboxStderr(data=data))
+
+
+def _exit(code: int):
+    return 
openshell_pb2.ExecSandboxEvent(exit=openshell_pb2.ExecSandboxExit(exit_code=code))
+
+
+def _result(code: int = 0, out: bytes = b"", err: bytes = b"") -> _Stream:
+    events = []
+    if out:
+        events.append(_stdout(out))
+    if err:
+        events.append(_stderr(err))
+    events.append(_exit(code))
+    return _Stream(events)
+
+
+def _sandbox(phase: openshell_pb2.SandboxPhase = 
openshell_pb2.SANDBOX_PHASE_READY, *conditions):
+    response = openshell_pb2.SandboxResponse()
+    response.sandbox.status.phase = phase
+    for reason, message in conditions:
+        response.sandbox.status.conditions.add(type="Ready", status="False", 
reason=reason, message=message)
+    return response
+
+
+def _config(
+    hosts=(),
+    *,
+    source: sandbox_pb2.PolicySource = sandbox_pb2.POLICY_SOURCE_SANDBOX,
+    admitted: bool = True,
+    landlock: str = "hard_requirement",
+    approval_mode: str | None = None,
+    proposals: bool = False,
+    extra_rule: tuple[str, int] | None = None,
+    allowed_ips=(),
+    binaries=("/**",),
+    middlewares=(),
+):
+    config = sandbox_pb2.GetSandboxConfigResponse(
+        version=1, policy_source=source, configuration_admitted=admitted
+    )
+    config.policy.landlock.compatibility = landlock
+    for name in middlewares:
+        # Reading a missing key of a message map adds it.
+        config.policy.network_middlewares[name]
+    if hosts or allowed_ips:
+        rule = config.policy.network_policies["airflow-egress"]
+        rule.name = "airflow-egress"
+        for host in hosts:
+            rule.endpoints.add(host=host, ports=[443], 
allowed_ips=list(allowed_ips))
+        if not hosts:
+            rule.endpoints.add(ports=[443], allowed_ips=list(allowed_ips))
+        rule.binaries.extend(sandbox_pb2.NetworkBinary(path=path) for path in 
binaries)
+    if extra_rule is not None:
+        rule = config.policy.network_policies["allow_extra"]
+        rule.endpoints.add(host=extra_rule[0], port=extra_rule[1])
+        rule.binaries.add(path="/usr/local/bin/python3.12")
+    if approval_mode is not None:
+        config.settings["proposal_approval_mode"].value.string_value = 
approval_mode
+        config.settings["proposal_approval_mode"].scope = 
sandbox_pb2.SETTING_SCOPE_GLOBAL
+    if proposals:
+        config.settings["agent_policy_proposals_enabled"].value.bool_value = 
True
+    return config
+
+
+def _backend(**kwargs) -> tuple[OpenShellSandboxBackend, mock.MagicMock]:
+    backend = OpenShellSandboxBackend(gateway="test-gateway", **kwargs)
+    client = mock.MagicMock(spec=SandboxClient)
+    client._stub = mock.MagicMock(spec=["GetSandbox", "GetSandboxConfig", 
"ExecSandbox"])
+    client._stub.GetSandbox.return_value = _sandbox()
+    client._stub.GetSandboxConfig.return_value = _config()
+    client._stub.ExecSandbox.return_value = _result()
+    backend._client = client
+    return backend, client
+
+
+def _exec_request(client, call: int = -1):
+    return client._stub.ExecSandbox.call_args_list[call].args[0]
+
+
[email protected](autouse=True)
+def _no_sleep():
+    with mock.patch(f"{_MODULE}.time.sleep", autospec=True):
+        yield
+
+
+def test_missing_sdk_error_is_actionable():
+    real_import = builtins.__import__
+
+    def blocked_import(name, *args, **kwargs):
+        if name.startswith("openshell"):
+            raise ImportError("blocked for test")
+        return real_import(name, *args, **kwargs)
+
+    backend = OpenShellSandboxBackend()
+    with mock.patch("builtins.__import__", side_effect=blocked_import):
+        with pytest.raises(SandboxTerminalError, match=r"\[openshell\]"):
+            backend.create()
+
+
[email protected](
+    ("kwargs", "message"),
+    [
+        ({"gateway": ""}, "gateway"),
+        ({"workspace": ""}, "workspace"),
+        ({"image": ""}, "image"),
+        ({"cpu": ""}, "cpu"),
+        ({"memory": ""}, "memory"),
+        ({"ready_timeout": 0}, "ready_timeout"),
+        ({"request_timeout": float("inf")}, "request_timeout"),
+    ],
+)
+def test_constructor_rejects_invalid_values(kwargs, message):
+    with pytest.raises(ValueError, match=message):
+        OpenShellSandboxBackend(**kwargs)
+
+
+class TestClient:
+    @mock.patch("openshell.SandboxClient.from_active_cluster", autospec=True)
+    def test_construction_reads_no_gateway_registration(self, 
from_active_cluster):
+        OpenShellSandboxBackend(gateway="prod")
+
+        from_active_cluster.assert_not_called()
+
+    def test_the_backend_can_be_deep_copied(self):
+        backend = OpenShellSandboxBackend(gateway="prod")
+
+        assert copy.deepcopy(backend)._gateway == "prod"
+
+    @mock.patch("openshell.SandboxClient.from_active_cluster", autospec=True)
+    def test_the_cli_gateway_registration_is_loaded_once(self, 
from_active_cluster):
+        backend = OpenShellSandboxBackend(gateway="prod", request_timeout=12)
+
+        assert backend._get_client() is backend._get_client()
+
+        from_active_cluster.assert_called_once_with(cluster="prod", timeout=12)
+
+    @mock.patch("openshell.SandboxClient.from_active_cluster", autospec=True)
+    def test_a_missing_registration_is_terminal(self, from_active_cluster):
+        from openshell import SandboxError as SdkError
+
+        from_active_cluster.side_effect = SdkError("gateway 'prod' not found")
+
+        with pytest.raises(SandboxTerminalError, match="gateway 'prod' not 
found"):
+            OpenShellSandboxBackend(gateway="prod")._get_client()
+
+
+class TestSpecRefusals:
+    @pytest.mark.parametrize(
+        ("spec", "message"),
+        [
+            (SandboxSpec(owner="dag/run"), "owner"),
+            (SandboxSpec(allow_egress_to_cidrs=["203.0.113.0/24"]), 
"allow_egress_to_cidrs"),
+            (SandboxSpec(block_network=False), "open network"),
+            (SandboxSpec(block_network=False, allow_egress_to=["pypi.org"]), 
"open network"),
+            (SandboxSpec(allow_egress_to="pypi.org"), "not one string"),
+            (SandboxSpec(env={"FOO": 1}), "strings to strings"),  # type: 
ignore[dict-item]
+        ],
+    )
+    def test_unenforceable_specs_are_refused_before_the_gateway_is_asked(self, 
spec, message):
+        backend, client = _backend()
+
+        with pytest.raises(SandboxTerminalError, match=message):
+            backend.create(spec=spec)
+
+        client.create.assert_not_called()
+
+    @pytest.mark.parametrize(
+        "host",
+        ["https://pypi.org";, "pypi.org:443", "203.0.113.7", "localhost", 
"*.com", "*", "pypi.org/simple", ""],
+    )
+    def test_entries_that_are_not_hostnames_are_refused(self, host):
+        backend, client = _backend()
+
+        with pytest.raises(SandboxTerminalError, match="bare hostnames"):
+            backend.create(spec=SandboxSpec(allow_egress_to=[host]))
+
+        client.create.assert_not_called()
+
+    @pytest.mark.parametrize(
+        "key", ["HTTPS_PROXY", "no_proxy", "SSL_CERT_FILE", 
"REQUESTS_CA_BUNDLE", "OPENSHELL_SANDBOX"]
+    )
+    def test_env_the_supervisor_would_silently_change_is_refused(self, key):
+        backend, client = _backend()
+
+        with pytest.raises(SandboxTerminalError, match=f"sets {key}"):
+            backend.create(spec=SandboxSpec(env={key: "value"}))
+
+        client.create.assert_not_called()
+
+
+class TestCreate:
+    def test_default_spec_is_deny_all_under_hard_landlock(self):
+        backend, client = _backend(image="python:3.13-slim", cpu="2", 
memory="4Gi")
+
+        with mock.patch(f"{_MODULE}.time.time", return_value=1_700_000_000.4):
+            name = backend.create(spec=SandboxSpec(env={"TOKEN": "value"}))
+
+        kwargs = client.create.call_args.kwargs
+        assert kwargs["name"] == name
+        assert name.startswith("airflow-")
+        assert len(name) <= 19
+        assert kwargs["workspace"] == "default"
+        assert kwargs["labels"] == {"created-by": "airflow", 
"airflow-created-at": "1700000000"}
+        spec = kwargs["spec"]
+        assert dict(spec.environment) == {"TOKEN": "value"}
+        assert spec.template.image == "python:3.13-slim"
+        assert dict(spec.template.resources)["limits"] == {"cpu": "2", 
"memory": "4Gi"}
+        assert len(spec.policy.network_policies) == 0
+        assert spec.policy.landlock.compatibility == "hard_requirement"
+        assert tuple(spec.policy.filesystem.read_write) == _READ_WRITE_PATHS
+        assert "/dev/shm" in spec.policy.filesystem.read_write
+        client.delete.assert_not_called()
+
+    def test_none_spec_still_gets_an_explicit_deny_all_policy(self):
+        # Without a policy OpenShell falls back to one baked into the image, 
which may be open.
+        backend, client = _backend()
+
+        backend.create()
+
+        spec = client.create.call_args.kwargs["spec"]
+        assert spec.HasField("policy")
+        assert len(spec.policy.network_policies) == 0
+
+    def test_hostnames_become_one_rule_on_port_443_for_any_binary(self):
+        backend, client = _backend()
+        client._stub.GetSandboxConfig.return_value = 
_config(["files.pythonhosted.org", "pypi.org"])
+
+        backend.create(spec=SandboxSpec(allow_egress_to=["PyPI.org", 
"files.pythonhosted.org", "pypi.org"]))
+
+        rules = client.create.call_args.kwargs["spec"].policy.network_policies
+        assert list(rules) == ["airflow-egress"]
+        assert [(endpoint.host, list(endpoint.ports)) for endpoint in 
rules["airflow-egress"].endpoints] == [
+            ("files.pythonhosted.org", [443]),
+            ("pypi.org", [443]),
+        ]
+        assert [binary.path for binary in rules["airflow-egress"].binaries] == 
["/**"]
+        client.delete.assert_not_called()
+
+    def test_waits_for_the_sandbox_to_become_ready(self):
+        backend, client = _backend()
+        client._stub.GetSandbox.side_effect = [
+            _sandbox(openshell_pb2.SANDBOX_PHASE_PROVISIONING),
+            _RpcError(grpc.StatusCode.UNAVAILABLE, "restarting"),
+            _sandbox(),
+        ]
+
+        backend.create(spec=SandboxSpec())
+
+        assert client._stub.GetSandbox.call_count == 3
+
+    def test_a_sandbox_that_fails_to_start_is_deleted_and_terminal(self):
+        backend, client = _backend()
+        client._stub.GetSandbox.return_value = _sandbox(
+            openshell_pb2.SANDBOX_PHASE_ERROR, ("IdentityResolutionFailed", 
"no passwd entry")
+        )
+
+        with pytest.raises(SandboxTerminalError, 
match="IdentityResolutionFailed: no passwd entry"):
+            backend.create(spec=SandboxSpec())
+
+        name = client.create.call_args.kwargs["name"]
+        client.delete.assert_called_once_with(name, workspace="default", 
allow_missing=True)
+
+    def test_a_sandbox_not_ready_in_time_is_deleted_and_terminal(self):
+        backend, client = _backend(ready_timeout=5)
+        client._stub.GetSandbox.return_value = 
_sandbox(openshell_pb2.SANDBOX_PHASE_PROVISIONING)
+
+        with mock.patch(_MONOTONIC_PATH, side_effect=[0.0, 1.0, 6.0]):
+            with pytest.raises(SandboxTerminalError, match="not ready: it is 
SANDBOX_PHASE_PROVISIONING"):
+                backend.create(spec=SandboxSpec())
+
+        client.delete.assert_called_once()
+
+    def test_a_create_the_gateway_rejects_is_terminal_and_cleaned_up(self):
+        backend, client = _backend()
+        client.create.side_effect = 
_RpcError(grpc.StatusCode.INVALID_ARGUMENT, "name exceeds maximum length")
+
+        with pytest.raises(SandboxTerminalError, match="INVALID_ARGUMENT: name 
exceeds maximum length"):
+            backend.create(spec=SandboxSpec())
+
+        client.delete.assert_called_once()
+
+    @pytest.mark.parametrize(
+        ("config", "message"),
+        [
+            (_config(source=sandbox_pb2.POLICY_SOURCE_GLOBAL), "gateway-wide 
policy"),
+            (_config(["pypi.org", "example.com"]), "admits"),
+            (_config(), "admits nothing where"),
+            (_config(extra_rule=("www.google.com", 443)), 
"www.google.com:443"),
+            (_config(["pypi.org"], binaries=("/usr/bin/curl",)), "admits"),
+            (_config(allowed_ips=["0.0.0.0/1"]), "admits addresses"),
+            (_config(["pypi.org"], approval_mode="auto"), 
"proposal_approval_mode is 'auto'"),
+            (_config(["pypi.org"], proposals=True), 
"agent_policy_proposals_enabled"),
+            (_config(["pypi.org"], landlock="best_effort"), "Landlock"),
+            (_config(["pypi.org"], admitted=False), "not admitted"),
+            (_config(["pypi.org"], middlewares=["inspect"]), "network 
middlewares"),
+            (_RpcError(grpc.StatusCode.PERMISSION_DENIED, "config:read"), 
"PERMISSION_DENIED"),
+        ],
+        ids=[
+            "global-override",
+            "wider-allowlist",
+            "missing-rule",
+            "auto-approved-rule",
+            "narrower-binaries",
+            "address-rule",
+            "auto-approval",
+            "agent-proposals",
+            "landlock-best-effort",
+            "not-admitted",
+            "network-middlewares",
+            "unreadable",
+        ],
+    )
+    def test_a_policy_other_than_the_requested_one_destroys_the_sandbox(self, 
config, message):
+        backend, client = _backend()
+        if isinstance(config, Exception):
+            client._stub.GetSandboxConfig.side_effect = config
+        else:
+            client._stub.GetSandboxConfig.return_value = config
+
+        with pytest.raises(SandboxTerminalError, match=message):
+            backend.create(spec=SandboxSpec(allow_egress_to=["pypi.org"]))
+
+        client.delete.assert_called_once()
+        assert backend._egress == {}
+
+    @pytest.mark.parametrize("code", [grpc.StatusCode.UNAVAILABLE, 
grpc.StatusCode.DEADLINE_EXCEEDED])
+    def test_a_policy_read_rides_out_a_restarting_gateway(self, code):
+        backend, client = _backend()
+        client._stub.GetSandboxConfig.side_effect = [_RpcError(code, 
"restarting"), _config()]
+
+        backend.create(spec=SandboxSpec())
+
+        assert client._stub.GetSandboxConfig.call_count == 2
+        client.delete.assert_not_called()
+
+    def test_manual_approval_mode_is_accepted(self):
+        backend, client = _backend()
+        client._stub.GetSandboxConfig.return_value = 
_config(approval_mode="manual")
+
+        backend.create(spec=SandboxSpec())
+
+        client.delete.assert_not_called()
+
+
+class TestRunCommand:
+    def test_the_command_travels_on_stdin_to_the_guest_wrapper(self):
+        backend, client = _backend()
+        client._stub.ExecSandbox.return_value = _result(3, b"out\n", b"err\n")
+
+        result = backend.run_command(
+            "box", "echo out; echo err >&2; exit 3", timeout=2.5, 
max_output_bytes=100
+        )
+
+        request = _exec_request(client)
+        assert request.sandbox == "box"
+        assert list(request.command) == ["/bin/sh", "-c", _RUN_WRAPPER, 
"airflow-exec", "3", "100"]
+        assert request.stdin == b"echo out; echo err >&2; exit 3"
+        assert request.no_login_shell is True
+        assert request.request_id
+        assert not request.HasField("execution_timeout")
+        assert client._stub.ExecSandbox.call_args.kwargs["timeout"] == 3 + 
_EXEC_GRACE
+        assert (result.exit_code, result.stdout, result.stderr) == (3, 
"out\n", "err\n")
+        assert not result.timed_out
+        assert not result.sandbox_terminated
+
+    @pytest.mark.parametrize(
+        ("exit_code", "elapsed", "timed_out"),
+        [(124, 3.1, True), (124, 0.2, False), (137, 3.1, False), (0, 3.1, 
False)],
+    )
+    def test_only_the_wrappers_status_after_the_budget_is_a_timeout(self, 
exit_code, elapsed, timed_out):
+        backend, client = _backend()
+        client._stub.ExecSandbox.return_value = _result(exit_code)
+
+        # The policy read before and after the command each take a deadline 
reading too.
+        with mock.patch(_MONOTONIC_PATH, side_effect=[0.0, 100.0, 100.0 + 
elapsed, 0.0]):
+            result = backend.run_command("box", "sleep 300", timeout=3, 
max_output_bytes=100)
+
+        assert result.timed_out is timed_out
+
+    @pytest.mark.parametrize(
+        ("exit_code", "stderr"),
+        [
+            pytest.param(125, b"docker: invalid reference format\n", 
id="a-command-exiting-125"),
+            pytest.param(1, f"{_STAGING_FAILED}\n".encode(), 
id="the-staging-line-with-another-status"),
+            pytest.param(125, f"{_STAGING_FAILED}\nmore\n".encode(), 
id="the-staging-line-not-last"),
+        ],
+    )
+    def 
test_a_result_unlike_the_wrappers_staging_failure_is_the_commands_own(self, 
exit_code, stderr):
+        backend, client = _backend()
+        client._stub.ExecSandbox.return_value = _result(exit_code, err=stderr)
+
+        result = backend.run_command("box", "cmd", timeout=5, 
max_output_bytes=100)
+
+        assert (result.exit_code, result.stderr) == (exit_code, 
stderr.decode())
+
+    @pytest.mark.parametrize(
+        ("max_output_bytes", "message"),
+        [
+            pytest.param(
+                100,
+                "The command was not run: the sandbox could not write it to 
/tmp "
+                "(cat: write error: No space left on device).",
+                id="with-the-cause",
+            ),
+            pytest.param(
+                10,
+                "The command was not run: the sandbox could not write it to 
/tmp.",
+                id="under-a-cap-shorter-than-the-staging-line",
+            ),
+        ],
+    )
+    def test_the_wrappers_own_staging_failure_is_a_recoverable_error(self, 
max_output_bytes, message):
+        backend, client = _backend()
+        client._stub.ExecSandbox.return_value = _result(
+            125, err=f"cat: write error: No space left on 
device\n{_STAGING_FAILED}\n".encode()
+        )
+
+        with pytest.raises(SandboxError) as error:
+            backend.run_command("box", "make", timeout=5, 
max_output_bytes=max_output_bytes)
+
+        assert str(error.value) == message
+        assert not isinstance(error.value, SandboxTerminalError)
+
+    def test_stderr_stays_capped_below_the_length_of_the_staging_line(self):
+        backend, client = _backend()
+        client._stub.ExecSandbox.return_value = _result(1, err=b"0123456789" * 
3)
+
+        result = backend.run_command("box", "cmd", timeout=5, 
max_output_bytes=20)
+
+        assert (result.stderr, result.stderr_truncated) == ("0123456789" * 2, 
True)
+
+    def test_each_stream_is_capped_and_flagged_on_its_own(self):
+        backend, client = _backend()
+        # The wrapper sends one byte past the cap for a stream it had to cut.
+        client._stub.ExecSandbox.return_value = _Stream(
+            [_stdout(b"partial line\nkept 1\n"), _stdout(b"kept 2\n"), 
_stderr(b"short\n"), _exit(0)]
+        )
+
+        result = backend.run_command("box", "cmd", timeout=5, 
max_output_bytes=20)
+
+        assert result.stdout == "kept 1\nkept 2\n"
+        assert result.stdout_truncated
+        assert result.stderr == "short\n"
+        assert not result.stderr_truncated
+
+    def test_one_line_longer_than_the_budget_still_reaches_the_model(self):
+        """Dropping the leading partial line must not empty the window: the 
toolset prints "(no output)" for it."""
+        backend, client = _backend()
+        client._stub.ExecSandbox.return_value = _Stream([_stdout(b"x" * 60000 
+ b"\n"), _exit(0)])
+
+        result = backend.run_command("box", "cat big.json", timeout=5, 
max_output_bytes=51200)
+
+        assert result.stdout != ""
+        assert result.stdout_truncated
+        assert len(result.stdout.encode()) <= 51200
+
+    def test_a_long_line_followed_by_a_short_one_keeps_the_window(self):
+        backend, client = _backend()
+        client._stub.ExecSandbox.return_value = _Stream(
+            [_stdout(b"L" * 200_000 + b"\nshort tail\n"), _exit(0)]
+        )
+
+        result = backend.run_command("box", "spew", timeout=5, 
max_output_bytes=51200)
+
+        assert len(result.stdout.encode()) > 51200 // 2
+        assert result.stdout.endswith("short tail\n")
+
+    def test_the_whole_second_deadline_the_command_got_is_reported(self):
+        backend, client = _backend()
+
+        result = backend.run_command("box", "true", timeout=2.2, 
max_output_bytes=100)
+
+        assert result.applied_timeout == 3.0
+
+    def 
test_the_whole_second_deadline_is_reported_for_an_abandoned_command_too(self):
+        backend, client = _backend()
+        client._stub.ExecSandbox.return_value = _Stream([], 
_RpcError(grpc.StatusCode.DEADLINE_EXCEEDED))
+
+        result = backend.run_command("box", "sleep 300", timeout=2.2, 
max_output_bytes=100)
+
+        assert (result.sandbox_terminated, result.applied_timeout) == (True, 
3.0)
+
+    def 
test_output_injected_past_the_wrapper_stays_bounded_in_worker_memory(self):
+        backend, client = _backend()
+        chunk = b"y\n" * 32768
+        client._stub.ExecSandbox.return_value = _Stream([_stdout(chunk) for _ 
in range(200)] + [_exit(0)])
+
+        result = backend.run_command("box", "yes > /proc/$PPID/fd/1", 
timeout=5, max_output_bytes=1024)
+
+        assert len(result.stdout.encode()) <= 1024
+        assert result.stdout_truncated
+
+    def test_undecodable_output_is_replaced_not_raised(self):
+        backend, client = _backend()
+        client._stub.ExecSandbox.return_value = _result(0, b"\xff\xfeok\n")
+
+        assert backend.run_command("box", "cmd", timeout=5, 
max_output_bytes=100).stdout == "��ok\n"
+
+    def 
test_a_command_over_the_request_limit_is_recoverable_and_not_sent(self):
+        backend, client = _backend()
+
+        with pytest.raises(SandboxError, match="write_file") as error:
+            backend.run_command("box", "x" * 1_000_001, timeout=5, 
max_output_bytes=100)
+
+        assert not isinstance(error.value, SandboxTerminalError)
+        client._stub.ExecSandbox.assert_not_called()
+
+    def test_a_hung_exec_destroys_the_sandbox(self):
+        backend, client = _backend()
+        client._stub.ExecSandbox.return_value = _Stream([], 
_RpcError(grpc.StatusCode.DEADLINE_EXCEEDED))
+
+        result = backend.run_command("box", "sleep 300", timeout=3, 
max_output_bytes=100)
+
+        assert (result.exit_code, result.timed_out, result.sandbox_terminated) 
== (-1, True, True)
+        client.delete.assert_called_once_with("box", workspace="default", 
allow_missing=True)
+
+    def 
test_a_sandbox_not_ready_after_a_gateway_restart_is_waited_for_and_retried(self):
+        backend, client = _backend()
+        client._stub.ExecSandbox.side_effect = [
+            _RpcError(grpc.StatusCode.FAILED_PRECONDITION, "sandbox is not 
ready"),
+            _result(0, b"ok\n"),
+        ]
+        client._stub.GetSandbox.side_effect = 
[_sandbox(openshell_pb2.SANDBOX_PHASE_PROVISIONING), _sandbox()]
+
+        result = backend.run_command("box", "echo ok", timeout=5, 
max_output_bytes=100)
+
+        assert result.stdout == "ok\n"
+        assert client._stub.ExecSandbox.call_count == 2
+        assert _exec_request(client, 0).request_id != _exec_request(client, 
1).request_id
+
+    def 
test_a_gateway_drop_during_the_not_ready_retry_is_reported_not_leaked(self):
+        backend, client = _backend()
+        client._stub.ExecSandbox.side_effect = [
+            _RpcError(grpc.StatusCode.FAILED_PRECONDITION, "sandbox is not 
ready"),
+            _Stream([], _RpcError(grpc.StatusCode.UNAVAILABLE, "exec relay 
closed")),
+        ]
+
+        with pytest.raises(SandboxError, match="may or may not have run") as 
error:
+            backend.run_command("box", "make install", timeout=5, 
max_output_bytes=100)
+
+        assert not isinstance(error.value, SandboxTerminalError)
+        assert client._stub.ExecSandbox.call_count == 2
+
+    def test_an_exec_the_gateway_drops_is_reported_not_retried(self):
+        backend, client = _backend()
+        client._stub.ExecSandbox.return_value = _Stream(
+            [],
+            _RpcError(
+                grpc.StatusCode.UNAVAILABLE, "exec relay closed before the 
command reported an exit status"
+            ),
+        )
+        client._stub.GetSandbox.side_effect = 
[_RpcError(grpc.StatusCode.UNAVAILABLE), _sandbox()]
+
+        with pytest.raises(SandboxError, match="may or may not have run") as 
error:
+            backend.run_command("box", "make install", timeout=5, 
max_output_bytes=100)
+
+        assert not isinstance(error.value, SandboxTerminalError)
+        client._stub.ExecSandbox.assert_called_once()
+
+    def 
test_a_dropped_exec_on_a_sandbox_that_does_not_come_back_is_terminal(self):
+        backend, client = _backend()
+        client._stub.ExecSandbox.return_value = _Stream([_stdout(b"x")])
+        client._stub.GetSandbox.return_value = 
_sandbox(openshell_pb2.SANDBOX_PHASE_ERROR)
+
+        with pytest.raises(SandboxTerminalError, match="SANDBOX_PHASE_ERROR"):
+            backend.run_command("box", "cmd", timeout=5, max_output_bytes=100)
+
+    @pytest.mark.parametrize(
+        ("code", "terminal"),
+        [
+            (grpc.StatusCode.NOT_FOUND, True),
+            (grpc.StatusCode.UNAUTHENTICATED, True),
+            (grpc.StatusCode.PERMISSION_DENIED, True),
+            (grpc.StatusCode.OUT_OF_RANGE, False),
+            (grpc.StatusCode.RESOURCE_EXHAUSTED, False),
+        ],
+    )
+    def test_gateway_errors_are_classified(self, code, terminal):
+        backend, client = _backend()
+        client._stub.ExecSandbox.side_effect = _RpcError(code, "detail")

Review Comment:
   Nothing here asserts `ExecSandbox` was called once. If `_exec_with_recovery` 
retried on every terminal status rather than only on "sandbox is not ready", 
this would still pass, since the second call raises the same error. `assert 
client._stub.ExecSandbox.call_count == 1` would catch that.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to