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


##########
providers/common/ai/src/airflow/providers/common/ai/sandbox/islo.py:
##########
@@ -0,0 +1,682 @@
+# 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.
+"""islo.dev microVM backend for 
:class:`~airflow.providers.common.ai.toolsets.sandbox.SandboxToolset`."""
+
+from __future__ import annotations
+
+import logging
+import math
+import shlex
+import time
+from contextlib import closing, contextmanager, suppress
+from typing import TYPE_CHECKING, Any, Literal, NoReturn
+
+from airflow.providers.common.ai.sandbox.base import (
+    SandboxBackend,
+    SandboxError,
+    SandboxExecResult,
+    SandboxFileTooLargeError,
+    SandboxTerminalError,
+    _check_export_deadline,
+    _export_deadline,
+    _new_sandbox_name,
+    _validate_positive_finite,
+)
+
+if TYPE_CHECKING:
+    from collections.abc import Generator, Iterator
+    from typing import BinaryIO
+
+    from islo import Islo
+
+    from airflow.providers.common.ai.sandbox.base import SandboxSpec
+
+log = logging.getLogger(__name__)
+
+# The vendor types ``status`` as a plain string on a model that allows extra
+# fields, so the vocabulary is open-ended. Track the statuses that mean "not
+# finished yet" instead of the ones that mean "finished": an unrecognised 
status
+# then reads as terminal, which surfaces a failure, rather than as 
still-running,
+# which would poll to the deadline and cost the agent its sandbox.
+_RUNNING_EXEC_STATUSES = frozenset({"pending", "queued", "starting", 
"running"})
+# Sandbox statuses that cannot serve a request. Same open vocabulary, opposite
+# bias: after a missing-file response the sandbox was reachable a moment ago, 
so
+# an unknown status reads as usable and only a known-dead one fails the task.
+_UNUSABLE_SANDBOX_STATUSES = frozenset({"stopping", "stopped", "deleting", 
"deleted", "error", "failed"})
+_AUTO_RESUME_POLICIES = frozenset({"never", "on_activity"})
+_POLL_INITIAL = 0.2
+_POLL_MAX = 2.0
+_POLL_BACKOFF = 1.5
+# The HTTP timeout of the requests that start and poll a command never drops
+# below this: a small command budget, or the last poll before the deadline, 
must
+# not become a one-second request budget that a slow API -- the first call 
after
+# an auto-resume, say -- turns into a task failure instead of a command result.
+_POLL_HTTP_TIMEOUT_MIN = 5.0
+_FILE_OP_TIMEOUT = 120.0
+# File-transfer responses that say nothing the model can act on: bad 
credentials,
+# a rate limit, or an overloaded API. Every other status is checked against the
+# sandbox before it is handed to the model; see ``_raise_file_op_error``.
+_TERMINAL_FILE_OP_STATUSES = frozenset({401, 403, 429})
+_API_MESSAGE_MAX_CHARS = 200
+# Measured against the compute API: each stream is capped at exactly this many
+# bytes with the tail kept, and one ``truncated`` flag covers both streams. The
+# wrapper never asks for more than this per stream, so the server's cap is not
+# reached by wrapper output and its flag stays a fallback.
+_SERVER_STREAM_CAP = 1024 * 1024
+_HELPER_OUTPUT_CAP = _SERVER_STREAM_CAP - 1
+# Runs the agent's command with each stream captured to a scratch file, then
+# emits only the last ``$2`` bytes of each. The tail is what the model needs (a
+# traceback and the exit status live at the end), and bounding inside the guest
+# keeps the transfer and the worker's copy at the caller's budget rather than
+# the server's 1 MiB.
+#
+# ``sh -c`` rather than a login shell: the spec's variables are the process
+# environment of every exec, and ``/etc/profile`` would run after them and win
+# for anything it also exports.
+#
+# Deliberately free of fifos, background jobs and ``wait``: a command that
+# backgrounds a process hands it the capture descriptor, so anything waiting 
for
+# end-of-input would block until that process exits -- ``sleep 20 & echo
+# started`` took 20s in a real microVM before this. Redirecting to files means
+# only the foreground command is waited on. The cost is that the scratch file
+# grows with total output, on the sandbox's own ephemeral disk.
+_COMMAND_WRAPPER = """\
+dir="${TMPDIR:-/tmp}/airflow-sandbox-$$"
+mkdir -m 700 "$dir" || exit 70
+trap 'rm -rf "$dir"' EXIT
+trap 'rm -rf "$dir"; exit 143' HUP INT TERM
+sh -c "$1" >"$dir/out" 2>"$dir/err"
+status=$?
+tail -c "$2" <"$dir/out"
+tail -c "$2" <"$dir/err" >&2
+exit "$status"
+"""
+
+
+@contextmanager
+def _translate_islo_errors(operation: str) -> Iterator[None]:
+    try:
+        yield
+    except SandboxError:
+        raise
+    except Exception as e:
+        try:
+            from islo.core.api_error import ApiError
+        except ImportError:
+            raise SandboxTerminalError(
+                'The Islo SDK is not installed. Install 
"apache-airflow-providers-common-ai[islo]".'
+            ) from e
+        if isinstance(e, ApiError):
+            status = f" (HTTP {e.status_code})" if e.status_code is not None 
else ""
+            raise SandboxTerminalError(f"Islo could not {operation}{status}.") 
from e
+        raise SandboxTerminalError(f"Islo could not {operation}: 
{type(e).__name__}.") from e
+
+
+def _raise_translated(error: Exception, operation: str) -> NoReturn:
+    """Raise ``error`` the way :func:`_translate_islo_errors` would have, for 
an error caught elsewhere."""
+    with _translate_islo_errors(operation):
+        raise error
+
+
+def _api_error_message(error: Exception) -> str:
+    """Return the ``message`` of an Islo error body, trimmed for a prompt, or 
``""``."""
+    body = getattr(error, "body", None)
+    message = body.get("message") if isinstance(body, dict) else None
+    if not isinstance(message, str):
+        return ""
+    message = message.strip().rstrip(".")
+    if len(message) > _API_MESSAGE_MAX_CHARS:
+        message = message[:_API_MESSAGE_MAX_CHARS] + "..."
+    return message
+
+
+def _is_transient_error(error: Exception) -> bool:
+    """Whether a failed call says nothing about the command: a 5xx, a 429, or 
no response at all."""
+    import httpx
+    from islo.core.api_error import ApiError
+
+    if isinstance(error, ApiError):
+        return error.status_code is None or error.status_code == 429 or 
error.status_code >= 500
+    return isinstance(error, httpx.TransportError)
+
+
+def _bound_result_stream(text: str, max_bytes: int, *, server_truncated: bool) 
-> tuple[str, bool]:
+    """
+    Trim one stream to ``max_bytes``, keeping the tail, and report whether 
bytes were dropped.
+
+    The sandbox is asked for one byte more than the budget, so a stream that
+    comes back over budget is the signal that the guest had more to give. The
+    server's own flag covers both streams at once, so it is only attributed to 
a
+    stream that sits at the server's cap.
+    """
+    encoded = text.encode("utf-8", errors="surrogatepass")
+    truncated = server_truncated and len(encoded) >= _SERVER_STREAM_CAP
+    if len(encoded) > max_bytes:
+        encoded = encoded[-max_bytes:]
+        truncated = True
+        # A byte-aligned cut usually lands mid-record, so drop the leading 
partial
+        # line -- but only while half the window survives, as Modal's 
``_drain``
+        # does. One line longer than the budget ends in its only newline, and
+        # dropping through it would leave nothing, which the model reads as 
"(no
+        # output)"; a partial line kept instead is still marked as cut.
+        newline = encoded.find(b"\n")
+        if newline != -1 and len(encoded) - (newline + 1) >= max_bytes // 2:
+            encoded = encoded[newline + 1 :]
+    return encoded.decode("utf-8", errors="replace"), truncated
+
+
+class IsloSandboxBackend(SandboxBackend):
+    """
+    Sandbox backend that runs agent commands in an `islo.dev 
<https://islo.dev>`__ microVM.
+
+    .. note::
+
+        Experimental: this can change or be removed in a minor release of this 
provider.
+        See :ref:`howto/stability`.
+
+    Islo is a hosted API with no local daemon or host-virtualization 
requirement,
+    so this backend works from an Airflow worker running in a container.
+
+    **Credentials are ambient.** The SDK reads ``ISLO_API_KEY``, and optionally
+    ``ISLO_BASE_URL`` and ``ISLO_COMPUTE_URL``, from the worker environment on 
first
+    use. Modal reads a ``modal`` connection first, but that connection type is 
owned by
+    the Modal provider; an ``islo`` connection type belongs in a future Islo 
provider,
+    not in this one.
+
+    File reads, writes and exports move file contents through Islo's native
+    streaming APIs. Everything else runs through a shell wrapper in the 
sandbox,
+    so the image needs ``sh``, ``mkdir`` and ``rm`` for the wrapper, ``tail`` 
to
+    bound command output, ``dirname`` to create a written file's parent
+    directory, ``stat`` to size an over-budget read and to check a file before 
it
+    is exported, and GNU ``find`` for directory listings. The server default
+    image and any Debian or Ubuntu based image provide them. Each command's
+    output is captured to a scratch file in the sandbox and only its last
+    ``max_output_bytes`` are returned, so the worker sees a bounded tail while
+    the sandbox's own ephemeral disk absorbs the rest.
+
+    Islo sets ``PATH`` for every command itself and drops a ``PATH`` given at
+    creation, so a spec that names it is refused rather than silently ignored.
+    ``SandboxSpec.owner`` is refused for the same reason: this backend keeps no
+    per-sandbox metadata, so a sandbox created here cannot be attached to 
later.
+
+    :param image: Sandbox image. ``None`` (default) uses the server default.
+    :param vcpus: Number of virtual CPUs. ``None`` uses the server default.
+    :param memory_mb: Memory in MB. ``None`` uses the server default.
+    :param pause_after_idle: Seconds without a command or file operation after
+        which the server pauses the microVM and releases its compute. ``None``
+        disables it. Default ``600``.
+    :param auto_resume: ``"on_activity"`` (default) resumes a paused sandbox on
+        the next command or file operation; ``"never"`` leaves it paused, and 
the
+        backend then treats a paused sandbox as unusable.
+    :param delete_after: Seconds after *creation* at which the server deletes 
the
+        sandbox whether or not it is in use. ``None`` disables it. Default
+        ``86400``.
+    """
+
+    name = "islo"
+
+    def __init__(
+        self,
+        *,
+        image: str | None = None,
+        vcpus: int | None = None,
+        memory_mb: int | None = None,
+        pause_after_idle: int | None = 600,
+        auto_resume: Literal["never", "on_activity"] = "on_activity",
+        delete_after: int | None = 86400,
+    ) -> None:
+        if pause_after_idle is not None:
+            _validate_positive_finite(pause_after_idle, "pause_after_idle")
+        if delete_after is not None:
+            _validate_positive_finite(delete_after, "delete_after")
+        if auto_resume not in _AUTO_RESUME_POLICIES:
+            raise ValueError(
+                f"auto_resume must be one of {sorted(_AUTO_RESUME_POLICIES)}, 
got {auto_resume!r}."
+            )
+        if vcpus is not None:
+            _validate_positive_finite(vcpus, "vcpus")
+        if memory_mb is not None:
+            _validate_positive_finite(memory_mb, "memory_mb")
+        if image == "":
+            raise ValueError("image must not be empty.")
+        self._image = image
+        self._vcpus = vcpus
+        self._memory_mb = memory_mb
+        self._pause_after_idle = pause_after_idle
+        self._auto_resume = auto_resume
+        self._delete_after = delete_after
+        self._client: Islo | None = None
+
+    def _get_client(self) -> Islo:
+        if self._client is not None:
+            return self._client
+        with _translate_islo_errors("initialize its client"):
+            from islo import Islo
+
+            self._client = Islo()

Review Comment:
   `Islo()` doesn't raise when `ISLO_API_KEY` is unset. It builds a client that 
sends no `Authorization` header, so the first sign of a missing key is `create` 
failing with a bare `Islo could not create a sandbox (HTTP ...).` Now that the 
environment is the only credential source, that's the most likely first-run 
mistake. Checking for the variable here and raising a `SandboxTerminalError` 
that names it would save a trip to the SDK source.



##########
providers/common/ai/src/airflow/providers/common/ai/sandbox/islo.py:
##########
@@ -0,0 +1,682 @@
+# 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.
+"""islo.dev microVM backend for 
:class:`~airflow.providers.common.ai.toolsets.sandbox.SandboxToolset`."""
+
+from __future__ import annotations
+
+import logging
+import math
+import shlex
+import time
+from contextlib import closing, contextmanager, suppress
+from typing import TYPE_CHECKING, Any, Literal, NoReturn
+
+from airflow.providers.common.ai.sandbox.base import (
+    SandboxBackend,
+    SandboxError,
+    SandboxExecResult,
+    SandboxFileTooLargeError,
+    SandboxTerminalError,
+    _check_export_deadline,
+    _export_deadline,
+    _new_sandbox_name,
+    _validate_positive_finite,
+)
+
+if TYPE_CHECKING:
+    from collections.abc import Generator, Iterator
+    from typing import BinaryIO
+
+    from islo import Islo
+
+    from airflow.providers.common.ai.sandbox.base import SandboxSpec
+
+log = logging.getLogger(__name__)
+
+# The vendor types ``status`` as a plain string on a model that allows extra
+# fields, so the vocabulary is open-ended. Track the statuses that mean "not
+# finished yet" instead of the ones that mean "finished": an unrecognised 
status
+# then reads as terminal, which surfaces a failure, rather than as 
still-running,
+# which would poll to the deadline and cost the agent its sandbox.
+_RUNNING_EXEC_STATUSES = frozenset({"pending", "queued", "starting", 
"running"})
+# Sandbox statuses that cannot serve a request. Same open vocabulary, opposite
+# bias: after a missing-file response the sandbox was reachable a moment ago, 
so
+# an unknown status reads as usable and only a known-dead one fails the task.
+_UNUSABLE_SANDBOX_STATUSES = frozenset({"stopping", "stopped", "deleting", 
"deleted", "error", "failed"})
+_AUTO_RESUME_POLICIES = frozenset({"never", "on_activity"})
+_POLL_INITIAL = 0.2
+_POLL_MAX = 2.0
+_POLL_BACKOFF = 1.5
+# The HTTP timeout of the requests that start and poll a command never drops
+# below this: a small command budget, or the last poll before the deadline, 
must
+# not become a one-second request budget that a slow API -- the first call 
after
+# an auto-resume, say -- turns into a task failure instead of a command result.
+_POLL_HTTP_TIMEOUT_MIN = 5.0
+_FILE_OP_TIMEOUT = 120.0
+# File-transfer responses that say nothing the model can act on: bad 
credentials,
+# a rate limit, or an overloaded API. Every other status is checked against the
+# sandbox before it is handed to the model; see ``_raise_file_op_error``.
+_TERMINAL_FILE_OP_STATUSES = frozenset({401, 403, 429})
+_API_MESSAGE_MAX_CHARS = 200
+# Measured against the compute API: each stream is capped at exactly this many
+# bytes with the tail kept, and one ``truncated`` flag covers both streams. The
+# wrapper never asks for more than this per stream, so the server's cap is not
+# reached by wrapper output and its flag stays a fallback.
+_SERVER_STREAM_CAP = 1024 * 1024
+_HELPER_OUTPUT_CAP = _SERVER_STREAM_CAP - 1
+# Runs the agent's command with each stream captured to a scratch file, then
+# emits only the last ``$2`` bytes of each. The tail is what the model needs (a
+# traceback and the exit status live at the end), and bounding inside the guest
+# keeps the transfer and the worker's copy at the caller's budget rather than
+# the server's 1 MiB.
+#
+# ``sh -c`` rather than a login shell: the spec's variables are the process
+# environment of every exec, and ``/etc/profile`` would run after them and win
+# for anything it also exports.
+#
+# Deliberately free of fifos, background jobs and ``wait``: a command that
+# backgrounds a process hands it the capture descriptor, so anything waiting 
for
+# end-of-input would block until that process exits -- ``sleep 20 & echo
+# started`` took 20s in a real microVM before this. Redirecting to files means
+# only the foreground command is waited on. The cost is that the scratch file
+# grows with total output, on the sandbox's own ephemeral disk.
+_COMMAND_WRAPPER = """\
+dir="${TMPDIR:-/tmp}/airflow-sandbox-$$"
+mkdir -m 700 "$dir" || exit 70
+trap 'rm -rf "$dir"' EXIT
+trap 'rm -rf "$dir"; exit 143' HUP INT TERM
+sh -c "$1" >"$dir/out" 2>"$dir/err"
+status=$?
+tail -c "$2" <"$dir/out"
+tail -c "$2" <"$dir/err" >&2
+exit "$status"
+"""
+
+
+@contextmanager
+def _translate_islo_errors(operation: str) -> Iterator[None]:
+    try:
+        yield
+    except SandboxError:
+        raise
+    except Exception as e:
+        try:
+            from islo.core.api_error import ApiError
+        except ImportError:
+            raise SandboxTerminalError(
+                'The Islo SDK is not installed. Install 
"apache-airflow-providers-common-ai[islo]".'
+            ) from e
+        if isinstance(e, ApiError):
+            status = f" (HTTP {e.status_code})" if e.status_code is not None 
else ""
+            raise SandboxTerminalError(f"Islo could not {operation}{status}.") 
from e
+        raise SandboxTerminalError(f"Islo could not {operation}: 
{type(e).__name__}.") from e
+
+
+def _raise_translated(error: Exception, operation: str) -> NoReturn:
+    """Raise ``error`` the way :func:`_translate_islo_errors` would have, for 
an error caught elsewhere."""
+    with _translate_islo_errors(operation):
+        raise error
+
+
+def _api_error_message(error: Exception) -> str:
+    """Return the ``message`` of an Islo error body, trimmed for a prompt, or 
``""``."""
+    body = getattr(error, "body", None)
+    message = body.get("message") if isinstance(body, dict) else None
+    if not isinstance(message, str):
+        return ""
+    message = message.strip().rstrip(".")
+    if len(message) > _API_MESSAGE_MAX_CHARS:
+        message = message[:_API_MESSAGE_MAX_CHARS] + "..."
+    return message
+
+
+def _is_transient_error(error: Exception) -> bool:
+    """Whether a failed call says nothing about the command: a 5xx, a 429, or 
no response at all."""
+    import httpx
+    from islo.core.api_error import ApiError
+
+    if isinstance(error, ApiError):
+        return error.status_code is None or error.status_code == 429 or 
error.status_code >= 500
+    return isinstance(error, httpx.TransportError)
+
+
+def _bound_result_stream(text: str, max_bytes: int, *, server_truncated: bool) 
-> tuple[str, bool]:
+    """
+    Trim one stream to ``max_bytes``, keeping the tail, and report whether 
bytes were dropped.
+
+    The sandbox is asked for one byte more than the budget, so a stream that
+    comes back over budget is the signal that the guest had more to give. The
+    server's own flag covers both streams at once, so it is only attributed to 
a
+    stream that sits at the server's cap.
+    """
+    encoded = text.encode("utf-8", errors="surrogatepass")
+    truncated = server_truncated and len(encoded) >= _SERVER_STREAM_CAP

Review Comment:
   Can `server_truncated` change the outcome here? `run_command` clamps the 
budget to `_SERVER_STREAM_CAP - 1`, so any stream long enough to satisfy 
`len(encoded) >= _SERVER_STREAM_CAP` is also over `max_bytes` and gets 
`truncated = True` from the branch below anyway. Dropping the flag term leaves 
`test_the_server_flag_marks_only_a_stream_at_the_server_cap` passing, so 
nothing pins it. Either remove the plumbing, or test `_bound_result_stream` 
directly with a stream at the cap and a budget above it.



##########
providers/common/ai/docs/sandbox/backends.rst:
##########
@@ -301,25 +385,28 @@ do not change. Four behaviours do, so read them before 
assuming the same Dag
 behaves identically everywhere:
 
 - **CPU.** ``sbx`` gives a sandbox every host CPU; Modal defaults to a request 
of
-  0.125 of one, so set ``cpu``; OpenSandbox takes ``cpu`` as a limit the 
server enforces.
+  0.125 of one, so set ``cpu``; OpenSandbox takes ``cpu`` as a limit the server
+  enforces; Islo uses the server default unless you set ``vcpus``.
 - **Egress allowlists.** ``sbx`` enforces ``allow_egress_to`` at the host 
policy
   layer; Modal matches TLS handshake names, which is weaker and has to be opted
   into; OpenSandbox enforces it in an egress sidecar, and the backend reads the
-  enforced policy back rather than trusting the create request. 
``allow_egress_to_cidrs``
-  is enforced at the address layer on Modal, refused on ``sbx``, and refused by
-  OpenSandbox because its SDK cannot prove that the sidecar is running in the
+  enforced policy back rather than trusting the create request; Islo has no
+  per-host form at all and refuses the spec rather than provisioning something
+  weaker. ``allow_egress_to_cidrs`` is enforced at the address layer on Modal,
+  refused on ``sbx`` and Islo, which have no per-sandbox address rule, and 
refused
+  by OpenSandbox because its SDK cannot prove that the sidecar is running in 
the
   ``dns+nft`` mode required for CIDR enforcement.
-- **Command timeouts.** A timeout destroys an ``sbx`` sandbox and its files;
-  Modal and a server-enforced OpenSandbox timeout preserve the sandbox and 
files.
-  OpenSandbox destroys it only if the command event stream itself stalls past 
the
-  client-side grace period.
+- **Command timeouts.** A timeout destroys an ``sbx`` or Islo sandbox and its
+  files; Modal and a server-enforced OpenSandbox timeout preserve the sandbox 
and
+  files. OpenSandbox destroys it only if the command event stream itself stalls
+  past the client-side grace period.
 - **Symlinks.** ``write_file`` through a symlink follows the link on ``sbx`` 
and
-  replaces it on Modal and OpenSandbox.
+  replaces it on Modal, OpenSandbox and Islo.

Review Comment:
   Was the Islo half of this measured? `write_file` hands the bytes to Islo's 
upload API, so following or replacing a symlink is the server's choice, and 
nothing in the unit tests, the system test or the threads here exercises it. If 
not, one probe on a live microVM, or dropping Islo from this bullet, would keep 
it in line with the other Islo claims, which all have a measurement behind them.



##########
providers/common/ai/tests/unit/common/ai/sandbox/test_islo.py:
##########
@@ -0,0 +1,1097 @@
+# 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 io
+import os
+import signal
+import subprocess
+import time
+from pathlib import Path
+from types import SimpleNamespace
+from unittest import mock
+
+import pytest
+
+pytest.importorskip("islo")
+
+import httpx
+from islo.core.api_error import ApiError
+from islo.errors import NotFoundError
+from islo.sandboxes.client import SandboxesClient
+
+from airflow.providers.common.ai.sandbox.base import (
+    SandboxError,
+    SandboxExecResult,
+    SandboxFileTooLargeError,
+    SandboxSpec,
+    SandboxTerminalError,
+)
+from airflow.providers.common.ai.sandbox.islo import (
+    _COMMAND_WRAPPER,
+    _SERVER_STREAM_CAP,
+    IsloSandboxBackend,
+    _bound_result_stream,
+)
+
+_MODULE = "airflow.providers.common.ai.sandbox.islo"
+_ISLO_PATH = "islo.Islo"
+
+
+def _exec_result(status="completed", exit_code=0, stdout="", stderr="", 
truncated=False):
+    return SimpleNamespace(
+        status=status, exit_code=exit_code, stdout=stdout, stderr=stderr, 
truncated=truncated
+    )
+
+
+def _sandbox_info(name="box-1", status="running", deleted_at=None):
+    return SimpleNamespace(name=name, status=status, deleted_at=deleted_at)
+
+
+def _export_check(size):
+    """What the guest's export check prints for a regular file of ``size`` 
bytes."""
+    return _exec_result(stdout=f"\nairflow-export-size:{size}\n")
+
+
+def _backend_with_client(**kwargs) -> tuple[IsloSandboxBackend, 
mock.MagicMock]:
+    backend = IsloSandboxBackend(**kwargs)
+    client = mock.MagicMock(spec=["sandboxes"])
+    client.sandboxes = mock.create_autospec(SandboxesClient, instance=True)
+    client.sandboxes.exec_in_sandbox.return_value = 
SimpleNamespace(exec_id="exec-1")
+    client.sandboxes.create_sandbox.return_value = _sandbox_info()
+    client.sandboxes.get_sandbox.return_value = _sandbox_info()
+    client.sandboxes.get_exec_result.return_value = _exec_result()
+    backend._client = client
+    return backend, client
+
+
+class TestCredentials:
+    @mock.patch(_ISLO_PATH, autospec=True)
+    def test_client_comes_from_the_sdk_environment(self, islo):
+        client = IsloSandboxBackend()._get_client()
+
+        islo.assert_called_once_with()
+        assert client is islo.return_value
+
+    @mock.patch(_ISLO_PATH, autospec=True)
+    def test_client_is_resolved_once_and_cached(self, islo):
+        backend = IsloSandboxBackend()
+
+        backend._get_client()
+        backend._get_client()
+
+        islo.assert_called_once_with()
+
+    def test_missing_sdk_error_is_actionable(self):
+        real_import = builtins.__import__
+
+        def blocked_import(name, *args, **kwargs):
+            if name.startswith("islo"):
+                raise ImportError("blocked for test")
+            return real_import(name, *args, **kwargs)
+
+        backend = IsloSandboxBackend()
+        with mock.patch("builtins.__import__", side_effect=blocked_import):
+            with pytest.raises(SandboxTerminalError, match=r"\[islo\]"):
+                backend.create()
+
+
[email protected](
+    ("kwargs", "message"),
+    [
+        ({"image": ""}, "image"),
+        ({"vcpus": 0}, "vcpus"),
+        ({"memory_mb": 0}, "memory_mb"),
+        ({"delete_after": 0}, "delete_after"),
+        ({"pause_after_idle": 0}, "pause_after_idle"),
+        ({"auto_resume": "sometimes"}, "auto_resume"),
+    ],
+)
+def test_constructor_rejects_invalid_values(kwargs, message):
+    with pytest.raises(ValueError, match=message):
+        IsloSandboxBackend(**kwargs)
+
+
+class TestCreate:
+    def test_refuses_an_owner_it_could_not_be_attached_by(self):
+        backend, client = _backend_with_client()
+
+        with pytest.raises(SandboxTerminalError, match="owner"):
+            backend.create(spec=SandboxSpec(owner="example_dag/manual__1"))
+
+        client.sandboxes.create_sandbox.assert_not_called()
+
+    def test_refuses_a_per_domain_egress_allowlist(self):

Review Comment:
   This is the only refusal test without 
`client.sandboxes.create_sandbox.assert_not_called()`, so a check that ran 
after provisioning would still pass. The four refusal tests differ only in the 
spec and the `match`, so one parametrized test with that assertion would cover 
them all.



##########
providers/common/ai/src/airflow/providers/common/ai/sandbox/islo.py:
##########
@@ -0,0 +1,682 @@
+# 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.
+"""islo.dev microVM backend for 
:class:`~airflow.providers.common.ai.toolsets.sandbox.SandboxToolset`."""
+
+from __future__ import annotations
+
+import logging
+import math
+import shlex
+import time
+from contextlib import closing, contextmanager, suppress
+from typing import TYPE_CHECKING, Any, Literal, NoReturn
+
+from airflow.providers.common.ai.sandbox.base import (
+    SandboxBackend,
+    SandboxError,
+    SandboxExecResult,
+    SandboxFileTooLargeError,
+    SandboxTerminalError,
+    _check_export_deadline,
+    _export_deadline,
+    _new_sandbox_name,
+    _validate_positive_finite,
+)
+
+if TYPE_CHECKING:
+    from collections.abc import Generator, Iterator
+    from typing import BinaryIO
+
+    from islo import Islo
+
+    from airflow.providers.common.ai.sandbox.base import SandboxSpec
+
+log = logging.getLogger(__name__)
+
+# The vendor types ``status`` as a plain string on a model that allows extra
+# fields, so the vocabulary is open-ended. Track the statuses that mean "not
+# finished yet" instead of the ones that mean "finished": an unrecognised 
status
+# then reads as terminal, which surfaces a failure, rather than as 
still-running,
+# which would poll to the deadline and cost the agent its sandbox.
+_RUNNING_EXEC_STATUSES = frozenset({"pending", "queued", "starting", 
"running"})
+# Sandbox statuses that cannot serve a request. Same open vocabulary, opposite
+# bias: after a missing-file response the sandbox was reachable a moment ago, 
so
+# an unknown status reads as usable and only a known-dead one fails the task.
+_UNUSABLE_SANDBOX_STATUSES = frozenset({"stopping", "stopped", "deleting", 
"deleted", "error", "failed"})
+_AUTO_RESUME_POLICIES = frozenset({"never", "on_activity"})
+_POLL_INITIAL = 0.2
+_POLL_MAX = 2.0
+_POLL_BACKOFF = 1.5
+# The HTTP timeout of the requests that start and poll a command never drops
+# below this: a small command budget, or the last poll before the deadline, 
must
+# not become a one-second request budget that a slow API -- the first call 
after
+# an auto-resume, say -- turns into a task failure instead of a command result.
+_POLL_HTTP_TIMEOUT_MIN = 5.0
+_FILE_OP_TIMEOUT = 120.0
+# File-transfer responses that say nothing the model can act on: bad 
credentials,
+# a rate limit, or an overloaded API. Every other status is checked against the
+# sandbox before it is handed to the model; see ``_raise_file_op_error``.
+_TERMINAL_FILE_OP_STATUSES = frozenset({401, 403, 429})
+_API_MESSAGE_MAX_CHARS = 200
+# Measured against the compute API: each stream is capped at exactly this many
+# bytes with the tail kept, and one ``truncated`` flag covers both streams. The
+# wrapper never asks for more than this per stream, so the server's cap is not
+# reached by wrapper output and its flag stays a fallback.
+_SERVER_STREAM_CAP = 1024 * 1024
+_HELPER_OUTPUT_CAP = _SERVER_STREAM_CAP - 1
+# Runs the agent's command with each stream captured to a scratch file, then
+# emits only the last ``$2`` bytes of each. The tail is what the model needs (a
+# traceback and the exit status live at the end), and bounding inside the guest
+# keeps the transfer and the worker's copy at the caller's budget rather than
+# the server's 1 MiB.
+#
+# ``sh -c`` rather than a login shell: the spec's variables are the process
+# environment of every exec, and ``/etc/profile`` would run after them and win
+# for anything it also exports.
+#
+# Deliberately free of fifos, background jobs and ``wait``: a command that
+# backgrounds a process hands it the capture descriptor, so anything waiting 
for
+# end-of-input would block until that process exits -- ``sleep 20 & echo
+# started`` took 20s in a real microVM before this. Redirecting to files means
+# only the foreground command is waited on. The cost is that the scratch file
+# grows with total output, on the sandbox's own ephemeral disk.
+_COMMAND_WRAPPER = """\
+dir="${TMPDIR:-/tmp}/airflow-sandbox-$$"
+mkdir -m 700 "$dir" || exit 70
+trap 'rm -rf "$dir"' EXIT
+trap 'rm -rf "$dir"; exit 143' HUP INT TERM
+sh -c "$1" >"$dir/out" 2>"$dir/err"
+status=$?
+tail -c "$2" <"$dir/out"
+tail -c "$2" <"$dir/err" >&2
+exit "$status"
+"""
+
+
+@contextmanager
+def _translate_islo_errors(operation: str) -> Iterator[None]:
+    try:
+        yield
+    except SandboxError:
+        raise
+    except Exception as e:
+        try:
+            from islo.core.api_error import ApiError
+        except ImportError:
+            raise SandboxTerminalError(
+                'The Islo SDK is not installed. Install 
"apache-airflow-providers-common-ai[islo]".'
+            ) from e
+        if isinstance(e, ApiError):
+            status = f" (HTTP {e.status_code})" if e.status_code is not None 
else ""
+            raise SandboxTerminalError(f"Islo could not {operation}{status}.") 
from e
+        raise SandboxTerminalError(f"Islo could not {operation}: 
{type(e).__name__}.") from e
+
+
+def _raise_translated(error: Exception, operation: str) -> NoReturn:
+    """Raise ``error`` the way :func:`_translate_islo_errors` would have, for 
an error caught elsewhere."""
+    with _translate_islo_errors(operation):
+        raise error
+
+
+def _api_error_message(error: Exception) -> str:
+    """Return the ``message`` of an Islo error body, trimmed for a prompt, or 
``""``."""
+    body = getattr(error, "body", None)
+    message = body.get("message") if isinstance(body, dict) else None
+    if not isinstance(message, str):
+        return ""
+    message = message.strip().rstrip(".")
+    if len(message) > _API_MESSAGE_MAX_CHARS:
+        message = message[:_API_MESSAGE_MAX_CHARS] + "..."
+    return message
+
+
+def _is_transient_error(error: Exception) -> bool:
+    """Whether a failed call says nothing about the command: a 5xx, a 429, or 
no response at all."""
+    import httpx
+    from islo.core.api_error import ApiError
+
+    if isinstance(error, ApiError):
+        return error.status_code is None or error.status_code == 429 or 
error.status_code >= 500
+    return isinstance(error, httpx.TransportError)
+
+
+def _bound_result_stream(text: str, max_bytes: int, *, server_truncated: bool) 
-> tuple[str, bool]:
+    """
+    Trim one stream to ``max_bytes``, keeping the tail, and report whether 
bytes were dropped.
+
+    The sandbox is asked for one byte more than the budget, so a stream that
+    comes back over budget is the signal that the guest had more to give. The
+    server's own flag covers both streams at once, so it is only attributed to 
a
+    stream that sits at the server's cap.
+    """
+    encoded = text.encode("utf-8", errors="surrogatepass")
+    truncated = server_truncated and len(encoded) >= _SERVER_STREAM_CAP
+    if len(encoded) > max_bytes:
+        encoded = encoded[-max_bytes:]
+        truncated = True
+        # A byte-aligned cut usually lands mid-record, so drop the leading 
partial
+        # line -- but only while half the window survives, as Modal's 
``_drain``
+        # does. One line longer than the budget ends in its only newline, and
+        # dropping through it would leave nothing, which the model reads as 
"(no
+        # output)"; a partial line kept instead is still marked as cut.
+        newline = encoded.find(b"\n")
+        if newline != -1 and len(encoded) - (newline + 1) >= max_bytes // 2:
+            encoded = encoded[newline + 1 :]
+    return encoded.decode("utf-8", errors="replace"), truncated
+
+
+class IsloSandboxBackend(SandboxBackend):
+    """
+    Sandbox backend that runs agent commands in an `islo.dev 
<https://islo.dev>`__ microVM.
+
+    .. note::
+
+        Experimental: this can change or be removed in a minor release of this 
provider.
+        See :ref:`howto/stability`.
+
+    Islo is a hosted API with no local daemon or host-virtualization 
requirement,
+    so this backend works from an Airflow worker running in a container.
+
+    **Credentials are ambient.** The SDK reads ``ISLO_API_KEY``, and optionally
+    ``ISLO_BASE_URL`` and ``ISLO_COMPUTE_URL``, from the worker environment on 
first
+    use. Modal reads a ``modal`` connection first, but that connection type is 
owned by
+    the Modal provider; an ``islo`` connection type belongs in a future Islo 
provider,
+    not in this one.
+
+    File reads, writes and exports move file contents through Islo's native
+    streaming APIs. Everything else runs through a shell wrapper in the 
sandbox,
+    so the image needs ``sh``, ``mkdir`` and ``rm`` for the wrapper, ``tail`` 
to
+    bound command output, ``dirname`` to create a written file's parent
+    directory, ``stat`` to size an over-budget read and to check a file before 
it
+    is exported, and GNU ``find`` for directory listings. The server default
+    image and any Debian or Ubuntu based image provide them. Each command's
+    output is captured to a scratch file in the sandbox and only its last
+    ``max_output_bytes`` are returned, so the worker sees a bounded tail while
+    the sandbox's own ephemeral disk absorbs the rest.
+
+    Islo sets ``PATH`` for every command itself and drops a ``PATH`` given at
+    creation, so a spec that names it is refused rather than silently ignored.
+    ``SandboxSpec.owner`` is refused for the same reason: this backend keeps no
+    per-sandbox metadata, so a sandbox created here cannot be attached to 
later.
+
+    :param image: Sandbox image. ``None`` (default) uses the server default.
+    :param vcpus: Number of virtual CPUs. ``None`` uses the server default.
+    :param memory_mb: Memory in MB. ``None`` uses the server default.
+    :param pause_after_idle: Seconds without a command or file operation after
+        which the server pauses the microVM and releases its compute. ``None``
+        disables it. Default ``600``.
+    :param auto_resume: ``"on_activity"`` (default) resumes a paused sandbox on
+        the next command or file operation; ``"never"`` leaves it paused, and 
the
+        backend then treats a paused sandbox as unusable.
+    :param delete_after: Seconds after *creation* at which the server deletes 
the
+        sandbox whether or not it is in use. ``None`` disables it. Default
+        ``86400``.
+    """
+
+    name = "islo"
+
+    def __init__(
+        self,
+        *,
+        image: str | None = None,
+        vcpus: int | None = None,
+        memory_mb: int | None = None,
+        pause_after_idle: int | None = 600,
+        auto_resume: Literal["never", "on_activity"] = "on_activity",
+        delete_after: int | None = 86400,
+    ) -> None:
+        if pause_after_idle is not None:
+            _validate_positive_finite(pause_after_idle, "pause_after_idle")
+        if delete_after is not None:
+            _validate_positive_finite(delete_after, "delete_after")
+        if auto_resume not in _AUTO_RESUME_POLICIES:
+            raise ValueError(
+                f"auto_resume must be one of {sorted(_AUTO_RESUME_POLICIES)}, 
got {auto_resume!r}."
+            )
+        if vcpus is not None:
+            _validate_positive_finite(vcpus, "vcpus")
+        if memory_mb is not None:
+            _validate_positive_finite(memory_mb, "memory_mb")
+        if image == "":
+            raise ValueError("image must not be empty.")
+        self._image = image
+        self._vcpus = vcpus
+        self._memory_mb = memory_mb
+        self._pause_after_idle = pause_after_idle
+        self._auto_resume = auto_resume
+        self._delete_after = delete_after
+        self._client: Islo | None = None
+
+    def _get_client(self) -> Islo:
+        if self._client is not None:
+            return self._client
+        with _translate_islo_errors("initialize its client"):
+            from islo import Islo
+
+            self._client = Islo()
+        return self._client
+
+    @staticmethod
+    def _request_options(
+        *, timeout: float, chunk_size: int | None = None, max_retries: int | 
None = None
+    ) -> dict[str, int]:
+        # ``max_retries`` is left to the SDK's default unless asked for: 
passing
+        # 0 would switch off the two transport retries it does on its own.
+        options = {"timeout_in_seconds": max(1, math.ceil(timeout))}
+        if chunk_size is not None:
+            options["chunk_size"] = chunk_size
+        if max_retries is not None:
+            options["max_retries"] = max_retries
+        return options
+
+    def _ensure_sandbox_usable(self, info: Any) -> None:

Review Comment:
   `SandboxResponse` declares `status` and `name` as required, and 
`ExecResultResponse` does the same for `status`, `stdout`, `stderr` and 
`truncated`, so `info: Any` with `getattr(..., None)` here (and `_await_exec -> 
Any`, `getattr(result, "truncated", False)` at line 458) hides those reads from 
mypy for no gain. Importing the two models under `TYPE_CHECKING` and reading 
the attributes directly would turn a field rename in the SDK into a type error 
rather than a silent `None`.



##########
providers/common/ai/src/airflow/providers/common/ai/sandbox/islo.py:
##########
@@ -0,0 +1,682 @@
+# 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.
+"""islo.dev microVM backend for 
:class:`~airflow.providers.common.ai.toolsets.sandbox.SandboxToolset`."""
+
+from __future__ import annotations
+
+import logging
+import math
+import shlex
+import time
+from contextlib import closing, contextmanager, suppress
+from typing import TYPE_CHECKING, Any, Literal, NoReturn
+
+from airflow.providers.common.ai.sandbox.base import (
+    SandboxBackend,
+    SandboxError,
+    SandboxExecResult,
+    SandboxFileTooLargeError,
+    SandboxTerminalError,
+    _check_export_deadline,
+    _export_deadline,
+    _new_sandbox_name,
+    _validate_positive_finite,
+)
+
+if TYPE_CHECKING:
+    from collections.abc import Generator, Iterator
+    from typing import BinaryIO
+
+    from islo import Islo
+
+    from airflow.providers.common.ai.sandbox.base import SandboxSpec
+
+log = logging.getLogger(__name__)
+
+# The vendor types ``status`` as a plain string on a model that allows extra
+# fields, so the vocabulary is open-ended. Track the statuses that mean "not
+# finished yet" instead of the ones that mean "finished": an unrecognised 
status
+# then reads as terminal, which surfaces a failure, rather than as 
still-running,
+# which would poll to the deadline and cost the agent its sandbox.
+_RUNNING_EXEC_STATUSES = frozenset({"pending", "queued", "starting", 
"running"})
+# Sandbox statuses that cannot serve a request. Same open vocabulary, opposite
+# bias: after a missing-file response the sandbox was reachable a moment ago, 
so
+# an unknown status reads as usable and only a known-dead one fails the task.
+_UNUSABLE_SANDBOX_STATUSES = frozenset({"stopping", "stopped", "deleting", 
"deleted", "error", "failed"})
+_AUTO_RESUME_POLICIES = frozenset({"never", "on_activity"})
+_POLL_INITIAL = 0.2
+_POLL_MAX = 2.0
+_POLL_BACKOFF = 1.5
+# The HTTP timeout of the requests that start and poll a command never drops
+# below this: a small command budget, or the last poll before the deadline, 
must
+# not become a one-second request budget that a slow API -- the first call 
after
+# an auto-resume, say -- turns into a task failure instead of a command result.
+_POLL_HTTP_TIMEOUT_MIN = 5.0
+_FILE_OP_TIMEOUT = 120.0
+# File-transfer responses that say nothing the model can act on: bad 
credentials,
+# a rate limit, or an overloaded API. Every other status is checked against the
+# sandbox before it is handed to the model; see ``_raise_file_op_error``.
+_TERMINAL_FILE_OP_STATUSES = frozenset({401, 403, 429})
+_API_MESSAGE_MAX_CHARS = 200
+# Measured against the compute API: each stream is capped at exactly this many
+# bytes with the tail kept, and one ``truncated`` flag covers both streams. The
+# wrapper never asks for more than this per stream, so the server's cap is not
+# reached by wrapper output and its flag stays a fallback.
+_SERVER_STREAM_CAP = 1024 * 1024
+_HELPER_OUTPUT_CAP = _SERVER_STREAM_CAP - 1
+# Runs the agent's command with each stream captured to a scratch file, then
+# emits only the last ``$2`` bytes of each. The tail is what the model needs (a
+# traceback and the exit status live at the end), and bounding inside the guest
+# keeps the transfer and the worker's copy at the caller's budget rather than
+# the server's 1 MiB.
+#
+# ``sh -c`` rather than a login shell: the spec's variables are the process
+# environment of every exec, and ``/etc/profile`` would run after them and win
+# for anything it also exports.
+#
+# Deliberately free of fifos, background jobs and ``wait``: a command that
+# backgrounds a process hands it the capture descriptor, so anything waiting 
for
+# end-of-input would block until that process exits -- ``sleep 20 & echo
+# started`` took 20s in a real microVM before this. Redirecting to files means
+# only the foreground command is waited on. The cost is that the scratch file
+# grows with total output, on the sandbox's own ephemeral disk.
+_COMMAND_WRAPPER = """\
+dir="${TMPDIR:-/tmp}/airflow-sandbox-$$"
+mkdir -m 700 "$dir" || exit 70
+trap 'rm -rf "$dir"' EXIT
+trap 'rm -rf "$dir"; exit 143' HUP INT TERM
+sh -c "$1" >"$dir/out" 2>"$dir/err"
+status=$?
+tail -c "$2" <"$dir/out"
+tail -c "$2" <"$dir/err" >&2
+exit "$status"
+"""
+
+
+@contextmanager
+def _translate_islo_errors(operation: str) -> Iterator[None]:
+    try:
+        yield
+    except SandboxError:
+        raise
+    except Exception as e:
+        try:
+            from islo.core.api_error import ApiError
+        except ImportError:
+            raise SandboxTerminalError(
+                'The Islo SDK is not installed. Install 
"apache-airflow-providers-common-ai[islo]".'
+            ) from e
+        if isinstance(e, ApiError):
+            status = f" (HTTP {e.status_code})" if e.status_code is not None 
else ""
+            raise SandboxTerminalError(f"Islo could not {operation}{status}.") 
from e
+        raise SandboxTerminalError(f"Islo could not {operation}: 
{type(e).__name__}.") from e
+
+
+def _raise_translated(error: Exception, operation: str) -> NoReturn:
+    """Raise ``error`` the way :func:`_translate_islo_errors` would have, for 
an error caught elsewhere."""
+    with _translate_islo_errors(operation):
+        raise error
+
+
+def _api_error_message(error: Exception) -> str:
+    """Return the ``message`` of an Islo error body, trimmed for a prompt, or 
``""``."""
+    body = getattr(error, "body", None)
+    message = body.get("message") if isinstance(body, dict) else None
+    if not isinstance(message, str):
+        return ""
+    message = message.strip().rstrip(".")
+    if len(message) > _API_MESSAGE_MAX_CHARS:
+        message = message[:_API_MESSAGE_MAX_CHARS] + "..."
+    return message
+
+
+def _is_transient_error(error: Exception) -> bool:
+    """Whether a failed call says nothing about the command: a 5xx, a 429, or 
no response at all."""
+    import httpx
+    from islo.core.api_error import ApiError
+
+    if isinstance(error, ApiError):
+        return error.status_code is None or error.status_code == 429 or 
error.status_code >= 500
+    return isinstance(error, httpx.TransportError)
+
+
+def _bound_result_stream(text: str, max_bytes: int, *, server_truncated: bool) 
-> tuple[str, bool]:
+    """
+    Trim one stream to ``max_bytes``, keeping the tail, and report whether 
bytes were dropped.
+
+    The sandbox is asked for one byte more than the budget, so a stream that
+    comes back over budget is the signal that the guest had more to give. The
+    server's own flag covers both streams at once, so it is only attributed to 
a
+    stream that sits at the server's cap.
+    """
+    encoded = text.encode("utf-8", errors="surrogatepass")
+    truncated = server_truncated and len(encoded) >= _SERVER_STREAM_CAP
+    if len(encoded) > max_bytes:
+        encoded = encoded[-max_bytes:]
+        truncated = True
+        # A byte-aligned cut usually lands mid-record, so drop the leading 
partial
+        # line -- but only while half the window survives, as Modal's 
``_drain``
+        # does. One line longer than the budget ends in its only newline, and
+        # dropping through it would leave nothing, which the model reads as 
"(no
+        # output)"; a partial line kept instead is still marked as cut.
+        newline = encoded.find(b"\n")
+        if newline != -1 and len(encoded) - (newline + 1) >= max_bytes // 2:
+            encoded = encoded[newline + 1 :]
+    return encoded.decode("utf-8", errors="replace"), truncated
+
+
+class IsloSandboxBackend(SandboxBackend):
+    """
+    Sandbox backend that runs agent commands in an `islo.dev 
<https://islo.dev>`__ microVM.
+
+    .. note::
+
+        Experimental: this can change or be removed in a minor release of this 
provider.
+        See :ref:`howto/stability`.
+
+    Islo is a hosted API with no local daemon or host-virtualization 
requirement,
+    so this backend works from an Airflow worker running in a container.
+
+    **Credentials are ambient.** The SDK reads ``ISLO_API_KEY``, and optionally
+    ``ISLO_BASE_URL`` and ``ISLO_COMPUTE_URL``, from the worker environment on 
first
+    use. Modal reads a ``modal`` connection first, but that connection type is 
owned by
+    the Modal provider; an ``islo`` connection type belongs in a future Islo 
provider,
+    not in this one.
+
+    File reads, writes and exports move file contents through Islo's native
+    streaming APIs. Everything else runs through a shell wrapper in the 
sandbox,
+    so the image needs ``sh``, ``mkdir`` and ``rm`` for the wrapper, ``tail`` 
to
+    bound command output, ``dirname`` to create a written file's parent
+    directory, ``stat`` to size an over-budget read and to check a file before 
it
+    is exported, and GNU ``find`` for directory listings. The server default
+    image and any Debian or Ubuntu based image provide them. Each command's
+    output is captured to a scratch file in the sandbox and only its last
+    ``max_output_bytes`` are returned, so the worker sees a bounded tail while
+    the sandbox's own ephemeral disk absorbs the rest.
+
+    Islo sets ``PATH`` for every command itself and drops a ``PATH`` given at
+    creation, so a spec that names it is refused rather than silently ignored.
+    ``SandboxSpec.owner`` is refused for the same reason: this backend keeps no
+    per-sandbox metadata, so a sandbox created here cannot be attached to 
later.
+
+    :param image: Sandbox image. ``None`` (default) uses the server default.
+    :param vcpus: Number of virtual CPUs. ``None`` uses the server default.
+    :param memory_mb: Memory in MB. ``None`` uses the server default.
+    :param pause_after_idle: Seconds without a command or file operation after
+        which the server pauses the microVM and releases its compute. ``None``
+        disables it. Default ``600``.
+    :param auto_resume: ``"on_activity"`` (default) resumes a paused sandbox on
+        the next command or file operation; ``"never"`` leaves it paused, and 
the
+        backend then treats a paused sandbox as unusable.
+    :param delete_after: Seconds after *creation* at which the server deletes 
the
+        sandbox whether or not it is in use. ``None`` disables it. Default
+        ``86400``.
+    """
+
+    name = "islo"
+
+    def __init__(
+        self,
+        *,
+        image: str | None = None,
+        vcpus: int | None = None,
+        memory_mb: int | None = None,
+        pause_after_idle: int | None = 600,
+        auto_resume: Literal["never", "on_activity"] = "on_activity",
+        delete_after: int | None = 86400,
+    ) -> None:
+        if pause_after_idle is not None:
+            _validate_positive_finite(pause_after_idle, "pause_after_idle")
+        if delete_after is not None:
+            _validate_positive_finite(delete_after, "delete_after")
+        if auto_resume not in _AUTO_RESUME_POLICIES:
+            raise ValueError(
+                f"auto_resume must be one of {sorted(_AUTO_RESUME_POLICIES)}, 
got {auto_resume!r}."
+            )
+        if vcpus is not None:
+            _validate_positive_finite(vcpus, "vcpus")
+        if memory_mb is not None:
+            _validate_positive_finite(memory_mb, "memory_mb")
+        if image == "":
+            raise ValueError("image must not be empty.")
+        self._image = image
+        self._vcpus = vcpus
+        self._memory_mb = memory_mb
+        self._pause_after_idle = pause_after_idle
+        self._auto_resume = auto_resume
+        self._delete_after = delete_after
+        self._client: Islo | None = None
+
+    def _get_client(self) -> Islo:
+        if self._client is not None:
+            return self._client
+        with _translate_islo_errors("initialize its client"):
+            from islo import Islo
+
+            self._client = Islo()
+        return self._client
+
+    @staticmethod
+    def _request_options(
+        *, timeout: float, chunk_size: int | None = None, max_retries: int | 
None = None
+    ) -> dict[str, int]:
+        # ``max_retries`` is left to the SDK's default unless asked for: 
passing
+        # 0 would switch off the two transport retries it does on its own.
+        options = {"timeout_in_seconds": max(1, math.ceil(timeout))}
+        if chunk_size is not None:
+            options["chunk_size"] = chunk_size
+        if max_retries is not None:
+            options["max_retries"] = max_retries
+        return options
+
+    def _ensure_sandbox_usable(self, info: Any) -> None:
+        status = getattr(info, "status", None)
+        unusable = getattr(info, "deleted_at", None) is not None or status in 
_UNUSABLE_SANDBOX_STATUSES
+        if status == "paused" and self._auto_resume != "on_activity":
+            unusable = True
+        if unusable:
+            raise SandboxTerminalError(
+                f"Islo sandbox {getattr(info, 'name', '?')!r} cannot serve 
requests (status={status!r})."
+            )
+
+    @staticmethod
+    def _check_spec(spec: SandboxSpec | None) -> None:
+        """Refuse a spec this backend cannot carry faithfully, before anything 
is provisioned."""
+        if spec is None:
+            return
+        if spec.owner is not None:
+            # An owner exists so that a later task can attach to the sandbox, 
and the
+            # ownership rules live in per-sandbox metadata this backend keeps 
none of,
+            # so recording one would promise an attach that cannot be checked.
+            raise SandboxTerminalError(
+                "SandboxSpec names an owner, but this backend keeps no 
per-sandbox metadata the "
+                "ownership rules could be read back from, 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:
+            raise SandboxTerminalError(
+                "The Islo backend cannot apply a per-domain egress allowlist; 
it can only turn "
+                "outbound access on or off. Drop allow_egress_to, or use a 
backend with "
+                "per-domain network rules."
+            )
+        if spec.allow_egress_to_cidrs:
+            raise SandboxTerminalError(
+                "SandboxSpec names allow_egress_to_cidrs, which the Islo 
backend cannot enforce: "
+                "it can only turn outbound access on or off. Drop 
allow_egress_to_cidrs, or use "
+                "a backend with an address-layer allowlist."
+            )
+        if spec.env and "PATH" in spec.env:
+            raise SandboxTerminalError(
+                "Islo sets PATH for every command itself and drops a PATH 
given at creation; "
+                "remove PATH from SandboxSpec.env."
+            )
+
+    def create(self, *, spec: SandboxSpec | None = None) -> str:
+        self._check_spec(spec)
+        with _translate_islo_errors("create a sandbox"):
+            from islo.types import AutoResumePolicy, LifecyclePolicy
+
+            kwargs: dict[str, Any] = {
+                "internet_enabled": False if spec is None else not 
spec.block_network,
+                "lifecycle": LifecyclePolicy(
+                    pause_after_idle=self._pause_after_idle,
+                    auto_resume=AutoResumePolicy(self._auto_resume),
+                    delete_after=self._delete_after,
+                ),
+            }
+            if self._image is not None:
+                kwargs["image"] = self._image
+            if self._vcpus is not None:
+                kwargs["vcpus"] = self._vcpus
+            if self._memory_mb is not None:
+                kwargs["memory_mb"] = self._memory_mb
+            if spec is not None and spec.env:
+                # Verified against a live microVM: variables set here are the
+                # process environment of every later exec.
+                kwargs["env"] = dict(spec.env)
+            # Bind the name before the call. If creation fails after the server
+            # provisioned the microVM -- a response timeout, a reset, a 5xx --
+            # this is the only handle that can still delete it, and without it
+            # the leak is neither cleanable nor traceable to a run.
+            name = _new_sandbox_name()
+            try:
+                sandbox = self._get_client().sandboxes.create_sandbox(
+                    name=name,
+                    
request_options=self._request_options(timeout=_FILE_OP_TIMEOUT),
+                    **kwargs,
+                )
+            except BaseException:
+                with suppress(Exception):
+                    self.destroy(name)
+                raise
+        try:
+            self._ensure_sandbox_usable(sandbox)
+        except SandboxTerminalError:
+            with suppress(Exception):
+                self.destroy(name)
+            raise
+        return sandbox.name
+
+    def _await_exec(self, sandbox: str, exec_id: str, *, deadline: float) -> 
Any:
+        client = self._get_client()
+        interval = _POLL_INITIAL
+        last_error: Exception | None = None
+        while time.monotonic() < deadline:
+            remaining = deadline - time.monotonic()
+            try:
+                result = client.sandboxes.get_exec_result(
+                    sandbox,
+                    exec_id,
+                    
request_options=self._request_options(timeout=max(_POLL_HTTP_TIMEOUT_MIN, 
remaining)),
+                )
+            except Exception as e:
+                if not _is_transient_error(e):
+                    _raise_translated(e, "poll a sandbox command")
+                # One failed poll says nothing about the command; the deadline 
decides.
+                last_error = e
+            else:
+                last_error = None
+                if result.status not in _RUNNING_EXEC_STATUSES:
+                    return result
+            time.sleep(min(interval, max(0.0, deadline - time.monotonic())))
+            interval = min(interval * _POLL_BACKOFF, _POLL_MAX)
+        if last_error is not None:
+            # Nothing was heard after the last failure, so "timed out" would 
be a guess.
+            _raise_translated(last_error, "poll a sandbox command")
+        return None
+
+    def _destroy_after_timeout(self, sandbox: str) -> None:
+        try:
+            self.destroy(sandbox)
+        except SandboxError:
+            # Warn rather than fail the task: the command merely ran long, and
+            # failing here would turn a timeout the model can react to into a
+            # task failure over a transient error. The warning says whether the
+            # lifecycle policy will reclaim the microVM, since with
+            # ``delete_after=None`` nothing will.
+            if self._delete_after is not None:
+                reclaim = f"its delete_after policy removes it 
{self._delete_after}s after creation"
+            else:
+                reclaim = "delete_after is disabled, so it persists until 
deleted by hand"
+            log.warning(
+                "Timed out running a command in Islo sandbox %s and could not 
confirm its deletion; %s.",
+                sandbox,
+                reclaim,
+                exc_info=True,
+            )
+
+    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")
+        client = self._get_client()
+        # The server returns at most _SERVER_STREAM_CAP bytes per stream, so a
+        # larger budget cannot be honoured and is clamped below it.
+        budget = min(max_output_bytes, _SERVER_STREAM_CAP - 1)
+        with _translate_islo_errors("start a sandbox command"):
+            response = client.sandboxes.exec_in_sandbox(
+                sandbox,
+                # One byte over the budget, so a stream that comes back over it
+                # is proof the guest had more to give.
+                command=[
+                    "sh",
+                    "-c",
+                    _COMMAND_WRAPPER,
+                    "airflow-sandbox",
+                    command,
+                    str(budget + 1),
+                ],
+                timeout_secs=max(1, math.ceil(timeout)),
+                
request_options=self._request_options(timeout=max(_POLL_HTTP_TIMEOUT_MIN, 
timeout)),

Review Comment:
   `exec_in_sandbox` is a POST with no idempotency key, and with `max_retries` 
left to the SDK, islo 0.3.19 resends it on any 5xx, 408, 409 or 429 and on 
`RemoteProtocolError` (`_should_retry` and the `except` above it in 
`core/http_client.py`). If Islo started the command and only the response was 
lost, the agent's command runs a second time, and `_await_exec` only ever sees 
the second run's `exec_id`. With an `httpx.MockTransport` that drops the first 
response, one `run_command` sent two POSTs to `/exec`. Keeping the SDK retries 
for the polls was right, but this one call wants `max_retries=0`: a failed 
start then fails the task, which beats running a `git push` or an `>> file` 
twice. The comment on `_request_options` could say why this call is the 
exception.



##########
providers/common/ai/tests/unit/common/ai/sandbox/test_islo.py:
##########
@@ -0,0 +1,1097 @@
+# 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 io
+import os
+import signal
+import subprocess
+import time
+from pathlib import Path
+from types import SimpleNamespace
+from unittest import mock
+
+import pytest
+
+pytest.importorskip("islo")
+
+import httpx
+from islo.core.api_error import ApiError
+from islo.errors import NotFoundError
+from islo.sandboxes.client import SandboxesClient
+
+from airflow.providers.common.ai.sandbox.base import (
+    SandboxError,
+    SandboxExecResult,
+    SandboxFileTooLargeError,
+    SandboxSpec,
+    SandboxTerminalError,
+)
+from airflow.providers.common.ai.sandbox.islo import (
+    _COMMAND_WRAPPER,
+    _SERVER_STREAM_CAP,
+    IsloSandboxBackend,
+    _bound_result_stream,
+)
+
+_MODULE = "airflow.providers.common.ai.sandbox.islo"
+_ISLO_PATH = "islo.Islo"
+
+
+def _exec_result(status="completed", exit_code=0, stdout="", stderr="", 
truncated=False):
+    return SimpleNamespace(
+        status=status, exit_code=exit_code, stdout=stdout, stderr=stderr, 
truncated=truncated
+    )
+
+
+def _sandbox_info(name="box-1", status="running", deleted_at=None):
+    return SimpleNamespace(name=name, status=status, deleted_at=deleted_at)
+
+
+def _export_check(size):
+    """What the guest's export check prints for a regular file of ``size`` 
bytes."""
+    return _exec_result(stdout=f"\nairflow-export-size:{size}\n")
+
+
+def _backend_with_client(**kwargs) -> tuple[IsloSandboxBackend, 
mock.MagicMock]:
+    backend = IsloSandboxBackend(**kwargs)
+    client = mock.MagicMock(spec=["sandboxes"])
+    client.sandboxes = mock.create_autospec(SandboxesClient, instance=True)
+    client.sandboxes.exec_in_sandbox.return_value = 
SimpleNamespace(exec_id="exec-1")
+    client.sandboxes.create_sandbox.return_value = _sandbox_info()
+    client.sandboxes.get_sandbox.return_value = _sandbox_info()
+    client.sandboxes.get_exec_result.return_value = _exec_result()
+    backend._client = client
+    return backend, client
+
+
+class TestCredentials:
+    @mock.patch(_ISLO_PATH, autospec=True)
+    def test_client_comes_from_the_sdk_environment(self, islo):
+        client = IsloSandboxBackend()._get_client()
+
+        islo.assert_called_once_with()
+        assert client is islo.return_value
+
+    @mock.patch(_ISLO_PATH, autospec=True)
+    def test_client_is_resolved_once_and_cached(self, islo):
+        backend = IsloSandboxBackend()
+
+        backend._get_client()
+        backend._get_client()
+
+        islo.assert_called_once_with()
+
+    def test_missing_sdk_error_is_actionable(self):
+        real_import = builtins.__import__
+
+        def blocked_import(name, *args, **kwargs):
+            if name.startswith("islo"):
+                raise ImportError("blocked for test")
+            return real_import(name, *args, **kwargs)
+
+        backend = IsloSandboxBackend()
+        with mock.patch("builtins.__import__", side_effect=blocked_import):
+            with pytest.raises(SandboxTerminalError, match=r"\[islo\]"):
+                backend.create()
+
+
[email protected](
+    ("kwargs", "message"),
+    [
+        ({"image": ""}, "image"),
+        ({"vcpus": 0}, "vcpus"),
+        ({"memory_mb": 0}, "memory_mb"),
+        ({"delete_after": 0}, "delete_after"),
+        ({"pause_after_idle": 0}, "pause_after_idle"),
+        ({"auto_resume": "sometimes"}, "auto_resume"),
+    ],
+)
+def test_constructor_rejects_invalid_values(kwargs, message):
+    with pytest.raises(ValueError, match=message):
+        IsloSandboxBackend(**kwargs)
+
+
+class TestCreate:
+    def test_refuses_an_owner_it_could_not_be_attached_by(self):
+        backend, client = _backend_with_client()
+
+        with pytest.raises(SandboxTerminalError, match="owner"):
+            backend.create(spec=SandboxSpec(owner="example_dag/manual__1"))
+
+        client.sandboxes.create_sandbox.assert_not_called()
+
+    def test_refuses_a_per_domain_egress_allowlist(self):
+        backend, _ = _backend_with_client()
+
+        with pytest.raises(SandboxTerminalError, match="per-domain egress 
allowlist"):
+            backend.create(spec=SandboxSpec(allow_egress_to=["example.com"]))
+
+    def test_refuses_an_address_egress_allowlist(self):
+        backend, client = _backend_with_client()
+
+        # internet_enabled is all-or-nothing, so honouring block_network=True 
alone
+        # would silently drop the ranges the Dag author asked to reach.
+        with pytest.raises(SandboxTerminalError, 
match="allow_egress_to_cidrs"):
+            
backend.create(spec=SandboxSpec(allow_egress_to_cidrs=["203.0.113.0/24"]))
+
+        client.sandboxes.create_sandbox.assert_not_called()
+
+    def test_refuses_a_path_the_runner_would_drop(self):
+        backend, client = _backend_with_client()
+
+        with pytest.raises(SandboxTerminalError, match="PATH"):
+            backend.create(spec=SandboxSpec(env={"PATH": "/opt/tool/bin"}))
+
+        client.sandboxes.create_sandbox.assert_not_called()
+
+    @pytest.mark.parametrize(
+        ("spec", "expected"),
+        [
+            (None, False),
+            (SandboxSpec(), False),
+            (SandboxSpec(block_network=True), False),
+            (SandboxSpec(block_network=False), True),
+        ],
+    )
+    def test_block_network_maps_to_internet_enabled(self, spec, expected):
+        backend, client = _backend_with_client()
+
+        backend.create(spec=spec)
+
+        assert 
client.sandboxes.create_sandbox.call_args.kwargs["internet_enabled"] is expected
+
+    def test_spec_and_sizing_are_passed_at_creation(self):
+        backend, client = _backend_with_client(
+            image="python",
+            vcpus=2,
+            memory_mb=1024,
+            pause_after_idle=300,
+            auto_resume="never",
+            delete_after=120,
+        )
+
+        name = backend.create(spec=SandboxSpec(env={"TOKEN": "value"}))
+
+        assert name == "box-1"
+        kwargs = client.sandboxes.create_sandbox.call_args.kwargs
+        assert kwargs["image"] == "python"
+        assert kwargs["vcpus"] == 2
+        assert kwargs["memory_mb"] == 1024
+        assert kwargs["env"] == {"TOKEN": "value"}
+        assert kwargs["lifecycle"].pause_after_idle == 300
+        assert kwargs["lifecycle"].auto_resume == "never"
+        assert kwargs["lifecycle"].delete_after == 120
+        assert kwargs["request_options"] == {"timeout_in_seconds": 120}
+
+    def 
test_default_lifecycle_pauses_idle_sandboxes_and_deletes_after_a_day(self):
+        backend, client = _backend_with_client()
+
+        backend.create()
+
+        lifecycle = 
client.sandboxes.create_sandbox.call_args.kwargs["lifecycle"]
+        assert lifecycle.pause_after_idle == 600
+        assert lifecycle.auto_resume == "on_activity"
+        assert lifecycle.delete_after == 86400
+
+    def test_disabled_lifecycle_timers_are_sent_as_unset(self):
+        backend, client = _backend_with_client(pause_after_idle=None, 
delete_after=None)
+
+        backend.create()
+
+        lifecycle = 
client.sandboxes.create_sandbox.call_args.kwargs["lifecycle"]
+        assert lifecycle.pause_after_idle is None
+        assert lifecycle.delete_after is None
+
+    def test_omitted_sizing_is_left_to_the_server(self):
+        backend, client = _backend_with_client()
+
+        backend.create()
+
+        assert not {"image", "vcpus", "memory_mb"} & 
client.sandboxes.create_sandbox.call_args.kwargs.keys()
+
+    def test_api_failure_is_terminal(self):
+        backend, client = _backend_with_client()
+        client.sandboxes.create_sandbox.side_effect = ApiError(status_code=503)
+
+        with pytest.raises(SandboxTerminalError, match="HTTP 503"):
+            backend.create()
+
+    def test_a_failed_create_deletes_the_name_it_had_already_bound(self):
+        backend, client = _backend_with_client()
+        client.sandboxes.create_sandbox.side_effect = ApiError(status_code=503)
+
+        with pytest.raises(SandboxTerminalError):
+            backend.create()
+
+        # The server may have provisioned the microVM before failing to answer,
+        # and this name is the only handle that can still reclaim it.
+        requested = client.sandboxes.create_sandbox.call_args.kwargs["name"]
+        client.sandboxes.delete_sandbox.assert_called_once()
+        assert 
client.sandboxes.delete_sandbox.call_args.kwargs["sandbox_name"] == requested
+
+    def test_a_cleanup_failure_does_not_mask_the_original_create_error(self):
+        backend, client = _backend_with_client()
+        client.sandboxes.create_sandbox.side_effect = ApiError(status_code=503)
+        client.sandboxes.delete_sandbox.side_effect = ApiError(status_code=500)
+
+        with pytest.raises(SandboxTerminalError, match="HTTP 503"):
+            backend.create()
+
+    def 
test_a_sandbox_that_cannot_serve_after_creation_is_destroyed_and_terminal(self):
+        backend, client = _backend_with_client()
+        client.sandboxes.create_sandbox.return_value = 
_sandbox_info(status="stopped")
+
+        with pytest.raises(SandboxTerminalError, match="cannot serve 
requests"):
+            backend.create()
+
+        client.sandboxes.delete_sandbox.assert_called_once()
+
+    def test_spec_env_is_forwarded_so_it_is_never_silently_dropped(self):
+        backend, client = _backend_with_client()
+
+        backend.create(spec=SandboxSpec(env={"TOKEN": "value", "OTHER": "2"}))
+
+        # base.py treats dropping a SandboxSpec field as a contract violation.
+        assert client.sandboxes.create_sandbox.call_args.kwargs["env"] == {
+            "TOKEN": "value",
+            "OTHER": "2",
+        }
+
+
+class TestRunCommand:
+    def test_polls_with_backoff(self):
+        backend, client = _backend_with_client()
+        client.sandboxes.get_exec_result.side_effect = [
+            _exec_result(status="running"),
+            _exec_result(status="running"),
+            _exec_result(status="running"),
+            _exec_result(stdout="done"),
+        ]
+
+        with mock.patch("time.sleep", autospec=True) as sleep:
+            result = backend.run_command("box", "x", timeout=60, 
max_output_bytes=1024)
+
+        intervals = [call.args[0] for call in sleep.call_args_list]
+        assert result.stdout == "done"
+        assert intervals == sorted(intervals)
+        assert intervals[-1] > intervals[0]
+
+    def test_user_command_is_an_argument_to_the_bounding_wrapper(self):
+        backend, client = _backend_with_client()
+        user_command = "echo '$HOME'; rm -f /tmp/nope"
+
+        backend.run_command("box", user_command, timeout=5, 
max_output_bytes=1024)
+
+        command = client.sandboxes.exec_in_sandbox.call_args.kwargs["command"]
+        assert command[:2] == ["sh", "-c"]
+        assert user_command not in command[2]
+        # One byte over the budget, so an over-budget stream proves truncation.
+        assert command[4:] == [user_command, "1025"]
+
+    def test_budget_above_the_server_cap_is_clamped_to_it(self):
+        backend, client = _backend_with_client()
+
+        backend.run_command("box", "x", timeout=5, max_output_bytes=5 * 
_SERVER_STREAM_CAP)
+
+        command = client.sandboxes.exec_in_sandbox.call_args.kwargs["command"]
+        assert command[-1] == str(_SERVER_STREAM_CAP)
+
+    def test_requests_leave_the_sdk_retries_in_place(self):
+        backend, client = _backend_with_client()
+
+        backend.run_command("box", "x", timeout=5, max_output_bytes=1024)
+
+        for call in (client.sandboxes.exec_in_sandbox, 
client.sandboxes.get_exec_result):
+            assert "max_retries" not in 
call.call_args.kwargs["request_options"]
+
+    def test_a_stream_within_budget_is_passed_through_untouched(self):
+        backend, client = _backend_with_client()
+        client.sandboxes.get_exec_result.return_value = 
_exec_result(stdout="a\nb\n", stderr="err\n")
+
+        result = backend.run_command("box", "x", timeout=5, 
max_output_bytes=1024)
+
+        assert result.stdout == "a\nb\n"
+        assert result.stderr == "err\n"
+        assert not result.stdout_truncated
+        assert not result.stderr_truncated
+
+    def test_an_over_budget_stream_keeps_the_tail_and_reports_truncation(self):
+        backend, client = _backend_with_client()
+        client.sandboxes.get_exec_result.return_value = _exec_result(
+            stdout="line1\nline2\nline3\n", stderr="e1\ne2\ne3\n"
+        )
+
+        result = backend.run_command("box", "x", timeout=5, max_output_bytes=8)
+
+        assert result.stdout == "line3\n"
+        assert result.stderr == "e2\ne3\n"
+        assert result.stdout_truncated
+        assert result.stderr_truncated
+
+    def test_truncation_drops_a_short_partial_leading_line(self):
+        backend, client = _backend_with_client()
+        # A byte-aligned cut of the last 12 bytes would land inside "line988".
+        client.sandboxes.get_exec_result.return_value = 
_exec_result(stdout="line988\nline989\n")
+
+        result = backend.run_command("box", "x", timeout=5, 
max_output_bytes=12)
+
+        assert result.stdout == "line989\n"
+        assert result.stdout_truncated
+
+    @pytest.mark.parametrize(
+        ("stdout", "expected"),
+        [
+            ("abcdefghij", "cdefghij"),
+            # The line's only newline is its last byte, so dropping through it
+            # would leave nothing and the model would read "(no output)".
+            ("x" * 20 + "\n", "xxxxxxx\n"),
+            ("z" * 20 + "\nok\n", "zzzz\nok\n"),
+        ],
+        ids=["no-newline", "trailing-newline", "short-line-after"],
+    )
+    def test_a_single_line_over_budget_is_cut_rather_than_dropped(self, 
stdout, expected):
+        backend, client = _backend_with_client()
+        client.sandboxes.get_exec_result.return_value = 
_exec_result(stdout=stdout)
+
+        result = backend.run_command("box", "x", timeout=5, max_output_bytes=8)
+
+        assert result.stdout == expected
+        assert result.stdout_truncated
+
+    def test_applies_the_byte_cap_on_utf8_boundaries(self):
+        backend, client = _backend_with_client()
+        client.sandboxes.get_exec_result.return_value = 
_exec_result(stdout="ééé")
+
+        result = backend.run_command("box", "x", timeout=5, max_output_bytes=4)

Review Comment:
   `ééé` with a 4-byte budget always cuts between whole characters, so this 
passes even with `errors="strict"` in `_bound_result_stream` and never reaches 
the replacement path. A budget of 3 keeps half of the first `é` and gives 
`"\ufffdé"`, which is the case worth pinning.



-- 
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