Files
2026-07-15 11:46:27 +08:00

404 lines
14 KiB
Python

"""Optional supervisor for Appium in the macOS single-machine real-device
workflow. Appium gates real driver connections; task inspection is provided by
the Host Agent's own console.
Lives in ``host_agent`` because it spawns and monitors host-level processes
alongside the heartbeat/claim loop. Off by default; see ``HostAgentConfig``.
Design (``openspec/changes/host-agent-dependency-supervisor/design.md``):
- Default off. Each dependency also has its own opt-in flag.
- Before spawn: TCP connect + dependency-specific HTTP health probe. A healthy
listener is *adopted* (never killed or restarted). An unhealthy listener is
a port conflict — logged and skipped. No listener → spawn.
- Only supervisor-spawned processes are restarted on unexpected exit, with
capped exponential backoff and a per-process-lifetime attempt ceiling.
- Adopted processes are never touched by ``stop()``.
- Spawns use stdlib ``subprocess.Popen`` (no new dependency). Child stdout/stderr
is forwarded into this module's logger, tagged per dependency.
"""
from __future__ import annotations
import asyncio
import logging
import socket
import subprocess
import threading
from collections.abc import Awaitable, Callable, Sequence
from dataclasses import dataclass
from enum import Enum
from typing import TYPE_CHECKING
import httpx
if TYPE_CHECKING:
from host_agent.config import HostAgentConfig
_LOGGER = logging.getLogger("host_agent.dependency_supervisor")
_PROBE_TCP_TIMEOUT_SECONDS = 1.0
_PROBE_HTTP_TIMEOUT_SECONDS = 2.0
class ProbeResult(str, Enum):
"""Outcome of probing a dependency's configured address."""
HEALTHY = "healthy"
UNHEALTHY_LISTENER = "unhealthy_listener"
NO_LISTENER = "no_listener"
def probe_appium(host: str, port: int) -> ProbeResult:
"""Probe Appium at ``host:port``. Healthy iff ``GET /status`` returns 200
with a JSON body (the documented shape per ``docs/MACOS_IPHONE_SETUP.md``
§6 — contains ``ready``/build info)."""
return _probe_http(host, port, path="/status")
def _probe_http(host: str, port: int, *, path: str) -> ProbeResult:
# Step 1: plain TCP connect — distinguish "nothing listening" (→ spawn)
# from "something is there but wrong" (→ port conflict, skip).
try:
with socket.create_connection((host, port), timeout=_PROBE_TCP_TIMEOUT_SECONDS):
pass
except OSError:
return ProbeResult.NO_LISTENER
# Step 2: dependency-specific HTTP health check.
try:
response = httpx.get(
f"http://{host}:{port}{path}",
timeout=_PROBE_HTTP_TIMEOUT_SECONDS,
)
except httpx.HTTPError:
return ProbeResult.UNHEALTHY_LISTENER
if response.status_code != 200:
return ProbeResult.UNHEALTHY_LISTENER
try:
response.json()
except ValueError:
return ProbeResult.UNHEALTHY_LISTENER
return ProbeResult.HEALTHY
def appium_argv_factory(host: str, port: int) -> list[str]:
return ["appium", "--address", host, "--port", str(port)]
@dataclass
class SupervisedDependency:
"""Config + mutable runtime state for one supervised external process."""
name: str
host: str
port: int
argv_factory: Callable[[str, int], Sequence[str]]
probe: Callable[[str, int], ProbeResult]
# Mutable state:
process: subprocess.Popen | None = None
adopted: bool = False
given_up: bool = False
restart_attempts: int = 0
ready: bool = False
@dataclass
class _SupervisorKnobs:
"""Tunables exposed for tests; production uses the defaults."""
startup_timeout_seconds: float = 30.0
readiness_poll_interval_seconds: float = 0.5
crash_poll_interval_seconds: float = 1.0
initial_backoff_seconds: float = 1.0
max_backoff_seconds: float = 30.0
terminate_grace_period_seconds: float = 5.0
class DependencySupervisor:
"""Manage zero or more supervised local dependencies.
Each dependency is either *adopted* (an existing healthy instance was
detected — never killed or restarted) or *spawned* (a child ``Popen``
handle the supervisor owns — restart-on-crash with capped exponential
backoff up to ``max_attempts`` per process lifetime).
"""
def __init__(
self,
dependencies: list[SupervisedDependency],
*,
max_attempts: int,
knobs: _SupervisorKnobs | None = None,
sleep: Callable[[float], Awaitable[None]] = asyncio.sleep,
popen_factory: Callable[[Sequence[str]], subprocess.Popen] = subprocess.Popen,
logger: logging.Logger = _LOGGER,
) -> None:
self._dependencies = dependencies
self._max_attempts = max_attempts
self._knobs = knobs or _SupervisorKnobs()
self._sleep = sleep
self._popen_factory = popen_factory
self._logger = logger
self._reader_threads: list[threading.Thread] = []
@classmethod
def from_host_agent_config(
cls,
ha_config: HostAgentConfig,
**kwargs: object,
) -> DependencySupervisor:
"""Build a supervisor reflecting ``HostAgentConfig`` flags.
Caller is responsible for only invoking this when
``dependency_supervisor_enabled`` is true; the resulting supervisor
will contain only the dependencies whose individual ``*_supervised``
flag is also true (possibly an empty list).
"""
deps: list[SupervisedDependency] = []
if ha_config.appium_supervised:
deps.append(
SupervisedDependency(
name="appium",
host=ha_config.appium_host,
port=ha_config.appium_port,
argv_factory=appium_argv_factory,
probe=probe_appium,
)
)
return cls(
deps,
max_attempts=ha_config.dependency_restart_max_attempts,
**kwargs, # type: ignore[arg-type]
)
@property
def dependencies(self) -> list[SupervisedDependency]:
return list(self._dependencies)
async def start(self) -> None:
"""Probe + spawn/adopt + readiness wait.
Call before the heartbeat loop's first ``connect_devices()`` pass so
a supervised Appium is up before any ``Driver.connect()`` attempt.
"""
for dep in self._dependencies:
await self._start_one(dep)
async def _start_one(self, dep: SupervisedDependency) -> None:
result = await asyncio.to_thread(dep.probe, dep.host, dep.port)
if result is ProbeResult.HEALTHY:
dep.adopted = True
dep.ready = True
self._logger.info(
"dependency-supervisor: %s adopted existing instance at %s:%s",
dep.name,
dep.host,
dep.port,
)
return
if result is ProbeResult.UNHEALTHY_LISTENER:
self._logger.error(
"dependency-supervisor: %s port %s occupied by an unhealthy "
"listener; leaving un-started to avoid a port conflict",
dep.name,
dep.port,
)
dep.given_up = True
return
await self._spawn_with_readiness_wait(dep)
async def _spawn_with_readiness_wait(
self,
dep: SupervisedDependency,
) -> bool:
"""Spawn ``dep`` and wait for it to become healthy.
Returns True on readiness, False otherwise. On spawn failure (missing
executable) or readiness timeout while the process is still running,
marks the dependency as ``given_up`` (no crash-restart loop entered).
If the process exits during startup, returns False without giving up
— the crash-restart loop in ``run()`` handles subsequent restarts.
"""
argv = list(dep.argv_factory(dep.host, dep.port))
try:
proc = self._popen_factory(
argv,
stdout=subprocess.PIPE,
stderr=subprocess.STDOUT,
bufsize=1,
text=True,
)
except FileNotFoundError:
self._logger.error(
"dependency-supervisor: %s spawn failed — executable not "
"found on PATH (tried: %s)",
dep.name,
argv[0],
)
dep.given_up = True
return False
except OSError as exc:
self._logger.error(
"dependency-supervisor: %s spawn failed: %s",
dep.name,
exc,
)
dep.given_up = True
return False
dep.process = proc
dep.adopted = False
dep.ready = False
self._start_reader_thread(dep, proc)
self._logger.info(
"dependency-supervisor: %s spawned (pid %s) at %s:%s",
dep.name,
proc.pid,
dep.host,
dep.port,
)
loop = asyncio.get_running_loop()
deadline = loop.time() + self._knobs.startup_timeout_seconds
while True:
if proc.poll() is not None:
self._logger.error(
"dependency-supervisor: %s exited during startup with code %s",
dep.name,
proc.returncode,
)
dep.ready = False
return False
probe_result = await asyncio.to_thread(dep.probe, dep.host, dep.port)
if probe_result is ProbeResult.HEALTHY:
dep.ready = True
self._logger.info(
"dependency-supervisor: %s ready at %s:%s",
dep.name,
dep.host,
dep.port,
)
return True
if loop.time() >= deadline:
self._logger.error(
"dependency-supervisor: %s started but did not become "
"healthy within %ss; process still running, not entering "
"crash-restart loop",
dep.name,
self._knobs.startup_timeout_seconds,
)
dep.ready = False
return False
await self._sleep(self._knobs.readiness_poll_interval_seconds)
def _start_reader_thread(
self,
dep: SupervisedDependency,
proc: subprocess.Popen,
) -> None:
def reader_loop() -> None:
stdout = proc.stdout
if stdout is None:
return
for raw_line in stdout:
line = raw_line.rstrip()
if line:
self._logger.info("%s: %s", dep.name, line)
thread = threading.Thread(
target=reader_loop,
name=f"dep-sup-{dep.name}",
daemon=True,
)
thread.start()
self._reader_threads.append(thread)
async def run(self, stop: asyncio.Event) -> None:
"""Background monitor: detect crashes and apply backoff restart.
Returns when ``stop`` is set. Crashes are only detected for
spawned processes (not adopted ones).
"""
if not self._dependencies:
await stop.wait()
return
while not stop.is_set():
for dep in self._dependencies:
await self._check_one(dep)
try:
await asyncio.wait_for(
stop.wait(),
timeout=self._knobs.crash_poll_interval_seconds,
)
except TimeoutError:
continue
async def _check_one(self, dep: SupervisedDependency) -> None:
if dep.given_up or dep.adopted or dep.process is None:
return
if dep.process.poll() is None:
return
self._logger.warning(
"dependency-supervisor: %s exited unexpectedly with code %s",
dep.name,
dep.process.returncode,
)
dep.process = None
dep.ready = False
dep.restart_attempts += 1
if dep.restart_attempts > self._max_attempts:
self._logger.error(
"dependency-supervisor: %s restart attempts exhausted "
"(crashes: %s, limit: %s); giving up for the rest of this "
"process lifetime",
dep.name,
dep.restart_attempts,
self._max_attempts,
)
dep.given_up = True
return
backoff = min(
self._knobs.initial_backoff_seconds * (2 ** (dep.restart_attempts - 1)),
self._knobs.max_backoff_seconds,
)
self._logger.info(
"dependency-supervisor: %s restarting in %.1fs (crash %s/%s)",
dep.name,
backoff,
dep.restart_attempts,
self._max_attempts,
)
await self._sleep(backoff)
await self._spawn_with_readiness_wait(dep)
async def stop(self) -> None:
"""Terminate spawned children. Adopted processes are left untouched."""
for dep in self._dependencies:
if dep.adopted or dep.process is None:
continue
proc = dep.process
self._logger.info(
"dependency-supervisor: stopping %s (pid %s)",
dep.name,
proc.pid,
)
proc.terminate()
try:
await asyncio.to_thread(
proc.wait, self._knobs.terminate_grace_period_seconds
)
except subprocess.TimeoutExpired:
self._logger.warning(
"dependency-supervisor: %s did not exit within %ss, killing",
dep.name,
self._knobs.terminate_grace_period_seconds,
)
proc.kill()
await asyncio.to_thread(proc.wait)
except OSError:
pass
finally:
dep.process = None