refactor(gateway): wait on the target user's bus socket, not the generic control-socket predicate

_user_systemd_socket_ready() accepts systemd/private alone, which is enough for
systemctl --user but not for the systemd-run --user that restart-safe workers
need; systemd_user_bus_env() requires the bus socket. Replace the uid threading
through five helpers with one _wait_for_target_user_bus(uid) that polls
/run/user/<uid>/bus, and move the post-enable wait + restart hint out of
_ensure_linger_enabled into _ensure_system_service_linger so the activity probe
runs only when linger was actually just enabled. Kanban applies the bus env
unconditionally like the cron sibling. Refs #104893.
This commit is contained in:
kshitijk4poor
2026-09-08 21:02:29 +05:30
committed by kshitij
parent 4c4845f0af
commit 6f11a3296e
5 changed files with 110 additions and 97 deletions
+46 -38
View File
@@ -2000,22 +2000,19 @@ class SystemScopeRequiresRootError(RuntimeError):
return self.args[0] if self.args else ""
def _user_runtime_dir(uid: int | None = None) -> Path:
"""``$XDG_RUNTIME_DIR`` or ``/run/user/<uid>`` (regardless of existence). An explicit *uid* — the
``User=`` of a system unit while root installs it — ignores the caller's env."""
if uid is not None:
return Path(f"/run/user/{uid}")
def _user_runtime_dir() -> Path:
"""``$XDG_RUNTIME_DIR`` or ``/run/user/<uid>`` (regardless of existence)."""
return Path(os.environ.get("XDG_RUNTIME_DIR") or f"/run/user/{os.getuid()}") # windows-footgun: ok — POSIX systemd helper, never invoked on Windows
def _user_dbus_socket_path(uid: int | None = None) -> Path:
def _user_dbus_socket_path() -> Path:
"""Return the expected per-user D-Bus socket path (regardless of existence)."""
return _user_runtime_dir(uid) / "bus"
return _user_runtime_dir() / "bus"
def _user_systemd_private_socket_path(uid: int | None = None) -> Path:
def _user_systemd_private_socket_path() -> Path:
"""Return the per-user systemd private socket path (regardless of existence)."""
return _user_runtime_dir(uid) / "systemd" / "private"
return _user_runtime_dir() / "systemd" / "private"
def _path_exists_safe(path: Path) -> bool:
@@ -2042,10 +2039,10 @@ def _runtime_dir_is_ours(runtime_dir: str) -> bool:
return False
def _user_systemd_socket_ready(uid: int | None = None) -> bool:
def _user_systemd_socket_ready() -> bool:
"""True when the user D-Bus socket OR the per-user systemd private socket exists (some distros
expose only the latter and ``systemctl --user`` still works). Inaccessible counts as not-ready."""
return _path_exists_safe(_user_dbus_socket_path(uid)) or _path_exists_safe(_user_systemd_private_socket_path(uid))
return _path_exists_safe(_user_dbus_socket_path()) or _path_exists_safe(_user_systemd_private_socket_path())
def _ensure_user_systemd_env() -> None:
@@ -2067,17 +2064,28 @@ def _ensure_user_systemd_env() -> None:
os.environ["DBUS_SESSION_BUS_ADDRESS"] = f"unix:path={bus_path}"
def _wait_for_user_dbus_socket(timeout: float = 3.0, *, uid: int | None = None) -> bool:
"""Poll up to ``timeout`` s for a user systemd control socket (user@.service takes a moment after
enable-linger). Only our own bus (``uid`` omitted) is adopted into the environment."""
def _wait_for_user_dbus_socket(timeout: float = 3.0) -> bool:
"""Poll up to ``timeout`` s for a user systemd control socket (user@.service takes a moment after enable-linger)."""
deadline = time.monotonic() + timeout
while time.monotonic() < deadline:
if _user_systemd_socket_ready(uid):
if uid is None:
_ensure_user_systemd_env()
if _user_systemd_socket_ready():
_ensure_user_systemd_env()
return True
time.sleep(0.2)
return _user_systemd_socket_ready(uid)
return _user_systemd_socket_ready()
def _wait_for_target_user_bus(uid: int, timeout: float = 5.0) -> bool:
"""Poll for ``/run/user/<uid>/bus`` of ANOTHER account (the system unit's ``User=`` while root installs).
Only the D-Bus socket counts — ``systemd/private`` alone is enough for ``systemctl --user`` but not for
the ``systemd-run --user`` that restart-safe workers need. Never adopts anything into our env."""
bus = Path(f"/run/user/{uid}/bus")
deadline = time.monotonic() + timeout
while time.monotonic() < deadline:
if _path_exists_safe(bus):
return True
time.sleep(0.2)
return _path_exists_safe(bus)
def _loginctl_enable_linger(username: str) -> subprocess.CompletedProcess:
@@ -3013,30 +3021,30 @@ def _ensure_linger_enabled(username: str | None = None, *, system: bool = False)
if result.returncode != 0:
_print_linger_enable_warning(username, _completed_process_detail(result) or linger_detail, system=system)
return False
if not system:
print("✓ Linger enabled — gateway will persist after logout")
return True
# logind starts user@<uid>.service asynchronously; a gateway started right after install must
# find the bus, so wait for the TARGET user's socket (root's own env says nothing about it).
import pwd
uid = pwd.getpwnam(username).pw_uid # windows-footgun: ok — POSIX systemd helper, never invoked on Windows
if _wait_for_user_dbus_socket(timeout=5.0, uid=uid):
print(f"✓ Enabled linger for {username} — user D-Bus now available")
else:
print(f"⚠ Linger enabled for {username}, but /run/user/{uid}/bus did not appear within 5s.")
print(f" Start the user manager: sudo systemctl start user@{uid}.service")
print(f"✓ Enabled linger for {username}" if system else "✓ Linger enabled — gateway will persist after logout")
return True
def _ensure_system_service_linger(username: str, *, running: bool = False) -> None:
def _ensure_system_service_linger(username: str) -> None:
"""Enable linger for the installed unit's ``User=`` (root included: restart-safe workers always cross
``systemd-run --user``, so a root gateway needs ``user@0.service`` just the same). A gateway that
was already ``running`` keeps its bus-less environment until restarted, and ``systemctl start`` on
an active unit is a no-op — say so rather than let the repair silently not take."""
if not _ensure_linger_enabled(username, system=True) or not running:
``systemd-run --user``, so a root gateway needs ``user@0.service`` just the same).
After a fresh enable, wait for the TARGET user's bus: logind starts ``user@<uid>.service``
asynchronously and ``--start-now`` boots the gateway immediately. A gateway that was already running
keeps its bus-less environment and ``systemctl start`` on an active unit is a no-op — say so rather
than let the repair silently not take."""
if not _ensure_linger_enabled(username, system=True):
return
print(" The running gateway was started without a user D-Bus; restart it to pick one up:")
print(f" sudo systemctl restart {get_service_name()}.service")
import pwd
uid = pwd.getpwnam(username).pw_uid # windows-footgun: ok — POSIX systemd helper, never invoked on Windows
if _wait_for_target_user_bus(uid):
print(f"✓ /run/user/{uid}/bus is up — cron and Kanban workers can use systemd-run --user")
else:
print(f"⚠ /run/user/{uid}/bus did not appear within 5s.")
print(f" Start the user manager: sudo systemctl start user@{uid}.service")
if _systemd_unit_is_active(system=True):
print(" The running gateway was started without a user D-Bus; restart it to pick one up:")
print(f" sudo systemctl restart {get_service_name()}.service")
def _select_systemd_scope(system: bool = False) -> bool:
@@ -3146,7 +3154,7 @@ def systemd_install(
print("Use --force to reinstall")
configured_user = _read_systemd_user_from_unit(unit_path) if system else None
if configured_user:
_ensure_system_service_linger(configured_user, running=_systemd_unit_is_active(system=True))
_ensure_system_service_linger(configured_user)
return
unit_path.parent.mkdir(parents=True, exist_ok=True)
+3 -5
View File
@@ -2257,11 +2257,9 @@ def _default_spawn(task: Task, workspace: str, *, board: Optional[str] = None) -
# A worker spawned by a managed systemd gateway must leave the gateway's
# cgroup before startup; otherwise restarting the service kills the worker
# that is performing the handoff.
scoped_cmd = _restart_safe_worker_argv(task, cmd)
if scoped_cmd != cmd:
from tools.process_registry import systemd_user_bus_env
env = systemd_user_bus_env(env)
cmd = scoped_cmd
cmd = _restart_safe_worker_argv(task, cmd)
from tools.process_registry import systemd_user_bus_env
env = systemd_user_bus_env(env)
log_f = _open_worker_log(task, board)
try:
proc = subprocess.Popen( # noqa: S603 -- argv is a fixed list built above
+48 -42
View File
@@ -67,15 +67,12 @@ class TestEnsureLingerEnabled:
assert "sudo loginctl enable-linger testuser" in out
assert "Permission denied" in out
def test_system_target_user_waits_for_that_users_bus(self, monkeypatch, capsys):
"""Root installing a system unit enables linger for User= and waits for THAT uid's bus, not its own."""
def test_system_scope_enables_target_user_without_logout_messaging(self, monkeypatch, capsys):
monkeypatch.setattr(gateway, "is_linux", lambda: True)
monkeypatch.setattr(gateway, "is_termux", lambda: False)
monkeypatch.setattr(gateway, "Path", lambda _path: SimpleNamespace(exists=lambda: False))
monkeypatch.setattr(gateway, "get_systemd_linger_status", lambda username=None: (False, ""))
monkeypatch.setattr("shutil.which", lambda name: "/usr/bin/loginctl")
monkeypatch.setattr("pwd.getpwnam", lambda name: SimpleNamespace(pw_uid=1001))
run_calls = []
def fake_run(cmd, **kwargs):
@@ -83,36 +80,58 @@ class TestEnsureLingerEnabled:
return SimpleNamespace(returncode=0, stdout="", stderr="")
monkeypatch.setattr(gateway.subprocess, "run", fake_run)
waited = []
monkeypatch.setattr(
gateway, "_wait_for_user_dbus_socket", lambda timeout=3.0, uid=None: waited.append(uid) or True
)
assert gateway._ensure_linger_enabled("alice", system=True) is True
assert run_calls == [["loginctl", "enable-linger", "alice"]]
assert waited == [1001]
out = capsys.readouterr().out
assert "alice" in out and "logout" not in out
assert f"sudo systemctl restart {gateway.get_service_name()}.service" not in out
def test_system_scope_warning_uses_system_restart(self, monkeypatch, capsys):
monkeypatch.setattr(gateway, "is_linux", lambda: True)
monkeypatch.setattr(gateway, "is_termux", lambda: False)
monkeypatch.setattr(gateway, "Path", lambda _path: SimpleNamespace(exists=lambda: False))
monkeypatch.setattr(gateway, "get_systemd_linger_status", lambda username=None: (False, ""))
monkeypatch.setattr("shutil.which", lambda name: "/usr/bin/loginctl")
monkeypatch.setattr(
gateway.subprocess, "run",
lambda *a, **k: SimpleNamespace(returncode=1, stdout="", stderr="Permission denied"),
)
assert gateway._ensure_linger_enabled("alice", system=True) is False
out = capsys.readouterr().out
assert f"sudo systemctl restart {gateway.get_service_name()}.service" in out
assert "systemctl --user" not in out
class TestTargetUidSocketPaths:
def test_explicit_uid_ignores_callers_runtime_env(self, monkeypatch):
monkeypatch.setenv("XDG_RUNTIME_DIR", "/run/user/0")
assert gateway._user_dbus_socket_path(1001) == gateway.Path("/run/user/1001/bus")
assert gateway._user_systemd_private_socket_path(1001) == gateway.Path("/run/user/1001/systemd/private")
def test_wait_for_foreign_uid_does_not_adopt_env(self, monkeypatch):
monkeypatch.setattr(gateway, "_user_systemd_socket_ready", lambda uid=None: True)
class TestEnsureSystemServiceLinger:
@pytest.mark.parametrize("running", [True, False])
def test_fresh_enable_waits_on_target_uid_and_hints_restart_only_when_running(
self, monkeypatch, capsys, running
):
"""Root installing a system unit waits for the TARGET user's bus (never adopting its own env) and
tells the operator to restart only when a bus-less gateway is already active."""
monkeypatch.setattr(gateway, "_ensure_linger_enabled", lambda username, system=False: True)
monkeypatch.setattr("pwd.getpwnam", lambda name: SimpleNamespace(pw_uid=1001))
waited = []
monkeypatch.setattr(gateway, "_wait_for_target_user_bus", lambda uid, timeout=5.0: waited.append(uid) or True)
adopted = []
monkeypatch.setattr(gateway, "_ensure_user_systemd_env", lambda: adopted.append(True))
monkeypatch.setattr(gateway, "_systemd_unit_is_active", lambda system=False: running)
assert gateway._wait_for_user_dbus_socket(timeout=0.1, uid=1001) is True
assert adopted == []
assert gateway._wait_for_user_dbus_socket(timeout=0.1) is True
assert adopted == [True]
gateway._ensure_system_service_linger("alice")
assert waited == [1001] and adopted == []
out = capsys.readouterr().out
assert (f"sudo systemctl restart {gateway.get_service_name()}.service" in out) is running
def test_already_enabled_skips_wait_and_activity_probe(self, monkeypatch):
monkeypatch.setattr(gateway, "_ensure_linger_enabled", lambda username, system=False: False)
monkeypatch.setattr(gateway, "_wait_for_target_user_bus", lambda uid, timeout=5.0: pytest.fail("waited"))
monkeypatch.setattr(gateway, "_systemd_unit_is_active", lambda system=False: pytest.fail("probed"))
gateway._ensure_system_service_linger("alice")
def test_systemd_install_calls_linger_helper(monkeypatch, tmp_path, capsys):
@@ -167,25 +186,18 @@ def test_systemd_install_targets_linger_at_system_service_user(monkeypatch, tmp_
lambda system=False, run_as_user=None: f"[Service]\nUser={run_as_user}\n",
)
monkeypatch.setattr(gateway, "_run_systemctl", lambda *args, **kwargs: None)
monkeypatch.setattr(
gateway, "_ensure_linger_enabled",
lambda username=None, system=False: helper_calls.append((username, system)) or False,
)
monkeypatch.setattr(gateway, "_ensure_system_service_linger", lambda username: helper_calls.append(username))
monkeypatch.setattr(gateway, "print_systemd_scope_conflict_warning", lambda: None)
monkeypatch.setattr(gateway, "print_legacy_unit_warning", lambda: None)
gateway.systemd_install(system=True, run_as_user=user)
assert helper_calls == [(user, True)]
assert helper_calls == [user]
@pytest.mark.parametrize("unit_is_current", [True, False])
@pytest.mark.parametrize("running", [True, False])
def test_existing_system_install_repairs_linger_and_says_restart(
monkeypatch, tmp_path, capsys, unit_is_current, running
):
"""Re-running install on an affected system service enables linger for User=; when the gateway is
already running it must be told to restart (systemctl start on an active unit is a no-op)."""
def test_existing_system_install_repairs_linger_for_configured_user(monkeypatch, tmp_path, unit_is_current):
"""Re-running install on an affected system service provisions linger for the unit's User=."""
unit_path = tmp_path / "systemd" / "hermes-gateway.service"
unit_path.parent.mkdir(parents=True)
unit_path.write_text("[Service]\nUser=alice\n", encoding="utf-8")
@@ -198,14 +210,8 @@ def test_existing_system_install_repairs_linger_and_says_restart(
monkeypatch.setattr(gateway, "systemd_unit_is_current", lambda system=False: unit_is_current)
monkeypatch.setattr(gateway, "refresh_systemd_unit_if_needed", lambda system=False: None)
monkeypatch.setattr(gateway, "_run_systemctl", lambda *args, **kwargs: None)
monkeypatch.setattr(gateway, "_systemd_unit_is_active", lambda system=False: running)
monkeypatch.setattr(
gateway, "_ensure_linger_enabled",
lambda username=None, system=False: helper_calls.append((username, system)) or True,
)
monkeypatch.setattr(gateway, "_ensure_system_service_linger", lambda username: helper_calls.append(username))
gateway.systemd_install(system=True, run_as_user="alice")
assert helper_calls == [("alice", True)]
restart_hint = f"sudo systemctl restart {gateway.get_service_name()}.service"
assert (restart_hint in capsys.readouterr().out) is running
assert helper_calls == ["alice"]
+9 -9
View File
@@ -25,8 +25,8 @@ from gateway.restart import (
class TestUserSystemdPrivateSocketPreflight:
def test_preflight_accepts_private_socket_without_dbus_bus(self, monkeypatch):
monkeypatch.setattr(gateway_cli, "_ensure_user_systemd_env", lambda: None)
monkeypatch.setattr(gateway_cli, "_user_dbus_socket_path", lambda uid=None: Path("/tmp/missing-bus"))
monkeypatch.setattr(gateway_cli, "_user_systemd_private_socket_path", lambda uid=None: Path("/tmp/private-socket"))
monkeypatch.setattr(gateway_cli, "_user_dbus_socket_path", lambda: Path("/tmp/missing-bus"))
monkeypatch.setattr(gateway_cli, "_user_systemd_private_socket_path", lambda: Path("/tmp/private-socket"))
monkeypatch.setattr(Path, "exists", lambda self: str(self) == "/tmp/private-socket")
gateway_cli._preflight_user_systemd(auto_enable_linger=False)
@@ -34,8 +34,8 @@ class TestUserSystemdPrivateSocketPreflight:
def test_wait_for_user_dbus_socket_accepts_private_socket(self, monkeypatch):
calls = []
monkeypatch.setattr(gateway_cli, "_ensure_user_systemd_env", lambda: calls.append("env"))
monkeypatch.setattr(gateway_cli, "_user_dbus_socket_path", lambda uid=None: Path("/tmp/missing-bus"))
monkeypatch.setattr(gateway_cli, "_user_systemd_private_socket_path", lambda uid=None: Path("/tmp/private-socket"))
monkeypatch.setattr(gateway_cli, "_user_dbus_socket_path", lambda: Path("/tmp/missing-bus"))
monkeypatch.setattr(gateway_cli, "_user_systemd_private_socket_path", lambda: Path("/tmp/private-socket"))
monkeypatch.setattr(Path, "exists", lambda self: str(self) == "/tmp/private-socket")
assert gateway_cli._wait_for_user_dbus_socket(timeout=0.1) is True
@@ -1595,11 +1595,11 @@ class TestPreflightUserSystemd:
"""Rick's scenario: no D-Bus, no linger, non-root SSH → clear error."""
monkeypatch.setattr(
gateway_cli, "_user_dbus_socket_path",
lambda uid=None: type("P", (), {"exists": lambda self: False})(),
lambda: type("P", (), {"exists": lambda self: False})(),
)
monkeypatch.setattr(
gateway_cli, "_user_systemd_private_socket_path",
lambda uid=None: type("P", (), {"exists": lambda self: False})(),
lambda: type("P", (), {"exists": lambda self: False})(),
)
monkeypatch.setattr(
gateway_cli, "get_systemd_linger_status", lambda username=None: (False, ""),
@@ -1629,11 +1629,11 @@ class TestPreflightUserSystemd:
"""Happy remediation path: polkit allows enable-linger, socket spawns."""
monkeypatch.setattr(
gateway_cli, "_user_dbus_socket_path",
lambda uid=None: type("P", (), {"exists": lambda self: False})(),
lambda: type("P", (), {"exists": lambda self: False})(),
)
monkeypatch.setattr(
gateway_cli, "_user_systemd_private_socket_path",
lambda uid=None: type("P", (), {"exists": lambda self: False})(),
lambda: type("P", (), {"exists": lambda self: False})(),
)
monkeypatch.setattr(
gateway_cli, "get_systemd_linger_status", lambda username=None: (False, ""),
@@ -1650,7 +1650,7 @@ class TestPreflightUserSystemd:
)
monkeypatch.setattr(
gateway_cli, "_wait_for_user_dbus_socket",
lambda timeout=5.0, uid=None: True,
lambda timeout=5.0: True,
)
# Should not raise.
+4 -3
View File
@@ -163,9 +163,10 @@ def systemd_user_bus_env(base_env: Optional[Dict[str, str]] = None) -> Dict[str,
System-level gateway units run as an unprivileged ``User=`` but normally do
not inherit login-session variables. When the conventional runtime
directory is owned by this uid and its bus exists, derive the two standard
variables. Derived fresh on every call rather than adopted once at boot: a
system unit has no ordering against ``user@<uid>.service``, and linger may be
enabled after the gateway started, so the bus can appear later (#104893).
variables. Derived fresh on every call rather than adopted once at boot:
linger may be enabled after the gateway started (existing installs), so
the bus can appear later and the probe's failure TTL must be able to
recover (#104893).
The returned copy is passed explicitly to the probe and every scoped spawn;
``os.environ`` is left unchanged.
"""