zozo123 commented on code in PR #71672: URL: https://github.com/apache/airflow/pull/71672#discussion_r4211150633
########## 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: Parametrized with a budget of 3, which keeps half of the first `é` and expects `"\ufffdé"`. It fails with `errors="strict"`. _🤖 Addressed by [Claude Code](https://claude.com/claude-code)_ ########## 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: Merged the four into `test_refuses_a_spec_it_cannot_carry_before_provisioning`, parametrized over the spec and the match, and each case asserts `create_sandbox.assert_not_called()`. _🤖 Addressed by [Claude Code](https://claude.com/claude-code)_ -- 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]
