fix(cron): a scoped worker whose user bus vanished names the cause and re-probes
One TTL for both probe verdicts (the success-TTL constant collapses into `_SYSTEMD_SCOPE_PROBE_TTL_SECONDS`), and the ack-wait loop consults `scoped_spawn_lost_user_bus()` when a scoped dispatch exits before the worker acknowledged: with `/run/user/<uid>/bus` gone the job error names the missing bus and the enable-linger remedy instead of the wrapper's bare `exit 1`, and the cached True flips so the next fire degrades to a direct external subprocess rather than consuming another occurrence on a dead wrapper. Contributor test trimmed to the revalidation invariant.
This commit is contained in:
@@ -3118,6 +3118,7 @@ def _launch_external_cron_worker(job: dict) -> bool:
|
||||
from tools.environments.local import build_subprocess_env, strip_launch_profile_env
|
||||
from tools.process_registry import (
|
||||
restart_safe_gateway_child_argv,
|
||||
scoped_spawn_lost_user_bus,
|
||||
systemd_user_bus_env,
|
||||
)
|
||||
|
||||
@@ -3260,6 +3261,14 @@ def _launch_external_cron_worker(job: dict) -> bool:
|
||||
with _running_lock:
|
||||
_restart_safe_waiter_job_ids.discard(job_id)
|
||||
payload_path.unlink(missing_ok=True)
|
||||
if dispatch.mode == "scoped" and scoped_spawn_lost_user_bus():
|
||||
# systemd-run itself failed (stderr is DEVNULL): name the cause, not the exit code.
|
||||
raise RuntimeError(
|
||||
"restart-safe systemd scope could not be created: the user D-Bus session at "
|
||||
f"/run/user/{os.getuid()}/bus disappeared after the gateway started. On a " # windows-footgun: ok — scoped dispatch exists only on Linux
|
||||
"system-level service install, run `sudo loginctl enable-linger <gateway-user>`; "
|
||||
"the next fire dispatches without scope isolation."
|
||||
)
|
||||
raise RuntimeError(
|
||||
f"cron external worker exited before ownership acknowledgement "
|
||||
f"(exit {returncode})"
|
||||
|
||||
@@ -258,6 +258,42 @@ def _stub_external_worker_launch(scheduler, monkeypatch):
|
||||
return spawned, payloads, handoff, get
|
||||
|
||||
|
||||
def test_scoped_wrapper_exit_without_user_bus_names_the_cause_and_invalidates_probe(
|
||||
tmp_path, monkeypatch
|
||||
):
|
||||
"""#110803: a stale True scope verdict wraps the worker in ``systemd-run --user --scope``
|
||||
after the user bus vanished; the wrapper exits 1 with no child. The job error must name the
|
||||
missing bus (not a bare exit code) and the cached verdict must flip so the next fire re-probes."""
|
||||
import cron.scheduler as scheduler
|
||||
import tools.process_registry as pr
|
||||
from tools.process_registry import GatewayChildDispatch
|
||||
|
||||
job = {"id": "job-bus", "execution_id": "exec-1", "prompt": "work"}
|
||||
monkeypatch.setattr(scheduler, "_get_hermes_home", lambda: tmp_path)
|
||||
monkeypatch.setattr(
|
||||
"tools.process_registry.restart_safe_gateway_child_argv",
|
||||
lambda command, **_: GatewayChildDispatch("scoped", ["systemd-run", "--", *command]),
|
||||
)
|
||||
monkeypatch.setattr(scheduler, "mark_execution_handoff_pending",
|
||||
lambda _eid: {"id": "exec-1", "handoff_pending": 1})
|
||||
|
||||
class DeadWrapper:
|
||||
returncode = 1
|
||||
|
||||
def poll(self):
|
||||
return 1
|
||||
|
||||
monkeypatch.setattr(scheduler.subprocess, "Popen", lambda *a, **k: DeadWrapper())
|
||||
# Bus gone: systemd_user_bus_env derives nothing.
|
||||
monkeypatch.setattr(pr, "systemd_user_bus_env", lambda base_env=None: dict(base_env or {}))
|
||||
monkeypatch.setattr(pr, "_SYSTEMD_SCOPE_AVAILABLE", True)
|
||||
monkeypatch.setattr(pr, "_SYSTEMD_SCOPE_PROBED_AT", pr.time.monotonic())
|
||||
|
||||
with pytest.raises(RuntimeError, match="user D-Bus session .* disappeared"):
|
||||
scheduler._launch_external_cron_worker(job)
|
||||
assert pr._SYSTEMD_SCOPE_AVAILABLE is False
|
||||
|
||||
|
||||
def test_launch_external_worker_uses_restart_safe_scope_and_acknowledges(
|
||||
tmp_path, monkeypatch
|
||||
):
|
||||
|
||||
@@ -2559,29 +2559,6 @@ class TestSystemdCgroupIsolation:
|
||||
value.startswith("OOMPolicy=") for value in probe_argv if isinstance(value, str)
|
||||
), probe_argv
|
||||
|
||||
def test_successful_systemd_probe_uses_cached_verdict_before_ttl(self, monkeypatch):
|
||||
"""A healthy user bus does not cause a probe for every worker spawn."""
|
||||
import tools.process_registry as pr
|
||||
|
||||
monkeypatch.setattr(pr, "_IS_LINUX", True)
|
||||
monkeypatch.setattr(pr, "_SYSTEMD_SCOPE_AVAILABLE", None)
|
||||
monkeypatch.setattr(pr, "_SYSTEMD_SCOPE_PROBED_AT", 0.0)
|
||||
clock = [100.0]
|
||||
probe_calls = []
|
||||
|
||||
def fake_run(*args, **kwargs):
|
||||
probe_calls.append(args)
|
||||
return subprocess.CompletedProcess(args=args[0], returncode=0)
|
||||
|
||||
monkeypatch.setattr("shutil.which", lambda name: "/usr/bin/systemd-run")
|
||||
monkeypatch.setattr("tools.process_registry.time.monotonic", lambda: clock[0])
|
||||
monkeypatch.setattr("subprocess.run", fake_run)
|
||||
|
||||
assert pr._systemd_run_user_scope_available() is True
|
||||
clock[0] += 30
|
||||
assert pr._systemd_run_user_scope_available() is True
|
||||
assert len(probe_calls) == 1
|
||||
|
||||
def test_successful_systemd_probe_revalidates_after_cache_ttl(self, monkeypatch):
|
||||
"""A vanished user bus invalidates a formerly successful scope verdict."""
|
||||
import tools.process_registry as pr
|
||||
|
||||
+22
-15
@@ -81,17 +81,18 @@ WATCH_GLOBAL_COOLDOWN_SECONDS = 30
|
||||
# Under a systemd gateway with MemoryMax, local background commands inherit the gateway's
|
||||
# cgroup, so a memory-heavy executor can get the ENTIRE gateway killed by systemd-oomd;
|
||||
# ``systemd-run --user --scope`` gives the worker its own transient cgroup. Usability is
|
||||
# probed with a bounded cache (binary present but user D-Bus absent in system services/containers).
|
||||
# probed and cached for a bounded TTL (binary present but user D-Bus absent in system services/containers).
|
||||
# A memory-heavy executor (Codex, tests, Node) can push the whole cgroup past MemoryMax and trigger
|
||||
# systemd-oomd to kill the ENTIRE gateway — taking down the messaging control plane and silently losing the
|
||||
# active turn. We probe whether ``systemd-run --user --scope`` is actually usable (the binary can
|
||||
# exist on the PATH while the user D-Bus session is unavailable — common for system services and
|
||||
# containers), and cache the result briefly. See #70716.
|
||||
# containers), and cache the verdict for a bounded TTL. See #70716.
|
||||
_SYSTEMD_SCOPE_AVAILABLE: Optional[bool] = None
|
||||
_SYSTEMD_SCOPE_PROBE_LOCK = threading.Lock()
|
||||
_SYSTEMD_SCOPE_PROBED_AT = 0.0
|
||||
_SYSTEMD_SCOPE_FAILURE_TTL_SECONDS = 60.0
|
||||
_SYSTEMD_SCOPE_SUCCESS_TTL_SECONDS = 60.0
|
||||
# Both verdicts expire: the user bus can vanish after a True (session logout without linger,
|
||||
# #110803) and reappear after a False (linger enabled later, #104893).
|
||||
_SYSTEMD_SCOPE_PROBE_TTL_SECONDS = 60.0
|
||||
_MIN_WORKER_MEMORY_MAX_BYTES = 64 * 1024 * 1024
|
||||
_DEFAULT_WORKER_MEMORY_MAX_BYTES = 1024 * 1024 * 1024
|
||||
_WORKER_MEMORY_MAX_CAP_BYTES = 4 * 1024 * 1024 * 1024
|
||||
@@ -205,19 +206,11 @@ def systemd_user_bus_env(base_env: Optional[Dict[str, str]] = None) -> Dict[str,
|
||||
|
||||
|
||||
def _systemd_scope_cached() -> Optional[bool]:
|
||||
"""Cached probe verdict, or None when a (re)probe is due.
|
||||
|
||||
Both verdicts expire: a user D-Bus can disappear after a successful probe,
|
||||
while a failed probe can recover after linger or a login session starts.
|
||||
"""
|
||||
"""Cached probe verdict, or None when a (re)probe is due."""
|
||||
if _SYSTEMD_SCOPE_AVAILABLE is None:
|
||||
return None
|
||||
ttl = (
|
||||
_SYSTEMD_SCOPE_SUCCESS_TTL_SECONDS
|
||||
if _SYSTEMD_SCOPE_AVAILABLE
|
||||
else _SYSTEMD_SCOPE_FAILURE_TTL_SECONDS
|
||||
)
|
||||
return None if time.monotonic() - _SYSTEMD_SCOPE_PROBED_AT >= ttl else _SYSTEMD_SCOPE_AVAILABLE
|
||||
stale = time.monotonic() - _SYSTEMD_SCOPE_PROBED_AT >= _SYSTEMD_SCOPE_PROBE_TTL_SECONDS
|
||||
return None if stale else _SYSTEMD_SCOPE_AVAILABLE
|
||||
|
||||
|
||||
def _systemd_run_user_scope_available() -> bool:
|
||||
@@ -332,6 +325,20 @@ class GatewayChildDispatch(NamedTuple):
|
||||
argv: List[str]
|
||||
|
||||
|
||||
def scoped_spawn_lost_user_bus() -> bool:
|
||||
"""After a ``systemd-run --user --scope`` wrapper exits before its child could start: True
|
||||
when the user bus is gone (:func:`systemd_user_bus_env` derives nothing), in which case the
|
||||
cached True verdict is replaced so the next dispatch re-probes and degrades instead of
|
||||
consuming another occurrence on the same dead wrapper (#110803)."""
|
||||
global _SYSTEMD_SCOPE_AVAILABLE, _SYSTEMD_SCOPE_PROBED_AT
|
||||
if systemd_user_bus_env({}):
|
||||
return False
|
||||
with _SYSTEMD_SCOPE_PROBE_LOCK:
|
||||
_SYSTEMD_SCOPE_AVAILABLE = False
|
||||
_SYSTEMD_SCOPE_PROBED_AT = time.monotonic()
|
||||
return True
|
||||
|
||||
|
||||
def restart_safe_gateway_child_argv(
|
||||
command: List[str], *, unit_suffix: str, require_restart_safe_scope: bool,
|
||||
) -> GatewayChildDispatch:
|
||||
|
||||
Reference in New Issue
Block a user