fix(cloakserve): add idle cleanup for seeded profiles (#352)
* fix(cloakserve): add idle cleanup for seeded profiles * docs(cloakserve): document idle process cleanup
This commit is contained in:
@@ -181,6 +181,7 @@ class ChromePool:
|
||||
default_seed: str | None = None,
|
||||
default_locale: str | None = None,
|
||||
default_timezone: str | None = None,
|
||||
idle_timeout: float = 0.0,
|
||||
):
|
||||
self._binary = binary
|
||||
self._global_args = global_args
|
||||
@@ -189,12 +190,14 @@ class ChromePool:
|
||||
self._default_seed = default_seed
|
||||
self._default_locale = default_locale
|
||||
self._default_timezone = default_timezone
|
||||
self._idle_timeout = idle_timeout
|
||||
self._processes: dict[str, ChromeProcess] = {}
|
||||
self._default: ChromeProcess | None = None
|
||||
self._locks: dict[str, asyncio.Lock] = {}
|
||||
self._next_port = BASE_CDP_PORT
|
||||
# Connection refcounting for status reporting
|
||||
self._connections: dict[str, int] = {}
|
||||
self._idle_tasks: dict[str, asyncio.Task] = {}
|
||||
|
||||
def _get_lock(self, seed: str) -> asyncio.Lock:
|
||||
if seed not in self._locks:
|
||||
@@ -224,6 +227,7 @@ class ChromePool:
|
||||
|
||||
def connect(self, seed_key: str) -> None:
|
||||
"""Increment connection refcount for a seed."""
|
||||
self._cancel_idle_cleanup(seed_key)
|
||||
self._connections[seed_key] = self._connections.get(seed_key, 0) + 1
|
||||
|
||||
def disconnect(self, seed_key: str) -> None:
|
||||
@@ -231,9 +235,54 @@ class ChromePool:
|
||||
count = self._connections.get(seed_key, 0) - 1
|
||||
if count <= 0:
|
||||
self._connections.pop(seed_key, None)
|
||||
self._schedule_idle_cleanup(seed_key)
|
||||
else:
|
||||
self._connections[seed_key] = count
|
||||
|
||||
def _cancel_idle_cleanup(self, seed_key: str) -> None:
|
||||
task = self._idle_tasks.pop(seed_key, None)
|
||||
if task is None or task.done():
|
||||
return
|
||||
try:
|
||||
current_task = asyncio.current_task()
|
||||
except RuntimeError:
|
||||
current_task = None
|
||||
if task is not current_task:
|
||||
task.cancel()
|
||||
|
||||
def _discard_idle_task(self, seed_key: str, task: asyncio.Task) -> None:
|
||||
if self._idle_tasks.get(seed_key) is task:
|
||||
self._idle_tasks.pop(seed_key, None)
|
||||
|
||||
def _schedule_idle_cleanup(self, seed_key: str) -> None:
|
||||
if self._idle_timeout <= 0 or seed_key not in self._processes:
|
||||
return
|
||||
|
||||
self._cancel_idle_cleanup(seed_key)
|
||||
try:
|
||||
loop = asyncio.get_running_loop()
|
||||
except RuntimeError:
|
||||
return
|
||||
|
||||
task = loop.create_task(
|
||||
self._cleanup_after_idle(seed_key, self._idle_timeout),
|
||||
name=f"cloakserve-idle-cleanup-{seed_key}",
|
||||
)
|
||||
self._idle_tasks[seed_key] = task
|
||||
task.add_done_callback(lambda done_task: self._discard_idle_task(seed_key, done_task))
|
||||
|
||||
async def _cleanup_after_idle(self, seed_key: str, timeout: float) -> None:
|
||||
try:
|
||||
await asyncio.sleep(timeout)
|
||||
if self._connections.get(seed_key, 0) > 0 or seed_key not in self._processes:
|
||||
return
|
||||
logger.info("Cleaning up idle Chrome process (seed=%s)", seed_key)
|
||||
await self._cleanup_process(seed_key)
|
||||
except asyncio.CancelledError:
|
||||
raise
|
||||
except Exception:
|
||||
logger.exception("Idle cleanup failed for seed=%s", seed_key)
|
||||
|
||||
async def get_or_launch(
|
||||
self,
|
||||
seed: str | None,
|
||||
@@ -271,6 +320,8 @@ class ChromePool:
|
||||
if seed_key in self._processes:
|
||||
proc = self._processes[seed_key]
|
||||
if proc.process.poll() is None:
|
||||
if seed_key in self._idle_tasks:
|
||||
self._schedule_idle_cleanup(seed_key)
|
||||
if any([extra_args, timezone, locale, proxy, geoip]):
|
||||
logger.warning(
|
||||
"Seed %s already running (port %d, tz=%s, locale=%s, proxy=%s) — "
|
||||
@@ -360,6 +411,7 @@ class ChromePool:
|
||||
|
||||
async def _cleanup_process(self, key: str) -> None:
|
||||
"""Terminate a Chrome process and clean up."""
|
||||
self._cancel_idle_cleanup(key)
|
||||
proc = self._processes.pop(key, None)
|
||||
if not proc:
|
||||
return
|
||||
@@ -377,6 +429,13 @@ class ChromePool:
|
||||
|
||||
async def shutdown(self) -> None:
|
||||
"""Terminate all Chrome processes."""
|
||||
idle_tasks = list(self._idle_tasks.values())
|
||||
self._idle_tasks.clear()
|
||||
for task in idle_tasks:
|
||||
if not task.done():
|
||||
task.cancel()
|
||||
if idle_tasks:
|
||||
await asyncio.gather(*idle_tasks, return_exceptions=True)
|
||||
for key in list(self._processes.keys()):
|
||||
await self._cleanup_process(key)
|
||||
logger.info("All Chrome processes terminated")
|
||||
@@ -479,6 +538,7 @@ async def handle_root(request: web.Request) -> web.Response:
|
||||
"port": proc.cdp_port,
|
||||
"seed": proc.seed,
|
||||
"connections": pool._connections.get(key, 0),
|
||||
"idle_cleanup_pending": key in pool._idle_tasks,
|
||||
"timezone": proc.timezone,
|
||||
"locale": proc.locale,
|
||||
"proxy": proc.proxy,
|
||||
@@ -486,6 +546,7 @@ async def handle_root(request: web.Request) -> web.Response:
|
||||
return web.json_response({
|
||||
"status": "ok",
|
||||
"active": len(processes),
|
||||
"idle_timeout": pool._idle_timeout,
|
||||
"processes": processes,
|
||||
})
|
||||
|
||||
@@ -687,6 +748,23 @@ def _default_data_dir() -> str:
|
||||
return str(Path.home() / ".cloakbrowser" / "cloakserve")
|
||||
|
||||
|
||||
def _parse_idle_timeout(value: str) -> float:
|
||||
value = value.strip()
|
||||
if value.lower() in {"0", "false", "off", "none", "disabled"}:
|
||||
return 0.0
|
||||
timeout = float(value)
|
||||
if timeout < 0:
|
||||
raise ValueError("--idle-timeout must be greater than or equal to 0")
|
||||
return timeout
|
||||
|
||||
|
||||
def _default_idle_timeout() -> float:
|
||||
value = os.environ.get("CLOAKSERVE_IDLE_TIMEOUT")
|
||||
if value is None:
|
||||
return 0.0
|
||||
return _parse_idle_timeout(value)
|
||||
|
||||
|
||||
def parse_cli_args(argv: list[str]) -> tuple[dict, list[str]]:
|
||||
"""Parse cloakserve-specific args, return (config, passthrough_args).
|
||||
|
||||
@@ -702,12 +780,14 @@ def parse_cli_args(argv: list[str]) -> tuple[dict, list[str]]:
|
||||
"default_seed": None,
|
||||
"default_locale": None,
|
||||
"default_timezone": None,
|
||||
"idle_timeout": _default_idle_timeout(),
|
||||
}
|
||||
passthrough = []
|
||||
# Flags consumed by cloakserve (not passed to Chrome)
|
||||
consumed_prefixes = (
|
||||
"--port=",
|
||||
"--data-dir=",
|
||||
"--idle-timeout=",
|
||||
"--remote-debugging-port=",
|
||||
"--remote-debugging-address=",
|
||||
)
|
||||
@@ -717,6 +797,8 @@ def parse_cli_args(argv: list[str]) -> tuple[dict, list[str]]:
|
||||
config["port"] = int(arg.split("=", 1)[1])
|
||||
elif arg.startswith("--data-dir="):
|
||||
config["data_dir"] = arg.split("=", 1)[1]
|
||||
elif arg.startswith("--idle-timeout="):
|
||||
config["idle_timeout"] = _parse_idle_timeout(arg.split("=", 1)[1])
|
||||
elif arg == "--headless=false" or arg == "--headless=False":
|
||||
config["headless"] = False
|
||||
passthrough.append(arg)
|
||||
@@ -761,6 +843,7 @@ def main() -> None:
|
||||
default_seed=config["default_seed"],
|
||||
default_locale=config["default_locale"],
|
||||
default_timezone=config["default_timezone"],
|
||||
idle_timeout=config["idle_timeout"],
|
||||
)
|
||||
|
||||
app = web.Application()
|
||||
|
||||
Reference in New Issue
Block a user