Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
65 changes: 48 additions & 17 deletions cron/scheduler_provider.py
Original file line number Diff line number Diff line change
Expand Up @@ -710,13 +710,22 @@ def _start_multiplex(
[p[0] if isinstance(p, tuple) else p for p in profile_homes],
)

# Recovery + initial heartbeat for every profile.
# A profile may have been deleted since this snapshot was taken;
# never recreate a deleted home's cron workspace via the heartbeat
# below (#47368).
# One profile's broken store (corrupt executions.db, unreadable
# cron dir) must not abort startup for every other profile (#74878).
for entry in _existing_profile_homes(profile_homes):
def eligible_profile_homes():
homes = _existing_profile_homes(profile_homes)
if profile_gate is None:
return homes
return [
entry
for entry in homes
if profile_gate(
entry[0] if isinstance(entry, tuple) else None,
entry[1] if isinstance(entry, tuple) else entry,
)
]

initialized_homes: set[str] = set()

def initialize_profile(entry) -> bool:
home = entry[1] if isinstance(entry, tuple) else entry
home_token = set_hermes_home_override(str(home))
try:
Expand All @@ -729,16 +738,47 @@ def _start_multiplex(
home,
)
record_ticker_heartbeat()
return True
except BaseException as e:
logger.error(
"Cron startup recovery error for profile at %s: %s",
home,
e,
exc_info=True,
)
return False
finally:
reset_hermes_home_override(home_token)

def initialize_eligible_profiles(entries):
eligible_keys = {
str(entry[1] if isinstance(entry, tuple) else entry)
for entry in entries
}
# Losing eligibility ends this process's ownership epoch. If the
# profile becomes eligible again, recover work abandoned by the
# previous owner before this scheduler resumes ticking it.
initialized_homes.intersection_update(eligible_keys)
for entry in entries:
home = entry[1] if isinstance(entry, tuple) else entry
key = str(home)
if key not in initialized_homes and initialize_profile(entry):
initialized_homes.add(key)
return [
entry
for entry in entries
if str(entry[1] if isinstance(entry, tuple) else entry)
in initialized_homes
]

# Recovery + initial heartbeat for every initially eligible profile.
# A profile may have been deleted since this snapshot was taken;
# never recreate a deleted home's cron workspace via the heartbeat
# below (#47368). One profile's broken store (corrupt executions.db,
# unreadable cron dir) must not abort startup for every other profile
# (#74878).
initialize_eligible_profiles(eligible_profile_homes())

consecutive_failures = 0
while not stop_event.is_set():
ok = False
Expand All @@ -747,16 +787,7 @@ def _start_multiplex(
# Worst per-profile failure this cycle (fd exhaustion wins) so the
# #87644 backoff/reclaim is applied once per cycle, not per profile.
_cycle_exc: BaseException | None = None
cycle_homes = _existing_profile_homes(profile_homes)
if profile_gate is not None:
cycle_homes = [
entry
for entry in cycle_homes
if profile_gate(
entry[0] if isinstance(entry, tuple) else None,
entry[1] if isinstance(entry, tuple) else entry,
)
]
cycle_homes = initialize_eligible_profiles(eligible_profile_homes())
try:
if can_dispatch is not None and not can_dispatch():
logger.debug("Cron dispatch paused while gateway drains existing work")
Expand Down
22 changes: 16 additions & 6 deletions hermes_cli/web_server.py
Original file line number Diff line number Diff line change
Expand Up @@ -303,16 +303,26 @@ def _start_desktop_cron_ticker(stop_event: "threading.Event", interval: int = 60
profile_homes = list(profiles_to_serve(multiplex=True))
if len(profile_homes) > 1:
start_kwargs["profile_homes"] = profile_homes
# Stand down, per tick, for any profile whose OWN gateway is
# running: that gateway ticks it with live adapters, and the
# tick-lock race otherwise lets this adapter-less ticker win
# and deliver the job through the standalone path (#100489).
# Stand down for any profile already owned by a gateway: either
# its own process or the live default-profile multiplexer. That
# gateway ticks it with live adapters, and the tick-lock race
# otherwise lets this adapter-less ticker win and deliver the
# job through the standalone path (#100489, #102156).
# Evaluated every cycle so a gateway starting/stopping later
# is picked up without a dashboard restart.
from hermes_cli.profiles import _check_gateway_running
from hermes_cli.profiles import (
_check_gateway_running,
_served_by_running_multiplexer,
)

start_kwargs["profile_gate"] = (
lambda _name, home: not _check_gateway_running(Path(home))
lambda name, home: not (
_check_gateway_running(Path(home))
or (
name is not None
and _served_by_running_multiplexer(name)
)
)
)
from hermes_logging import enable_profile_log_routing

Expand Down
120 changes: 115 additions & 5 deletions tests/cron/test_cron_multiplex_desktop_ticker_scope.py
Original file line number Diff line number Diff line change
Expand Up @@ -99,23 +99,132 @@ def _tick(*args, **kwargs):

assert not thread.is_alive()
assert set(ticked) == {str(orphan)}
# The gated profile gets no tick-loop success marker either: its own
# gateway owns that status surface.
# The gated profile gets no startup heartbeat or tick-loop success marker:
# its gateway owns that status surface and the scheduler never enters it.
assert not (own_gateway / "cron" / "ticker_heartbeat").exists()
assert not (own_gateway / "cron" / "ticker_last_success").exists()
assert (orphan / "cron" / "ticker_last_success").exists()


def test_desktop_ticker_gates_on_profile_gateway_running(tmp_path, monkeypatch):
"""The desktop ticker wires the gate to ``_check_gateway_running``."""
def test_multiplex_ticker_recovers_each_new_ownership_epoch(tmp_path):
from cron.scheduler_provider import InProcessCronScheduler
from hermes_constants import get_hermes_home

home = tmp_path / "handoff"
(home / "cron").mkdir(parents=True)
stop = threading.Event()
events: list[tuple[str, str]] = []
gate_results = iter((False, True, False, True))

def _recover():
events.append(("recover", str(get_hermes_home())))
return 0

def _tick(*args, **kwargs):
events.append(("tick", str(get_hermes_home())))
if sum(kind == "tick" for kind, _ in events) == 2:
stop.set()
return 0

provider = InProcessCronScheduler()
with (
patch.object(provider, "recover_interrupted", side_effect=_recover),
patch("cron.scheduler.tick", side_effect=_tick),
):
thread = threading.Thread(
target=provider.start,
args=(stop,),
kwargs={
"interval": 0,
"profile_homes": [("handoff", home)],
"profile_gate": lambda name, candidate: next(gate_results),
},
daemon=True,
)
thread.start()
thread.join(timeout=5)
stop.set()
thread.join(timeout=5)

assert not thread.is_alive()
assert events == [
("recover", str(home)),
("tick", str(home)),
("recover", str(home)),
("tick", str(home)),
]


def test_multiplex_ticker_retries_failed_ownership_initialization(tmp_path):
from cron.scheduler_provider import InProcessCronScheduler
from hermes_constants import get_hermes_home

home = tmp_path / "handoff"
(home / "cron").mkdir(parents=True)
stop = threading.Event()
events: list[tuple[str, str]] = []
attempts = 0

def _recover():
nonlocal attempts
attempts += 1
events.append((f"recover-{attempts}", str(get_hermes_home())))
if attempts == 1:
raise OSError("transient ledger lock")
return 0

def _tick(*args, **kwargs):
events.append(("tick", str(get_hermes_home())))
stop.set()
return 0

provider = InProcessCronScheduler()
with (
patch.object(provider, "recover_interrupted", side_effect=_recover),
patch("cron.scheduler.tick", side_effect=_tick),
):
thread = threading.Thread(
target=provider.start,
args=(stop,),
kwargs={
"interval": 0,
"profile_homes": [("handoff", home)],
"profile_gate": lambda name, candidate: True,
},
daemon=True,
)
thread.start()
thread.join(timeout=5)
stop.set()
thread.join(timeout=5)

assert not thread.is_alive()
assert events == [
("recover-1", str(home)),
("recover-2", str(home)),
("tick", str(home)),
]


def test_desktop_ticker_gates_on_every_profile_gateway_owner(tmp_path, monkeypatch):
"""The desktop ticker stands down for direct and multiplex gateway owners."""
from hermes_cli import web_server

homes = [("default", tmp_path / "default"), ("ops", tmp_path / "ops")]
homes = [
("default", tmp_path / "default"),
("ops", tmp_path / "ops"),
("mux", tmp_path / "mux"),
]
monkeypatch.setattr(
"hermes_cli.profiles.profiles_to_serve", lambda multiplex=False: list(homes)
)
monkeypatch.setattr(
"hermes_cli.profiles._check_gateway_running", lambda home: home.name == "ops"
)
monkeypatch.setattr(
"hermes_cli.profiles._served_by_running_multiplexer",
lambda name: name == "mux",
)
captured = {}

class _Provider:
Expand All @@ -137,3 +246,4 @@ def start(self, stop_event, **kwargs):
assert gate is not None, "desktop ticker did not install a profile gate"
assert gate("default", tmp_path / "default") is True
assert gate("ops", tmp_path / "ops") is False
assert gate("mux", tmp_path / "mux") is False
9 changes: 6 additions & 3 deletions tests/cron/test_scheduler_provider.py
Original file line number Diff line number Diff line change
Expand Up @@ -889,6 +889,9 @@ def _tick(*args, **kwargs):
thread.join(timeout=5)

assert not thread.is_alive()
assert recovery_homes == [str(failing_home), str(healthy_home)]
# The failing profile stays in rotation: its ledger may still hold jobs.
assert set(tick_homes) == {str(failing_home), str(healthy_home)}
assert recovery_homes[:2] == [str(failing_home), str(healthy_home)]
assert recovery_homes.count(str(failing_home)) >= 2
assert recovery_homes.count(str(healthy_home)) == 1
# Recovery failures are retried without blocking a healthy sibling, but
# the failed profile must not tick until its ledger can be recovered.
assert set(tick_homes) == {str(healthy_home)}
Loading