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