Files
k7/tests/unit/test_core_mocked.py
G 13fafe0997 Release 0.4.0
Per-key node pins, k7d-fc pause/resume/exec, and HA-soak fixes. Playbook
pins k7d 0.7.0. GitHub .deb, Launchpad PPA, and PyPI k7-sdk are 0.4.0.
2026-09-19 23:22:14 +02:00

1835 lines
71 KiB
Python

"""Mock-based unit tests for K7Core lifecycle methods (no real k8s needed)."""
from types import SimpleNamespace
from unittest.mock import AsyncMock, MagicMock, patch
from kubernetes_asyncio.client.exceptions import ApiException
from k7.core.core import K7_TENANT_LABEL, K7Core
from k7.core.models import ExecResult, OperationResult, SandboxConfig
from tests.unit.conftest import mock_deployment, mock_pod
# ---------------------------------------------------------------------------
# Helpers
# ---------------------------------------------------------------------------
def _setup_clients(core, apps=None, v1=None, net=None, metrics=None, custom=None, batch=None, apiext=None):
"""Wire MagicMock clients into a K7Core instance."""
core._config_loaded = True
if apps is not None:
core._apps_v1_client = apps
if v1 is not None:
core._core_v1_client = v1
if net is not None:
core._networking_v1_client = net
if metrics is not None:
core._metrics_client = metrics
if custom is not None:
core._custom_objects_client = custom
if batch is not None:
core._batch_v1_client = batch
if apiext is not None:
core._apiextensions_v1_client = apiext
# --- _detect_backend ---
class TestDetectBackend:
async def test_annotation_takes_priority(self, core: K7Core):
dep = mock_deployment(annotations={"k7.katakate.org/backend": "kata-qemu-longhorn"})
mock_apps = AsyncMock()
mock_apps.read_namespaced_deployment.return_value = dep
core._apps_v1_client = mock_apps
core._config_loaded = True
assert await core._detect_backend("test-sb", "default") == "kata-qemu-longhorn"
async def test_host_file_when_no_sandbox(self, core: K7Core, tmp_path):
backend_file = tmp_path / "backend"
backend_file.write_text("kata-qemu-longhorn\n")
with patch("os.path.exists", return_value=True), patch("builtins.open", return_value=open(backend_file)):
assert await core._detect_backend(None) == "kata-qemu-longhorn"
async def test_default_when_nothing_found(self, core: K7Core):
with patch("os.path.exists", return_value=False):
assert await core._detect_backend(None) == "kata-firecracker-devmapper"
async def test_default_when_annotation_missing(self, core: K7Core):
dep = mock_deployment(annotations={})
mock_apps = AsyncMock()
mock_apps.read_namespaced_deployment.return_value = dep
core._apps_v1_client = mock_apps
core._config_loaded = True
with patch("os.path.exists", return_value=False):
assert await core._detect_backend("test-sb", "default") == "kata-firecracker-devmapper"
async def test_annotation_api_error_falls_through(self, core: K7Core):
mock_apps = AsyncMock()
mock_apps.read_namespaced_deployment.side_effect = ApiException(status=404)
core._apps_v1_client = mock_apps
core._config_loaded = True
with patch("os.path.exists", return_value=False):
assert await core._detect_backend("missing", "default") == "kata-firecracker-devmapper"
async def test_none_file_is_not_kfd(self, core: K7Core, tmp_path):
backend_file = tmp_path / "backend"
backend_file.write_text("none\n")
with patch("os.path.exists", return_value=True), patch("builtins.open", return_value=open(backend_file)):
assert await core._detect_backend(None) == ""
# --- delete_sandbox ---
class TestDeleteSandbox:
async def test_successful_delete(self, core: K7Core):
mock_apps = AsyncMock()
mock_v1 = AsyncMock()
mock_net = AsyncMock()
mock_custom = AsyncMock()
mock_custom.list_namespaced_custom_object.return_value = {"items": []}
core._apps_v1_client = mock_apps
core._core_v1_client = mock_v1
core._networking_v1_client = mock_net
core._custom_objects_client = mock_custom
core._config_loaded = True
result = await core.delete_sandbox("my-sb", "default")
assert result.success is True
mock_apps.delete_namespaced_deployment.assert_called_once_with(name="my-sb", namespace="default")
mock_v1.delete_namespaced_secret.assert_called_once_with(name="my-sb-env", namespace="default")
mock_v1.delete_namespaced_config_map.assert_any_call(name="my-sb-persist-wrapper", namespace="default")
mock_v1.delete_namespaced_config_map.assert_any_call(name="my-sb-docker-vehicle", namespace="default")
mock_v1.delete_namespaced_persistent_volume_claim.assert_any_call(name="my-sb-root-lh", namespace="default")
mock_v1.delete_namespaced_persistent_volume_claim.assert_any_call(name="my-sb-docker-lh", namespace="default")
async def test_delete_ignores_404s(self, core: K7Core):
mock_apps = AsyncMock()
mock_apps.delete_namespaced_deployment.side_effect = ApiException(status=404)
mock_v1 = AsyncMock()
mock_v1.delete_namespaced_secret.side_effect = ApiException(status=404)
mock_v1.delete_namespaced_config_map.side_effect = ApiException(status=404)
mock_v1.delete_namespaced_persistent_volume_claim.side_effect = ApiException(status=404)
mock_v1.patch_namespaced_persistent_volume_claim.side_effect = ApiException(status=404)
mock_net = AsyncMock()
mock_net.delete_namespaced_network_policy.side_effect = ApiException(status=404)
mock_custom = AsyncMock()
mock_custom.list_namespaced_custom_object.return_value = {"items": []}
core._apps_v1_client = mock_apps
core._core_v1_client = mock_v1
core._networking_v1_client = mock_net
core._custom_objects_client = mock_custom
core._config_loaded = True
result = await core.delete_sandbox("gone-sb", "default")
assert result.success is True
async def test_delete_reports_non_404_errors(self, core: K7Core):
mock_apps = AsyncMock()
mock_apps.delete_namespaced_deployment.side_effect = ApiException(status=500, reason="Server Error")
mock_v1 = AsyncMock()
mock_net = AsyncMock()
mock_custom = AsyncMock()
mock_custom.list_namespaced_custom_object.return_value = {"items": []}
core._apps_v1_client = mock_apps
core._core_v1_client = mock_v1
core._networking_v1_client = mock_net
core._custom_objects_client = mock_custom
core._config_loaded = True
result = await core.delete_sandbox("err-sb", "default")
assert result.success is False
assert "deployment" in result.error
# --- list_sandboxes ---
class TestListSandboxes:
async def test_returns_sandbox_info_list(self, core: K7Core):
dep = mock_deployment(name="sb1", annotations={"k7.katakate.org/backend": "kata-firecracker-devmapper"})
pod = mock_pod(name="sb1-abc", image="ubuntu:22.04")
mock_apps = AsyncMock()
core._apps_v1_client = mock_apps
core._config_loaded = True
mock_v1 = AsyncMock()
mock_v1.list_namespaced_pod.return_value = MagicMock(items=[pod])
core._core_v1_client = mock_v1
with patch.object(core, "_get_kata_sandboxes", new_callable=AsyncMock, return_value=[dep]):
result = await core.list_sandboxes(namespace="default")
assert len(result) == 1
assert result[0].name == "sb1"
assert result[0].image == "ubuntu:22.04"
assert result[0].backend == "kata-firecracker-devmapper"
assert result[0].ready == "True"
async def test_returns_empty_on_no_sandboxes(self, core: K7Core):
mock_apps = AsyncMock()
core._apps_v1_client = mock_apps
core._config_loaded = True
with patch.object(core, "_get_kata_sandboxes", new_callable=AsyncMock, return_value=[]):
result = await core.list_sandboxes(namespace="default")
assert result == []
async def test_handles_pod_with_no_items(self, core: K7Core):
dep = mock_deployment(name="sb-no-pod", annotations={"k7.katakate.org/backend": "kata-qemu-longhorn"})
mock_apps = AsyncMock()
core._apps_v1_client = mock_apps
core._config_loaded = True
mock_v1 = AsyncMock()
mock_v1.list_namespaced_pod.return_value = MagicMock(items=[])
core._core_v1_client = mock_v1
with patch.object(core, "_get_kata_sandboxes", new_callable=AsyncMock, return_value=[dep]):
result = await core.list_sandboxes(namespace="default")
assert len(result) == 1
assert result[0].status == "No Pods"
assert result[0].ready == "False"
async def test_exception_returns_empty_list(self, core: K7Core):
core._config_loaded = True
with patch.object(core, "_get_kata_sandboxes", new_callable=AsyncMock, side_effect=Exception("k8s down")):
result = await core.list_sandboxes(namespace="default")
assert result == []
# --- _get_kata_sandboxes ---
class TestGetKataSandboxes:
async def test_filters_by_runtime_class(self, core: K7Core):
kata_dep = mock_deployment(name="kata-sb", runtime_class="kata")
non_kata = mock_deployment(name="regular", runtime_class="runc")
non_kata.metadata.labels = {}
mock_apps = AsyncMock()
mock_apps.list_namespaced_deployment.return_value = MagicMock(items=[kata_dep, non_kata])
core._apps_v1_client = mock_apps
core._config_loaded = True
result = await core._get_kata_sandboxes(namespace="default")
assert len(result) == 1
assert result[0].metadata.name == "kata-sb"
async def test_filters_by_label_fallback(self, core: K7Core):
dep = mock_deployment(name="label-sb", runtime_class="runc", labels={"runtime": "kata"})
mock_apps = AsyncMock()
mock_apps.list_namespaced_deployment.return_value = MagicMock(items=[dep])
core._apps_v1_client = mock_apps
core._config_loaded = True
result = await core._get_kata_sandboxes(namespace="default")
assert len(result) == 1
async def test_all_namespaces_when_none(self, core: K7Core):
mock_apps = AsyncMock()
mock_apps.list_deployment_for_all_namespaces.return_value = MagicMock(items=[])
core._apps_v1_client = mock_apps
core._config_loaded = True
await core._get_kata_sandboxes(namespace=None)
mock_apps.list_deployment_for_all_namespaces.assert_called_once()
# --- _ensure_root_pvc ---
class TestEnsureRootPvc:
async def test_exists_with_filesystem_mode(self, core: K7Core):
mock_v1 = AsyncMock()
pvc = MagicMock()
pvc.spec.volume_mode = "Filesystem"
mock_v1.read_namespaced_persistent_volume_claim.return_value = pvc
_setup_clients(core, v1=mock_v1)
result = await core._ensure_root_pvc("sb1", "default", "10Gi")
assert result.success is True
assert result.data["created"] is False
async def test_exists_with_wrong_volume_mode(self, core: K7Core):
mock_v1 = AsyncMock()
pvc = MagicMock()
pvc.spec.volume_mode = "Block"
mock_v1.read_namespaced_persistent_volume_claim.return_value = pvc
_setup_clients(core, v1=mock_v1)
result = await core._ensure_root_pvc("sb1", "default", "10Gi")
assert result.success is False
assert "Block" in result.error
async def test_creates_when_not_found(self, core: K7Core):
mock_v1 = AsyncMock()
mock_v1.read_namespaced_persistent_volume_claim.side_effect = ApiException(status=404)
_setup_clients(core, v1=mock_v1)
result = await core._ensure_root_pvc("sb1", "default", "10Gi")
assert result.success is True
assert result.data["created"] is True
mock_v1.create_namespaced_persistent_volume_claim.assert_called_once()
async def test_create_fails(self, core: K7Core):
mock_v1 = AsyncMock()
mock_v1.read_namespaced_persistent_volume_claim.side_effect = ApiException(status=404)
mock_v1.create_namespaced_persistent_volume_claim.side_effect = ApiException(status=500)
_setup_clients(core, v1=mock_v1)
result = await core._ensure_root_pvc("sb1", "default", "10Gi")
assert result.success is False
async def test_non_404_read_error(self, core: K7Core):
mock_v1 = AsyncMock()
mock_v1.read_namespaced_persistent_volume_claim.side_effect = ApiException(status=500)
_setup_clients(core, v1=mock_v1)
result = await core._ensure_root_pvc("sb1", "default", "10Gi")
assert result.success is False
assert "lookup error" in result.error
# --- _ensure_persist_wrapper_configmap ---
class TestEnsurePersistWrapperConfigmap:
async def test_exists(self, core: K7Core):
mock_v1 = AsyncMock()
_setup_clients(core, v1=mock_v1)
result = await core._ensure_persist_wrapper_configmap("default", "sb1", "#!/bin/bash")
assert result.success is True
assert result.data["created"] is False
async def test_creates_when_not_found(self, core: K7Core):
mock_v1 = AsyncMock()
mock_v1.read_namespaced_config_map.side_effect = ApiException(status=404)
_setup_clients(core, v1=mock_v1)
result = await core._ensure_persist_wrapper_configmap("default", "sb1", "#!/bin/bash")
assert result.success is True
assert result.data["created"] is True
mock_v1.create_namespaced_config_map.assert_called_once()
async def test_non_404_read_error(self, core: K7Core):
mock_v1 = AsyncMock()
mock_v1.read_namespaced_config_map.side_effect = ApiException(status=500)
_setup_clients(core, v1=mock_v1)
result = await core._ensure_persist_wrapper_configmap("default", "sb1", "#!/bin/bash")
assert result.success is False
assert "lookup error" in result.error
async def test_create_fails(self, core: K7Core):
mock_v1 = AsyncMock()
mock_v1.read_namespaced_config_map.side_effect = ApiException(status=404)
mock_v1.create_namespaced_config_map.side_effect = ApiException(status=500)
_setup_clients(core, v1=mock_v1)
result = await core._ensure_persist_wrapper_configmap("default", "sb1", "#!/bin/bash")
assert result.success is False
# --- _update_deployment_annotation ---
class TestUpdateDeploymentAnnotation:
async def test_patches_deployment(self, core: K7Core):
mock_apps = AsyncMock()
_setup_clients(core, apps=mock_apps)
await core._update_deployment_annotation("sb1", "default", "key", "val")
mock_apps.patch_namespaced_deployment.assert_called_once_with(
name="sb1",
namespace="default",
body={"metadata": {"annotations": {"key": "val"}}},
)
# --- pause_sandbox ---
class TestPauseSandbox:
async def test_scales_to_zero(self, core: K7Core):
mock_apps = AsyncMock()
_setup_clients(core, apps=mock_apps)
result = await core.pause_sandbox("sb1")
assert result.success is True
assert "paused" in result.message
mock_apps.patch_namespaced_deployment_scale.assert_called_once()
async def test_with_snapshot_explicit_pvc(self, core: K7Core):
mock_apps = AsyncMock()
_setup_clients(core, apps=mock_apps)
with patch.object(
core,
"_create_volume_snapshot",
new_callable=AsyncMock,
return_value=OperationResult(success=True, message="snap ok"),
) as snap_mock:
result = await core.pause_sandbox("sb1", pvc_name="pvc1", snapshot_name="snap1")
assert result.success is True
assert "Snapshot snap1 created" in result.message
snap_mock.assert_called_once()
assert snap_mock.call_args.kwargs["pvc_name"] == "pvc1"
assert snap_mock.call_args.kwargs["snapshot_name"] == "snap1"
async def test_snapshot_only_defaults_pvc_to_root(self, core: K7Core):
"""Regression: snapshot_name alone triggers snapshot.
The snapshot used to fire only when ``pvc_name AND snapshot_name`` were
both set, so ``k7 pause foo --snapshot bar`` silently did nothing.
Now ``snapshot_name`` alone is sufficient and ``pvc_name``
defaults to ``_root_pvc_name(name)``.
"""
mock_apps = AsyncMock()
_setup_clients(core, apps=mock_apps)
with patch.object(
core,
"_create_volume_snapshot",
new_callable=AsyncMock,
return_value=OperationResult(success=True, message="snap ok"),
) as snap_mock:
result = await core.pause_sandbox("sb1", snapshot_name="snap1")
assert result.success is True
assert "Snapshot snap1 created" in result.message
snap_mock.assert_called_once()
assert snap_mock.call_args.kwargs["pvc_name"] == "sb1-root-lh"
assert snap_mock.call_args.kwargs["snapshot_name"] == "snap1"
async def test_no_snapshot_when_name_unset(self, core: K7Core):
"""Without ``snapshot_name`` we must NOT call _create_volume_snapshot."""
mock_apps = AsyncMock()
_setup_clients(core, apps=mock_apps)
with patch.object(
core,
"_create_volume_snapshot",
new_callable=AsyncMock,
) as snap_mock:
result = await core.pause_sandbox("sb1", pvc_name="ignored-without-name")
assert result.success is True
assert "Snapshot" not in result.message
snap_mock.assert_not_called()
async def test_snapshot_fails(self, core: K7Core):
mock_apps = AsyncMock()
_setup_clients(core, apps=mock_apps)
with patch.object(
core,
"_create_volume_snapshot",
new_callable=AsyncMock,
return_value=OperationResult(success=False, error="snap err"),
):
result = await core.pause_sandbox("sb1", snapshot_name="snap1")
assert result.success is False
assert "snap err" in result.error
async def test_api_exception(self, core: K7Core):
mock_apps = AsyncMock()
mock_apps.patch_namespaced_deployment_scale.side_effect = ApiException(status=500)
_setup_clients(core, apps=mock_apps)
result = await core.pause_sandbox("sb1")
assert result.success is False
# --- resume_sandbox ---
class TestResumeSandbox:
async def test_scales_to_one(self, core: K7Core):
mock_apps = AsyncMock()
_setup_clients(core, apps=mock_apps)
result = await core.resume_sandbox("sb1")
assert result.success is True
assert "resumed" in result.message
async def test_api_exception(self, core: K7Core):
mock_apps = AsyncMock()
mock_apps.patch_namespaced_deployment_scale.side_effect = ApiException(status=500)
_setup_clients(core, apps=mock_apps)
result = await core.resume_sandbox("sb1")
assert result.success is False
# --- delete_all_sandboxes ---
class TestDeleteAllSandboxes:
async def test_deletes_multiple(self, core: K7Core):
mock_apps = AsyncMock()
mock_v1 = AsyncMock()
mock_net = AsyncMock()
mock_custom = AsyncMock()
mock_custom.list_namespaced_custom_object.return_value = {"items": []}
_setup_clients(core, apps=mock_apps, v1=mock_v1, net=mock_net, custom=mock_custom)
deps = [mock_deployment(name="a"), mock_deployment(name="b")]
with patch.object(core, "_get_kata_sandboxes", new_callable=AsyncMock, return_value=deps):
result = await core.delete_all_sandboxes()
assert result.success is True
assert "2" in result.message
async def test_one_fails(self, core: K7Core):
mock_apps = AsyncMock()
mock_apps.delete_namespaced_deployment.side_effect = [
None,
ApiException(status=500, reason="fail"),
]
mock_v1 = AsyncMock()
mock_net = AsyncMock()
mock_custom = AsyncMock()
mock_custom.list_namespaced_custom_object.return_value = {"items": []}
_setup_clients(core, apps=mock_apps, v1=mock_v1, net=mock_net, custom=mock_custom)
deps = [mock_deployment(name="ok"), mock_deployment(name="fail")]
with patch.object(core, "_get_kata_sandboxes", new_callable=AsyncMock, return_value=deps):
result = await core.delete_all_sandboxes()
assert result.success is False
async def test_no_sandboxes(self, core: K7Core):
mock_apps = AsyncMock()
mock_v1 = AsyncMock()
mock_net = AsyncMock()
mock_custom = AsyncMock()
mock_custom.list_namespaced_custom_object.return_value = {"items": []}
_setup_clients(core, apps=mock_apps, v1=mock_v1, net=mock_net, custom=mock_custom)
with patch.object(core, "_get_kata_sandboxes", new_callable=AsyncMock, return_value=[]):
result = await core.delete_all_sandboxes()
assert result.success is True
assert result.data == []
async def test_exception_from_listing(self, core: K7Core):
_setup_clients(core)
with patch.object(core, "_get_kata_sandboxes", new_callable=AsyncMock, side_effect=Exception("boom")):
result = await core.delete_all_sandboxes()
assert result.success is False
# --- fork_sandbox ---
class TestForkSandbox:
def _setup_fork(self, core):
"""Set up mock clients for fork tests and return them."""
mock_apps = AsyncMock()
mock_v1 = AsyncMock()
# Source has no egress policy at all → the fork inherits open egress.
mock_net = AsyncMock()
mock_net.read_namespaced_network_policy.side_effect = ApiException(status=404)
mock_custom = AsyncMock()
mock_custom.get_namespaced_custom_object.side_effect = ApiException(status=404)
_setup_clients(core, apps=mock_apps, v1=mock_v1, net=mock_net, custom=mock_custom)
core._detect_backend = AsyncMock(return_value="kata-qemu-longhorn")
src_dep = MagicMock()
src_dep.metadata.name = "src"
src_dep.metadata.labels = {"app": "src", "katakate.org/sandbox": "src"}
src_dep.spec.selector.match_labels = {"app": "src"}
src_dep.spec.template.metadata.labels = {"app": "src", "katakate.org/sandbox": "src"}
src_dep.spec.template.spec.volumes = []
src_dep.spec.replicas = 1
# 0 ready replicas ⇒ fork skips the pre-snapshot guest `sync`
# (a paused/quiesced source needs no flush) — keeps these mocked
# tests off the exec_command path.
src_dep.status.ready_replicas = 0
mock_apps.read_namespaced_deployment.return_value = src_dep
src_pvc = MagicMock()
src_pvc.spec.access_modes = ["ReadWriteOnce"]
src_pvc.spec.resources = MagicMock()
src_pvc.spec.storage_class_name = "longhorn"
src_pvc.spec.volume_mode = "Filesystem"
mock_v1.read_namespaced_persistent_volume_claim.return_value = src_pvc
return mock_apps, mock_v1
async def test_source_deployment_not_found(self, core: K7Core):
mock_apps = AsyncMock()
mock_apps.read_namespaced_deployment.side_effect = ApiException(status=404)
_setup_clients(core, apps=mock_apps, v1=AsyncMock())
core._detect_backend = AsyncMock(return_value="kata-qemu-longhorn")
result = await core.fork_sandbox("src", "dst")
assert result.success is False
async def test_source_pvc_not_found(self, core: K7Core):
mock_apps = AsyncMock()
mock_v1 = AsyncMock()
mock_v1.read_namespaced_persistent_volume_claim.side_effect = ApiException(status=404)
_setup_clients(core, apps=mock_apps, v1=mock_v1)
core._detect_backend = AsyncMock(return_value="kata-qemu-longhorn")
result = await core.fork_sandbox("src", "dst")
assert result.success is False
assert "not found" in result.error
async def test_snapshot_creation_fails(self, core: K7Core):
self._setup_fork(core)
with patch.object(
core,
"_create_volume_snapshot",
new_callable=AsyncMock,
return_value=OperationResult(success=False, error="snap fail"),
):
result = await core.fork_sandbox("src", "dst")
assert result.success is False
assert "snap fail" in result.error
async def test_snapshot_not_ready(self, core: K7Core):
self._setup_fork(core)
with (
patch.object(
core,
"_create_volume_snapshot",
new_callable=AsyncMock,
return_value=OperationResult(success=True),
),
patch.object(
core,
"_wait_for_snapshot_ready",
new_callable=AsyncMock,
return_value=OperationResult(success=False, error="timeout"),
),
):
result = await core.fork_sandbox("src", "dst")
assert result.success is False
async def test_target_already_exists(self, core: K7Core):
mock_apps, mock_v1 = self._setup_fork(core)
mock_apps.create_namespaced_deployment.side_effect = ApiException(status=409)
with (
patch.object(
core,
"_create_volume_snapshot",
new_callable=AsyncMock,
return_value=OperationResult(success=True),
),
patch.object(
core,
"_wait_for_snapshot_ready",
new_callable=AsyncMock,
return_value=OperationResult(success=True),
),
patch.object(
core,
"_wait_for_pvc_bound",
new_callable=AsyncMock,
return_value=OperationResult(success=True),
),
):
result = await core.fork_sandbox("src", "dst")
assert result.success is False
assert "already exists" in result.error
async def test_happy_path(self, core: K7Core):
mock_apps, _ = self._setup_fork(core)
with (
patch.object(
core,
"_create_volume_snapshot",
new_callable=AsyncMock,
return_value=OperationResult(success=True),
),
patch.object(
core,
"_wait_for_snapshot_ready",
new_callable=AsyncMock,
return_value=OperationResult(success=True),
),
patch.object(
core,
"_wait_for_pvc_bound",
new_callable=AsyncMock,
return_value=OperationResult(success=True),
),
):
result = await core.fork_sandbox("src", "dst")
assert result.success is True
assert "Forked" in result.message
mock_apps.create_namespaced_deployment.assert_called_once()
def _patch_fork_storage(self, core: K7Core):
return (
patch.object(core, "_create_volume_snapshot", new_callable=AsyncMock, return_value=OperationResult(True)),
patch.object(core, "_wait_for_snapshot_ready", new_callable=AsyncMock, return_value=OperationResult(True)),
patch.object(core, "_wait_for_pvc_bound", new_callable=AsyncMock, return_value=OperationResult(True)),
)
async def test_fork_applies_network_policies(self, core: K7Core):
"""a fork used to come up with no policy at all."""
self._setup_fork(core)
core._custom_objects_client.get_namespaced_custom_object.side_effect = None
core._custom_objects_client.get_namespaced_custom_object.return_value = {
"spec": {"egress": [{"toFQDNs": [{"matchName": "example.com"}]}, {"toCIDR": ["10.0.0.0/8"]}]}
}
snap, ready, bound = self._patch_fork_storage(core)
with (
snap,
ready,
bound,
patch.object(
core,
"_apply_sandbox_network_policies",
new_callable=AsyncMock,
return_value=OperationResult(True),
) as apply_policies,
):
result = await core.fork_sandbox("src", "dst")
assert result.success is True
assert apply_policies.call_args.kwargs["name"] == "dst"
assert apply_policies.call_args.kwargs["egress_whitelist"] == ["example.com", "10.0.0.0/8"]
async def test_fork_of_open_egress_source_inherits_open_egress(self, core: K7Core):
self._setup_fork(core)
snap, ready, bound = self._patch_fork_storage(core)
with (
snap,
ready,
bound,
patch.object(
core,
"_apply_sandbox_network_policies",
new_callable=AsyncMock,
return_value=OperationResult(True),
) as apply_policies,
):
result = await core.fork_sandbox("src", "dst")
assert result.success is True
assert apply_policies.call_args.kwargs["egress_whitelist"] is None
async def test_fork_policy_failure_rolls_back_deployment(self, core: K7Core):
self._setup_fork(core)
snap, ready, bound = self._patch_fork_storage(core)
with (
snap,
ready,
bound,
patch.object(
core,
"_apply_sandbox_network_policies",
new_callable=AsyncMock,
return_value=OperationResult(False, error="boom"),
),
patch.object(
core,
"_delete_sandbox_resources",
new_callable=AsyncMock,
return_value=OperationResult(True),
) as rollback,
):
result = await core.fork_sandbox("src", "dst")
assert result.success is False
assert "rolled back" in result.error and "boom" in result.error
rollback.assert_awaited_once_with("dst", "default")
class TestApplySandboxNetworkPolicies:
"""opt-in ingress, deny-all by default."""
def _net(self, core: K7Core):
mock_net = AsyncMock()
_setup_clients(core, net=mock_net)
return mock_net
def _ingress_body(self, mock_net):
for call in mock_net.create_namespaced_network_policy.call_args_list:
body = call.kwargs["body"]
if body.metadata.name.endswith("-deny-ingress"):
return body
raise AssertionError("no deny-ingress policy created")
async def test_default_is_deny_all(self, core: K7Core):
mock_net = self._net(core)
result = await core._apply_sandbox_network_policies(name="sb", namespace="default", egress_whitelist=None)
assert result.success is True
body = self._ingress_body(mock_net)
assert body.spec.policy_types == ["Ingress"]
assert body.spec.ingress == []
assert body.spec.pod_selector.match_labels == {"katakate.org/sandbox": "sb"}
async def test_ports_without_sources_allow_same_namespace_sandboxes(self, core: K7Core):
mock_net = self._net(core)
await core._apply_sandbox_network_policies(
name="sb", namespace="default", egress_whitelist=None, ingress_ports=[8000, 8001]
)
rule = self._ingress_body(mock_net).spec.ingress[0]
assert [(p.protocol, p.port) for p in rule.ports] == [("TCP", 8000), ("TCP", 8001)]
expr = rule._from[0].pod_selector.match_expressions[0]
assert (expr.key, expr.operator) == ("katakate.org/sandbox", "Exists")
async def test_explicit_sources_are_scoped(self, core: K7Core):
mock_net = self._net(core)
await core._apply_sandbox_network_policies(
name="sb",
namespace="default",
egress_whitelist=None,
ingress_ports=[8000],
ingress_from=["sandbox:alice", "cidr:203.0.113.0/24"],
)
peers = self._ingress_body(mock_net).spec.ingress[0]._from
assert peers[0].pod_selector.match_labels == {"katakate.org/sandbox": "alice"}
assert peers[1].ip_block.cidr == "203.0.113.0/24"
async def test_bad_source_fails_without_creating_the_policy(self, core: K7Core):
mock_net = self._net(core)
result = await core._apply_sandbox_network_policies(
name="sb", namespace="default", egress_whitelist=None, ingress_ports=[8000], ingress_from=["nonsense"]
)
assert result.success is False
assert "nonsense" in result.error
mock_net.create_namespaced_network_policy.assert_not_called()
async def test_policy_name_is_unchanged_when_it_carries_allow_rules(self, core: K7Core):
"""Renaming it would strand policies created by older k7 versions,
which `_delete_sandbox_resources` deletes by name."""
mock_net = self._net(core)
await core._apply_sandbox_network_policies(
name="sb", namespace="default", egress_whitelist=None, ingress_ports=[8000]
)
assert self._ingress_body(mock_net).metadata.name == "sb-deny-ingress"
class TestExposeSandboxPorts:
"""NodePort Service in front of matching ingress rules."""
async def test_service_uses_local_external_traffic_policy(self, core: K7Core):
mock_v1 = AsyncMock()
created = MagicMock()
port = MagicMock()
port.port = 8000
port.node_port = 31234
created.spec.ports = [port]
mock_v1.create_namespaced_service.return_value = created
_setup_clients(core, v1=mock_v1)
result = await core._apply_sandbox_expose_service("sb", "default", [8000])
assert result.success is True
assert result.data == {"node_ports": {8000: 31234}}
body = mock_v1.create_namespaced_service.call_args.kwargs["body"]
assert body.metadata.name == "sb-expose"
assert body.spec.type == "NodePort"
# Without Local the client's source IP is SNAT'd and cidr: rules are theatre.
assert body.spec.external_traffic_policy == "Local"
assert body.spec.selector == {"app": "sb"}
async def test_unallocated_node_port_is_a_loud_failure(self, core: K7Core):
mock_v1 = AsyncMock()
created = MagicMock()
port = MagicMock()
port.port = 8000
port.node_port = None
created.spec.ports = [port]
mock_v1.create_namespaced_service.return_value = created
_setup_clients(core, v1=mock_v1)
result = await core._apply_sandbox_expose_service("sb", "default", [8000])
assert result.success is False
assert "no node port" in result.error
async def test_expose_without_matching_ingress_port_is_rejected(self, core: K7Core):
mock_apps = AsyncMock()
_setup_clients(core, apps=mock_apps, v1=AsyncMock(), net=AsyncMock())
cfg = SandboxConfig(name="sb1", image="alpine:3.20", ingress_ports=[8000], expose_ports=[9000])
result = await core.create_sandbox(cfg)
assert result.success is False
assert "[9000]" in result.error
mock_apps.create_namespaced_deployment.assert_not_called()
class TestReadSandboxIngressRules:
async def test_deny_all_reads_back_as_none(self, core: K7Core):
np = MagicMock()
np.spec.ingress = []
mock_net = AsyncMock()
mock_net.read_namespaced_network_policy.return_value = np
_setup_clients(core, net=mock_net)
assert await core._read_sandbox_ingress_rules("sb", "default") == (None, None)
async def test_round_trips_ports_and_sources(self, core: K7Core):
"""A fork must re-derive the source's policy from what it reads back."""
core_net = AsyncMock()
_setup_clients(core, net=core_net)
await core._apply_sandbox_network_policies(
name="sb",
namespace="default",
egress_whitelist=None,
ingress_ports=[8000],
ingress_from=["sandbox:alice", "namespace:team-a", "cidr:203.0.113.0/24"],
)
applied = [
c.kwargs["body"]
for c in core_net.create_namespaced_network_policy.call_args_list
if c.kwargs["body"].metadata.name == "sb-deny-ingress"
][0]
core_net.read_namespaced_network_policy.return_value = applied
ports, sources = await core._read_sandbox_ingress_rules("sb", "default")
assert ports == [8000]
assert sources == ["sandbox:alice", "namespace:team-a", "cidr:203.0.113.0/24"]
async def test_default_source_reads_back_as_empty_list(self, core: K7Core):
core_net = AsyncMock()
_setup_clients(core, net=core_net)
await core._apply_sandbox_network_policies(
name="sb", namespace="default", egress_whitelist=None, ingress_ports=[8000]
)
applied = [
c.kwargs["body"]
for c in core_net.create_namespaced_network_policy.call_args_list
if c.kwargs["body"].metadata.name == "sb-deny-ingress"
][0]
core_net.read_namespaced_network_policy.return_value = applied
assert await core._read_sandbox_ingress_rules("sb", "default") == ([8000], [])
class TestReadSandboxEgressWhitelist:
async def test_cilium_wildcard_round_trips(self, core: K7Core):
mock_custom = AsyncMock()
mock_custom.get_namespaced_custom_object.return_value = {
"spec": {"egress": [{"toFQDNs": [{"matchPattern": "**.docker.com"}, {"matchName": "docker.io"}]}]}
}
_setup_clients(core, custom=mock_custom)
assert await core._read_sandbox_egress_whitelist("sb", "default") == ["*.docker.com", "docker.io"]
async def test_v1_netpol_cidrs(self, core: K7Core):
np = MagicMock()
peer = MagicMock()
peer.ip_block.cidr = "10.0.0.0/8"
rule = MagicMock()
rule.to = [peer]
np.spec.egress = [rule]
mock_net = AsyncMock()
mock_net.read_namespaced_network_policy.return_value = np
mock_custom = AsyncMock()
mock_custom.get_namespaced_custom_object.side_effect = ApiException(status=404)
_setup_clients(core, net=mock_net, custom=mock_custom)
assert await core._read_sandbox_egress_whitelist("sb", "default") == ["10.0.0.0/8"]
async def test_block_all_egress_round_trips_as_empty_list(self, core: K7Core):
"""`egress_whitelist=[]` (block everything) must not degrade to None
(open) when a fork inherits it."""
np = MagicMock()
np.spec.egress = []
mock_net = AsyncMock()
mock_net.read_namespaced_network_policy.return_value = np
mock_custom = AsyncMock()
mock_custom.get_namespaced_custom_object.side_effect = ApiException(status=404)
_setup_clients(core, net=mock_net, custom=mock_custom)
assert await core._read_sandbox_egress_whitelist("sb", "default") == []
async def test_no_policy_means_open_egress(self, core: K7Core):
mock_net = AsyncMock()
mock_net.read_namespaced_network_policy.side_effect = ApiException(status=404)
mock_custom = AsyncMock()
mock_custom.get_namespaced_custom_object.side_effect = ApiException(status=404)
_setup_clients(core, net=mock_net, custom=mock_custom)
assert await core._read_sandbox_egress_whitelist("sb", "default") is None
# --- per-node agent forwarding ---
class TestK7dAgentForwarding:
def _remote_pod(self, node: str = "node-remote", name: str = "sb1-pod"):
pod = mock_pod(name=name, phase="Running")
pod.spec.node_name = node
return pod
async def test_pause_forwards_to_remote_agent(self, core: K7Core):
"""k7d pause of a sandbox on another node is forwarded, not failed."""
dep = mock_deployment(annotations={"k7.katakate.org/backend": "k7d"})
mock_apps = AsyncMock()
mock_apps.read_namespaced_deployment.return_value = dep
mock_v1 = AsyncMock()
mock_v1.list_namespaced_pod.return_value = MagicMock(items=[self._remote_pod()])
_setup_clients(core, apps=mock_apps, v1=mock_v1)
with (
patch.dict("os.environ", {"K7_NODE_NAME": "node-local"}),
patch.object(
core,
"_k7d_forward_vm_op",
new_callable=AsyncMock,
return_value=OperationResult(success=True, message="forwarded"),
) as fwd,
):
result = await core.pause_sandbox("sb1")
assert result.success is True
fwd.assert_called_once()
assert fwd.call_args.args[0] == "node-remote"
assert fwd.call_args.args[1] == "pause"
assert fwd.call_args.args[2]["pod_name"] == "sb1-pod"
mock_apps.patch_namespaced_deployment.assert_awaited()
async def test_fork_looks_up_on_remote_creates_locally(self, core: K7Core):
"""k7d fork Kubernetes writes stay on k7-api; agent only looks up the VM."""
dep = mock_deployment(annotations={"k7.katakate.org/backend": "k7d"})
dep.spec.template.metadata.labels = {"app": "sb1"}
dep.metadata.labels = {"app": "sb1"}
dep.spec.selector.match_labels = {"app": "sb1"}
mock_apps = AsyncMock()
mock_apps.read_namespaced_deployment.return_value = dep
mock_v1 = AsyncMock()
mock_v1.list_namespaced_pod.return_value = MagicMock(items=[self._remote_pod()])
_setup_clients(core, apps=mock_apps, v1=mock_v1)
vm = {"vm_id": "vm-1-0", "sandbox_id": "cri-src", "cluster_id": "cri-src"}
with (
patch.dict("os.environ", {"K7_NODE_NAME": "node-local"}),
patch.object(core, "lookup_k7d_vm", new=AsyncMock(return_value=vm)),
patch.object(core, "_wait_for_pod_container_started", new=AsyncMock(return_value="fork-pod")),
patch.object(core, "_read_sandbox_egress_whitelist", new=AsyncMock(return_value=[])),
patch.object(core, "_read_sandbox_ingress_rules", new=AsyncMock(return_value=([], []))),
patch.object(core, "_apply_sandbox_network_policies", new=AsyncMock(return_value=OperationResult(True))),
patch.object(core, "_k7d_forward_vm_op", new_callable=AsyncMock) as fwd,
):
result = await core.fork_sandbox("sb1", "sb1-fork")
assert result.success, result.error
fwd.assert_not_called()
mock_apps.create_namespaced_deployment.assert_awaited()
new_dep = mock_apps.create_namespaced_deployment.call_args.kwargs["body"]
assert new_dep.spec.template.spec.node_name == "node-remote"
async def test_agent_refuses_to_reforward(self, core: K7Core):
"""An agent that resolves a sandbox to yet another node must fail
loudly instead of forwarding again (no proxy loops)."""
with patch.dict("os.environ", {"K7_AGENT": "1", "K7_NODE_NAME": "node-a"}):
try:
await core._k7d_forward_vm_op("node-b", "pause", {}, timeout=5)
raise AssertionError("expected RuntimeError")
except RuntimeError as e:
assert "refusing to re-forward" in str(e)
async def test_missing_agent_pod_fails_loudly(self, core: K7Core):
mock_v1 = AsyncMock()
mock_v1.list_namespaced_pod.return_value = MagicMock(items=[])
_setup_clients(core, v1=mock_v1)
try:
await core._k7d_agent_base_url("node-remote")
raise AssertionError("expected RuntimeError")
except RuntimeError as e:
assert "k7-agent" in str(e)
async def test_missing_token_fails_loudly(self, core: K7Core, tmp_path):
with patch.dict("os.environ", {"K7_AGENT_TOKEN_FILE": str(tmp_path / "absent")}):
try:
core._k7d_agent_token()
raise AssertionError("expected RuntimeError")
except RuntimeError as e:
assert "agent token" in str(e)
async def test_nodes_storage_reports_unreachable_agent(self, core: K7Core, tmp_path):
"""A node with no Ready agent gets a loud error entry, not omission."""
token_file = tmp_path / "agent_token"
token_file.write_text("tok\n")
mock_v1 = AsyncMock()
mock_v1.list_node.return_value = SimpleNamespace(items=[_node_named("n1")])
mock_v1.list_namespaced_pod.return_value = MagicMock(items=[])
_setup_clients(core, v1=mock_v1)
with patch.dict("os.environ", {"K7_AGENT_TOKEN_FILE": str(token_file)}):
result = await core.nodes_storage()
assert "n1" in result
assert "k7-agent" in result["n1"]["error"]
def _node_named(name: str):
n = MagicMock()
n.metadata.name = name
return n
# --- exec_command ---
class TestExecCommand:
async def test_sandbox_not_found(self, core: K7Core):
mock_apps = AsyncMock()
mock_apps.read_namespaced_deployment.side_effect = ApiException(status=404)
_setup_clients(core, apps=mock_apps, v1=AsyncMock())
result = await core.exec_command("missing", "echo hi")
assert result.exit_code == 1
assert "not found" in result.stderr
async def test_no_pods(self, core: K7Core):
mock_apps = AsyncMock()
mock_v1 = AsyncMock()
mock_v1.list_namespaced_pod.return_value = MagicMock(items=[])
_setup_clients(core, apps=mock_apps, v1=mock_v1)
result = await core.exec_command("sb1", "echo hi")
assert result.exit_code == 1
assert "No pods" in result.stderr
async def test_pod_not_running(self, core: K7Core):
mock_apps = AsyncMock()
mock_v1 = AsyncMock()
pod = mock_pod(phase="Pending")
mock_v1.list_namespaced_pod.return_value = MagicMock(items=[pod])
_setup_clients(core, apps=mock_apps, v1=mock_v1)
result = await core.exec_command("sb1", "echo hi")
assert result.exit_code == 1
assert "not running" in result.stderr
async def test_successful_exec(self, core: K7Core):
mock_apps = AsyncMock()
mock_v1 = AsyncMock()
pod = mock_pod(phase="Running")
mock_v1.list_namespaced_pod.return_value = MagicMock(items=[pod])
_setup_clients(core, apps=mock_apps, v1=mock_v1)
with patch("k7.core.core.WsApiClient") as mock_ws_cls:
mock_ws_ctx = AsyncMock()
mock_ws_cls.return_value.__aenter__ = AsyncMock(return_value=mock_ws_ctx)
mock_ws_cls.return_value.__aexit__ = AsyncMock(return_value=False)
mock_v1_ws = AsyncMock()
mock_v1_ws.connect_get_namespaced_pod_exec.return_value = "hello\n"
with patch("k7.core.core.client.CoreV1Api", return_value=mock_v1_ws):
result = await core.exec_command("sb1", "echo hello")
assert result.exit_code == 0
assert result.stdout == "hello\n"
assert result.duration_ms >= 0
# --- get_sandbox_metrics ---
class TestGetSandboxMetrics:
async def test_running_pod_with_metrics(self, core: K7Core):
mock_apps = AsyncMock()
mock_v1 = AsyncMock()
mock_metrics = AsyncMock()
_setup_clients(core, apps=mock_apps, v1=mock_v1, metrics=mock_metrics)
dep = mock_deployment(name="sb1")
pod = mock_pod(name="sb1-abc", phase="Running")
mock_v1.list_namespaced_pod.return_value = MagicMock(items=[pod])
mock_metrics.get_namespaced_custom_object.return_value = {
"containers": [{"usage": {"cpu": "100m", "memory": "128Mi"}}]
}
with patch.object(core, "_get_kata_sandboxes", new_callable=AsyncMock, return_value=[dep]):
result = await core.get_sandbox_metrics(namespace="default")
assert len(result) == 1
assert result[0]["cpu_usage"] == "100m"
assert result[0]["memory_usage"] == "128Mi"
async def test_no_pods_skipped(self, core: K7Core):
mock_v1 = AsyncMock()
mock_v1.list_namespaced_pod.return_value = MagicMock(items=[])
mock_metrics = AsyncMock()
_setup_clients(core, v1=mock_v1, metrics=mock_metrics)
dep = mock_deployment(name="sb1")
with patch.object(core, "_get_kata_sandboxes", new_callable=AsyncMock, return_value=[dep]):
result = await core.get_sandbox_metrics(namespace="default")
assert result == []
async def test_pod_not_running_skipped(self, core: K7Core):
mock_v1 = AsyncMock()
pod = mock_pod(phase="Pending")
mock_v1.list_namespaced_pod.return_value = MagicMock(items=[pod])
mock_metrics = AsyncMock()
_setup_clients(core, v1=mock_v1, metrics=mock_metrics)
dep = mock_deployment(name="sb1")
with patch.object(core, "_get_kata_sandboxes", new_callable=AsyncMock, return_value=[dep]):
result = await core.get_sandbox_metrics(namespace="default")
assert result == []
async def test_metrics_api_error(self, core: K7Core):
_setup_clients(core)
with patch.object(core, "_get_kata_sandboxes", new_callable=AsyncMock, side_effect=Exception("boom")):
result = await core.get_sandbox_metrics(namespace="default")
assert result == []
# --- _wait_for_pvc_bound ---
class TestWaitForPvcBound:
async def test_bound_immediately(self, core: K7Core):
mock_v1 = AsyncMock()
pvc = MagicMock()
pvc.status.phase = "Bound"
mock_v1.read_namespaced_persistent_volume_claim.return_value = pvc
_setup_clients(core, v1=mock_v1)
result = await core._wait_for_pvc_bound("pvc1", timeout=5)
assert result.success is True
async def test_bound_after_retry(self, core: K7Core):
mock_v1 = AsyncMock()
pending = MagicMock()
pending.status.phase = "Pending"
bound = MagicMock()
bound.status.phase = "Bound"
mock_v1.read_namespaced_persistent_volume_claim.side_effect = [pending, bound]
_setup_clients(core, v1=mock_v1)
with patch("k7.core.core.asyncio.sleep", new_callable=AsyncMock):
result = await core._wait_for_pvc_bound("pvc1", timeout=10)
assert result.success is True
async def test_404_then_bound(self, core: K7Core):
mock_v1 = AsyncMock()
bound = MagicMock()
bound.status.phase = "Bound"
mock_v1.read_namespaced_persistent_volume_claim.side_effect = [
ApiException(status=404),
bound,
]
_setup_clients(core, v1=mock_v1)
with patch("k7.core.core.asyncio.sleep", new_callable=AsyncMock):
result = await core._wait_for_pvc_bound("pvc1", timeout=10)
assert result.success is True
async def test_non_404_error(self, core: K7Core):
mock_v1 = AsyncMock()
mock_v1.read_namespaced_persistent_volume_claim.side_effect = ApiException(status=500)
_setup_clients(core, v1=mock_v1)
result = await core._wait_for_pvc_bound("pvc1", timeout=5)
assert result.success is False
assert "check error" in result.error
async def test_timeout(self, core: K7Core):
mock_v1 = AsyncMock()
pending = MagicMock()
pending.status.phase = "Pending"
mock_v1.read_namespaced_persistent_volume_claim.return_value = pending
_setup_clients(core, v1=mock_v1)
elapsed = 0.0
def advancing_time():
nonlocal elapsed
elapsed += 1.0
return elapsed
with (
patch("k7.core.core.time.time", side_effect=advancing_time),
patch("k7.core.core.asyncio.sleep", new_callable=AsyncMock),
):
result = await core._wait_for_pvc_bound("pvc1", timeout=2)
assert result.success is False
assert "not Bound" in result.error
# --- _wait_for_job ---
class TestWaitForJob:
async def test_job_succeeds(self, core: K7Core):
mock_batch = AsyncMock()
job = MagicMock()
job.status.succeeded = 1
job.status.failed = 0
mock_batch.read_namespaced_job.return_value = job
_setup_clients(core, batch=mock_batch)
result = await core._wait_for_job("j1", timeout=5)
assert result.success is True
async def test_job_fails(self, core: K7Core):
mock_batch = AsyncMock()
job = MagicMock()
job.status.succeeded = 0
job.status.failed = 1
mock_batch.read_namespaced_job.return_value = job
_setup_clients(core, batch=mock_batch)
result = await core._wait_for_job("j1", timeout=5)
assert result.success is False
assert "failed" in result.error
async def test_404_then_succeeds(self, core: K7Core):
mock_batch = AsyncMock()
job_ok = MagicMock()
job_ok.status.succeeded = 1
job_ok.status.failed = 0
mock_batch.read_namespaced_job.side_effect = [ApiException(status=404), job_ok]
_setup_clients(core, batch=mock_batch)
with patch("k7.core.core.asyncio.sleep", new_callable=AsyncMock):
result = await core._wait_for_job("j1", timeout=10)
assert result.success is True
async def test_non_404_error(self, core: K7Core):
mock_batch = AsyncMock()
mock_batch.read_namespaced_job.side_effect = ApiException(status=500)
_setup_clients(core, batch=mock_batch)
result = await core._wait_for_job("j1", timeout=5)
assert result.success is False
assert "wait error" in result.error
async def test_timeout(self, core: K7Core):
mock_batch = AsyncMock()
job = MagicMock()
job.status.succeeded = 0
job.status.failed = 0
mock_batch.read_namespaced_job.return_value = job
_setup_clients(core, batch=mock_batch)
elapsed = 0.0
def advancing_time():
nonlocal elapsed
elapsed += 1.0
return elapsed
with (
patch("k7.core.core.time.time", side_effect=advancing_time),
patch("k7.core.core.asyncio.sleep", new_callable=AsyncMock),
):
result = await core._wait_for_job("j1", timeout=2)
assert result.success is False
assert "did not complete" in result.error
# --- create_sandbox validation ---
class TestCreateSandboxValidation:
async def test_invalid_limits_rejected(self, core: K7Core):
mock_apps = AsyncMock()
mock_v1 = AsyncMock()
mock_net = AsyncMock()
_setup_clients(core, apps=mock_apps, v1=mock_v1, net=mock_net)
cfg = SandboxConfig(name="sb1", image="alpine:3.18", limits={"cpu": "0m"})
result = await core.create_sandbox(cfg)
assert result.success is False
assert "Invalid resource limits" in result.error
async def test_empty_env_file_rejected(self, core: K7Core, tmp_path):
mock_apps = AsyncMock()
mock_v1 = AsyncMock()
mock_net = AsyncMock()
_setup_clients(core, apps=mock_apps, v1=mock_v1, net=mock_net)
env_file = tmp_path / ".env"
env_file.write_text("# only comments\n\n")
cfg = SandboxConfig(name="sb1", image="alpine:3.18", env_file=str(env_file))
result = await core.create_sandbox(cfg)
assert result.success is False
assert "empty or invalid" in result.error
async def test_ingress_from_without_ports_rejected(self, core: K7Core):
_setup_clients(core, apps=AsyncMock(), v1=AsyncMock(), net=AsyncMock())
cfg = SandboxConfig(name="sb1", image="alpine:3.18", ingress_from=["sandbox:alice"])
result = await core.create_sandbox(cfg)
assert result.success is False
assert "ingress_from requires ingress_ports" in result.error
async def test_unparseable_ingress_source_rejected_before_provisioning(self, core: K7Core):
mock_apps = AsyncMock()
_setup_clients(core, apps=mock_apps, v1=AsyncMock(), net=AsyncMock())
cfg = SandboxConfig(name="sb1", image="alpine:3.18", ingress_ports=[8000], ingress_from=["nonsense"])
result = await core.create_sandbox(cfg)
assert result.success is False
assert "nonsense" in result.error
mock_apps.create_namespaced_deployment.assert_not_called()
# --- _wait_for_pod_container_started ---
class TestWaitForPodContainerStarted:
async def test_returns_pod_name_when_running(self, core: K7Core):
mock_v1 = AsyncMock()
pod = mock_pod(name="sb1-abc", phase="Running")
cs = MagicMock()
cs.name = "sandbox"
cs.state.running = True
pod.status.container_statuses = [cs]
mock_v1.list_namespaced_pod.return_value = MagicMock(items=[pod])
_setup_clients(core, v1=mock_v1)
result = await core._wait_for_pod_container_started("sb1", "default", timeout_seconds=5)
assert result == "sb1-abc"
async def test_returns_none_on_timeout_no_pods(self, core: K7Core):
mock_v1 = AsyncMock()
mock_v1.list_namespaced_pod.return_value = MagicMock(items=[])
_setup_clients(core, v1=mock_v1)
elapsed = 0.0
def advancing_time():
nonlocal elapsed
elapsed += 1.0
return elapsed
with (
patch("k7.core.core.time.time", side_effect=advancing_time),
patch("k7.core.core.asyncio.sleep", new_callable=AsyncMock),
):
result = await core._wait_for_pod_container_started("sb1", "default", timeout_seconds=2)
assert result is None
async def test_waits_for_running_phase(self, core: K7Core):
mock_v1 = AsyncMock()
pending_pod = mock_pod(name="sb1-abc", phase="Pending")
running_pod = mock_pod(name="sb1-abc", phase="Running")
cs = MagicMock()
cs.name = "sandbox"
cs.state.running = True
running_pod.status.container_statuses = [cs]
mock_v1.list_namespaced_pod.side_effect = [
MagicMock(items=[pending_pod]),
MagicMock(items=[running_pod]),
]
_setup_clients(core, v1=mock_v1)
with patch("k7.core.core.asyncio.sleep", new_callable=AsyncMock):
result = await core._wait_for_pod_container_started("sb1", "default", timeout_seconds=60)
assert result == "sb1-abc"
async def test_returns_none_when_container_not_started(self, core: K7Core):
mock_v1 = AsyncMock()
pod = mock_pod(name="sb1-abc", phase="Running")
cs = MagicMock()
cs.name = "sandbox"
cs.state.running = None
cs.state.waiting = MagicMock()
pod.status.container_statuses = [cs]
mock_v1.list_namespaced_pod.return_value = MagicMock(items=[pod])
_setup_clients(core, v1=mock_v1)
elapsed = 0.0
def advancing_time():
nonlocal elapsed
elapsed += 1.0
return elapsed
with (
patch("k7.core.core.time.time", side_effect=advancing_time),
patch("k7.core.core.asyncio.sleep", new_callable=AsyncMock),
):
result = await core._wait_for_pod_container_started("sb1", "default", timeout_seconds=2)
assert result is None
# --- _list_backends_per_node ---
def _node(name: str, labels: dict[str, str], taints: list | None = None):
n = MagicMock()
n.metadata.name = name
n.metadata.labels = labels
n.spec.taints = taints or []
return n
class TestListBackendsPerNode:
async def test_extracts_backend_labels(self, core: K7Core):
mock_v1 = AsyncMock()
nodes = [
_node(
"node1",
{
"k7.katakate.org/backend-kata-qemu-longhorn": "true",
"k7.katakate.org/backend-kata-firecracker-devmapper": "true",
"kubernetes.io/hostname": "node1",
},
),
_node("node2", {"k7.katakate.org/backend-kata-qemu-longhorn": "true"}),
_node("node3", {}),
]
mock_v1.list_node.return_value = SimpleNamespace(items=nodes)
_setup_clients(core, v1=mock_v1)
result = await core._list_backends_per_node()
assert result["node1"] == ["kata-firecracker-devmapper", "kata-qemu-longhorn"]
assert result["node2"] == ["kata-qemu-longhorn"]
assert result["node3"] == []
async def test_skips_label_with_false_value(self, core: K7Core):
mock_v1 = AsyncMock()
mock_v1.list_node.return_value = SimpleNamespace(
items=[_node("n", {"k7.katakate.org/backend-kata-qemu-longhorn": "false"})]
)
_setup_clients(core, v1=mock_v1)
result = await core._list_backends_per_node()
assert result["n"] == []
class TestListClusterNodes:
async def test_hostname_backends_and_tenant(self, core: K7Core):
taint = SimpleNamespace(key=K7_TENANT_LABEL, value="acme", effect="NoSchedule")
mock_v1 = AsyncMock()
mock_v1.list_node.return_value = SimpleNamespace(
items=[
_node(
"k7-node-01",
{
"kubernetes.io/hostname": "k7-node-01",
"k7.katakate.org/backend-k7d": "true",
K7_TENANT_LABEL: "acme",
},
taints=[taint],
),
_node("k7-node-02", {"kubernetes.io/hostname": "k7-node-02"}),
]
)
_setup_clients(core, v1=mock_v1)
rows = await core.list_cluster_nodes()
by_name = {r["name"]: r for r in rows}
assert by_name["k7-node-01"]["hostname"] == "k7-node-01"
assert by_name["k7-node-01"]["backends"] == ["k7d"]
assert by_name["k7-node-01"]["tenant"] == "acme"
assert by_name["k7-node-02"]["hostname"] == "k7-node-02"
assert by_name["k7-node-02"]["backends"] == []
assert by_name["k7-node-02"]["tenant"] == ""
# --- _check_scheduling ---
def _event(reason: str, message: str):
e = MagicMock()
e.reason = reason
e.message = message
return e
class TestCheckScheduling:
async def test_running_pod_returns_success(self, core: K7Core):
mock_v1 = AsyncMock()
pod = mock_pod(name="sb1-x", phase="Running")
mock_v1.list_namespaced_pod.return_value = MagicMock(items=[pod])
_setup_clients(core, v1=mock_v1)
result = await core._check_scheduling("sb1", "default", "kata-qemu-longhorn", timeout_seconds=1)
assert result.success is True
async def test_failed_scheduling_returns_explanatory_error(self, core: K7Core):
mock_v1 = AsyncMock()
pod = mock_pod(name="sb1-x", phase="Pending")
mock_v1.list_namespaced_pod.return_value = MagicMock(items=[pod])
events = SimpleNamespace(
items=[
_event(
"FailedScheduling",
"0/2 nodes are available: 2 node(s) didn't match Pod's node affinity/selector.",
)
]
)
mock_v1.list_namespaced_event.return_value = events
# Two nodes labeled with FD only — no QL anywhere
mock_v1.list_node.return_value = SimpleNamespace(
items=[
_node("node1", {"k7.katakate.org/backend-kata-firecracker-devmapper": "true"}),
_node("node2", {"k7.katakate.org/backend-kata-firecracker-devmapper": "true"}),
]
)
_setup_clients(core, v1=mock_v1)
result = await core._check_scheduling("sb1", "default", "kata-qemu-longhorn", timeout_seconds=1)
assert result.success is False
assert "kata-qemu-longhorn" in result.error
assert "kata-firecracker-devmapper" in result.error
assert "node1" in result.error
assert "node2" in result.error
async def test_no_pod_yet_eventually_returns_success(self, core: K7Core):
mock_v1 = AsyncMock()
mock_v1.list_namespaced_pod.return_value = MagicMock(items=[])
_setup_clients(core, v1=mock_v1)
elapsed = 0.0
def advancing_time():
nonlocal elapsed
elapsed += 1.0
return elapsed
with (
patch("k7.core.core.time.time", side_effect=advancing_time),
patch("k7.core.core.asyncio.sleep", new_callable=AsyncMock),
):
result = await core._check_scheduling("sb1", "default", "kata-qemu-longhorn", timeout_seconds=2)
assert result.success is True
async def test_unrelated_event_does_not_trigger(self, core: K7Core):
mock_v1 = AsyncMock()
pod = mock_pod(name="sb1-x", phase="Pending")
mock_v1.list_namespaced_pod.return_value = MagicMock(items=[pod])
mock_v1.list_namespaced_event.return_value = SimpleNamespace(
items=[_event("Scheduled", "Successfully assigned default/sb1-x to node1")]
)
_setup_clients(core, v1=mock_v1)
elapsed = 0.0
def advancing_time():
nonlocal elapsed
elapsed += 1.0
return elapsed
with (
patch("k7.core.core.time.time", side_effect=advancing_time),
patch("k7.core.core.asyncio.sleep", new_callable=AsyncMock),
):
result = await core._check_scheduling("sb1", "default", "kata-qemu-longhorn", timeout_seconds=2)
assert result.success is True
# --- _cilium_available / _apply_cilium_egress_policy ---
class TestCiliumAvailable:
async def test_returns_true_when_crd_exists(self, core: K7Core):
fake_api_ext = AsyncMock()
fake_api_ext.read_custom_resource_definition.return_value = MagicMock()
_setup_clients(core, apiext=fake_api_ext)
assert await core._cilium_available() is True
async def test_returns_false_when_crd_missing(self, core: K7Core):
fake_api_ext = AsyncMock()
fake_api_ext.read_custom_resource_definition.side_effect = ApiException(status=404)
_setup_clients(core, apiext=fake_api_ext)
assert await core._cilium_available() is False
class TestApplyCiliumEgressPolicy:
async def test_builds_policy_with_fqdns_and_cidrs(self, core: K7Core):
mock_custom = AsyncMock()
_setup_clients(core, custom=mock_custom)
result = await core._apply_cilium_egress_policy(
sandbox_name="sb1",
namespace="default",
cidrs=["10.0.0.0/8"],
fqdns=["api.openai.com", "*.huggingface.co"],
)
assert result.success is True
args, kwargs = mock_custom.create_namespaced_custom_object.call_args
body = kwargs["body"]
assert body["kind"] == "CiliumNetworkPolicy"
assert body["metadata"]["name"] == "sb1-egress"
assert body["spec"]["endpointSelector"]["matchLabels"]["katakate.org/sandbox"] == "sb1"
egress = body["spec"]["egress"]
# DNS rule first, then FQDN rule, then CIDR rule
assert any("toFQDNs" in r for r in egress)
fqdn_rule = next(r for r in egress if "toFQDNs" in r)
assert {"matchName": "api.openai.com"} in fqdn_rule["toFQDNs"]
# Leading `*.` must be translated to Cilium's multi-label wildcard
# `**.` — a single `*` never crosses label boundaries, silently
# breaking nested subdomains like production.cloudfront.docker.com.
assert {"matchPattern": "**.huggingface.co"} in fqdn_rule["toFQDNs"]
assert any("toCIDR" in r and r["toCIDR"] == ["10.0.0.0/8"] for r in egress)
# DNS proxy allowance is present
assert any("toPorts" in r and any("rules" in p and "dns" in p["rules"] for p in r["toPorts"]) for r in egress)
async def test_leading_wildcard_becomes_multilabel_pattern(self, core: K7Core):
mock_custom = AsyncMock()
_setup_clients(core, custom=mock_custom)
result = await core._apply_cilium_egress_policy(
sandbox_name="sb1",
namespace="default",
cidrs=[],
fqdns=["*.docker.com", "**.already.com", "api-*.foo.com"],
)
assert result.success is True
body = mock_custom.create_namespaced_custom_object.call_args.kwargs["body"]
fqdn_rule = next(r for r in body["spec"]["egress"] if "toFQDNs" in r)
assert {"matchPattern": "**.docker.com"} in fqdn_rule["toFQDNs"]
# Explicit `**.` is passed through untouched; mid-label wildcards too.
assert {"matchPattern": "**.already.com"} in fqdn_rule["toFQDNs"]
assert {"matchPattern": "api-*.foo.com"} in fqdn_rule["toFQDNs"]
async def test_replaces_on_409(self, core: K7Core):
mock_custom = AsyncMock()
mock_custom.create_namespaced_custom_object.side_effect = ApiException(status=409)
_setup_clients(core, custom=mock_custom)
result = await core._apply_cilium_egress_policy("sb1", "default", [], ["api.openai.com"])
assert result.success is True
mock_custom.replace_namespaced_custom_object.assert_called_once()
async def test_non_conflict_error_propagates(self, core: K7Core):
mock_custom = AsyncMock()
mock_custom.create_namespaced_custom_object.side_effect = ApiException(status=500)
_setup_clients(core, custom=mock_custom)
result = await core._apply_cilium_egress_policy("sb1", "default", [], ["api.openai.com"])
assert result.success is False
assert "CiliumNetworkPolicy" in result.error
# --- kubernetes_asyncio ApiClient lifecycle ---
class TestApiClientLifecycle:
async def test_typed_clients_share_one_api_client(self, core: K7Core):
api = MagicMock()
api.close = AsyncMock()
with (
patch("k7.core.core.client.ApiClient", return_value=api) as api_cls,
patch("k7.core.core.client.CoreV1Api") as core_cls,
patch("k7.core.core.client.AppsV1Api") as apps_cls,
):
await core._get_core_v1_client()
await core._get_apps_v1_client()
api_cls.assert_called_once()
core_cls.assert_called_once_with(api)
apps_cls.assert_called_once_with(api)
await core.aclose()
api.close.assert_awaited_once()
assert core._api_client is None
assert core._core_v1_client is None
assert core._apps_v1_client is None
async def test_aclose_is_idempotent(self, core: K7Core):
await core.aclose()
await core.aclose()
async def test_aclose_allows_a_new_session(self, core: K7Core):
first = MagicMock()
first.close = AsyncMock()
second = MagicMock()
second.close = AsyncMock()
with (
patch("k7.core.core.client.ApiClient", side_effect=[first, second]),
patch("k7.core.core.client.CoreV1Api"),
):
await core._get_core_v1_client()
await core.aclose()
first.close.assert_awaited_once()
await core._get_core_v1_client()
await core.aclose()
second.close.assert_awaited_once()
class TestCreateSnapshotMessage:
async def test_named_snapshot_keeps_create_message(self, core: K7Core):
"""Regression: quiesced success used to drop the VolumeSnapshot name.
``k7 snapshot create`` echoes ``result.message``; empty stdout failed
``TestSnapshotCrud`` with ``assert snap in cp.stdout``.
"""
dep = MagicMock()
dep.metadata.annotations = {}
dep.status.ready_replicas = 1
mock_apps = AsyncMock()
mock_apps.read_namespaced_deployment.return_value = dep
_setup_clients(core, apps=mock_apps)
with (
patch.object(core, "_detect_backend", new=AsyncMock(return_value="kata-qemu-longhorn")),
patch.object(
core,
"_create_volume_snapshot",
new=AsyncMock(
return_value=OperationResult(success=True, message="Snapshot snap1 created for PVC sb-root-lh")
),
),
patch.object(
core,
"exec_command",
new=AsyncMock(return_value=ExecResult(exit_code=0, stdout="", stderr="", duration_ms=0)),
),
):
result = await core.create_snapshot("sb", "snap1")
assert result.success, result.error
assert "snap1" in result.message