Files

442 lines
18 KiB
Python

from __future__ import annotations
import json
import subprocess
import tempfile
import unittest
from pathlib import Path
from typing import Any
from .runtime import (
BrowserRuntimeError,
GenerationConflict,
NativeRuntimeManager,
RuntimeCleanupPending,
UnitStatus,
)
class FakeUnits:
def __init__(self) -> None:
self.units: dict[str, UnitStatus] = {}
self.commands: list[tuple[str, list[str]]] = []
self.next_pid = 1000
self.fail_stop: set[str] = set()
def start(self, unit: str, command: list[str], **kwargs: Any) -> UnitStatus:
del kwargs
self.commands.append((unit, list(command)))
self.next_pid += 1
status = UnitStatus(True, "active", self.next_pid, str(self.next_pid), "")
self.units[unit] = status
return status
def status(self, unit: str) -> UnitStatus:
return self.units.get(unit, UnitStatus(False, "inactive", 0, "", ""))
def stop(self, unit: str, timeout: float) -> None:
del timeout
if unit in self.fail_stop:
raise OSError("stop failed")
self.units[unit] = UnitStatus(False, "inactive", 0, "", "")
def reset(self, unit: str) -> None:
del unit
class SystemdUnitManagerTests(unittest.TestCase):
def test_start_status_stop_and_reset(self) -> None:
calls: list[list[str]] = []
stop_calls = 0
def runner(args: list[str], **kwargs: Any) -> subprocess.CompletedProcess[str]:
nonlocal stop_calls
del kwargs
calls.append(args)
if args[0] == "systemctl" and args[2] == "show":
return subprocess.CompletedProcess(
args,
0,
"ActiveState=active\nSubState=running\nMainPID=42\n"
"ExecMainStartTimestampMonotonic=1\nExecMainStatus=0\n",
"",
)
if args[0] == "systemctl" and args[2] == "stop":
stop_calls += 1
return subprocess.CompletedProcess(args, 1 if stop_calls == 1 else 0, "", "")
return subprocess.CompletedProcess(args, 0, "", "")
from .runtime import SystemdUnitManager
with tempfile.TemporaryDirectory() as directory:
manager = SystemdUnitManager("systemd-run", "systemctl", runner)
status = manager.start(
"creatorhub-test.service",
["/bin/true"],
environment={"B": "2", "A": "1"},
working_directory=Path(directory),
stdout_path=Path(directory) / "runtime.log",
limits={"MemoryMax": "1G"},
)
self.assertEqual(status.pid, 42)
self.assertTrue(any("--setenv" in item for item in calls[0]))
self.assertEqual(manager.status("creatorhub-test.service").state, "running")
manager.stop("creatorhub-test.service", 1)
manager.reset("creatorhub-test.service")
self.assertTrue(any(args[2] == "kill" for args in calls))
self.assertTrue(any(args[2] == "reset-failed" for args in calls))
def test_status_and_start_failures_are_explicit(self) -> None:
from .runtime import BrowserRuntimeError, SystemdUnitManager
with self.assertRaises(BrowserRuntimeError):
SystemdUnitManager(runner=lambda *args, **kwargs: (_ for _ in ()).throw(OSError("down"))).status("x")
def failed_start(args: list[str], **kwargs: Any) -> subprocess.CompletedProcess[str]:
del kwargs
return subprocess.CompletedProcess(args, 1, "", "rejected")
with tempfile.TemporaryDirectory() as directory, self.assertRaises(
BrowserRuntimeError
):
SystemdUnitManager("systemd-run", "systemctl", failed_start).start(
"x", ["/bin/true"], environment={}, working_directory=Path(directory),
stdout_path=Path(directory) / "x.log", limits={},
)
def missing_status(args: list[str], **kwargs: Any) -> subprocess.CompletedProcess[str]:
del kwargs
if args[2] == "show":
return subprocess.CompletedProcess(args, 1, "", "missing")
return subprocess.CompletedProcess(args, 0, "", "")
manager = SystemdUnitManager("systemd-run", "systemctl", missing_status)
self.assertFalse(manager.status("missing").active)
self.assertEqual(manager.status("missing").state, "not-found")
manager.reset("missing")
def test_stop_and_reset_failures_are_not_hidden(self) -> None:
from .runtime import BrowserRuntimeError, SystemdUnitManager
def kill_fails(args: list[str], **kwargs: Any) -> subprocess.CompletedProcess[str]:
del kwargs
if args[2] == "stop":
return subprocess.CompletedProcess(args, 1, "", "failed")
if args[2] == "kill":
return subprocess.CompletedProcess(args, 1, "", "failed")
return subprocess.CompletedProcess(args, 0, "", "")
with self.assertRaises(BrowserRuntimeError):
SystemdUnitManager("systemd-run", "systemctl", kill_fails).stop("x", 1)
def reset_fails(args: list[str], **kwargs: Any) -> subprocess.CompletedProcess[str]:
del kwargs
return subprocess.CompletedProcess(args, 2, "", "failed")
with self.assertRaises(BrowserRuntimeError):
SystemdUnitManager("systemd-run", "systemctl", reset_fails).reset("x")
def test_allocator_existing_and_file_lock_edges(self) -> None:
from .runtime import (
BrowserRuntimeError,
DisplayAllocator,
FileLock,
PortAllocator,
)
with tempfile.TemporaryDirectory() as directory:
lock_dir = Path(directory)
lock = FileLock(lock_dir / "resource.lock")
self.assertTrue(lock.acquire())
self.assertTrue(lock.acquire())
lock.release()
lock.release()
ports = PortAllocator(lock_dir)
lease = ports.reserve_existing("cdp", 19999)
with self.assertRaises(BrowserRuntimeError):
ports.reserve_existing("cdp", 19999)
lease.lock.release()
with self.assertRaises(BrowserRuntimeError):
ports.reserve_existing("cdp", 0)
displays = DisplayAllocator(lock_dir)
display = displays.reserve_existing(100)
with self.assertRaises(BrowserRuntimeError):
displays.reserve_existing(100)
display.lock.release()
with self.assertRaises(BrowserRuntimeError):
displays.reserve_existing(0)
class NativeRuntimeManagerTests(unittest.TestCase):
def setUp(self) -> None:
self.temp = tempfile.TemporaryDirectory()
root = Path(self.temp.name)
self.state = root / "state"
self.profiles = root / "profiles"
self.units = FakeUnits()
self.manager = NativeRuntimeManager(
state_dir=self.state,
profile_root=self.profiles,
node_id="node-a",
browser_path="/bin/true",
unit_manager=self.units,
min_free_bytes=0,
)
self.manager._wait_for_display = lambda record: None
self.manager._wait_for_cdp = lambda record: None
def tearDown(self) -> None:
self.manager.close()
self.temp.cleanup()
def payload(self, **extra: Any) -> dict[str, Any]:
value: dict[str, Any] = {
"alias": "account-a",
"name": "Account A",
"profile_id": "account-a-id",
"cmd": ["--fingerprint=1000", "about:blank"],
"binding_version": 1,
"network_exit_id": "",
"network_exit": {},
"stopped": False,
}
value.update(extra)
return value
def test_create_allocates_native_runtime_and_persists_generation(self) -> None:
result = self.manager.create(self.payload())
self.assertEqual(result["state"], "running")
self.assertEqual(result["node_id"], "node-a")
self.assertRegex(result["runtime_id"], r"^[a-f0-9]{64}$")
self.assertTrue(result["endpoint"].startswith("http://127.0.0.1:"))
runtime_file = next(self.state.glob("runtimes/*/runtime.json"))
stored = json.loads(runtime_file.read_text())
self.assertEqual(stored["runtime_id"], result["runtime_id"])
browser_command = self.units.commands[-1][1]
self.assertNotIn("--no-sandbox", browser_command)
self.assertNotIn("--", browser_command)
self.assertIn("--disable-gpu", browser_command)
self.assertIn("--disable-gpu-compositing", browser_command)
self.assertIn("--remote-debugging-address=127.0.0.1", browser_command)
self.assertIn("--user-data-dir=" + stored["profile_dir"], browser_command)
def test_insufficient_disk_is_rejected_before_process_side_effects(self) -> None:
limited = NativeRuntimeManager(
state_dir=self.state / "disk-state",
profile_root=self.profiles / "disk-profiles",
node_id="node-a",
browser_path="/bin/true",
unit_manager=self.units,
min_free_bytes=10**20,
)
with self.assertRaises(BrowserRuntimeError) as caught:
limited.create(self.payload(alias="disk-account"))
self.assertIn("free disk space", str(caught.exception))
self.assertEqual(self.units.commands, [])
limited.close()
def test_external_display_reuses_existing_xvfb_without_starting_or_stopping_it(self) -> None:
external = NativeRuntimeManager(
state_dir=self.state / "external-state",
profile_root=self.profiles / "external-profiles",
node_id="node-a",
browser_path="/bin/true",
unit_manager=self.units,
external_display=99,
min_free_bytes=0,
)
external._display_available = lambda display: display == 99
external._wait_for_cdp = lambda record: None
result = external.create(self.payload(alias="external-account"))
self.assertEqual(result["display"], 99)
self.assertEqual(result["display_mode"], "external")
self.assertEqual(len(self.units.commands), 1)
self.assertIn("browser", self.units.commands[0][0])
generation = {
"binding_version": 1,
"runtime_id": result["runtime_id"],
"network_id": result["network_id"],
}
external.change_state("external-account", "stop", generation)
self.assertEqual(external.list_public()[0]["state"], "stopped")
external.close()
def test_external_display_rejects_a_second_active_runtime(self) -> None:
external = NativeRuntimeManager(
state_dir=self.state / "external-state-2",
profile_root=self.profiles / "external-profiles-2",
node_id="node-a",
browser_path="/bin/true",
unit_manager=self.units,
external_display=99,
min_free_bytes=0,
)
external._display_available = lambda display: display == 99
external._wait_for_cdp = lambda record: None
external.create(self.payload(alias="external-account"))
with self.assertRaises(BrowserRuntimeError) as caught:
external.create(self.payload(alias="external-account-2", profile_id="other"))
self.assertIn("already in use", str(caught.exception))
external.close()
def test_generation_mismatch_cannot_stop_or_remove_new_runtime(self) -> None:
result = self.manager.create(self.payload())
wrong = {
"binding_version": 1,
"runtime_id": "0" * 64,
"network_id": result["network_id"],
}
with self.assertRaises(GenerationConflict):
self.manager.change_state("account-a", "stop", wrong)
with self.assertRaises(GenerationConflict):
self.manager.remove("account-a", wrong)
self.assertEqual(self.manager.list_public()[0]["state"], "running")
def test_stop_and_repeated_remove_keep_profile(self) -> None:
result = self.manager.create(self.payload())
generation = {
"binding_version": 1,
"runtime_id": result["runtime_id"],
"network_id": result["network_id"],
}
self.manager.change_state("account-a", "stop", generation)
self.assertEqual(self.manager.list_public()[0]["state"], "stopped")
self.manager.remove("account-a", generation)
self.manager.remove("account-a", generation)
self.assertEqual(self.manager.list_public()[0]["state"], "released")
profile_dirs = list(self.profiles.iterdir())
self.assertEqual(len(profile_dirs), 1)
self.assertTrue(profile_dirs[0].is_dir())
def test_same_profile_is_exclusive_across_manager_instances(self) -> None:
first = self.manager.create(self.payload())
del first
other = NativeRuntimeManager(
state_dir=self.state,
profile_root=self.profiles,
node_id="node-a",
browser_path="/bin/true",
unit_manager=self.units,
min_free_bytes=0,
)
other._wait_for_display = lambda record: None
other._wait_for_cdp = lambda record: None
with self.assertRaises(BrowserRuntimeError):
other.create(self.payload(alias="account-b"))
other.close()
def test_cleanup_failure_is_visible_and_retryable(self) -> None:
result = self.manager.create(self.payload())
generation = {
"binding_version": 1,
"runtime_id": result["runtime_id"],
"network_id": result["network_id"],
}
self.units.fail_stop.update(
{
f"creatorhub-{result['alias']}-{result['runtime_id'][:16]}-browser.service",
}
)
with self.assertRaises(RuntimeCleanupPending):
self.manager.remove("account-a", generation)
public = self.manager.list_public()[0]
self.assertEqual(public["cleanup_state"], "pending")
self.assertTrue(public["cleanup_error"])
self.units.fail_stop.clear()
self.manager.retry_cleanup("account-a", generation)
self.assertEqual(self.manager.list_public()[0]["state"], "released")
def test_reserved_browser_flags_are_rejected_before_side_effects(self) -> None:
with self.assertRaises(BrowserRuntimeError):
self.manager.create(
self.payload(cmd=["--no-sandbox", "about:blank"])
)
self.assertEqual(list(self.state.glob("runtimes/*/runtime.json")), [])
def test_stopped_runtime_can_be_created_without_processes(self) -> None:
result = self.manager.create(self.payload(stopped=True))
self.assertEqual(result["state"], "stopped")
self.assertEqual(self.units.commands, [])
def test_timeout_marks_runtime_cleanup_pending(self) -> None:
self.manager.ready_timeout = 0.001
self.manager._wait_for_display = NativeRuntimeManager._wait_for_display.__get__(self.manager)
self.manager._wait_for_cdp = NativeRuntimeManager._wait_for_cdp.__get__(self.manager)
with self.assertRaises(BrowserRuntimeError) as caught:
self.manager.create(self.payload())
self.assertIn("Xvfb", str(caught.exception))
public = self.manager.list_public()[0]
self.assertEqual(public["cleanup_state"], "pending")
self.assertEqual(public["state"], "failed")
def test_cancel_during_start_is_visible_and_does_not_start_browser(self) -> None:
started = __import__("threading").Event()
original_wait = self.manager._wait_for_display
def wait(record: object) -> None:
started.set()
while True:
self.manager._check_cancel(record) # type: ignore[arg-type]
__import__("time").sleep(0.001)
self.manager._wait_for_display = wait
result: list[object] = []
def create() -> None:
try:
result.append(self.manager.create(self.payload()))
except Exception as exc: # the assertion below checks the typed outcome
result.append(exc)
thread = __import__("threading").Thread(target=create)
thread.start()
self.assertTrue(started.wait(1))
runtime_file = next(self.state.glob("runtimes/*/runtime.json"))
stored = json.loads(runtime_file.read_text())
self.manager.cancel(
"account-a",
{
"binding_version": stored["binding_version"],
"runtime_id": stored["runtime_id"],
"network_id": stored["network_id"],
},
)
thread.join(1)
self.assertFalse(thread.is_alive())
self.assertTrue(result and isinstance(result[0], BrowserRuntimeError))
self.assertEqual(self.manager.list_public()[0]["cleanup_state"], "pending")
self.manager._wait_for_display = original_wait
def test_restart_recovers_only_matching_active_units(self) -> None:
result = self.manager.create(self.payload())
self.manager.close()
recovered = NativeRuntimeManager(
state_dir=self.state,
profile_root=self.profiles,
node_id="node-a",
browser_path="/bin/true",
unit_manager=self.units,
min_free_bytes=0,
)
self.assertEqual(recovered.list_public()[0]["runtime_id"], result["runtime_id"])
self.assertEqual(recovered.list_public()[0]["node_id"], "node-a")
recovered.close()
def test_old_released_generation_cannot_touch_new_runtime(self) -> None:
old = self.manager.create(self.payload(stopped=True))
old_generation = {
"binding_version": 1,
"runtime_id": old["runtime_id"],
"network_id": old["network_id"],
}
self.manager.remove("account-a", old_generation)
new = self.manager.create(self.payload(stopped=True))
self.assertNotEqual(old["runtime_id"], new["runtime_id"])
self.manager.remove("account-a", old_generation)
self.assertEqual(self.manager.list_public()[-1]["runtime_id"], new["runtime_id"])
self.assertEqual(self.manager.list_public()[-1]["state"], "stopped")
if __name__ == "__main__":
unittest.main()