diff --git a/cron/delivery_queue.py b/cron/delivery_queue.py new file mode 100644 index 0000000000..a7dbc8eea2 --- /dev/null +++ b/cron/delivery_queue.py @@ -0,0 +1,244 @@ +"""Profile-local durable handoff for cron delivery through live gateway adapters. + +A restart-safe cron worker executes outside the gateway cgroup. It cannot own +relay/E2EE adapter objects, so it queues the final send here. A gateway claims +each row at most once. If that gateway dies after claiming, the outcome is +marked unknown and never retried: losing a delivery is safer than duplicating a +possibly-completed send. +""" + +from __future__ import annotations + +import json +import os +import sqlite3 +import threading +import time +import uuid +from contextlib import contextmanager +from pathlib import Path +from typing import Any, Callable, Iterator, Optional + +from hermes_constants import get_hermes_home +from hermes_time import now as _hermes_now + +DELIVERY_DB: Optional[Path] = None +_PROCESS_ID = uuid.uuid4().hex +_lock = threading.RLock() +_ACTIVE_DELIVERIES: set[str] = set() +_TERMINAL = ("delivered", "failed", "unknown") + + +def _path() -> Path: + return DELIVERY_DB or (get_hermes_home().resolve() / "cron" / "deliveries.db") + + +@contextmanager +def _transaction() -> Iterator[sqlite3.Connection]: + with _lock: + path = _path() + path.parent.mkdir(parents=True, exist_ok=True) + conn = sqlite3.connect(path, timeout=5) + try: + path.chmod(0o600) + except OSError: + pass + conn.row_factory = sqlite3.Row + try: + conn.execute("PRAGMA busy_timeout=5000") + conn.execute("PRAGMA journal_mode=WAL") + conn.execute("PRAGMA synchronous=FULL") + conn.execute( + """CREATE TABLE IF NOT EXISTS deliveries ( + execution_id TEXT PRIMARY KEY, + job_json TEXT NOT NULL, + content TEXT NOT NULL, + status TEXT NOT NULL CHECK(status IN + ('pending','delivering','delivered','failed','unknown')), + owner_process_id TEXT, + owner_pid INTEGER, + owner_started_at INTEGER, + created_at TEXT NOT NULL, + finished_at TEXT, + error TEXT + )""" + ) + with conn: + yield conn + finally: + conn.close() + + +def _process_start_time(pid: int) -> Optional[int]: + try: + from gateway.status import get_process_start_time + + return get_process_start_time(pid) + except Exception: + return None + + +def _owner_is_live(pid: int, started_at: Optional[int]) -> bool: + try: + from gateway.status import _pid_exists + + if not _pid_exists(pid): + return False + except Exception: + return True + if started_at is None: + return pid == os.getpid() + return _process_start_time(pid) == started_at + + +def enqueue(execution_id: str, job: dict, content: str) -> dict: + """Persist one idempotent delivery request before the worker waits.""" + with _transaction() as conn: + conn.execute( + """INSERT OR IGNORE INTO deliveries + (execution_id, job_json, content, status, created_at) + VALUES (?, ?, ?, 'pending', ?)""", + ( + str(execution_id), + json.dumps(job, ensure_ascii=False, sort_keys=True), + str(content), + _hermes_now().isoformat(), + ), + ) + row = conn.execute( + "SELECT * FROM deliveries WHERE execution_id=?", (str(execution_id),) + ).fetchone() + return dict(row) + + +def get_status(execution_id: str) -> Optional[dict]: + with _transaction() as conn: + row = conn.execute( + "SELECT * FROM deliveries WHERE execution_id=?", (str(execution_id),) + ).fetchone() + return dict(row) if row is not None else None + + +def claim_next() -> Optional[dict]: + """Atomically claim one pending send before touching the transport.""" + pid = os.getpid() + started = _process_start_time(pid) + with _transaction() as conn: + row = conn.execute( + "SELECT execution_id FROM deliveries WHERE status='pending' " + "ORDER BY created_at, execution_id LIMIT 1" + ).fetchone() + if row is None: + return None + cur = conn.execute( + """UPDATE deliveries SET status='delivering', owner_process_id=?, + owner_pid=?, owner_started_at=? + WHERE execution_id=? AND status='pending'""", + (_PROCESS_ID, pid, started, row["execution_id"]), + ) + if cur.rowcount != 1: + return None + claimed = conn.execute( + "SELECT * FROM deliveries WHERE execution_id=?", (row["execution_id"],) + ).fetchone() + _ACTIVE_DELIVERIES.add(row["execution_id"]) + result = dict(claimed) + result["job"] = json.loads(result.pop("job_json")) + return result + + +def _finish(execution_id: str, *, error: Optional[str]) -> bool: + status = "failed" if error else "delivered" + with _transaction() as conn: + cur = conn.execute( + """UPDATE deliveries SET status=?, finished_at=?, error=? + WHERE execution_id=? AND status='delivering' + AND owner_process_id=? AND owner_pid=?""", + ( + status, + _hermes_now().isoformat(), + error, + execution_id, + _PROCESS_ID, + os.getpid(), + ), + ) + return cur.rowcount == 1 + + +def recover_abandoned() -> int: + """Fence dead delivery owners as unknown; never replay uncertain sends.""" + changed = 0 + with _transaction() as conn: + rows = conn.execute( + "SELECT execution_id, owner_process_id, owner_pid, owner_started_at " + "FROM deliveries WHERE status='delivering'" + ).fetchall() + for row in rows: + same_process = row["owner_process_id"] == _PROCESS_ID + if same_process: + with _lock: + if row["execution_id"] in _ACTIVE_DELIVERIES: + continue + elif _owner_is_live(int(row["owner_pid"]), row["owner_started_at"]): + continue + error = ( + "Gateway finished delivery but could not persist its outcome; " + "send was not retried." + if same_process + else "Gateway exited during delivery; send outcome is unknown and was not retried." + ) + cur = conn.execute( + """UPDATE deliveries SET status='unknown', finished_at=?, error=? + WHERE execution_id=? AND status='delivering'""", + ( + _hermes_now().isoformat(), + error, + row["execution_id"], + ), + ) + changed += cur.rowcount + return changed + + +def drain(send: Callable[[dict, str], Optional[str]], *, limit: int = 20) -> int: + """Deliver pending rows through *send*, terminalizing every claimed row.""" + recover_abandoned() + processed = 0 + for _ in range(max(0, limit)): + row = claim_next() + if row is None: + break + with _lock: + _ACTIVE_DELIVERIES.add(row["execution_id"]) + try: + try: + error = send(row["job"], row["content"]) + except BaseException as exc: + error = f"{type(exc).__name__}: {exc}" + _finish(row["execution_id"], error=error) + finally: + with _lock: + _ACTIVE_DELIVERIES.discard(row["execution_id"]) + processed += 1 + return processed + + +def enqueue_and_wait( + execution_id: str, + job: dict, + content: str, + *, + timeout: Optional[float] = None, +) -> Optional[str]: + """Queue delivery and wait for a gateway's terminal at-most-once outcome.""" + enqueue(execution_id, job, content) + deadline = None if timeout is None else time.monotonic() + timeout + while deadline is None or time.monotonic() < deadline: + row = get_status(execution_id) + if row and row["status"] in _TERMINAL: + return None if row["status"] == "delivered" else str( + row.get("error") or f"delivery {row['status']}" + ) + time.sleep(0.25) + return "timed out waiting for live gateway delivery; request remains pending" diff --git a/cron/executions.py b/cron/executions.py index efadab464a..d05afc2dec 100644 --- a/cron/executions.py +++ b/cron/executions.py @@ -160,6 +160,33 @@ def create_execution(job_id: str, *, source: str) -> Dict[str, Any]: return record # type: ignore[return-value] +def adopt_claimed_execution(execution_id: str) -> Optional[Dict[str, Any]]: + """Atomically transfer and start an attempt in its worker process. + + The dispatching gateway creates the row before spawning a restart-safe + worker. Adoption is the single ``claimed`` → ``running`` gate: only the + winner may acknowledge ownership or run side effects. + """ + pid = os.getpid() + process_started_at = _process_start_time(pid) + now = _hermes_now().isoformat() + with _transaction() as conn: + cur = conn.execute( + """UPDATE executions + SET process_id=?, pid=?, process_started_at=?, + status='running', started_at=? + WHERE id=? AND status='claimed'""", + (_PROCESS_ID, pid, process_started_at, now, execution_id), + ) + if cur.rowcount != 1: + return None + record = _record(conn.execute( + "SELECT * FROM executions WHERE id=?", (execution_id,) + ).fetchone()) + _emit_execution_state(record) + return record + + def mark_execution_running(execution_id: str) -> Optional[Dict[str, Any]]: """Transition one claimed attempt to running exactly once.""" now = _hermes_now().isoformat() diff --git a/cron/scheduler.py b/cron/scheduler.py index ba2d0bde4b..aeda1b4196 100644 --- a/cron/scheduler.py +++ b/cron/scheduler.py @@ -3203,6 +3203,16 @@ def _deliver_result( logger.warning("Job '%s': %s", job["id"], msg) return msg + # Restart-safe workers intentionally have no live gateway adapter objects. + # Hand the send back through a durable queue so the current or replacement + # gateway performs it with relay/E2EE parity. The execution id is the + # idempotency key; the queue never retries an uncertain claimed send. + external_execution = os.environ.get("_HERMES_CRON_EXTERNAL_WORKER", "") + if external_execution and adapters is None: + from cron.delivery_queue import enqueue_and_wait + + return enqueue_and_wait(external_execution, job, content) + from tools.send_message_tool import _send_to_platform from gateway.config import load_gateway_config, Platform @@ -4118,6 +4128,20 @@ def _deliver_result( return None +def drain_delivery_queue(adapters, loop) -> int: + """Send queued worker results through this gateway's live adapters.""" + from cron.delivery_queue import drain + + return drain( + lambda queued_job, queued_content: _deliver_result( + queued_job, + queued_content, + adapters=adapters, + loop=loop, + ) + ) + + _DEFAULT_SCRIPT_TIMEOUT = 3600 # seconds (1 hour) # Backward-compatible module override used by tests and emergency monkeypatches. _SCRIPT_TIMEOUT = _DEFAULT_SCRIPT_TIMEOUT @@ -7368,6 +7392,34 @@ def run_one_job( run cooperatively — agent interruption AND script process-tree kill — through the single fenced completion path. """ + # Every gateway path (built-in scheduler, external providers, and direct + # API fires) crosses this seam. Ensure the detached worker has a durable + # attempt to adopt before any launch can occur. + if adapters is not None and not job.get("execution_id"): + execution = create_execution(job["id"], source="direct") + job["execution_id"] = execution["id"] + + if adapters is not None: + try: + if _launch_external_cron_worker(job): + return True + except Exception as handoff_error: + error = f"Restart-safe cron worker dispatch failed: {handoff_error}" + logger.error("Job '%s': %s", job["id"], error) + claim = job.get("fire_claim") + owner = str(claim.get("by") or "") if isinstance(claim, dict) else "" + try: + mark_job_run( + job["id"], + False, + error, + **({"expected_fire_owner": owner} if owner else {}), + ) + finally: + execution_id = job.get("execution_id") + if execution_id: + finish_execution(execution_id, success=False, error=error) + return True if extra_prompt is None: # A gateway-forwarded manual run (`hermes cron run --prompt` / # cronjob(action='run', prompt=...) on a relay-fronted target) stamps @@ -7491,7 +7543,16 @@ def _run_one_job_body( # The attempt is claimed durably before executor/provider dispatch and # becomes running only immediately before the actual run. - mark_execution_running(execution_id) + # Detached workers atomically transition the attempt to running while + # adopting it. In-process paths must win the claimed->running CAS + # here before any user script or agent side effect may begin. + external_owner = os.environ.get("_HERMES_CRON_EXTERNAL_WORKER") == execution_id + if not external_owner and mark_execution_running(execution_id) is None: + logger.warning( + "Cron job %s lost execution ownership before start; skipping", + job["id"], + ) + return True # Run and deliver under the profile's secret scope. get_secret() fails # closed outside a scope once profile isolation is active, and cron @@ -7967,6 +8028,210 @@ def _run_one_job_body( reset_terminal_scope(_terminal_scope_token) +def _launch_external_cron_worker(job: dict) -> bool: + """Launch *job* outside a managed gateway cgroup when required. + + Returns ``False`` when the caller is not a managed systemd gateway and the + existing in-process path should be used. In managed topology, failure to + establish the transient scope raises: falling back would recreate the + restart interruption this handoff exists to prevent. + """ + execution_id = str(job["execution_id"]) + job_id = str(job["id"]) + handoff_dir = _get_hermes_home() / "cron" / "external-workers" + payload_path = handoff_dir / f"{execution_id}.json" + ack_path = handoff_dir / f"{execution_id}.ready" + command = [ + sys.executable, + "-m", + "cron.scheduler", + "--external-worker-file", + str(payload_path), + "--ack-file", + str(ack_path), + ] + + from agent.secret_scope import is_multiplex_active + from tools.environments.local import build_subprocess_env + from tools.process_registry import restart_safe_gateway_child_argv + + multiplex_active = is_multiplex_active() + scoped_command = restart_safe_gateway_child_argv( + command, + unit_suffix=f"cron-{job_id}-exec-{execution_id}", + ) + if scoped_command == command: + return False + + _ensure_cron_dir(handoff_dir) + try: + handoff_dir.chmod(0o700) + except OSError: + pass + fd = os.open(payload_path, os.O_WRONLY | os.O_CREAT | os.O_EXCL, 0o600) + try: + with os.fdopen(fd, "w", encoding="utf-8") as payload_file: + json.dump( + { + "job": job, + "profile_home": str(_get_hermes_home().resolve()), + "multiplex_active": multiplex_active, + }, + payload_file, + ) + payload_file.flush() + os.fsync(payload_file.fileno()) + except BaseException: + payload_path.unlink(missing_ok=True) + raise + + worker_env = build_subprocess_env( + scrub_secrets=multiplex_active, + inherit_profile_home=True, + extra={"HERMES_HOME": str(_get_hermes_home().resolve())}, + ) + try: + process = subprocess.Popen( + scoped_command, + cwd=str(Path(__file__).resolve().parent.parent), + env=worker_env, + stdin=subprocess.DEVNULL, + stdout=subprocess.DEVNULL, + stderr=subprocess.DEVNULL, + start_new_session=True, + creationflags=windows_hide_flags(), + ) + except BaseException: + payload_path.unlink(missing_ok=True) + raise + + deadline = time.monotonic() + 5.0 + while time.monotonic() < deadline: + if ack_path.exists(): + try: + acknowledgement = json.loads(ack_path.read_text(encoding="utf-8")) + except Exception: + logger.exception( + "Cron external worker %s published an unreadable acknowledgement; " + "treating handoff as ownership-uncertain", + execution_id, + ) + return True + finally: + ack_path.unlink(missing_ok=True) + if acknowledgement.get("execution_id") != execution_id: + logger.error( + "Cron external worker acknowledgement mismatch for %s; " + "treating handoff as ownership-uncertain", + execution_id, + ) + return True + logger.info( + "Cron job '%s' handed to restart-safe worker pid=%s execution=%s", + job_id, + acknowledgement.get("pid"), + execution_id, + ) + return True + returncode = process.poll() + if returncode is not None: + payload_path.unlink(missing_ok=True) + raise RuntimeError( + f"cron external worker exited before ownership acknowledgement " + f"(exit {returncode})" + ) + time.sleep(0.05) + + # The child may have adopted the durable row just before publishing its + # acknowledgement. Never fall back to in-process execution on an uncertain + # handoff: that could duplicate side effects. The execution owner/dead-owner + # recovery ledger remains the authority. + logger.warning( + "Cron external worker for job '%s' did not acknowledge within 5s; " + "leaving the durable execution claim untouched", + job_id, + ) + return True + + +def _run_external_worker_payload(payload_path: Path, ack_path: Path) -> bool: + """Adopt and execute one gateway-dispatched cron payload. + + The execution row is created by the gateway before spawn, then transferred + here before the ready acknowledgement is published. No side effect runs + unless that durable ownership transfer succeeds. + """ + try: + payload = json.loads(payload_path.read_text(encoding="utf-8")) + job = payload["job"] + profile_home = Path(payload["profile_home"]).resolve() + execution_id = str(job["execution_id"]) + except Exception: + logger.exception("Cron external worker could not load payload %s", payload_path) + return False + finally: + try: + payload_path.unlink(missing_ok=True) + except OSError: + pass + + from agent.secret_scope import ( + build_profile_secret_scope, + is_multiplex_active, + reset_secret_scope, + set_multiplex_active, + set_secret_scope, + ) + from cron.executions import adopt_claimed_execution + from hermes_cli.env_loader import hydrate_profile_secret_sources + from hermes_constants import ( + reset_hermes_home_override, + set_hermes_home_override, + ) + + home_token = set_hermes_home_override(profile_home) + previous_multiplex = is_multiplex_active() + multiplex_active = bool(payload.get("multiplex_active", False)) + set_multiplex_active(multiplex_active) + hydrate_profile_secret_sources(profile_home) + secret_token = set_secret_scope(build_profile_secret_scope(profile_home)) + try: + with use_cron_store(profile_home): + if adopt_claimed_execution(execution_id) is None: + logger.error( + "Cron external worker refused execution %s: durable ownership " + "could not be established", + execution_id, + ) + return False + try: + ack_path.parent.mkdir(parents=True, exist_ok=True) + fd = os.open(ack_path, os.O_WRONLY | os.O_CREAT | os.O_EXCL, 0o600) + with os.fdopen(fd, "w", encoding="utf-8") as ack_file: + json.dump({"pid": os.getpid(), "execution_id": execution_id}, ack_file) + ack_file.flush() + os.fsync(ack_file.fileno()) + except Exception: + logger.exception( + "Cron external worker could not publish ready acknowledgement for %s", + execution_id, + ) + return False + old_external_execution = os.environ.get("_HERMES_CRON_EXTERNAL_WORKER") + os.environ["_HERMES_CRON_EXTERNAL_WORKER"] = execution_id + try: + return run_one_job(job, adapters=None, loop=None, verbose=False) + finally: + if old_external_execution is None: + os.environ.pop("_HERMES_CRON_EXTERNAL_WORKER", None) + else: + os.environ["_HERMES_CRON_EXTERNAL_WORKER"] = old_external_execution + finally: + reset_secret_scope(secret_token) + set_multiplex_active(previous_multiplex) + reset_hermes_home_override(home_token) + + def _notify_provider_jobs_changed() -> None: """Best-effort: tell the active scheduler provider the job set changed. @@ -8582,4 +8847,14 @@ def tick( if __name__ == "__main__": + if "--external-worker-file" in sys.argv: + import argparse + + parser = argparse.ArgumentParser(add_help=False) + parser.add_argument("--external-worker-file", type=Path, required=True) + parser.add_argument("--ack-file", type=Path, required=True) + args = parser.parse_args() + raise SystemExit( + 0 if _run_external_worker_payload(args.external_worker_file, args.ack_file) else 1 + ) tick(verbose=True) diff --git a/gateway/run.py b/gateway/run.py index 2602479525..27f88ea75f 100644 --- a/gateway/run.py +++ b/gateway/run.py @@ -33276,7 +33276,30 @@ def _run_planned_stop_watcher( stop_event.wait(poll_interval) -def _start_gateway_housekeeping(stop_event: threading.Event, adapters=None, loop=None, interval: int = 60, cron_provider=None): +def _drain_restart_safe_cron_deliveries(adapters, loop, runner=None) -> None: + """Drain each profile's worker queue through its matching live adapters.""" + from cron import scheduler as cron_scheduler + + if adapters: + cron_scheduler.drain_delivery_queue(adapters, loop) + if runner is None: + return + for profile_name, profile_home in _handoff_watch_scopes(runner)[1:]: + profile_adapters = getattr(runner, "_profile_adapters", {}).get(profile_name) + if not profile_adapters: + continue + with _profile_runtime_scope(profile_home): + cron_scheduler.drain_delivery_queue(profile_adapters, loop) + + +def _start_gateway_housekeeping( + stop_event: threading.Event, + adapters=None, + loop=None, + interval: int = 60, + cron_provider=None, + runner=None, +): """Background thread for gateway-only periodic chores (NOT cron). Split out of the historical ``_start_cron_ticker`` so the cron *trigger* @@ -33331,6 +33354,19 @@ def _start_gateway_housekeeping(stop_event: threading.Event, adapters=None, loop while not stop_event.is_set(): tick_count += 1 + # Restart-safe cron workers run outside the gateway cgroup and queue + # their final send for whichever gateway instance is live. Drain on + # the gateway-wide housekeeper rather than the built-in scheduler tick: + # external providers do not run that ticker. + profile_adapters = ( + getattr(runner, "_profile_adapters", {}) if runner is not None else {} + ) + if adapters or any(profile_adapters.values()): + try: + _drain_restart_safe_cron_deliveries(adapters, loop, runner) + except Exception as exc: + logger.debug("Cron durable delivery queue drain error: %s", exc) + if tick_count % CHANNEL_DIR_EVERY == 0 and adapters: try: from gateway.channel_directory import build_channel_directory @@ -34512,6 +34548,7 @@ async def start_gateway(config: Optional[GatewayConfig] = None, replace: bool = "adapters": runner.adapters, "loop": asyncio.get_running_loop(), "cron_provider": cron_provider, + "runner": runner, }, daemon=True, name="gateway-housekeeping", diff --git a/hermes_cli/kanban_db.py b/hermes_cli/kanban_db.py index 198669792e..d685853066 100644 --- a/hermes_cli/kanban_db.py +++ b/hermes_cli/kanban_db.py @@ -10717,6 +10717,32 @@ def _retag_legacy_worker_sessions(workspaces_root_path: str) -> None: _log.debug("kanban worker: legacy session retag skipped (%s)", exc) +def _restart_safe_worker_argv(task: Task, command: list[str]) -> list[str]: + """Wrap a managed-gateway worker in the shared restart-safe scope.""" + if task.current_run_id is None: + # Outside managed systemd this is harmless, but a managed dispatch must + # never mint an untraceable scope. Check topology through the shared + # helper first, using a placeholder suffix that cannot be launched. + from tools.process_registry import restart_safe_gateway_child_argv + + scoped = restart_safe_gateway_child_argv( + command, unit_suffix=f"kanban-{task.id}-run-missing" + ) + if scoped is not command: + raise RuntimeError( + "cannot create restart-safe systemd scope for Kanban worker: " + "the claimed task has no current run id" + ) + return command + + from tools.process_registry import restart_safe_gateway_child_argv + + return restart_safe_gateway_child_argv( + command, + unit_suffix=f"kanban-{task.id}-run-{task.current_run_id}", + ) + + def _default_spawn( task: Task, workspace: str, @@ -10744,7 +10770,13 @@ def _default_spawn( profile_arg = normalize_profile_name(task.assignee) prompt = f"work kanban task {task.id}" - env = dict(os.environ) + from agent.secret_scope import is_multiplex_active + from tools.environments.local import build_subprocess_env + + env = build_subprocess_env( + scrub_secrets=is_multiplex_active(), + inherit_profile_home=True, + ) # The dispatcher is detached from every conversation. Its worker must never # inherit routing mirrored by a previous gateway turn, even before the first # session binds ContextVars in this process. @@ -10895,6 +10927,12 @@ def _default_spawn( # turn, prints text, exits rc=0, and the dispatcher records a # protocol violation (incident 2026-06-09 t_d9cbe312). cmd.append("-Q") + + # 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. + cmd = _restart_safe_worker_argv(task, cmd) + # Redirect output to a per-task log under /logs/. # Anchored at the board root (not the shared kanban root), so # `hermes kanban log` on a specific board reads its own file and diff --git a/tests/cron/test_delivery_queue.py b/tests/cron/test_delivery_queue.py new file mode 100644 index 0000000000..7fc6517158 --- /dev/null +++ b/tests/cron/test_delivery_queue.py @@ -0,0 +1,94 @@ +"""Durable at-most-once delivery handoff for restart-safe cron workers.""" + +from __future__ import annotations + +from unittest.mock import Mock + +import pytest + + +def test_pending_delivery_is_claimed_and_sent_once(tmp_path, monkeypatch): + import cron.delivery_queue as queue + + monkeypatch.setattr(queue, "DELIVERY_DB", tmp_path / "deliveries.db") + queue.enqueue("exec-1", {"id": "job-1"}, "brief") + send = Mock(return_value=None) + + assert queue.drain(send) == 1 + assert queue.drain(send) == 0 + send.assert_called_once_with({"id": "job-1"}, "brief") + assert queue.get_status("exec-1")["status"] == "delivered" + + +def test_dead_delivery_owner_becomes_unknown_and_is_not_retried( + tmp_path, monkeypatch +): + import cron.delivery_queue as queue + + monkeypatch.setattr(queue, "DELIVERY_DB", tmp_path / "deliveries.db") + queue.enqueue("exec-1", {"id": "job-1"}, "brief") + assert queue.claim_next() is not None + monkeypatch.setattr(queue, "_PROCESS_ID", "replacement-gateway") + monkeypatch.setattr(queue, "_owner_is_live", lambda _pid, _started: False) + + assert queue.recover_abandoned() == 1 + send = Mock() + assert queue.drain(send) == 0 + send.assert_not_called() + assert queue.get_status("exec-1")["status"] == "unknown" + + +def test_delivery_failure_is_terminal_and_not_retried(tmp_path, monkeypatch): + import cron.delivery_queue as queue + + monkeypatch.setattr(queue, "DELIVERY_DB", tmp_path / "deliveries.db") + queue.enqueue("exec-1", {"id": "job-1"}, "brief") + send = Mock(return_value="transport failed") + + assert queue.drain(send) == 1 + assert queue.drain(send) == 0 + assert send.call_count == 1 + status = queue.get_status("exec-1") + assert status["status"] == "failed" + assert status["error"] == "transport failed" + + +def test_wait_timeout_cancels_unclaimed_delivery(tmp_path, monkeypatch): + import cron.delivery_queue as queue + + monkeypatch.setattr(queue, "DELIVERY_DB", tmp_path / "deliveries.db") + job = {"id": "job-3", "deliver": "origin"} + + error = queue.enqueue_and_wait("exec-3", job, "result", timeout=0) + + assert "timed out" in error + assert queue.get_status("exec-3")["status"] == "pending" + send = Mock(return_value=None) + assert queue.drain(send) == 1 + send.assert_called_once() + + +def test_same_gateway_recovers_terminalization_failure_without_resending( + tmp_path, monkeypatch +): + import cron.delivery_queue as queue + + monkeypatch.setattr(queue, "DELIVERY_DB", tmp_path / "deliveries.db") + queue.enqueue("exec-4", {"id": "job-4"}, "result") + send = Mock(return_value=None) + original_finish = queue._finish + monkeypatch.setattr( + queue, + "_finish", + Mock(side_effect=OSError("database temporarily unavailable")), + ) + + with pytest.raises(OSError, match="temporarily unavailable"): + queue.drain(send) + + monkeypatch.setattr(queue, "_finish", original_finish) + assert queue.drain(send) == 0 + send.assert_called_once() + status = queue.get_status("exec-4") + assert status["status"] == "unknown" + assert "not retried" in status["error"] diff --git a/tests/cron/test_execution_ledger.py b/tests/cron/test_execution_ledger.py index ffe39a5164..f2b8a9cacf 100644 --- a/tests/cron/test_execution_ledger.py +++ b/tests/cron/test_execution_ledger.py @@ -221,7 +221,7 @@ def test_run_one_job_records_running_then_terminal(monkeypatch): monkeypatch.setattr( scheduler, "mark_execution_running", - lambda execution_id: events.append(("running", execution_id)), + lambda execution_id: events.append(("running", execution_id)) or {}, raising=False, ) monkeypatch.setattr( diff --git a/tests/cron/test_parallel_pool.py b/tests/cron/test_parallel_pool.py index b159f265d7..9853dbc229 100644 --- a/tests/cron/test_parallel_pool.py +++ b/tests/cron/test_parallel_pool.py @@ -187,7 +187,7 @@ class TestRunningJobGuard: if job_id == "healthy-job" else None, ) - monkeypatch.setattr(sched, "mark_execution_running", lambda *_a, **_kw: None) + monkeypatch.setattr(sched, "mark_execution_running", lambda *_a, **_kw: {}) monkeypatch.setattr(sched, "heartbeat_fire_claim", lambda *_a, **_kw: True) n = sched.tick(verbose=False) diff --git a/tests/cron/test_restart_safe_worker.py b/tests/cron/test_restart_safe_worker.py new file mode 100644 index 0000000000..e0c0207f41 --- /dev/null +++ b/tests/cron/test_restart_safe_worker.py @@ -0,0 +1,245 @@ +"""Restart-safe cron worker handoff and ownership contracts.""" + +from __future__ import annotations + +import json +from pathlib import Path +from unittest.mock import Mock + +import pytest + + +@pytest.fixture +def execution_ledger(tmp_path, monkeypatch): + import cron.executions as executions + + monkeypatch.setattr(executions, "EXECUTIONS_FILE", tmp_path / "executions.db") + return executions + + +def test_execution_owner_moves_to_external_worker_before_running( + execution_ledger, monkeypatch +): + record = execution_ledger.create_execution("job-1", source="builtin") + monkeypatch.setattr(execution_ledger.os, "getpid", lambda: 4242) + monkeypatch.setattr(execution_ledger, "_process_start_time", lambda pid: 9876) + + adopted = execution_ledger.adopt_claimed_execution(record["id"]) + + assert adopted is not None + assert adopted["pid"] == 4242 + assert adopted["process_started_at"] == 9876 + assert adopted["status"] == "running" + assert execution_ledger.adopt_claimed_execution(record["id"]) is None + assert execution_ledger.mark_execution_running(record["id"]) is None + + +def test_restart_safe_gateway_child_fails_closed_without_scope(monkeypatch): + import tools.process_registry as process_registry + + monkeypatch.setattr(process_registry, "_IS_WINDOWS", False) + monkeypatch.setattr(process_registry, "_is_supervised_gateway_process", lambda: True) + monkeypatch.setenv("INVOCATION_ID", "managed-service") + monkeypatch.setattr(process_registry, "_systemd_run_user_scope_available", lambda: False) + + with pytest.raises(RuntimeError, match="systemd-run --user --scope is unavailable"): + process_registry.restart_safe_gateway_child_argv( + ["python", "worker.py"], unit_suffix="cron-job-1" + ) + + +def test_restart_safe_gateway_child_is_unchanged_outside_managed_gateway(monkeypatch): + import tools.process_registry as process_registry + + command = ["python", "worker.py"] + monkeypatch.setattr(process_registry, "_is_supervised_gateway_process", lambda: False) + + assert process_registry.restart_safe_gateway_child_argv( + command, unit_suffix="cron-job-1" + ) is command + + +def test_external_worker_adopts_execution_and_runs_payload_once( + tmp_path, monkeypatch +): + import cron.scheduler as scheduler + + payload = tmp_path / "payload.json" + ack = tmp_path / "ready.json" + payload.write_text( + json.dumps({ + "job": {"id": "job-1", "execution_id": "exec-1"}, + "profile_home": str(tmp_path / "profile"), + }), + encoding="utf-8", + ) + from hermes_constants import get_hermes_home + + observed_homes = [] + adopted = Mock( + side_effect=lambda execution_id: ( + observed_homes.append(get_hermes_home().resolve()) + or {"id": execution_id, "status": "running"} + ) + ) + run = Mock( + side_effect=lambda *_args, **_kwargs: ( + observed_homes.append(get_hermes_home().resolve()) or True + ) + ) + monkeypatch.setattr("cron.executions.adopt_claimed_execution", adopted) + monkeypatch.setattr(scheduler, "run_one_job", run) + + assert scheduler._run_external_worker_payload(payload, ack) is True + + adopted.assert_called_once_with("exec-1") + run.assert_called_once() + assert run.call_args.args[0]["id"] == "job-1" + expected_home = (tmp_path / "profile").resolve() + assert observed_homes == [expected_home, expected_home] + assert ack.exists() + assert not payload.exists() + + +def test_external_worker_refuses_to_run_without_durable_ownership( + tmp_path, monkeypatch +): + import cron.scheduler as scheduler + + payload = tmp_path / "payload.json" + ack = tmp_path / "ready.json" + payload.write_text( + json.dumps({ + "job": {"id": "job-1", "execution_id": "exec-1"}, + "profile_home": str(tmp_path / "profile"), + }), + encoding="utf-8", + ) + monkeypatch.setattr("cron.executions.adopt_claimed_execution", lambda _id: None) + run = Mock() + monkeypatch.setattr(scheduler, "run_one_job", run) + + assert scheduler._run_external_worker_payload(payload, ack) is False + + run.assert_not_called() + assert not ack.exists() + + +def test_launch_external_worker_uses_restart_safe_scope_and_acknowledges( + tmp_path, monkeypatch +): + import cron.scheduler as scheduler + + job = {"id": "job-1", "execution_id": "exec-1", "prompt": "work"} + monkeypatch.setattr(scheduler, "_get_hermes_home", lambda: tmp_path) + wrapped_commands = [] + + def wrap(command, *, unit_suffix): + wrapped_commands.append((command, unit_suffix)) + return ["scope", "--", *command] + + monkeypatch.setattr( + "tools.process_registry.restart_safe_gateway_child_argv", wrap + ) + + class FakeProcess: + returncode = None + + def poll(self): + return self.returncode + + spawned = [] + + def popen(command, **kwargs): + spawned.append((command, kwargs)) + ack_index = command.index("--ack-file") + 1 + Path(command[ack_index]).write_text( + json.dumps({"pid": 4321, "execution_id": "exec-1"}), + encoding="utf-8", + ) + return FakeProcess() + + monkeypatch.setattr(scheduler.subprocess, "Popen", popen) + monkeypatch.setenv("ANTHROPIC_API_KEY", "should-not-cross-profile") + from agent.secret_scope import set_multiplex_active + + set_multiplex_active(True) + try: + assert scheduler._launch_external_cron_worker(job) is True + finally: + set_multiplex_active(False) + assert wrapped_commands[0][1] == "cron-job-1-exec-exec-1" + assert spawned[0][0][0:2] == ["scope", "--"] + assert spawned[0][1]["start_new_session"] is True + assert "ANTHROPIC_API_KEY" not in spawned[0][1]["env"] + payload = json.loads((tmp_path / "cron/external-workers/exec-1.json").read_text()) + assert payload["multiplex_active"] is True + + +def test_launch_external_worker_stays_in_process_outside_managed_gateway( + monkeypatch, +): + import cron.scheduler as scheduler + + command_calls = [] + + def unchanged(command, *, unit_suffix): + command_calls.append((command, unit_suffix)) + return command + + monkeypatch.setattr( + "tools.process_registry.restart_safe_gateway_child_argv", unchanged + ) + popen = Mock() + monkeypatch.setattr(scheduler.subprocess, "Popen", popen) + + assert scheduler._launch_external_cron_worker( + {"id": "job-1", "execution_id": "exec-1"} + ) is False + assert command_calls + popen.assert_not_called() + + +def test_shared_run_path_hands_gateway_fire_to_external_worker(monkeypatch): + import cron.scheduler as scheduler + + launch = Mock(return_value=True) + run = Mock(side_effect=AssertionError("agent ran inside gateway")) + monkeypatch.setattr(scheduler, "_launch_external_cron_worker", launch) + monkeypatch.setattr(scheduler, "run_job", run) + job = {"id": "job-1", "execution_id": "exec-1"} + + assert scheduler.run_one_job(job, adapters={"discord": object()}) is True + + launch.assert_called_once_with(job) + run.assert_not_called() + + +def test_shared_run_path_creates_execution_before_managed_handoff(monkeypatch): + import cron.scheduler as scheduler + + created = Mock(return_value={"id": "exec-new"}) + launch = Mock(return_value=True) + monkeypatch.setattr(scheduler, "create_execution", created) + monkeypatch.setattr(scheduler, "_launch_external_cron_worker", launch) + job = {"id": "manual-job"} + + assert scheduler.run_one_job(job, adapters={"discord": object()}) is True + + created.assert_called_once_with("manual-job", source="direct") + assert job["execution_id"] == "exec-new" + launch.assert_called_once_with(job) + + +def test_lost_execution_start_cas_prevents_side_effects(monkeypatch): + import cron.scheduler as scheduler + + run = Mock(side_effect=AssertionError("side effect ran without ownership")) + monkeypatch.setattr(scheduler, "claim_dispatch", lambda _job_id: True) + monkeypatch.setattr(scheduler, "mark_execution_running", lambda _execution_id: None) + monkeypatch.setattr(scheduler, "run_job", run) + + assert scheduler.run_one_job( + {"id": "job-1", "execution_id": "exec-1"}, adapters=None + ) is True + run.assert_not_called() diff --git a/tests/cron/test_run_one_job.py b/tests/cron/test_run_one_job.py index 93bcef00a5..9b241de1b6 100644 --- a/tests/cron/test_run_one_job.py +++ b/tests/cron/test_run_one_job.py @@ -88,7 +88,7 @@ def test_run_one_job_exception_delivers_failure_alert(monkeypatch): s, "create_execution", lambda *_a, **_kw: {"id": "exec-j3"} ) monkeypatch.setattr(s, "claim_dispatch", lambda _job_id: True) - monkeypatch.setattr(s, "mark_execution_running", lambda _execution_id: None) + monkeypatch.setattr(s, "mark_execution_running", lambda _execution_id: {}) monkeypatch.setattr( s, "run_job", @@ -141,7 +141,7 @@ def test_run_one_job_exception_records_failure_alert_delivery_error(monkeypatch) s, "create_execution", lambda *_a, **_kw: {"id": "exec-j4"} ) monkeypatch.setattr(s, "claim_dispatch", lambda _job_id: True) - monkeypatch.setattr(s, "mark_execution_running", lambda _execution_id: None) + monkeypatch.setattr(s, "mark_execution_running", lambda _execution_id: {}) monkeypatch.setattr( s, "run_job", @@ -165,7 +165,7 @@ def _patch_escaped_failure(monkeypatch, delivered, *, exec_id, err): """Make run_job raise, and capture what the escape handler delivers.""" monkeypatch.setattr(s, "create_execution", lambda *_a, **_kw: {"id": exec_id}) monkeypatch.setattr(s, "claim_dispatch", lambda _job_id: True) - monkeypatch.setattr(s, "mark_execution_running", lambda _execution_id: None) + monkeypatch.setattr(s, "mark_execution_running", lambda _execution_id: {}) monkeypatch.setattr( s, "run_job", @@ -246,7 +246,7 @@ def test_run_one_job_exception_after_delivery_does_not_redeliver(monkeypatch): s, "create_execution", lambda *_a, **_kw: {"id": "exec-j5"} ) monkeypatch.setattr(s, "claim_dispatch", lambda _job_id: True) - monkeypatch.setattr(s, "mark_execution_running", lambda _execution_id: None) + monkeypatch.setattr(s, "mark_execution_running", lambda _execution_id: {}) monkeypatch.setattr( s, "run_job", @@ -288,7 +288,7 @@ def test_run_one_job_keyboard_interrupt_skips_delivery_and_reraises(monkeypatch) s, "create_execution", lambda *_a, **_kw: {"id": "exec-j6"} ) monkeypatch.setattr(s, "claim_dispatch", lambda _job_id: True) - monkeypatch.setattr(s, "mark_execution_running", lambda _execution_id: None) + monkeypatch.setattr(s, "mark_execution_running", lambda _execution_id: {}) monkeypatch.setattr( s, "run_job", diff --git a/tests/cron/test_script_claim_heartbeat.py b/tests/cron/test_script_claim_heartbeat.py index effb3f8dd4..a989b01457 100644 --- a/tests/cron/test_script_claim_heartbeat.py +++ b/tests/cron/test_script_claim_heartbeat.py @@ -564,7 +564,7 @@ def test_terminal_owner_cas_failure_marks_ledger_ownership_lost(monkeypatch): finish = MagicMock() monkeypatch.setattr(scheduler, "heartbeat_fire_claim", lambda *args, **kwargs: True) monkeypatch.setattr(scheduler, "claim_dispatch", lambda *_args, **_kwargs: True) - monkeypatch.setattr(scheduler, "mark_execution_running", lambda *_args: None) + monkeypatch.setattr(scheduler, "mark_execution_running", lambda *_args: {}) monkeypatch.setattr( scheduler, "run_job", diff --git a/tests/gateway/test_cron_delivery_housekeeping.py b/tests/gateway/test_cron_delivery_housekeeping.py new file mode 100644 index 0000000000..8e38071679 --- /dev/null +++ b/tests/gateway/test_cron_delivery_housekeeping.py @@ -0,0 +1,82 @@ +"""Gateway-independent draining of restart-safe cron deliveries.""" + +from contextlib import contextmanager +from types import SimpleNamespace + +import cron.scheduler as scheduler +import gateway.run as gateway_run + + +class _OneTickStopEvent: + def __init__(self): + self.waited = False + + def is_set(self): + return self.waited + + def wait(self, timeout=None): + self.waited = True + return True + + +def test_gateway_housekeeping_drains_cron_delivery_with_live_adapters(monkeypatch): + adapters = {"discord": object()} + loop = object() + calls = [] + monkeypatch.setattr( + scheduler, + "drain_delivery_queue", + lambda live_adapters, live_loop: calls.append((live_adapters, live_loop)), + raising=False, + ) + + gateway_run._start_gateway_housekeeping( + _OneTickStopEvent(), adapters=adapters, loop=loop, interval=0 + ) + + assert calls == [(adapters, loop)] + + +def test_multiplex_housekeeping_drains_each_profile_with_its_adapters( + tmp_path, monkeypatch +): + root_adapters = {} + secondary_adapters = {"telegram": "secondary"} + runner = SimpleNamespace( + config=SimpleNamespace(multiplex_profiles=True), + adapters=root_adapters, + _profile_adapters={"secondary": secondary_adapters}, + ) + secondary_home = tmp_path / "secondary" + calls = [] + + monkeypatch.setattr( + gateway_run, + "_handoff_watch_scopes", + lambda _runner: [(None, None), ("secondary", secondary_home)], + ) + + @contextmanager + def fake_scope(home): + calls.append(("scope", home)) + yield + + monkeypatch.setattr(gateway_run, "_profile_runtime_scope", fake_scope) + monkeypatch.setattr( + scheduler, + "drain_delivery_queue", + lambda adapters, loop: calls.append(("drain", adapters)), + ) + + gateway_run._start_gateway_housekeeping( + _OneTickStopEvent(), + adapters=root_adapters, + loop=object(), + interval=0, + runner=runner, + ) + + assert calls == [ + ("scope", secondary_home), + ("drain", secondary_adapters), + ] diff --git a/tests/hermes_cli/test_kanban_gateway_restart_handoff.py b/tests/hermes_cli/test_kanban_gateway_restart_handoff.py new file mode 100644 index 0000000000..d0efb6ed89 --- /dev/null +++ b/tests/hermes_cli/test_kanban_gateway_restart_handoff.py @@ -0,0 +1,178 @@ +"""Managed-gateway isolation for dispatcher-owned Kanban workers.""" + +from __future__ import annotations + +import json +import subprocess +import sys +import time +from pathlib import Path + +import pytest + +from hermes_cli import kanban_db as kb + + +@pytest.fixture +def worker_setup(tmp_path: Path, monkeypatch: pytest.MonkeyPatch) -> tuple[Path, kb.Task]: + root = tmp_path / ".hermes" + profile = root / "profiles" / "coder" + profile.mkdir(parents=True) + root.joinpath("config.yaml").write_text("{}\n", encoding="utf-8") + profile.joinpath("config.yaml").write_text("{}\n", encoding="utf-8") + monkeypatch.setenv("HERMES_HOME", str(root)) + monkeypatch.setattr(Path, "home", lambda: tmp_path) + monkeypatch.setattr(kb, "_resolve_hermes_argv", lambda: ["hermes"]) + + workspace = tmp_path / "candidate-worktree" + workspace.mkdir() + task = kb.Task( + id="t_candidate_restart", + title="activate candidate", + body=None, + assignee="coder", + status="running", + priority=0, + created_by="test", + created_at=1, + started_at=1, + completed_at=None, + workspace_kind="worktree", + workspace_path=str(workspace), + claim_lock="host:dispatcher", + claim_expires=999, + tenant=None, + branch_name="wt/t_candidate_restart", + current_run_id=23, + ) + return workspace, task + + +@pytest.mark.linux_only +def test_managed_gateway_worker_is_spawned_in_restart_safe_scope( + worker_setup: tuple[Path, kb.Task], monkeypatch: pytest.MonkeyPatch +) -> None: + workspace, task = worker_setup + captured_cmd: list[str] = [] + captured_env: dict[str, str] = {} + captured_cwd: str | None = None + + class FakeProc: + pid = 4242 + + def fake_popen(cmd, **kwargs): + nonlocal captured_cwd + captured_cmd.extend(cmd) + captured_env.update(kwargs.get("env") or {}) + captured_cwd = kwargs.get("cwd") + return FakeProc() + + monkeypatch.setenv("INVOCATION_ID", "managed-gateway-test") + monkeypatch.setenv("ANTHROPIC_API_KEY", "must-not-cross-profile") + monkeypatch.setattr("agent.secret_scope.is_multiplex_active", lambda: True) + monkeypatch.setattr(subprocess, "Popen", fake_popen) + monkeypatch.setattr("tools.process_registry._is_supervised_gateway_process", lambda: True) + monkeypatch.setattr("tools.process_registry._systemd_run_user_scope_available", lambda: True) + monkeypatch.setattr("tools.process_registry._worker_memory_max_bytes", lambda: 536_870_912) + monkeypatch.setattr("shutil.which", lambda name: "/usr/bin/systemd-run") + + assert kb._default_spawn(task, str(workspace)) == 4242 + assert captured_cmd[:4] == ["/usr/bin/systemd-run", "--user", "--scope", "--quiet"] + unit_index = captured_cmd.index("--unit") + assert captured_cmd[unit_index + 1] == "hermes-worker-kanban-t_candidate_restart-run-23" + assert "MemoryMax=536870912" in captured_cmd + separator = captured_cmd.index("--") + assert captured_cmd[separator + 1 : separator + 4] == ["hermes", "-p", "coder"] + assert captured_cwd == str(workspace) + assert captured_env["HERMES_KANBAN_TASK"] == task.id + assert captured_env["HERMES_KANBAN_RUN_ID"] == "23" + assert "ANTHROPIC_API_KEY" not in captured_env + + +@pytest.mark.linux_only +def test_managed_gateway_worker_spawn_fails_closed_without_scope( + worker_setup: tuple[Path, kb.Task], monkeypatch: pytest.MonkeyPatch +) -> None: + workspace, task = worker_setup + popen_calls: list[list[str]] = [] + monkeypatch.setenv("INVOCATION_ID", "managed-gateway-test") + monkeypatch.setattr(subprocess, "Popen", lambda cmd, **kwargs: popen_calls.append(list(cmd))) + monkeypatch.setattr("tools.process_registry._is_supervised_gateway_process", lambda: True) + monkeypatch.setattr("tools.process_registry._systemd_run_user_scope_available", lambda: False) + + with pytest.raises(RuntimeError, match="restart-safe systemd scope"): + kb._default_spawn(task, str(workspace)) + assert popen_calls == [] + + +@pytest.mark.linux_only +def test_managed_gateway_scope_builder_fails_closed_if_binary_disappears( + worker_setup: tuple[Path, kb.Task], monkeypatch: pytest.MonkeyPatch +) -> None: + workspace, task = worker_setup + monkeypatch.setenv("INVOCATION_ID", "managed-gateway-test") + monkeypatch.setattr("tools.process_registry._is_supervised_gateway_process", lambda: True) + monkeypatch.setattr("tools.process_registry._systemd_run_user_scope_available", lambda: True) + monkeypatch.setattr("shutil.which", lambda _name: None) + monkeypatch.setattr(subprocess, "Popen", lambda *_args, **_kwargs: pytest.fail("unsafe direct spawn")) + + with pytest.raises(RuntimeError, match="restart-safe systemd scope"): + kb._default_spawn(task, str(workspace)) + + +def test_standalone_dispatcher_keeps_direct_worker_spawn( + worker_setup: tuple[Path, kb.Task], monkeypatch: pytest.MonkeyPatch +) -> None: + workspace, task = worker_setup + captured_cmd: list[str] = [] + + class FakeProc: + pid = 4243 + + monkeypatch.setattr(subprocess, "Popen", lambda cmd, **kwargs: captured_cmd.extend(cmd) or FakeProc()) + monkeypatch.setattr("tools.process_registry._is_supervised_gateway_process", lambda: False) + monkeypatch.setattr( + "tools.process_registry._systemd_run_user_scope_available", + lambda: pytest.fail("scope probe must not run outside managed gateway"), + ) + + assert kb._default_spawn(task, str(workspace)) == 4243 + assert captured_cmd[:3] == ["hermes", "-p", "coder"] + + +@pytest.mark.linux_only +def test_real_user_systemd_scope_preserves_worker_context( + worker_setup: tuple[Path, kb.Task], monkeypatch: pytest.MonkeyPatch +) -> None: + from tools import process_registry + + if not process_registry._systemd_run_user_scope_available(): + pytest.skip("systemd-run --user --scope is unavailable on this host") + + workspace, task = worker_setup + receipt = workspace / "worker-receipt.json" + script = ( + "import json, os, pathlib, sys, time; " + "pathlib.Path(sys.argv[1]).write_text(json.dumps({" + "'pid': os.getpid(), 'cwd': os.getcwd(), " + "'task': os.environ.get('HERMES_KANBAN_TASK'), " + "'run': os.environ.get('HERMES_KANBAN_RUN_ID'), " + "'cgroup': pathlib.Path('/proc/self/cgroup').read_text()})); time.sleep(0.5)" + ) + monkeypatch.setattr(kb, "_resolve_hermes_argv", lambda: [sys.executable, "-c", script, str(receipt)]) + monkeypatch.setenv("INVOCATION_ID", "managed-gateway-test") + monkeypatch.setattr(process_registry, "_is_supervised_gateway_process", lambda: True) + + pid = kb._default_spawn(task, str(workspace)) + deadline = time.monotonic() + 5 + while not receipt.exists() and time.monotonic() < deadline: + time.sleep(0.05) + + assert receipt.exists() + payload = json.loads(receipt.read_text(encoding="utf-8")) + assert payload["pid"] == pid + assert payload["cwd"] == str(workspace) + assert payload["task"] == task.id + assert payload["run"] == "23" + assert ".scope" in payload["cgroup"] + assert "hermes-gateway.service" not in payload["cgroup"] diff --git a/tools/process_registry.py b/tools/process_registry.py index 8937c4cce8..4ec461df14 100644 --- a/tools/process_registry.py +++ b/tools/process_registry.py @@ -323,6 +323,35 @@ def _build_systemd_scope_argv( ] +def restart_safe_gateway_child_argv( + command: List[str], *, unit_suffix: str +) -> List[str]: + """Place a managed-systemd gateway child outside the gateway cgroup. + + Children that must survive an intentional gateway restart cannot rely on + ``start_new_session`` alone: systemd still kills every process in the + service cgroup. In that topology, require a transient user scope and fail + closed if it cannot be established. Standalone processes, non-systemd + supervisors, and non-Linux hosts retain the direct command. + """ + if _IS_WINDOWS: + return command + if not _is_supervised_gateway_process() or not os.environ.get("INVOCATION_ID"): + return command + if not _systemd_run_user_scope_available(): + raise RuntimeError( + "cannot create restart-safe systemd scope for gateway child: " + "systemd-run --user --scope is unavailable" + ) + scoped = _build_systemd_scope_argv(command, unit_suffix=unit_suffix) + if scoped == command: + raise RuntimeError( + "cannot create restart-safe systemd scope for gateway child: " + "systemd-run disappeared after the availability probe" + ) + return scoped + + def _stop_systemd_unit(unit_name: str) -> bool: """Stop a transient systemd user scope by unit name.