fix(bypass): close the gaps a helper that outlives its request opened (#1244)

assumptions the code around it still made were written for a helper that
was killed after every request.

A bypass that hits the child's deadline is cancelled from the calling
thread, which returns the moment the cancellation is scheduled - so the
helper went on to serve the next request while the abandoned one was
still closing its browser, on the same loop, sharing the DISPLAY globals
and one process group. The deadline now lives inside the loop, where
asyncio.wait_for() waits for the unwind before it raises, with the
calling thread keeping a bounded backstop in case the cleanup wedges
too. Both budgets are set so the child still answers before the parent
gives up on it.

The helper's cookie store survived the request as well, and the whole of
it is exported back to the parent on every answer - so clearance the
parent had purged for one host came back the next time some other host
was solved, the dead-cookie resurrection _redirect_loop_handoff purges
to avoid. The child starts each request from an empty store again; the
parent already runs the cached-cookie check against a superset of it.

DNS config is compared against what the helper is actually resolving
through rather than skipped whenever the parent reports "auto", so a
user flipping CUSTOM_DNS back to auto - which applies live - reaches a
warm helper instead of leaving it on an abandoned DoH resolver.

The 15s exit grace is now asked only of a helper that can still read its
stdin. One dropped mid-bypass never returns to that read, so the grace
could only end in the kill - while a user cancelling a download, and
every bypass queued behind them on LOCKED, waited it out.

Result files are cleaned on the timeout and cancellation paths too,
staging file included, rather than only when the answer was read.
This commit is contained in:
CaliBrain
2026-08-20 19:29:33 -04:00
committed by GitHub
parent bd21ec1257
commit 5247ec6124
3 changed files with 340 additions and 48 deletions
+115 -47
View File
@@ -58,14 +58,21 @@ _LOADING_BODY_LENGTH_MAX = 50
_PAGE_BODY_PREVIEW_CHARS = 500
_BROWSER_START_TIMEOUT_SECONDS = 45.0
_BYPASS_SUBPROCESS_TIMEOUT_SECONDS = 420.0
# Same budget as the Docker helper process, applied to the in-process CDP path so both
# branches of get() are bounded the same way.
_IN_PROCESS_BYPASS_TIMEOUT_SECONDS = _BYPASS_SUBPROCESS_TIMEOUT_SECONDS
# How long a cancelled bypass may take to close its browser before the calling thread
# stops waiting for it. Counted on top of the bypass deadline, so every budget below is
# set to leave room for it.
_CDP_UNWIND_GRACE_SECONDS = 15.0
# Same wall-clock budget as the Docker helper process, applied to the in-process CDP path
# so both branches of get() are bounded the same way: the deadline plus the unwind grace
# comes to _BYPASS_SUBPROCESS_TIMEOUT_SECONDS either way.
_IN_PROCESS_BYPASS_TIMEOUT_SECONDS = _BYPASS_SUBPROCESS_TIMEOUT_SECONDS - _CDP_UNWIND_GRACE_SECONDS
_BYPASS_CHILD_ENV = "SHELFMARK_INTERNAL_BYPASSER_CHILD"
# The helper bounds each bypass below the parent's deadline, so it is the side that gives
# up first: it still gets to report the timeout and close its browser, and stays available
# for the next request. A parent that hit its deadline first could only kill the helper,
# throwing away a process the next request would have to start again.
# throwing away a process the next request would have to start again. The 30s covers the
# unwind grace as well, so a helper that times out and closes its browser as slowly as it
# is allowed to still answers with 15s to spare.
_CHILD_BYPASS_TIMEOUT_SECONDS = _BYPASS_SUBPROCESS_TIMEOUT_SECONDS - 30.0
# The helper publishes its answer by writing the result file the request named, so the
# parent waits by watching for that file rather than by reading a stream it would have to
@@ -220,14 +227,33 @@ class _CdpWorker:
msg = "CDP worker loop failed to start"
raise RuntimeError(msg)
@staticmethod
async def _bounded(coro: Any, timeout: float | None) -> Any:
"""Run the coroutine under its deadline, on the loop that owns it.
The deadline has to be enforced from inside the loop rather than by the calling
thread: asyncio.wait_for() cancels the bypass and then *waits for it to unwind*,
so `finally: await _close_cdp_driver(driver)` has finished by the time this
raises. Cancelling from outside returns the moment the cancellation is scheduled,
which in a helper serving many requests let the abandoned bypass close its browser
while the next one was already opening its own - on the same loop, sharing the
DISPLAY globals and one process group.
"""
if timeout is None:
return await coro
return await asyncio.wait_for(coro, timeout)
def run(self, coro: Any, timeout: float | None = None) -> Any:
self.start()
if not self._loop or self._loop.is_closed():
msg = "CDP worker loop not available"
raise RuntimeError(msg)
future = asyncio.run_coroutine_threadsafe(coro, self._loop)
future = asyncio.run_coroutine_threadsafe(self._bounded(coro, timeout), self._loop)
# Backstop for an unwind that wedges too: _close_cdp_driver awaits websockets that
# a dead browser may never answer, and _bounded cannot outlive its own cleanup.
wait_for = None if timeout is None else timeout + _CDP_UNWIND_GRACE_SECONDS
try:
return future.result(timeout=timeout)
return future.result(timeout=wait_for)
except TimeoutError:
# Otherwise the coroutine keeps running in the worker loop after we stop
# waiting, holding the browser and racing the next bypass.
@@ -874,6 +900,11 @@ def _terminate_helper_session(proc: subprocess.Popen[str]) -> None:
proc.wait(timeout=5)
def _part_path(result_path: Path) -> Path:
"""Where the helper stages a result before renaming it into place."""
return result_path.with_name(result_path.name + ".part")
class _BypassHelper:
"""The helper subprocess that runs the bypasses, kept alive across them.
@@ -935,22 +966,31 @@ class _BypassHelper:
self._proc = self._spawn()
return self._proc
def _discard(self) -> None:
"""Stop the helper and forget it."""
def _discard(self, *, wait_for_exit: bool = True) -> None:
"""Stop the helper and forget it.
`wait_for_exit` belongs to a helper that could still act on the closed pipe: an
idle one is sitting in its stdin read, notices EOF and exits on its own. A helper
dropped mid-bypass is blocked inside the solve and will not return to that read,
so the grace cannot end in anything but the kill below - and the caller waiting it
out is a user cancelling a download, holding LOCKED while every other bypass in
the worker queues behind them.
"""
proc = self._proc
self._proc = None
if proc is None:
return
# Closing stdin ends the helper's request loop, so an idle helper gets to exit on
# its own. One mid-bypass cannot answer, and is killed below once the grace passes.
# its own. One mid-bypass cannot answer, and is killed below.
with suppress(OSError):
if proc.stdin is not None and not proc.stdin.closed:
proc.stdin.close()
try:
proc.wait(timeout=_HELPER_SHUTDOWN_GRACE_SECONDS)
except subprocess.TimeoutExpired:
logger.warning("Bypass helper did not exit on request, killing its session")
if wait_for_exit:
try:
proc.wait(timeout=_HELPER_SHUTDOWN_GRACE_SECONDS)
except subprocess.TimeoutExpired:
logger.warning("Bypass helper did not exit on request, killing its session")
# Tear the session down either way: a helper killed mid-bypass leaves its Chrome
# and Xvfb running, and those leftovers are what made the next worker's browser
@@ -1013,7 +1053,8 @@ class _BypassHelper:
# A live helper can die between the liveness check and the write, so one retry
# on a fresh process. A fresh one failing here is a real failure.
logger.info("Bypass helper closed its pipe (%s), retrying on a new one", exc)
self._discard()
# Nothing to ask of a helper we cannot write to: its read end is already gone.
self._discard(wait_for_exit=False)
proc = self._ensure_running()
self._write(proc, request_line)
@@ -1034,27 +1075,33 @@ class _BypassHelper:
cancel_flag: Event | None,
) -> dict[str, Any]:
deadline = time.monotonic() + timeout
while not result_path.exists():
if proc.poll() is not None:
returncode = proc.returncode
self._discard()
msg = f"Internal bypasser helper exited without a result (code {returncode})"
raise RuntimeError(msg)
if cancel_flag is not None and cancel_flag.is_set():
# The helper is mid-bypass and cannot be told to stop, so it goes.
self._discard()
_check_cancellation(cancel_flag, "Bypass cancelled while waiting for helper")
if time.monotonic() >= deadline:
self._discard()
msg = "Internal bypasser helper process timed out"
raise TimeoutError(msg)
time.sleep(_HELPER_RESULT_POLL_SECONDS)
try:
while not result_path.exists():
if proc.poll() is not None:
returncode = proc.returncode
self._discard(wait_for_exit=False)
msg = f"Internal bypasser helper exited without a result (code {returncode})"
raise RuntimeError(msg)
if cancel_flag is not None and cancel_flag.is_set():
# The helper is mid-bypass and cannot be told to stop, so it goes.
self._discard(wait_for_exit=False)
_check_cancellation(cancel_flag, "Bypass cancelled while waiting for helper")
if time.monotonic() >= deadline:
self._discard(wait_for_exit=False)
msg = "Internal bypasser helper process timed out"
raise TimeoutError(msg)
time.sleep(_HELPER_RESULT_POLL_SECONDS)
return json.loads(result_path.read_text(encoding="utf-8"))
finally:
with suppress(OSError):
result_path.unlink()
# Every way out of here is final for this request: either the answer has been
# read, or the helper that would have written it has just been killed. Nothing
# will write these paths afterwards and nothing will come looking for them, so
# they are cleaned on the failure paths too - otherwise every cancelled
# download and every wedged solve leaves one behind for the container's life.
for path in (result_path, _part_path(result_path)):
with suppress(OSError):
path.unlink()
_BYPASS_HELPER = _BypassHelper()
@@ -1415,26 +1462,38 @@ def get_bypassed_page(
return response_html
def _dns_fingerprint(dns_config: dict[str, Any]) -> tuple[str, tuple[str, ...], bool]:
"""Reduce a DNS config to what has to match for two of them to be the same one."""
provider = str(dns_config.get("provider") or "").strip().lower()
servers = dns_config.get("servers") if provider == "manual" else None
server_list = tuple(str(server) for server in servers) if isinstance(servers, list) else ()
return (provider, server_list, bool(dns_config.get("doh_enabled")))
def _apply_parent_dns_config(dns_config: dict[str, Any]) -> None:
"""Mirror the parent process's active DNS provider in this helper subprocess.
DNS state is in-memory only, so a fresh helper defaults to system DNS and would
pre-resolve AA hostnames (for Chrome's --host-resolver-rules) against a resolver
that may be blocked/hijacked. Re-applying the parent's provider keeps the helper on
the same DoH/custom resolver the parent already validated.
DNS state is in-memory only, so a helper left to itself would pre-resolve AA hostnames
(for Chrome's --host-resolver-rules) against a resolver that may be blocked or
hijacked. Re-applying the parent's provider keeps the helper on the same DoH/custom
resolver the parent already validated.
Compared against what this process is *actually* resolving through, rather than
against the last config it happened to be handed. The helper now outlives the request,
so it has to be able to travel back to auto as well as away from it - which a user
flipping CUSTOM_DNS in settings does live, without a restart - and asking the network
module what it is doing beats keeping a second, drifting copy of that answer here.
"""
provider = str(dns_config.get("provider") or "").strip().lower()
# "auto" means the parent has not rotated off system DNS yet, so the helper's own
# default initialization already matches it - nothing to override.
if not provider or provider == "auto":
wanted = _dns_fingerprint(dns_config)
provider, servers, use_doh = wanted
if not provider:
return
manual_servers = dns_config.get("servers") if provider == "manual" else None
# set_dns_provider() rebuilds the resolvers, so it is worth doing only on a real change.
if wanted == _dns_fingerprint(network.get_dns_config()):
return
try:
network.set_dns_provider(
provider,
manual_servers,
use_doh=bool(dns_config.get("doh_enabled")),
)
network.set_dns_provider(provider, list(servers) or None, use_doh=use_doh)
except (OSError, RuntimeError, ValueError) as exc:
logger.warning("Could not apply parent DNS config (%s): %s", provider, exc)
@@ -1477,7 +1536,7 @@ def _publish_result(result_path: Path, payload: dict[str, Any]) -> None:
The parent decides the request is answered the moment this path exists, so it must
never observe a half-written file. Rename within the same directory is atomic.
"""
tmp_path = result_path.with_name(result_path.name + ".part")
tmp_path = _part_path(result_path)
tmp_path.write_text(json.dumps(payload), encoding="utf-8")
tmp_path.replace(result_path)
@@ -1495,6 +1554,15 @@ def _handle_child_request(request_line: str) -> int:
if isinstance(dns_config, dict):
_apply_parent_dns_config(dns_config)
# The parent owns the cookie store; this process only solves. Starting each request
# from an empty store is what a helper spawned per request gave for free, and losing
# it is what let clearance the parent had deliberately purged for some *other* host
# survive here and get merged back over the parent's copy by the export below - the
# dead-cookie resurrection that http.py's _redirect_loop_handoff purges to avoid.
# Nothing is lost by dropping it: get() below re-checks cached cookies, and the
# parent already ran that same check against a store that is a superset of this one.
clear_cf_cookies()
try:
html = get(url, retry=retry)
cookies, user_agents = export_store()
-1
View File
@@ -1,6 +1,5 @@
import asyncio
import json
import subprocess
import threading
from pathlib import Path
+225
View File
@@ -35,6 +35,7 @@ class _FakeProc:
def __init__(self) -> None:
self.stdin = _FakeStdin()
self.returncode: int | None = None
self.waited = False
# A pid nothing may actually be signalled by: _terminate_helper_session is patched
# out in these tests, and a stray killpg on a live pid would take out the test run.
type(self)._next_pid += 1
@@ -44,6 +45,7 @@ class _FakeProc:
return self.returncode
def wait(self, timeout: float | None = None) -> int:
self.waited = True
if self.returncode is None:
self.returncode = 0
return self.returncode
@@ -403,8 +405,13 @@ def test_child_deadline_leaves_the_parent_room_to_hear_the_answer(monkeypatch):
"""If the parent gave up first it could only kill the helper, losing a warm process."""
import shelfmark.bypass.internal_bypasser as internal_bypasser
# The child's worst case is its deadline plus the grace it is given to close the
# browser after that deadline cancels the bypass, and all of it has to fit inside the
# parent's wait - otherwise the parent gives up first and kills a helper that was
# about to answer.
assert (
internal_bypasser._CHILD_BYPASS_TIMEOUT_SECONDS
+ internal_bypasser._CDP_UNWIND_GRACE_SECONDS
< internal_bypasser._BYPASS_SUBPROCESS_TIMEOUT_SECONDS
)
@@ -431,3 +438,221 @@ def test_in_process_bypass_keeps_the_parents_budget(monkeypatch):
internal_bypasser._run_bypass_in_current_process("https://example.com", 1)
assert timeouts == [internal_bypasser._IN_PROCESS_BYPASS_TIMEOUT_SECONDS]
def test_timed_out_bypass_finishes_unwinding_before_the_call_returns():
"""A helper serving the next request must not race the browser teardown of the last.
The deadline cancels the bypass, but cancelling from the calling thread only schedules
that - it returns while `finally: await _close_cdp_driver(driver)` is still running.
In a helper that now outlives the request, the next bypass would open its Chrome on the
same loop while the abandoned one was still closing its own, sharing the DISPLAY
globals and one process group.
"""
import shelfmark.bypass.internal_bypasser as internal_bypasser
events: list[str] = []
async def _wedged():
try:
await asyncio.sleep(30)
finally:
# Teardown that yields, the way closing websockets and Chrome does.
await asyncio.sleep(0.05)
events.append("browser closed")
with pytest.raises(TimeoutError):
internal_bypasser._CDP_WORKER.run(_wedged(), timeout=0.1)
assert events == ["browser closed"], "run() returned before the bypass had unwound"
def test_unwind_that_wedges_does_not_hold_the_caller_forever(monkeypatch):
"""The grace is a bound, not a promise: cleanup can hang on a dead browser too."""
import shelfmark.bypass.internal_bypasser as internal_bypasser
monkeypatch.setattr(internal_bypasser, "_CDP_UNWIND_GRACE_SECONDS", 0.1)
async def _wedged_on_both_ends():
try:
await asyncio.sleep(30)
finally:
await asyncio.sleep(30)
with pytest.raises(TimeoutError):
internal_bypasser._CDP_WORKER.run(_wedged_on_both_ends(), timeout=0.1)
def test_cancelling_does_not_wait_out_the_shutdown_grace(monkeypatch, tmp_path):
"""The grace only helps a helper that can still read its stdin.
One dropped mid-bypass is blocked inside the solve and will never reach its read loop,
so waiting it out cannot end in anything but the kill - while the user who asked to
cancel, and every bypass queued behind them on LOCKED, waits for it.
"""
import threading
import shelfmark.bypass.internal_bypasser as internal_bypasser
procs: list[_FakeProc] = []
helper = _helper_with_fake_spawn(monkeypatch, procs)
cancel_flag = threading.Event()
cancel_flag.set()
payload = {
"url": "https://example.com",
"retry": 1,
"result_path": str(tmp_path / "never-written.json"),
}
with pytest.raises(internal_bypasser.BypassCancelledError):
helper.run(payload, timeout=5, cancel_flag=cancel_flag)
assert not procs[0].waited, "a helper wedged mid-bypass was given the full exit grace"
assert procs[0].stdin.closed
def test_idle_helper_still_gets_its_grace(monkeypatch, tmp_path):
"""The reaper drops a helper that *is* in its read loop, and that one gets to exit."""
procs: list[_FakeProc] = []
helper = _helper_with_fake_spawn(monkeypatch, procs)
helper.run(_answered_payload(tmp_path), timeout=5, cancel_flag=None)
helper._discard()
assert procs[0].waited, "an idle helper should be asked to exit before being killed"
def test_failed_request_leaves_no_result_files_behind(monkeypatch, tmp_path):
"""Result paths are unique per request, so anything left is left for good."""
procs: list[_FakeProc] = []
helper = _helper_with_fake_spawn(monkeypatch, procs)
result_path = tmp_path / "result.json"
# A helper killed part-way through _publish_result leaves the staging file.
(tmp_path / "result.json.part").write_text('{"ok": tr', encoding="utf-8")
payload = {"url": "https://example.com", "retry": 1, "result_path": str(result_path)}
with pytest.raises(TimeoutError):
helper.run(payload, timeout=0.05, cancel_flag=None)
assert list(tmp_path.iterdir()) == []
def test_child_does_not_export_cookies_left_by_an_earlier_request(monkeypatch, tmp_path):
"""The parent owns the store; a warm helper must not push its own history back over it.
http.py purges a host's clearance the moment that host challenges a request carrying
it. A helper that kept its store across requests would still be holding the purged
cookies, and the next solve - for some entirely different host - would export them and
the parent would merge them straight back in.
"""
import shelfmark.bypass.internal_bypasser as internal_bypasser
def _solve_host(url, retry=None, cancel_flag=None):
# A solve fills the store for the host it solved, which is all it should report.
host = url.rsplit("/", 1)[-1]
internal_bypasser.import_store({host: {"cf_clearance": "fresh"}}, {host: "UA"})
return "<html>ok</html>"
monkeypatch.setattr(internal_bypasser, "get", _solve_host)
internal_bypasser.clear_cf_cookies()
for index, host in enumerate(("first.example", "second.example")):
internal_bypasser._handle_child_request(
json.dumps(
{
"url": f"https://example.com/{host}",
"retry": 1,
"result_path": str(tmp_path / f"{index}.json"),
}
)
)
second = json.loads((tmp_path / "1.json").read_text(encoding="utf-8"))
assert list(second["cookies"]) == ["second.example"], (
"the helper exported clearance won by an earlier request"
)
assert list(second["user_agents"]) == ["second.example"]
internal_bypasser.clear_cf_cookies()
def _record_dns_calls(monkeypatch):
"""Stand in for the network module: report a resolver state, record changes to it.
A helper starts on system DNS, which is what the parent reports as "auto".
"""
import shelfmark.bypass.internal_bypasser as internal_bypasser
calls: list[tuple] = []
state = {"provider": "auto", "servers": [], "doh_enabled": False}
def _set(provider, servers=None, use_doh=None):
calls.append((provider, servers, use_doh))
state.update({"provider": provider, "servers": servers or [], "doh_enabled": bool(use_doh)})
monkeypatch.setattr(internal_bypasser.network, "set_dns_provider", _set)
monkeypatch.setattr(internal_bypasser.network, "get_dns_config", lambda: dict(state))
return internal_bypasser, calls
def test_helper_follows_the_parent_back_to_auto_dns(monkeypatch):
"""A user flipping CUSTOM_DNS back to auto applies live - the helper has to hear it.
The old early-return on "auto" was correct only because a fresh helper had never been
told anything else. One that outlives the request has, and would go on resolving AA
through a resolver the parent has already abandoned.
"""
internal_bypasser, calls = _record_dns_calls(monkeypatch)
internal_bypasser._apply_parent_dns_config(
{"provider": "cloudflare", "servers": [], "doh_enabled": True}
)
internal_bypasser._apply_parent_dns_config(
{"provider": "auto", "servers": [], "doh_enabled": False}
)
assert calls == [("cloudflare", None, True), ("auto", None, False)]
def test_helper_does_not_reinitialize_dns_for_an_unchanged_config(monkeypatch):
"""set_dns_provider() rebuilds resolvers; every request would pay for it otherwise."""
internal_bypasser, calls = _record_dns_calls(monkeypatch)
for _ in range(3):
internal_bypasser._apply_parent_dns_config(
{"provider": "quad9", "servers": [], "doh_enabled": True}
)
assert calls == [("quad9", None, True)]
def test_fresh_helper_leaves_auto_dns_alone(monkeypatch):
"""A helper starts on system DNS, which is what the parent reports as auto."""
internal_bypasser, calls = _record_dns_calls(monkeypatch)
internal_bypasser._apply_parent_dns_config(
{"provider": "auto", "servers": [], "doh_enabled": False}
)
assert calls == []
def test_failed_dns_apply_is_retried_on_the_next_request(monkeypatch):
"""A provider that did not land leaves the resolver where it was, so the next request
sees the same mismatch and tries again."""
internal_bypasser, calls = _record_dns_calls(monkeypatch)
def _explode(provider, servers=None, use_doh=None):
calls.append((provider, servers, use_doh))
msg = "resolver unreachable"
raise OSError(msg)
monkeypatch.setattr(internal_bypasser.network, "set_dns_provider", _explode)
config = {"provider": "google", "servers": [], "doh_enabled": True}
internal_bypasser._apply_parent_dns_config(config)
internal_bypasser._apply_parent_dns_config(config)
assert calls == [("google", None, True), ("google", None, True)]